From f637433d01ea4d72fb7409a61337c77404fb4ea6 Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 7 Aug 2026 13:52:59 +0000 Subject: [PATCH 1/2] docs: document how a task result is derived A task function may return a plain value, a Promise or an Observable, but the way an Observable is consumed was not documented anywhere. The first emitted value becomes the task's result and the subscription is closed right after it, so later emissions never occur and the source is cancelled. The cancellation is the surprising part: a request that reports progress resolves with its first progress event and is then aborted. Document this on the Task type, on register() and in the README, together with the last() and toArray() alternatives and the caveat that last() never emits for a source that does not complete. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_018q1BpbY6mjVrbgaPBQDA5Z --- .changeset/document-task-results.md | 8 +++++ README.md | 47 +++++++++++++++++++++++++++++ src/resolver.interface.ts | 33 ++++++++++++++++++++ src/resolver.ts | 6 +++- 4 files changed, 93 insertions(+), 1 deletion(-) create mode 100644 .changeset/document-task-results.md diff --git a/.changeset/document-task-results.md b/.changeset/document-task-results.md new file mode 100644 index 0000000..c27bf1d --- /dev/null +++ b/.changeset/document-task-results.md @@ -0,0 +1,8 @@ +--- +'@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 +never occur and the source is cancelled. The `Task` type, `register()` and the README now describe +this, including how to use `last()` or `toArray()` when a different value is needed. diff --git a/README.md b/README.md index 6150b3e..417bf83 100644 --- a/README.md +++ b/README.md @@ -158,6 +158,51 @@ 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 +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 never occur and the source is cancelled: + +```typescript +// Only the first value is used, the source is cancelled afterwards +fn: () => of(1, 2, 3); // { data: 1 } + +// A request that reports progress resolves with its first progress event, +// and the request itself is cancelled +fn: () => http.get('/user', { 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()` only emits once the source completes. Applying it 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 { EmptyTaskError, isError } from '@robinw151/resolver'; + +const result = await lastValueFrom(new Resolver().register({ id: 'empty', fn: () => EMPTY }).resolve()); + +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: @@ -165,6 +210,8 @@ Tasks can return either successful data or errors. The resolver handles both cas - 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. diff --git a/src/resolver.interface.ts b/src/resolver.interface.ts index 2e362a7..09b55a5 100644 --- a/src/resolver.interface.ts +++ b/src/resolver.interface.ts @@ -24,8 +24,41 @@ export type ResolverResultWithLoadingState = 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 { + /** + * 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 never occur and the source is 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()` never emits for a source that does not complete, which leaves the task, + * and therefore the whole resolution, pending indefinitely. + */ fn: (args: TArgs, globalArgs: TGlobalArgs) => TResult | Promise | Observable; } diff --git a/src/resolver.ts b/src/resolver.ts index e67c2d3..ea37eae 100644 --- a/src/resolver.ts +++ b/src/resolver.ts @@ -143,7 +143,11 @@ export class Resolver { * * @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 never occur and the source is cancelled. Pipe the source with `last()` or + * `toArray()` when a different value is needed. * @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. * From d6f7ffe0a47570d34a41788cdb290478f78efb7c Mon Sep 17 00:00:00 2001 From: Claude Date: Fri, 7 Aug 2026 15:05:19 +0000 Subject: [PATCH 2/2] docs: correct observable contract wording and examples Address review feedback: - Fix the HttpClient progress example. With reportProgress alone the request is still observed as a body, so it emits only the response and the example contradicted the point it was making. Observed as events it emits HttpEventType.Sent first, which is the behaviour worth warning about. - Say that later emissions are never observed and that sources honoring unsubscription are cancelled, instead of claiming the source is always cancelled. Unsubscribing only stops upstream work when the source implements teardown. - Extend the completion caveat to toArray(), which carries the same risk as last(). - Complete the imports in the README examples. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_018q1BpbY6mjVrbgaPBQDA5Z --- .changeset/document-task-results.md | 5 +++-- README.md | 17 ++++++++++------- src/resolver.interface.ts | 10 ++++++---- src/resolver.ts | 5 +++-- 4 files changed, 22 insertions(+), 15 deletions(-) diff --git a/.changeset/document-task-results.md b/.changeset/document-task-results.md index c27bf1d..c28b3e7 100644 --- a/.changeset/document-task-results.md +++ b/.changeset/document-task-results.md @@ -4,5 +4,6 @@ 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 -never occur and the source is cancelled. The `Task` type, `register()` and the README now describe -this, including how to use `last()` or `toArray()` when a different value is needed. +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. diff --git a/README.md b/README.md index 417bf83..104466c 100644 --- a/README.md +++ b/README.md @@ -163,21 +163,23 @@ const resolver = new Resolver() 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 never occur and the source is cancelled: +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 source is cancelled afterwards +// Only the first value is used, the subscription is closed afterwards fn: () => of(1, 2, 3); // { data: 1 } -// A request that reports progress resolves with its first progress event, -// and the request itself is cancelled -fn: () => http.get('/user', { reportProgress: true }); +// 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: @@ -189,12 +191,13 @@ 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()` only emits once the source completes. Applying it 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. +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 { EmptyTaskError, isError } from '@robinw151/resolver'; +import { EMPTY, lastValueFrom } from 'rxjs'; +import { EmptyTaskError, isError, Resolver } from '@robinw151/resolver'; const result = await lastValueFrom(new Resolver().register({ id: 'empty', fn: () => EMPTY }).resolve()); diff --git a/src/resolver.interface.ts b/src/resolver.interface.ts index 09b55a5..635c146 100644 --- a/src/resolver.interface.ts +++ b/src/resolver.interface.ts @@ -44,8 +44,9 @@ export interface Task of(1, 2, 3).pipe(toArray()); // { data: [1, 2, 3] } * ``` * - * Note that `last()` never emits for a source that does not complete, which leaves the task, - * and therefore the whole resolution, pending indefinitely. + * 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 | Observable; } diff --git a/src/resolver.ts b/src/resolver.ts index ea37eae..1eefc8b 100644 --- a/src/resolver.ts +++ b/src/resolver.ts @@ -146,8 +146,9 @@ export class Resolver { * - `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 never occur and the source is cancelled. Pipe the source with `last()` or - * `toArray()` when a different value is needed. + * 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. *