Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions .changeset/document-task-results.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
---
'@robinw151/resolver': patch
---

Document how a task's result is derived from what its function returns. For an `Observable` the
first emitted value becomes the result and the subscription is closed afterwards, so later emissions
are never observed and sources that honor unsubscription are cancelled. The `Task` type,
`register()` and the README now describe this, including how to use `last()` or `toArray()` when a
different value is needed, and that both only emit once the source completes.
50 changes: 50 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -158,13 +158,63 @@ const resolver = new Resolver()
);
```

## Task Results

Every task produces exactly one result. A task function may return a plain value, a `Promise` or an `Observable`:

```typescript
import { of } from 'rxjs';

const resolver = new Resolver()
.register({ id: 'value', fn: () => 1 })
.register({ id: 'promise', fn: () => Promise.resolve(2) })
.register({ id: 'observable', fn: () => of(3) });
```

For an `Observable` the **first emitted value** becomes the task's result. The subscription is closed right after that value, so later emissions are never observed, and sources that honor unsubscription are cancelled:

```typescript
// Only the first value is used, the subscription is closed afterwards
fn: () => of(1, 2, 3); // { data: 1 }

// Observed as events, an HttpClient request emits `HttpEventType.Sent` first,
// so the task resolves with that event and the request is cancelled
fn: () => http.get('/user', { observe: 'events', reportProgress: true });
```

Pipe the source when a different value is needed:

```typescript
import { last, toArray } from 'rxjs';

fn: () => of(1, 2, 3).pipe(last()); // { data: 3 }
fn: () => of(1, 2, 3).pipe(toArray()); // { data: [1, 2, 3] }
```

Keep in mind that `last()` and `toArray()` only emit once the source completes. Applying either to a source that never completes leaves the task, and therefore the whole resolution, pending indefinitely. The default behavior has no such risk, which is why an infinite source such as `interval(1000)` resolves with its first value instead of hanging.

A source that completes without emitting any value cannot produce a result. Such a task resolves with an [`EmptyTaskError`](#error-handling) instead of blocking the resolution:

```typescript
import { EMPTY, lastValueFrom } from 'rxjs';
import { EmptyTaskError, isError, Resolver } from '@robinw151/resolver';

const result = await lastValueFrom(new Resolver().register({ id: 'empty', fn: () => EMPTY }).resolve());
Comment thread
coderabbitai[bot] marked this conversation as resolved.

if (isError(result.tasks.empty)) {
console.log(result.tasks.empty.error instanceof EmptyTaskError); // true
}
```

## Error Handling

Tasks can return either successful data or errors. The resolver handles both cases gracefully:

- Successful tasks return `{ data: TResult }`
- Failed tasks return `{ error: unknown }`

A task whose `Observable` completes without emitting a value fails with an `EmptyTaskError`, which is exported from the package and carries the `taskId` of the task that produced it.

## Global Arguments

The resolver supports global arguments that are passed to all task functions during execution. This is useful for sharing configuration, API keys, or other context across all tasks.
Expand Down
35 changes: 35 additions & 0 deletions src/resolver.interface.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,43 @@ export type ResolverResultWithLoadingState<TGlobalArgs, TResult> = Observable<
>;

// Task
/**
* A single unit of work that can be registered with a resolver.
*
* Every task produces exactly one result. How that result is obtained depends on what `fn`
* returns, but the outcome is always a single `TaskResult`, either `{ data }` or `{ error }`.
*/
export interface Task<TId extends string, TArgs = unknown, TGlobalArgs = unknown, TResult = unknown> {
/**
* Unique identifier of the task within a resolver.
*/
readonly id: TId;

/**
* The function that performs the work of the task.
*
* It receives the results of the task's dependencies as its first argument and the resolver's
* global arguments as its second argument, and may return a plain value, a `Promise` or an
* `Observable`.
*
* For an `Observable` the **first emitted value** becomes the task's result and the subscription
* is then closed, so later emissions are never observed and sources that honor unsubscription
* are cancelled. A source that completes without emitting resolves the task with an
* `EmptyTaskError`.
*
* Pipe the source when a different value is needed, for example `last()` to wait for the final
* value of a completing source, or `toArray()` to collect every value:
*
* ```typescript
* fn: () => of(1, 2, 3); // { data: 1 }
* fn: () => of(1, 2, 3).pipe(last()); // { data: 3 }
* fn: () => of(1, 2, 3).pipe(toArray()); // { data: [1, 2, 3] }
* ```
*
* Note that `last()` and `toArray()` only emit once the source completes. Applying either to a
* source that does not complete leaves the task, and therefore the whole resolution, pending
* indefinitely.
*/
fn: (args: TArgs, globalArgs: TGlobalArgs) => TResult | Promise<TResult> | Observable<TResult>;
}

Expand Down
7 changes: 6 additions & 1 deletion src/resolver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -143,7 +143,12 @@ export class Resolver<TGlobalArgs = unknown, TResult = object> {
*
* @param task - The task configuration object containing:
* - `id`: Unique identifier for the task
* - `fn`: Function that executes the task, receiving resolved dependencies and global args
* - `fn`: Function that executes the task, receiving resolved dependencies and global args.
* It may return a plain value, a `Promise` or an `Observable`. For an `Observable` the first
* emitted value becomes the task's result and the subscription is then closed, so later
* emissions are never observed and sources that honor unsubscription are cancelled. Pipe the
* source with `last()` or `toArray()` when a different value is needed, keeping in mind that
* both only emit once the source completes.
* @param dependencies - Optional array of task IDs that this task depends on. These tasks
* must be registered before this task and will be resolved before this task executes.
*
Expand Down