feat(durable): add durable task runtime

Adds defineTask, registry-resolved phase handlers, and a scheduler that
decides every task transition in one callback on the Session line:
reservation with migration, dependencies, precedence rules in a step
before each phase, handover to replaced definitions, direct-task abort
(mark, signal, join, fresh abort invocation), orphaning of blocked tasks,
open reconciliation, and close that joins invocations without writing
outcomes. runtime.commit() callbacks return the typed next state; Tx no
longer exposes setTask. Adds Harness.resume, getTask, waitForTask,
abortTask, and task-aware waitForIdle; RegistryReader.subscribe wakes
the scheduler.
This commit is contained in:
Mario Zechner
2026-09-28 23:55:24 +02:00
parent 540e174c72
commit cb7969d212
22 changed files with 3316 additions and 134 deletions
+5
View File
@@ -11,6 +11,10 @@
- Made task conversation membership immutable after task creation.
- Added `ConversationQuery` to Storage and transaction conversation scans.
- Added required `StoredDocument.deltasSinceBase` to Storage document reads.
- Task definitions now require an exhaustive `phases` map and an `abort` handler; define them with `defineTask()`.
- `RegistryReader` now requires `subscribe()`.
- Removed `Tx.setTask()`; a task changes its own state by returning the next state from its `runtime.commit()` callback.
- `Session.subscribeClose()` listeners now run synchronously when close begins, after admission is sealed.
### Added
@@ -23,6 +27,7 @@
- Added `Harness.open()` with lazy root creation, atomic conversation creation and forks with `init`, conversation-bound commits, fork-aware entry pagination, model context derivation, and the built-in `ConversationConfig` document with model, thinking level, and active tool accessors.
- Added `createRegistry()` for tools, tool wrappers, hooks, tasks, and system prompt sections with batched publication and stable keyed ordering.
- Added `defineEntry()` typed entry kinds.
- Added the durable task runtime: `defineTask()`, registry-resolved phase handlers with checkpoint progress rules decided on the Session line, migration at reservation, typed runtime commits, memos, `sleep()`, and invocation-owned watches, plus `Harness.resume()`, `getTask()`, `waitForTask()`, `abortTask()`, and task-aware `waitForIdle()` on the Harness and conversations. Open reconciles running tasks to pending; tasks without a fitting definition stay blocked until registration, and aborting them settles them as `orphaned`.
### Fixed
+37 -10
View File
@@ -11,7 +11,7 @@ facades, membranes, document routing, view projection, events, or clone chains.
- Obsolete `pico` and `pico4` prototypes were removed.
- `pico3` remains.
- Packages 1–13 are implemented in `packages/durable`; Package 10 was already satisfied by Chord's canonical structural diff implementation.
- Packages 1–14 are implemented in `packages/durable`; Package 10 was already satisfied by Chord's canonical structural diff implementation.
## 1. Records, cursors, and memory tables
@@ -334,9 +334,8 @@ and reopen.
Implement `defineTask`, exhaustive phase maps, full checkpoint replacement,
kind migration at reservation, runtime commits and memos, invocation lifetime
gates, scheduler reservation, dependencies, terminal outcomes, typed waits,
holds, quiescence, and joins. Complete the task-facing §2.2 methods: `resume`,
`suspend`, `hold`, `getTask`, `waitForTask`, `markTask`, `abortTask`, and the
task-aware portion of idle waits. Task definitions come from the registry snapshot;
and joins. Complete the task-facing §2.2 methods: `resume`, `getTask`,
`waitForTask`, `abortTask`, and the task-aware portion of idle waits. Task definitions come from the registry snapshot;
add `RegistryReader.subscribe()` so registry changes wake the scheduler.
Include the execution-critical abort core: durable direct-task marks,
@@ -350,6 +349,20 @@ orphans a task. Implement per-phase registry snapshots, refresh at
every normal phase boundary, and hand over when the task definition object changed
and the new definition can reserve the task.
Deferred to later packages; do not stub them in Package 14:
- Orphaning writes the terminal `orphaned` record and retires task documents.
Unanswered input submissions and clearing turn control are added in Package
15, when submissions and turn control exist.
- `TaskRuntime.hooks` (the hook runner) is added in Package 16 with hook
dispatch. Package 14 adds `runtime.registry` (the phase's registry snapshot)
and `runtime.models`.
- `TaskRuntime.conversation()` returns a `ConversationHandle`, whose `submit`
needs Package 15 and whose `abort`/`waitForIdle` need ownership traversal; it
is added with the invocation-bound owned APIs in Package 18.
- Idle waits count live non-background tasks directly (Harness-wide or in the
addressed conversation). Traversal through owned conversations and background
boundaries replaces this in Package 18.
Use a fake two-phase external effect to test the real runtime. Acceptance is an
intent/effect/outcome task interrupted after intent, closed, reopened on the same
storage, and safely resumed to a durable terminal receipt. Also test unchanged-
@@ -357,16 +370,24 @@ checkpoint faulting, same-phase progress, cancellation precedence, thrown
handlers, dependencies, result values and entry IDs, first-writer-wins memos,
terminal checkpoint/memo removal, task-document retirement, close/reopen without
abort marks or fabricated outcomes, no fresh phase/abort dispatch while closing,
watch cleanup, holds and quiescence with eligible work, blocked tasks unblocked
watch cleanup, blocked tasks unblocked
by later registration, migration failure leaving the record unchanged, handover
after a same-name task replacement, no handover to a missing or incompatible
definition, no overlap between predecessor and successor
invocations, abort of a blocked task settling as `orphaned` with full cleanup,
mark-only versus signalling abort, and crashes at every direct-task abort stage.
run commits rejected after the mark, and crashes at every direct-task abort stage.
Implemented design notes: every task transition is decided by one callback
serialized on the Session line (reservation, abort marks, runtime commits, and the
step before each phase). `Tx` has no task replacement; `runtime.commit()`
callbacks return the typed next state. `suspend()`, `hold()`, `quiescent()`, and
a public `markTask()` were dropped from §2.2.
## 15. First runnable no-tool chat turn
Implement the smallest real input-to-answer vertical path. Define the final
Implement the smallest real input-to-answer vertical path. Extend Package 14's
`orphaned` settlement to mark affected input submissions unanswered and clear
matching turn control in the same commit. Define the final
inbox, turn-control, and generation presentation documents needed by this path;
do not use provisional kinds or schemas. Add input `Submission` admission,
request-ID deduplication, reacquisition and waiting, idle placement, active-turn
@@ -406,8 +427,8 @@ awaits its own input submission rather than global idle.
## 16. First coding-agent tool turn
Implement hook dispatch (Session-wide and scoped to a conversation or its owned
subtree) and wire the real generation → tool tasks → post-tools → generation
Implement hook dispatch and `TaskRuntime.hooks`, deferred from Package 14
(Session-wide and scoped to a conversation or its owned subtree) and wire the real generation → tool tasks → post-tools → generation
chain. Implement offered-set checks, tool pinning from the phase snapshot,
declaration and argument validation against both the offered declaration and
the pinned implementation, hook composition, durable execution intent, stored
@@ -441,9 +462,15 @@ self-head cuts, successor turns, queued reset/handoff, and every terminal cleanu
Successful inputs still require an answer; writes settle on placement and never
start generation.
Add `harness.blockedTasks()`, deferred from Package 14: the pending tasks this
process cannot run and why (`missing_task`, `task_too_old`, `migration_failed`
with the stored and registered versions and the migration error). It is derived
from task records, the current registry snapshot, and scheduler memory, and is
never persisted.
Specify and implement the Harness activity view:
the active conversations, notifications when a conversation becomes active or
idle, and Harness quiescence, all derived from committed turn-control and task
idle, all derived from committed turn-control and task
state.
Define any remaining built-in preference/presentation documents once with final
+97 -40
View File
@@ -411,9 +411,6 @@ interface Conversation {
interface Harness extends Session {
resume(): void;
suspend(context: Context): Promise<void>;
quiescent(): boolean;
hold(): () => void;
root(
context: Context,
@@ -433,7 +430,6 @@ interface Harness extends Session {
conversationId?: ConversationId,
): Promise<"aborted" | "already_placed" | "settled" | "not_found">;
abortTask(id: TaskId, context: Context): Promise<"marked" | "terminal">;
markTask(id: TaskId, context: Context): Promise<"marked" | "terminal">;
waitForTask<R>(id: TaskId<R>, context: Context): Promise<SettledTask<R>>;
waitForIdle(context: Context): Promise<void>;
}
@@ -473,11 +469,9 @@ write occurs. Reopen finds the same root by its reserved ID. A conversation with
no configured model produces a durable `no_model` generation failure.
`resume()` is idempotent while running and only enables scheduling. It does not
repeat open-time reconciliation. `suspend()` is terminal for that Harness
instance and follows the close semantics below; `resume()` after suspend/close
rejects. `quiescent()` means no task invocation is currently executing; eligible
or delayed durable tasks may still exist. `hold()` is available only while
quiescent and pauses reservation until its idempotent release function runs.
repeat open-time reconciliation, and it throws after close. Work that must happen
before any task runs, such as registration or seeding, happens before
`resume()`.
`createConversation({ ownership, init, input })` and
`fork(at, { ownership, init })` commit atomically: the conversation, its
@@ -567,10 +561,10 @@ handoff write and then resolves; while busy, placement follows section 6 and may
occur later. Observe its placement through the conversation watch. An idle wait
does not guarantee placement of queued passive writes.
`markTask()` commits `abortRequested` and the durable foreground-subtree
cascade; it neither signals nor joins active invocations. The scheduler notices
the marks on its next drain. `abortTask()` also signals and joins the active run
before starting the abort invocation. `Conversation.abort()` withdraws queued
`abortTask()` commits `abortRequested` and the durable foreground-subtree
cascade, then signals and joins the active run; the scheduler then starts the
abort invocation. Marking without signalling is internal: the cascade marks
owned tasks in the same commit. `Conversation.abort()` withdraws queued
input submissions, marks non-background tasks selected by ordinary ownership
traversal, signals them, and resolves only after that scope is ordinarily idle.
Passive writes and background subtrees survive. Conversation idle means no live
@@ -581,9 +575,8 @@ wait aborts only that waiter.
Conversation handles are stateless; compare them by `id`. Hosts discover
conversations through lookups and scans. An activity view that lists active
conversations and reports conversations becoming active or idle, plus Harness
quiescence, is specified with the task runtime and turn control (Packages 14
and 17); there is no creation listener.
conversations and reports conversations becoming active or idle is specified
with turn control (Package 17); there is no creation listener.
`submit()` returns after durable admission, not settlement. An input submission
creates a user message with the admission timestamp; `whenBusy` defaults to
@@ -610,7 +603,7 @@ uses the same serialized asynchronous, bounded-buffer contract as `watchDoc()`
in section 9.2. Neither carries semantic events or owns a second persistence
authority.
`close()` is equivalent to `suspend()` for v1. It seals mutation admission and
`close()` seals mutation admission and
task reservation, signals invocations, and stops future watch deliveries. Outside
the Session line it lets already-admitted storage commits settle, joins
task/tool/hook invocations, then closes states and storage. Already-running watch
@@ -869,7 +862,6 @@ interface Tx {
createTask<I, S extends { phase: string }, R, H extends object>(
task: Task<I, S, R, H>, input: I, options?: TaskOptions,
): Promise<TaskId<R>>;
setTask(value: TaskRecord<JsonValue, JsonValue, JsonValue>): void;
doc<T extends JsonObject>(token: SessionDocToken<T>): Promise<Draft<T>>;
doc<T extends JsonObject>(token: ConversationDocToken<T>, conversationId: ConversationId): Promise<Draft<T>>;
@@ -1205,23 +1197,24 @@ Session-line operation so it never exposes a mixture from one commit.
A Session commit callback may be asynchronous. It owns the Session mutation
line through callback execution, preparation, storage settlement, committed
baseline adoption, and publication enqueue. Commit/close observers run on the
baseline adoption, and publication enqueue. Commit observers run on the
line; document-state and watch user callbacks run later.
External model, process, tool, network, and human effects run outside it.
`subscribeCommits()` observes complete immutable publications synchronously
after adoption, and `subscribeClose()` observes Session closure synchronously.
Both run on the Session line and return idempotent disposers; their listeners
must not throw, block, or call Session APIs. Document-state subscribers and watch
`subscribeCommits()` observes complete immutable publications synchronously on
the line after adoption. `subscribeClose()` observes close synchronously when it
begins, after admission is sealed; the Harness stops watches and signals task
invocations there. Both return idempotent disposers; their listeners must not
throw, block, or call Session APIs. Document-state subscribers and watch
listeners still run later, off the line.
```ts
await session.commit(async tx => {
const task = await tx.task(taskId); // table read
const conversation = await tx.conversation(conversationId); // table read
const live = await tx.doc(LiveDoc, conversationId);
await tx.appendEntry(conversationId, message); // first table write
delete live.message; // document mutation remains valid
tx.setTask(nextTask(task));
await tx.createTask(Follow, { after: message.id }); // further table writes are fine
}, context);
```
@@ -1288,6 +1281,9 @@ type RunningTask<I, S, R> = TaskRecord<I, S, R> & {
readonly state: Extract<TaskState<S, R>, { status: "running" }>;
};
/** State a task commits for itself: a replacement checkpoint or its terminal outcome. */
type NextTaskState<S, R> = Extract<TaskState<S, R>, { status: "running" | "terminal" }>;
interface HookRunner<H extends object> {
each<K extends keyof H>(name: K, invoke: (handler: H[K]) => void | Promise<void>): Promise<void>;
}
@@ -1302,10 +1298,16 @@ interface TaskRuntime<I, S, R, H extends object> extends DocumentObserver {
readonly taskId: TaskId<R>;
readonly conversationId: ConversationId;
readonly signal: AbortSignal;
/** Registry snapshot of the current phase; refreshed at every phase boundary. */
readonly registry: RegistrySnapshot<ToolRegistration>;
readonly models: Models;
readonly hooks: HookRunner<H>;
commit(
change: (tx: Tx, current: RunningTask<I, S, R>) => void | Promise<void>,
change: (
tx: Tx,
current: RunningTask<I, S, R>,
) => NextTaskState<S, R> | undefined | Promise<NextTaskState<S, R> | undefined>,
context: Context,
): Promise<void>;
@@ -1322,7 +1324,7 @@ type TaskDefinition<I, S extends { phase: string }, R, H extends object> = {
readonly phases: {
[P in S["phase"]]: PhaseHandler<I, Extract<S, { phase: P }>, S, R, H>;
};
abort(task: TaskRecord<I, S, R>, runtime: TaskRuntime<I, S, R, H>, context: Context): Promise<void>;
abort(task: RunningTask<I, S, R>, runtime: TaskRuntime<I, S, R, H>, context: Context): Promise<void>;
migrate?(input: JsonValue, checkpoint: JsonValue, fromVersion: number): {
input: I;
checkpoint: S;
@@ -1355,23 +1357,71 @@ Storage admission so persisted conversation-owner edges cannot become stale.
The phase map is exhaustive and phase-narrowed. A handler may perform several
commits around one effect, but each durable checkpoint is a full replacement.
`TaskRuntime.commit()` rereads and gates the current durable task on the Session
line before invoking its callback. Transaction methods replace its checkpoint
or write its terminal outcome.
line before invoking its callback. It rejects when the invocation has ended, the
Harness is closing, the task is terminal, or a run invocation's task carries an
abort mark. Its `tx.createTask()` defaults to the task's conversation. When the
callback returns a state, the runtime replaces the task's state in the same
commit, so the checkpoint or outcome is atomic with the callback's entries,
documents, and child tasks and is type-checked against the task's checkpoint and
result types. Returning nothing leaves the state unchanged. A terminal state
drops the memos. `pending` is never returned; only reconciliation and handover
write it.
`Tx` has no task replacement operation. A task changes only its own state,
through its runtime. The scheduler owns reservation, reconciliation, handover,
faults, and orphaning; `abortTask()` owns abort marks. Other code
stops a task with `abortTask()` and reads its result with `waitForTask()`.
```ts
// Intent, effect, outcome.
prepare: async (task, runtime, context) => {
await runtime.commit(() => ({ status: "running", checkpoint: { phase: "charge", key: newKey() } }), context);
},
charge: async (task, runtime, context) => {
const receipt = await payments.charge(task.state.checkpoint.key); // idempotent by key
await runtime.commit(async (tx, current) => {
const entry = await tx.appendEntry(current.conversationId, receiptEntry(receipt));
return { status: "terminal", outcome: { status: "completed", result: { entryId: entry.id } } };
}, context);
},
// The abort handler decides the outcome; returning without one faults the task.
abort: async (task, runtime, context) => {
await payments.cancel(task.state.checkpoint);
await runtime.commit(() => ({ status: "terminal", outcome: { status: "aborted", reason: "user" } }), context);
},
```
`memo(name, candidate)` is one gated commit; `memo(name)` reads the committed
record. `sleep(until)` compares against the Harness `now` clock and rejects when
the invocation is signalled or its context is cancelled. Watches acquired through
the runtime stop when the invocation ends.
Reservation durably changes `pending` to `running`. One invocation runs phase
handlers in sequence; checkpoint commits retain `running`. After a handler
settles, the scheduler rereads the task and applies the first matching rule:
handlers in sequence; checkpoint commits retain `running`. Before every phase,
including the first, the scheduler runs one synchronous step callback on the
Session line. It reads the committed task, applies the first matching rule
below to the phase that just returned, writes any fault or handover in the same
commit, and ends the invocation there when a rule stops it. A runtime commit the
invocation queued earlier therefore either lands before the step and counts, or
reaches the line after it and rejects. Before the first phase only rules 1–3
apply:
1. Terminal: stop.
2. Session closing: stop; preserve the checkpoint and any abort mark for reopen.
3. Run mode with a durable abort mark: end and join the run invocation, then
dispatch a fresh abort invocation.
4. Uncaught error: write terminal `faulted`.
5. Checkpoint changed, including progress within the same phase: invoke its
phase handler in the same task invocation.
5. Checkpoint changed, including progress within the same phase: refresh the
registry snapshot and either hand over (section 5.4) or invoke the phase
handler in the same task invocation.
6. Checkpoint unchanged: write terminal `faulted` because no durable progress
was made.
An abort invocation runs its abort handler once. After it settles, a step
applies rules 1, 2, and 4; a handler that returns without a terminal outcome
faults the task. When Storage rejects a step's fault or handover write, the
invocation still ends; the task stays `running` and the next reservation runs it
again.
On open, running-task reconciliation changes surviving `running` tasks back to
`pending`, preserving their checkpoint and abort mark. Task migration runs at
reservation, atomically with `pending -> running`. One callback handles every
@@ -1497,6 +1547,9 @@ An abort invocation resolves the definition the same way. When the definition
can take the task, its abort handler runs; otherwise the task is orphaned as
described below.
A failed migration is reported once through `onReport` and retried only after
the registry resolves a different definition object for the task's kind.
The Harness never terminalizes a task merely because registry code is missing or
incompatible, neither at open nor later. A blocked task without an abort mark
keeps its durable record unchanged, remains live, still blocks ordinary idle
@@ -1508,7 +1561,11 @@ Aborting a blocked task cannot run its abort handler, because that code is
missing or cannot take the task. The Harness therefore settles it as terminal
`orphaned` instead of `aborted`: `aborted` means the task's own abort handler
ran and decided the outcome, while `orphaned` means no task code ran, so external
effects the task started may remain uncleaned. Only an abort (direct, by
effects the task started may remain uncleaned. The orphaned `reason` is the
blocked reason (`missing_task`, `task_too_old`, or `migration_failed`). When
`abortTask()` finds no active invocation and the current snapshot
cannot take the task, the marking commit settles it directly; otherwise the
scheduler settles it when it would reserve the abort invocation. Only an abort (direct, by
conversation, or by cascade) orphans a task; a missing definition alone never
does. The orphaning commit performs the cleanup the task's code cannot: affected
input submissions become unanswered with the reason, any matching active turn
@@ -1517,16 +1574,16 @@ written; the terminal task record and unanswered submissions carry the reason.
Faulting a turn task performs the same control/submission cleanup with a
`faulted` outcome.
At every normal phase boundary (section 5.1, rule 5), after rereading the task
and applying the precedence rules, the scheduler refreshes the invocation's
registry snapshot. If the task definition resolved by name is a different object
At every normal phase boundary (section 5.1, rule 5), the step refreshes the
invocation's registry snapshot. If the task definition resolved by name is a different object
than the one the invocation started with and the new definition can reserve the
task (same version, or a higher version with `migrate`), the invocation hands
over: it commits the task back to `pending` with its checkpoint, memos, and abort
mark, ends, and the next reservation starts a fresh invocation under the new
over: the step commits the task back to `pending` with its checkpoint, memos,
and abort mark and ends the invocation, and the next reservation starts a fresh invocation under the new
definition, applying the reservation rules above. When the definition is missing
or cannot reserve the task, the invocation keeps running under its old
definition, reports once, and reconsiders at its next boundary. A handler that
definition, reports once through `onReport` per resolved definition, and
reconsiders at its next boundary. A handler that
never settles never hands over.
## 6. Submissions and inbox
+65 -4
View File
@@ -1,4 +1,5 @@
import type { Context, Draft } from "@earendil-works/chord";
import type { Context, Draft, JsonValue } from "@earendil-works/chord";
import { withoutAbortSignal } from "@earendil-works/chord/context";
import type { ModelThinkingLevel } from "@earendil-works/pi-ai";
import { SessionImpl } from "../session/session.ts";
import type {
@@ -11,11 +12,14 @@ import type {
EntryRecord,
Page,
Storage,
TaskId,
TaskRecord,
Tx,
} from "../types.ts";
import { ROOT_CONVERSATION_ID } from "../types.ts";
import { ConversationConfig, type ConversationConfigState } from "./config.ts";
import { captureContextBounds, deriveContext } from "./context.ts";
import { TaskScheduler } from "./scheduler.ts";
import type {
ContextView,
Conversation,
@@ -26,6 +30,7 @@ import type {
ModelRef,
RegistryReader,
RegistrySnapshot,
SettledTask,
ToolRegistration,
} from "./types.ts";
@@ -44,6 +49,7 @@ type ConversationHost<Tool extends ToolRegistration> = {
readonly harness: HarnessImpl<Tool>;
readonly storage: Storage;
readonly registry: RegistryReader<Tool>;
readonly tasks: TaskScheduler;
create(target: CreateTarget, init: ConversationInit | undefined, context: Context): Promise<Conversation>;
};
@@ -121,6 +127,10 @@ class ConversationImpl<Tool extends ToolRegistration> implements Conversation {
);
}
waitForIdle(context: Context): Promise<void> {
return this.#host.tasks.waitForIdle(this.id, context);
}
async #config(context: Context): Promise<Readonly<ConversationConfigState>> {
return (
(await this.#host.harness.snapshot(ConversationConfig, this.id, context)) ??
@@ -140,20 +150,59 @@ class HarnessImpl<Tool extends ToolRegistration> extends SessionImpl implements
readonly #storage: Storage;
readonly #registry: RegistryReader<Tool>;
readonly #host: ConversationHost<Tool>;
readonly #tasks: TaskScheduler;
#closed = false;
constructor(storage: Storage, options: HarnessOptions<Tool>) {
constructor(storage: Storage, options: HarnessOptions<Tool>, context: Context) {
super(storage);
this.#storage = storage;
this.#registry = options.registry;
this.#tasks = new TaskScheduler({
session: this,
storage,
registry: options.registry,
models: options.models,
now: options.now ?? Date.now,
report: options.onReport ?? (() => {}),
context: withoutAbortSignal(context),
});
this.#host = {
harness: this,
storage,
registry: options.registry,
tasks: this.#tasks,
create: (target, init, context) => this.#create(target, init, context),
};
}
/** Reconcile surviving `running` tasks to `pending`; part of open. */
openTasks(context: Context): Promise<void> {
return this.#tasks.open(context);
}
resume(): void {
this.#assertOpen();
this.#tasks.resume();
}
getTask<R>(id: TaskId<R>, context: Context): Promise<TaskRecord<JsonValue, JsonValue, R> | undefined> {
return this.readOnLine(() => this.#storage.task(id, context)) as Promise<
TaskRecord<JsonValue, JsonValue, R> | undefined
>;
}
abortTask(id: TaskId, context: Context): Promise<"marked" | "terminal"> {
return this.#tasks.abort(id, context);
}
waitForTask<R>(id: TaskId<R>, context: Context): Promise<SettledTask<R>> {
return this.#tasks.waitForTask(id, context) as Promise<SettledTask<R>>;
}
waitForIdle(context: Context): Promise<void> {
return this.#tasks.waitForIdle(undefined, context);
}
root(context: Context, options?: { readonly init?: ConversationInit }): Promise<Conversation> {
return this.#create({ kind: "root" }, options?.init, context);
}
@@ -173,11 +222,16 @@ class HarnessImpl<Tool extends ToolRegistration> extends SessionImpl implements
return super.close(context);
}
/** Join task invocations after admission is sealed and before Storage closes; writes no task outcome. */
protected override beforeClose(): Promise<void> {
return this.#tasks.join();
}
async #create(target: CreateTarget, init: ConversationInit | undefined, context: Context): Promise<Conversation> {
this.#assertOpen();
const id =
target.kind === "root"
? await this.commitRoot(async (tx) => {
? await this.commitWith(async (tx) => {
if ((await tx.conversation(ROOT_CONVERSATION_ID)) !== undefined) return ROOT_CONVERSATION_ID;
return this.#stageConfiguration(tx, await tx.createRootConversation(), false, init);
}, context)
@@ -239,6 +293,13 @@ export const Harness = {
context: Context,
): Promise<Harness> {
context.abortSignal?.throwIfAborted();
return new HarnessImpl(storage, options);
const harness = new HarnessImpl(storage, options, context);
try {
await harness.openTasks(context);
} catch (error) {
await harness.close(context);
throw error;
}
return harness;
},
};
+7
View File
@@ -158,6 +158,7 @@ class RegistryImpl<Tool extends ToolRegistration> implements Registry<Tool> {
#batch: Batch<Tool> | undefined;
/** First position of every key ever published; re-registered keys keep it. */
readonly #positions = new Map<string, number>();
readonly #listeners = new Set<() => void>();
#nextPosition = 0;
readonly tools: Registry<Tool>["tools"];
@@ -222,6 +223,11 @@ class RegistryImpl<Tool extends ToolRegistration> implements Registry<Tool> {
return this.#current;
}
subscribe(listener: () => void): () => void {
this.#listeners.add(listener);
return () => this.#listeners.delete(listener);
}
batch(register: () => void): Registration {
if (this.#batch !== undefined) throw new Error("Registry batches cannot be nested");
const batch: Batch<Tool> = { added: [], disposed: new Set() };
@@ -305,6 +311,7 @@ class RegistryImpl<Tool extends ToolRegistration> implements Registry<Tool> {
this.#current = new RegistryState(records);
for (const record of batch.added) record.status = "published";
for (const record of batch.disposed) record.status = "disposed";
for (const listener of [...this.#listeners]) listener();
}
}
+703
View File
@@ -0,0 +1,703 @@
import { type Context, copyJson, type JsonValue } from "@earendil-works/chord";
import { awaitWithContext, withAbortSignal } from "@earendil-works/chord/context";
import type { Models } from "@earendil-works/pi-ai";
import type { SessionImpl } from "../session/session.ts";
import type { Transaction } from "../session/transaction.ts";
import type {
CommitPublication,
ConversationId,
DocumentWatch,
JsonObject,
NextTaskState,
RunningTask,
Storage,
TaskDefinition,
TaskId,
TaskOutcome,
TaskRecord,
TaskRuntime,
TaskState,
Tx,
} from "../types.ts";
import type { AnyTask, RegistryReader, RegistrySnapshot, SettledTask } from "./types.ts";
type AnyTaskRecord = TaskRecord<JsonValue, JsonValue, JsonValue>;
type LiveTaskRecord = Extract<AnyTaskRecord, { readonly state: { readonly status: "pending" | "running" } }>;
type Checkpoint = { readonly phase: string };
type ErasedDefinition = TaskDefinition<JsonValue, Checkpoint, JsonValue, object>;
type ErasedRuntime = TaskRuntime<JsonValue, Checkpoint, JsonValue, object>;
type ErasedRunningTask = RunningTask<JsonValue, Checkpoint, JsonValue>;
const SCAN_PAGE_SIZE = 256;
/** Longest delay `setTimeout` supports; longer sleeps wait in several steps. */
const MAX_TIMER_DELAY = 2_147_483_647;
/** Why a pending task cannot be reserved under a registry snapshot. Derived, never persisted. */
type BlockedReason = "missing_task" | "task_too_old" | "migration_failed";
type Resolution =
| { readonly kind: "ready"; readonly task: AnyTask; readonly record: LiveTaskRecord }
| { readonly kind: "blocked"; readonly reason: BlockedReason };
/** One in-memory execution of a task in run or abort mode. */
type Invocation = {
readonly taskId: TaskId;
readonly conversationId: ConversationId;
readonly mode: "run" | "abort";
readonly controller: AbortController;
/** Context passed to handlers; cancelled by `controller`. */
readonly context: Context;
/** Watches acquired through the runtime; stopped at invocation end. */
readonly watches: Set<DocumentWatch<JsonObject>>;
ended: boolean;
readonly done: Promise<void>;
readonly finish: () => void;
};
type Reservation = {
readonly invocation: Invocation;
readonly task: AnyTask;
readonly snapshot: RegistrySnapshot;
};
/** Replacement definition already reported as unable to take over, per invocation. */
type ReportedTask = { readonly task: AnyTask | undefined } | undefined;
/** Outcome of the phase that just returned, judged by the next step. */
type PhaseResult = { readonly checkpoint: Checkpoint; readonly failure?: { readonly error: unknown } };
type Waiter<T> = { readonly resolve: (value: T) => void; readonly reject: (error: unknown) => void };
export type TaskSchedulerOptions = {
readonly session: SessionImpl;
readonly storage: Storage;
readonly registry: RegistryReader;
readonly models: Models;
readonly now: () => number;
readonly report: (error: unknown) => void;
/** Context for scheduler commits and invocations; carries no caller cancellation. */
readonly context: Context;
};
/**
* Durable task scheduler of one Harness.
*
* `#live` mirrors every committed pending/running task record. The synchronous commit listener updates it on the
* Session line, so code running on the line reads exactly the committed state from it.
*
* Invariant: every task transition is decided and written by one callback serialized on the Session line. That covers
* reservation, marks, runtime commits, and the synchronous step before each phase, which applies the precedence rules
* and writes a fault or handover. Handlers and joins run off the line. An invocation ends inside the step that decides its end,
* so a runtime commit it queued either lands before that decision or is rejected.
*/
export class TaskScheduler {
readonly #session: SessionImpl;
readonly #storage: Storage;
readonly #registry: RegistryReader;
readonly #models: Models;
readonly #now: () => number;
readonly #report: (error: unknown) => void;
readonly #context: Context;
readonly #live = new Map<TaskId, LiveTaskRecord>();
readonly #invocations = new Map<TaskId, Invocation>();
readonly #taskWaiters = new Map<TaskId, Set<Waiter<SettledTask<JsonValue>>>>();
/** Idle waiters by conversation; `undefined` waits for the whole Harness. */
readonly #idleWaiters = new Map<ConversationId | undefined, Set<Waiter<void>>>();
/** Definition whose migration failed per task; retried only once the registry resolves another definition. */
readonly #failedMigrations = new Map<TaskId, AnyTask>();
#unsubscribeRegistry: () => void = () => {};
#enabled = false;
#closing = false;
#dirty = false;
#draining = false;
constructor(options: TaskSchedulerOptions) {
this.#session = options.session;
this.#storage = options.storage;
this.#registry = options.registry;
this.#models = options.models;
this.#now = options.now;
this.#report = options.report;
this.#context = options.context;
}
/** Load live tasks and change surviving `running` tasks back to `pending`. Dispatches nothing. */
async open(context: Context): Promise<void> {
this.#session.subscribeCommits((publication) => this.#observe(publication));
this.#session.subscribeClose(() => this.#seal());
this.#unsubscribeRegistry = this.#registry.subscribe(() => this.#kick());
await this.#session.commitWith(async (tx) => {
const running: LiveTaskRecord[] = [];
for (const status of ["pending", "running"] as const) {
let cursor: Parameters<Tx["scanTasks"]>[2];
do {
const page = await tx.scanTasks({ status }, SCAN_PAGE_SIZE, cursor);
for (const record of page.items as readonly LiveTaskRecord[]) {
this.#live.set(record.id, record);
if (status === "running") running.push(record);
}
cursor = page.next;
} while (cursor !== undefined);
}
for (const record of running)
tx.setTask(withState(record, { status: "pending", checkpoint: record.state.checkpoint }));
}, context);
}
resume(): void {
this.#enabled = true;
this.#kick();
}
/** Wait for every invocation signalled by `#seal()`. Writes nothing. */
async join(): Promise<void> {
await Promise.allSettled([...this.#invocations.values()].map((invocation) => invocation.done));
}
/**
* Commit the abort mark, or settle a task that no registered definition can take as `orphaned`, then signal and
* join the run invocation seen on the line. The next drain starts the abort invocation.
*/
async abort(id: TaskId, context: Context): Promise<"marked" | "terminal"> {
const marked = await this.#session.commitWith(async (tx) => {
const current = await tx.task(id);
if (current === undefined) throw new Error(`Task ${id} does not exist`);
if (current.state.status === "terminal") return { result: "terminal" as const };
const live = current as LiveTaskRecord;
const invocation = this.#invocations.get(id);
if (invocation === undefined) {
const resolution = this.#resolve(live, this.#registry.snapshot());
if (resolution.kind === "blocked") {
tx.setTask(withState(live, terminal({ status: "orphaned", reason: resolution.reason })));
return { result: "marked" as const };
}
}
if (!live.abortRequested) tx.setTask({ ...live, abortRequested: true });
return { result: "marked" as const, run: invocation?.mode === "run" ? invocation : undefined };
}, context);
if (marked.run !== undefined) {
marked.run.controller.abort();
await awaitWithContext(marked.run.done, context);
}
return marked.result;
}
async waitForTask(id: TaskId, context: Context): Promise<SettledTask<JsonValue>> {
// Check and register on the line so no terminal publication falls between them.
const found = await this.#session.readOnLine(async () => {
if (this.#closing) throw closedError();
if (this.#live.has(id)) return { promise: addWaiter(this.#taskWaiters, id, context) };
const record = await this.#storage.task(id, context);
if (record === undefined) throw new Error(`Task ${id} does not exist`);
return { promise: Promise.resolve(record as SettledTask<JsonValue>) };
});
return found.promise;
}
/** Resolve when no live non-background task exists, optionally within one conversation. */
waitForIdle(conversationId: ConversationId | undefined, context: Context): Promise<void> {
if (this.#closing) return Promise.reject(closedError());
if (this.#idle(conversationId)) return Promise.resolve();
return addWaiter(this.#idleWaiters, conversationId, context);
}
// ─── Scheduling ────────────────────────────────────────────────────────
#observe(publication: CommitPublication): void {
let changed = false;
for (const change of publication.changes) {
if (change.type !== "task") continue;
changed = true;
const record = change.value;
if (record.state.status !== "terminal") {
this.#live.set(record.id, record as LiveTaskRecord);
continue;
}
this.#live.delete(record.id);
this.#failedMigrations.delete(record.id);
settleWaiters(this.#taskWaiters, record.id, (waiter) => waiter.resolve(record as SettledTask<JsonValue>));
}
if (!changed) return;
for (const conversationId of [...this.#idleWaiters.keys()]) {
if (this.#idle(conversationId)) settleWaiters(this.#idleWaiters, conversationId, (waiter) => waiter.resolve());
}
this.#kick();
}
/** Close listener: runs synchronously once admission is sealed, before `join()`. */
#seal(): void {
this.#closing = true;
this.#unsubscribeRegistry();
const error = closedError();
for (const id of [...this.#taskWaiters.keys()]) settleWaiters(this.#taskWaiters, id, (w) => w.reject(error));
for (const id of [...this.#idleWaiters.keys()]) settleWaiters(this.#idleWaiters, id, (w) => w.reject(error));
for (const invocation of this.#invocations.values()) invocation.controller.abort();
}
#kick(): void {
this.#dirty = true;
if (this.#draining || !this.#enabled || this.#closing) return;
this.#draining = true;
// Never commit synchronously from a commit or registry listener.
queueMicrotask(() => void this.#drain());
}
async #drain(): Promise<void> {
try {
while (this.#dirty && this.#enabled && !this.#closing) {
this.#dirty = false;
for (const reservation of await this.#reserve()) this.#start(reservation);
}
} catch (error) {
if (!this.#closing) this.#report(error);
} finally {
this.#draining = false;
// A wakeup that arrived during a failed pass still needs its pass.
if (this.#dirty) this.#kick();
}
}
/** Reserve every eligible task in one commit; orphan abort-marked tasks no definition can take. */
async #reserve(): Promise<Reservation[]> {
const reservations: Reservation[] = [];
try {
await this.#session.commitWith(async (tx) => {
if (!this.#enabled || this.#closing) return;
// Taken once per pass, and only when some task is a candidate.
let snapshot: RegistrySnapshot | undefined;
for (const record of [...this.#live.values()]) {
if (this.#invocations.has(record.id)) continue;
const mode = record.abortRequested ? "abort" : "run";
// An abort mark bypasses dependencies; a task already running has passed them.
if (
mode === "run" &&
record.state.status === "pending" &&
record.after.some((id) => this.#live.has(id))
) {
continue;
}
snapshot ??= this.#registry.snapshot();
const resolution = this.#resolve(record, snapshot);
if (resolution.kind === "blocked") {
if (mode === "abort") {
tx.setTask(withState(record, terminal({ status: "orphaned", reason: resolution.reason })));
}
continue;
}
if (resolution.record !== record || record.state.status !== "running") {
tx.setTask(
withState(resolution.record, {
status: "running",
checkpoint: resolution.record.state.checkpoint,
}),
);
}
// Registered on the line, so marks and later reservations see it and close joins it.
const invocation = this.#createInvocation(record, mode);
reservations.push({ invocation, task: resolution.task, snapshot });
}
}, this.#context);
} catch (error) {
for (const { invocation } of reservations) {
this.#invocations.delete(invocation.taskId);
invocation.finish();
}
throw error;
}
return reservations;
}
/** Resolve the record's definition by kind, migrating an older stored version. */
#resolve(record: LiveTaskRecord, snapshot: RegistrySnapshot): Resolution {
const task = snapshot.task(record.kind);
if (task === undefined) return { kind: "blocked", reason: "missing_task" };
const definition = erased(task);
if (definition.version === record.version) return { kind: "ready", task, record };
if (definition.version < record.version) return { kind: "blocked", reason: "task_too_old" };
if (this.#failedMigrations.get(record.id) === task) return { kind: "blocked", reason: "migration_failed" };
try {
if (definition.migrate === undefined) {
throw new Error(
`Task ${record.kind} version ${definition.version} has no migration from ${record.version}`,
);
}
const migrated = definition.migrate(record.input, record.state.checkpoint, record.version);
const migratedRecord = {
...record,
version: definition.version,
input: copyJson(migrated.input),
state: { status: record.state.status, checkpoint: copyJson(migrated.checkpoint) },
} as LiveTaskRecord;
return { kind: "ready", task, record: migratedRecord };
} catch (error) {
this.#failedMigrations.set(record.id, task);
this.#report(error);
return { kind: "blocked", reason: "migration_failed" };
}
}
#createInvocation(record: LiveTaskRecord, mode: "run" | "abort"): Invocation {
const controller = new AbortController();
let finish!: () => void;
const done = new Promise<void>((resolve) => {
finish = resolve;
});
const invocation: Invocation = {
taskId: record.id,
conversationId: record.conversationId,
mode,
controller,
context: withAbortSignal(controller.signal, this.#context),
watches: new Set(),
ended: false,
done,
finish,
};
this.#invocations.set(record.id, invocation);
return invocation;
}
#start(reservation: Reservation): void {
const invocation = reservation.invocation;
void (async () => {
try {
if (invocation.mode === "run") await this.#run(reservation);
else await this.#runAbort(reservation);
} catch (error) {
this.#report(error);
} finally {
this.#end(invocation);
invocation.finish();
this.#kick();
}
})();
}
/** Run phase handlers, each preceded by a step that decides on the line whether the invocation continues. */
async #run(reservation: Reservation): Promise<void> {
const invocation = reservation.invocation;
const state = { task: reservation.task, snapshot: reservation.snapshot, reported: undefined as ReportedTask };
const runtime = this.#runtime(invocation, () => state.snapshot);
let previous: PhaseResult | undefined;
for (;;) {
const current = await this.#step(invocation, (tx, current) => this.#decide(tx, current, previous, state));
// Close may seal between the decision and dispatch.
if (current === undefined || this.#closing) return;
const checkpoint = current.state.checkpoint;
try {
await erased(state.task).phases[checkpoint.phase]!(current, runtime, invocation.context);
previous = { checkpoint };
} catch (error) {
previous = { checkpoint, failure: { error } };
}
}
}
/**
* Precedence rules for a run invocation, on the line. Rules 1 (terminal) and 2 (closing) are applied by `#step`.
* Returns whether the invocation continues with the next phase.
*/
#decide(
tx: Transaction,
current: ErasedRunningTask,
previous: PhaseResult | undefined,
state: { task: AnyTask; snapshot: RegistrySnapshot; reported: ReportedTask },
): boolean {
// 3. abort mark: end; the next drain starts a fresh abort invocation.
if (current.abortRequested) return false;
if (previous === undefined) return true;
// 4. uncaught error.
if (previous.failure !== undefined) {
tx.setTask(faulted(current, previous.failure.error));
return false;
}
// 6. no durable progress.
if (jsonEqual(current.state.checkpoint, previous.checkpoint)) {
const message = `Task ${current.kind} phase ${previous.checkpoint.phase} returned without durable progress`;
tx.setTask(faulted(current, new Error(message)));
return false;
}
// 5. progress: refresh the snapshot; hand over to a replacement definition that can take the task.
state.snapshot = this.#registry.snapshot();
const next = state.snapshot.task(current.kind);
if (next !== state.task) {
if (next !== undefined && canReserve(next, current)) {
tx.setTask(withState(current, { status: "pending", checkpoint: current.state.checkpoint }));
return false;
}
if (state.reported === undefined || state.reported.task !== next) {
state.reported = { task: next };
const cause = next === undefined ? "missing_task" : "incompatible_task";
this.#report(
new Error(`Task ${current.id} keeps running under its old ${current.kind} definition`, { cause }),
);
}
}
return true;
}
/** Run the abort handler once; rules 1, 2, and 4 apply, and returning without an outcome faults. */
async #runAbort(reservation: Reservation): Promise<void> {
const invocation = reservation.invocation;
const current = this.#live.get(invocation.taskId) as ErasedRunningTask | undefined;
if (current === undefined || this.#closing) return;
let failure: { readonly error: unknown } | undefined;
try {
const runtime = this.#runtime(invocation, () => reservation.snapshot);
await erased(reservation.task).abort(current, runtime, invocation.context);
} catch (error) {
failure = { error };
}
const message = `Abort handler of task ${invocation.taskId} returned without a terminal outcome`;
await this.#step(invocation, (tx, current) => {
tx.setTask(faulted(current, failure?.error ?? new Error(message)));
return false;
});
}
/**
* One synchronous decision on the Session line. A terminal task (rule 1) or a closing Harness (rule 2) ends the
* invocation without a write; otherwise `decide` may stage a write and returns whether the invocation continues.
* Ending happens inside the callback. A rejected step, such as admission after close, also ends the invocation.
*/
async #step(
invocation: Invocation,
decide: (tx: Transaction, current: ErasedRunningTask) => boolean,
): Promise<ErasedRunningTask | undefined> {
try {
return await this.#session.commitWith((tx) => {
const current = this.#live.get(invocation.taskId) as ErasedRunningTask | undefined;
if (current !== undefined && !this.#closing && decide(tx, current)) return current;
this.#end(invocation);
return undefined;
}, this.#context);
} catch (error) {
this.#end(invocation);
if (!this.#closing) this.#report(error);
return undefined;
}
}
/** End an invocation: its runtime operations reject from now on, its watches stop, and its task is free. */
#end(invocation: Invocation): void {
if (invocation.ended) return;
invocation.ended = true;
if (this.#invocations.get(invocation.taskId) === invocation) this.#invocations.delete(invocation.taskId);
for (const watch of invocation.watches) void watch.stop();
}
#idle(conversationId: ConversationId | undefined): boolean {
for (const record of this.#live.values()) {
if (record.background) continue;
if (conversationId === undefined || record.conversationId === conversationId) return false;
}
return true;
}
// ─── Invocation runtime ──────────────────────────────────────────────────
#runtime(invocation: Invocation, snapshot: () => RegistrySnapshot): ErasedRuntime {
return {
taskId: invocation.taskId as TaskId<JsonValue>,
conversationId: invocation.conversationId,
signal: invocation.controller.signal,
models: this.#models,
get registry() {
return snapshot();
},
commit: (change, context) =>
this.#gated(
invocation,
async (tx, current) => {
const next = await change(tx, current);
if (next !== undefined) tx.setTask(withState(current, next));
},
context,
),
memo: ((name: string, ...rest: readonly unknown[]) =>
rest.length === 1
? this.#readMemo(invocation, name)
: this.#writeMemo(invocation, name, rest[0] as JsonValue, rest[1] as Context)) as ErasedRuntime["memo"],
sleep: (until, context) => this.#sleep(invocation, until, context),
watchDoc: ((...args: readonly unknown[]) => this.#watchDoc(invocation, args)) as ErasedRuntime["watchDoc"],
};
}
/** Commit after rereading the task on the line and gating the invocation. */
#gated<T>(
invocation: Invocation,
change: (tx: Transaction, current: ErasedRunningTask) => T | Promise<T>,
context: Context,
): Promise<T> {
if (invocation.ended) return Promise.reject(endedError(invocation));
return this.#session.commitWith(
async (tx) => {
if (invocation.ended) throw endedError(invocation);
if (this.#closing) throw closedError();
const current = this.#live.get(invocation.taskId) as ErasedRunningTask | undefined;
if (current === undefined) throw new Error(`Task ${invocation.taskId} is terminal`);
if (invocation.mode === "run" && current.abortRequested) {
throw new Error(`Task ${invocation.taskId} has a durable abort mark`);
}
return change(tx, current);
},
context,
invocation.conversationId,
);
}
#readMemo(invocation: Invocation, name: string): Promise<JsonValue | undefined> {
if (invocation.ended) return Promise.reject(endedError(invocation));
return Promise.resolve(memoOf(this.#live.get(invocation.taskId), name));
}
#writeMemo(invocation: Invocation, name: string, candidate: JsonValue, context: Context): Promise<JsonValue> {
return this.#gated(
invocation,
(tx, current) => {
const winner = memoOf(current, name);
if (winner !== undefined) return winner;
tx.setTask({ ...current, memos: { ...current.memos, [name]: candidate } } as AnyTaskRecord);
return candidate;
},
context,
);
}
/** Wait until the Harness clock reaches `until`, rechecking it after every timer. */
async #sleep(invocation: Invocation, until: number, context: Context): Promise<void> {
if (invocation.ended) throw endedError(invocation);
const signals = [invocation.controller.signal];
if (context.abortSignal !== undefined) signals.push(context.abortSignal);
const signal = AbortSignal.any(signals);
for (;;) {
signal.throwIfAborted();
const remaining = until - this.#now();
if (remaining <= 0) return;
await delay(Math.min(remaining, MAX_TIMER_DELAY), signal);
}
}
async #watchDoc(invocation: Invocation, args: readonly unknown[]): Promise<DocumentWatch<JsonObject> | undefined> {
if (invocation.ended) throw endedError(invocation);
const watchDoc = this.#session.watchDoc.bind(this.#session) as (
...args: readonly unknown[]
) => Promise<DocumentWatch<JsonObject> | undefined>;
const watch = await watchDoc(...args);
if (watch === undefined) return undefined;
if (invocation.ended) {
void watch.stop();
throw endedError(invocation);
}
invocation.watches.add(watch);
void watch.closed.then(() => invocation.watches.delete(watch));
return watch;
}
}
/** Own memo entry only; memo names such as `toString` must not resolve to inherited properties. */
function memoOf(record: LiveTaskRecord | undefined, name: string): JsonValue | undefined {
const memos = record?.memos;
return memos !== undefined && Object.hasOwn(memos, name) ? memos[name] : undefined;
}
function erased(task: AnyTask): ErasedDefinition {
return task.definition as unknown as ErasedDefinition;
}
function endedError(invocation: Invocation): Error {
return new Error(`Task ${invocation.taskId} invocation has ended`);
}
function closedError(): Error {
return new Error("Harness is closed");
}
function terminal(outcome: TaskOutcome<JsonValue>): NextTaskState<JsonValue, JsonValue> {
return { status: "terminal", outcome };
}
function faulted(record: LiveTaskRecord, error: unknown): AnyTaskRecord {
const message = error instanceof Error ? error.message : String(error);
return withState(record, terminal({ status: "faulted", error: { message } }));
}
/** Replace a live record's state; memos disappear in the terminal replacement. */
function withState(record: LiveTaskRecord, state: TaskState<JsonValue, JsonValue>): AnyTaskRecord {
if (state.status !== "terminal") return { ...record, state } as AnyTaskRecord;
const { memos: _memos, ...rest } = record;
return { ...rest, state };
}
/** Whether a definition can take the task at reservation: same version, or newer with a migration. */
function canReserve(task: AnyTask, record: LiveTaskRecord): boolean {
const definition = task.definition;
return (
definition.version === record.version || (definition.version > record.version && definition.migrate !== undefined)
);
}
/** Wait in `waiters[key]` until settled or `context` is cancelled. */
function addWaiter<K, T>(waiters: Map<K, Set<Waiter<T>>>, key: K, context: Context): Promise<T> {
return new Promise<T>((resolve, reject) => {
const signal = context.abortSignal;
if (signal?.aborted) return reject(signal.reason);
let set = waiters.get(key);
if (set === undefined) {
set = new Set();
waiters.set(key, set);
}
const own = set;
const remove = (): void => {
own.delete(waiter);
if (own.size === 0 && waiters.get(key) === own) waiters.delete(key);
signal?.removeEventListener("abort", onAbort);
};
const waiter: Waiter<T> = {
resolve: (value) => {
remove();
resolve(value);
},
reject: (error) => {
remove();
reject(error);
},
};
const onAbort = (): void => waiter.reject(signal!.reason);
signal?.addEventListener("abort", onAbort, { once: true });
own.add(waiter);
});
}
function settleWaiters<K, T>(waiters: Map<K, Set<Waiter<T>>>, key: K, settle: (waiter: Waiter<T>) => void): void {
const set = waiters.get(key);
if (set === undefined) return;
waiters.delete(key);
for (const waiter of [...set]) settle(waiter);
}
function delay(ms: number, signal: AbortSignal): Promise<void> {
return new Promise((resolve, reject) => {
const onAbort = (): void => {
clearTimeout(timer);
reject(signal.reason);
};
const timer = setTimeout(() => {
signal.removeEventListener("abort", onAbort);
resolve();
}, ms);
signal.addEventListener("abort", onAbort, { once: true });
});
}
function jsonEqual(left: JsonValue | undefined, right: JsonValue | undefined): boolean {
if (left === right) return true;
if (typeof left !== "object" || typeof right !== "object" || left === null || right === null) return false;
if (Array.isArray(left) || Array.isArray(right)) {
if (!Array.isArray(left) || !Array.isArray(right) || left.length !== right.length) return false;
return left.every((value, index) => jsonEqual(value, right[index]));
}
const keys = Object.keys(left);
if (keys.length !== Object.keys(right).length) return false;
return keys.every((key) => Object.hasOwn(right, key) && jsonEqual(left[key], right[key]));
}
+19 -2
View File
@@ -215,6 +215,8 @@ export interface RegistrySnapshot<Tool extends ToolRegistration = ToolRegistrati
export interface RegistryReader<Tool extends ToolRegistration = ToolRegistration> {
/** Immutable view of the whole current registry. */
snapshot(): RegistrySnapshot<Tool>;
/** Called synchronously after every publication; wakes the scheduler to reconsider blocked tasks. */
subscribe(listener: () => void): () => void;
}
/** Application-owned registry of tools, hooks, tasks, and the system prompt. */
@@ -305,16 +307,31 @@ export interface Conversation {
context: Context,
): Promise<Page<EntryRecord, Cursor>>;
fork(at: EntryId, options: ConversationCreateOptions, context: Context): Promise<Conversation>;
/** Resolve when no live non-background task belongs to this conversation. */
waitForIdle(context: Context): Promise<void>;
}
// TODO: decide how Harness exposes subscribeCommits() and subscribeClose(). Their listeners run on the Session line
// and must not throw or call Session APIs, and Harness close will also join task invocations.
/** Durable agent harness over one Session. */
export interface Harness extends Session {
/** Enable task scheduling. Idempotent; throws after close. */
resume(): void;
/** Return the reserved root conversation, creating it with `init` in one commit when absent. */
root(context: Context, options?: { readonly init?: ConversationInit }): Promise<Conversation>;
// Conversation activity (active/idle notifications, quiescence) is specified with the task runtime and turn
// control in Packages 14 and 17.
// Conversation activity (active/idle notifications) is specified with turn control in Package 17.
conversation(id: ConversationId, context: Context): Promise<Conversation | undefined>;
createConversation(options: ConversationCreateOptions, context: Context): Promise<Conversation>;
getTask<R>(id: TaskId<R>, context: Context): Promise<TaskRecord<JsonValue, JsonValue, R> | undefined>;
/**
* Commit the abort mark, signal and join an active run invocation, and schedule the abort invocation. A task whose
* definition cannot take it settles as `orphaned` instead.
*/
abortTask(id: TaskId, context: Context): Promise<"marked" | "terminal">;
/** Resolve with the terminal receipt; cancelling `context` cancels only this wait. */
waitForTask<R>(id: TaskId<R>, context: Context): Promise<SettledTask<R>>;
/** Resolve when no live non-background task exists. */
waitForIdle(context: Context): Promise<void>;
}
+5
View File
@@ -40,6 +40,7 @@ export type {
} from "./harness/types.ts";
export { createSession } from "./session/session.ts";
export { MemoryStorage } from "./storage/memory.ts";
export { defineTask } from "./tasks.ts";
export type {
CheckpointInfo,
CommitChange,
@@ -77,10 +78,13 @@ export type {
Id,
JsonObject,
LatestConversationSemantics,
NextTaskState,
Page,
PhaseHandler,
RewindableConversationDocFamilyToken,
RewindableConversationDocToken,
RewindableConversationSemantics,
RunningTask,
Seq,
Session,
SessionDocFamilyToken,
@@ -102,6 +106,7 @@ export type {
TaskOutcomeError,
TaskQuery,
TaskRecord,
TaskRuntime,
TaskState,
Tx,
WatchEnd,
+24 -20
View File
@@ -83,9 +83,12 @@ export class SessionImpl implements Session {
return this.#enqueue(() => this.#runCommit(change, context));
}
/** Internal commit whose `tx.createTask()` defaults to `defaultConversationId`; used by Conversation handles. */
/**
* Internal commit exposing the concrete transaction and its internal operations, such as the reserved-ID root
* bootstrap and task replacement. `tx.createTask()` defaults to `defaultConversationId`.
*/
commitWith<T>(
change: (tx: Tx) => T | Promise<T>,
change: (tx: Transaction) => T | Promise<T>,
context: Context,
defaultConversationId?: ConversationId,
): Promise<T> {
@@ -97,16 +100,6 @@ export class SessionImpl implements Session {
return this.#enqueue(() => this.#runCommit(change, context, defaultConversationId));
}
/** Internal commit exposing the concrete transaction, whose reserved-ID root bootstrap is not on the public `Tx`. */
commitRoot<T>(change: (tx: Transaction) => T | Promise<T>, context: Context): Promise<T> {
try {
this.#assertUsable();
} catch (error) {
return Promise.reject(error);
}
return this.#enqueue(() => this.#runCommit(change, context));
}
/** Internal: run a read-only job on the mutation line so multi-read derivations observe one committed state. */
readOnLine<T>(job: () => Promise<T>): Promise<T> {
try {
@@ -363,17 +356,28 @@ export class SessionImpl implements Session {
close(context: Context): Promise<void> {
if (this.#closing === undefined) {
const cleanup = withoutAbortSignal(context);
this.#closing = this.#enqueue(async () => {
for (const listener of [...this.#closeListeners]) listener();
this.#closeListeners.clear();
this.#commitListeners.clear();
this.#documents.clear();
await this.#storage.close(cleanup);
});
// Seal admission before anything else runs, then stop observers; admitted work settles before Storage closes.
this.#closing = Promise.resolve()
.then(() => this.beforeClose())
.then(() =>
this.#enqueue(async () => {
this.#commitListeners.clear();
this.#documents.clear();
await this.#storage.close(cleanup);
}),
);
const listeners = [...this.#closeListeners];
this.#closeListeners.clear();
for (const listener of listeners) listener();
}
return awaitWithContext(this.#closing, context);
}
/** Runs after close seals admission and before the line closes Storage; must not reject. */
protected beforeClose(): Promise<void> {
return Promise.resolve();
}
/** Register a synchronous post-adoption listener. It must not throw, block, or call Session operations. */
subscribeCommits(listener: (publication: CommitPublication, context: Context) => void): () => void {
this.#assertUsable();
@@ -381,7 +385,7 @@ export class SessionImpl implements Session {
return () => this.#commitListeners.delete(listener);
}
/** Register a synchronous close listener. It must not throw, block, or call Session operations. */
/** Register a listener called synchronously when close begins. It must not throw, block, or call Session operations. */
subscribeClose(listener: () => void): () => void {
this.#assertUsable();
this.#closeListeners.add(listener);
@@ -308,6 +308,7 @@ export class Transaction implements Tx {
});
}
/** Internal: replace one task record completely. Tasks change their own state through their runtime. */
setTask(value: AnyTaskRecord): void {
this.#assertOpen();
this.#hasTableWrite = true;
+8
View File
@@ -0,0 +1,8 @@
import type { Task, TaskDefinition } from "./types.ts";
/** Define an executable task. Register it in the registry so a Harness can run tasks of its kind. */
export function defineTask<I, S extends { phase: string }, R, H extends object = object>(
definition: TaskDefinition<I, S, R, H>,
): Task<I, S, R, H> {
return { definition };
}
+70 -8
View File
@@ -1,6 +1,7 @@
import type { AttachedReplicatedState, Context, Draft, JsonValue } from "@earendil-works/chord";
import type { Op } from "@earendil-works/chord/delta";
import type { Message } from "@earendil-works/pi-ai";
import type { Message, Models } from "@earendil-works/pi-ai";
import type { RegistrySnapshot } from "./harness/types.ts";
/** JSON object used as the root of every durable document. */
export type JsonObject = { [key: string]: JsonValue };
@@ -124,9 +125,59 @@ export type TaskDocFamilyToken<T extends JsonObject, I extends JsonValue> = DocF
DocFamilyDefinition<T, I> & { readonly scope: "task" }
>;
declare const taskResultType: unique symbol;
/** Live task record reserved by one invocation. */
export type RunningTask<I, S, R> = TaskRecord<I, S, R> & {
readonly state: Extract<TaskState<S, R>, { readonly status: "running" }>;
};
/** Task definition fields currently supported by Session task creation. */
/** Next state a task commits for itself: a replacement checkpoint or its terminal outcome. */
export type NextTaskState<S, R> = Extract<TaskState<S, R>, { readonly status: "running" | "terminal" }>;
/**
* Runs one checkpoint phase. It must commit a changed checkpoint or a terminal outcome through `runtime.commit()`;
* returning without durable progress faults the task.
*/
export type PhaseHandler<I, P, S, R, H extends object> = (
task: RunningTask<I, P, R>,
runtime: TaskRuntime<I, S, R, H>,
context: Context,
) => Promise<void>;
/**
* Operations of one task invocation. Every operation rejects after the invocation ends; watches acquired through it
* stop at invocation end. `_H` is the task's hook map, consumed once the runtime gains its hook runner.
*/
export interface TaskRuntime<I, S, R, _H extends object> extends DocumentObserver {
readonly taskId: TaskId<R>;
readonly conversationId: ConversationId;
/** Aborted when the run is signalled by `abortTask()` or the Harness closes. */
readonly signal: AbortSignal;
/** Registry snapshot of the current phase; refreshed at every phase boundary. */
readonly registry: RegistrySnapshot;
readonly models: Models;
/**
* Commit on the Session line after rereading the task. Rejects when the task is terminal, the invocation ended, the
* Harness is closing, or, in a run invocation, the task carries an abort mark. A returned state replaces the task's
* state in the same commit; returning nothing leaves it unchanged. `tx.createTask()` defaults to the task's
* conversation.
*/
commit(
change: (
tx: Tx,
current: RunningTask<I, S, R>,
) => NextTaskState<S, R> | undefined | Promise<NextTaskState<S, R> | undefined>,
context: Context,
): Promise<void>;
/** Read a durable memo of this task. */
memo<T extends JsonValue>(name: string, context: Context): Promise<T | undefined>;
/** Store `candidate` unless a memo already exists; return the durable winner. */
memo<T extends JsonValue>(name: string, candidate: T, context: Context): Promise<T>;
/** Resolve once the Harness clock reaches `until`; rejects when the invocation or `context` is cancelled. */
sleep(until: number, context: Context): Promise<void>;
}
/** Executable durable state machine definition, registered in the registry by `name`. */
export type TaskDefinition<I, S extends { phase: string }, R, H extends object> = {
/** Registered task kind persisted in `TaskRecord.kind`. */
readonly name: string;
@@ -134,9 +185,22 @@ export type TaskDefinition<I, S extends { phase: string }, R, H extends object>
readonly version: number;
/** First durable checkpoint for a newly created task. */
initial(input: I): S;
/** Exhaustive phase map; each handler receives the task narrowed to its phase. */
readonly phases: {
readonly [P in S["phase"]]: PhaseHandler<I, Extract<S, { phase: P }>, S, R, H>;
};
/** Runs in a fresh invocation after an abort mark and must commit a terminal outcome. */
abort(task: RunningTask<I, S, R>, runtime: TaskRuntime<I, S, R, H>, context: Context): Promise<void>;
/** Convert a record stored by any older supported version; runs at reservation. */
migrate?(
input: JsonValue,
checkpoint: JsonValue,
fromVersion: number,
): {
input: I;
checkpoint: S;
};
readonly hooks?: H;
/** Type-only result marker until phase handlers commit typed outcomes. */
readonly [taskResultType]?: R;
};
/** Typed executable task definition. */
@@ -608,8 +672,6 @@ export interface Tx {
input: I,
options?: TaskOptions,
): Promise<TaskId<R>>;
/** Replace one task record completely. */
setTask(value: TaskRecord<JsonValue, JsonValue, JsonValue>): void;
doc<T extends JsonObject>(token: SessionDocToken<T>): Promise<Draft<T>>;
doc<T extends JsonObject>(token: ConversationDocToken<T>, conversationId: ConversationId): Promise<Draft<T>>;
@@ -710,7 +772,7 @@ export interface Session extends DocumentObserver {
close(context: Context): Promise<void>;
/** Observe complete commits synchronously after adoption. The listener must not throw, block, or call Session APIs. */
subscribeCommits(listener: (publication: CommitPublication, context: Context) => void): () => void;
/** Observe close synchronously on the Session line. The listener must not throw, block, or call Session APIs. */
/** Observe close synchronously when it begins. The listener must not throw, block, or call Session APIs. */
subscribeClose(listener: () => void): () => void;
snapshot<T extends JsonObject>(token: SessionDocToken<T>, context: Context): Promise<Readonly<T> | undefined>;
@@ -7,10 +7,10 @@ import {
type ConversationId,
defineDoc,
defineEntry,
defineTask,
type EntryRecord,
MemoryStorage,
ROOT_CONVERSATION_ID,
type Task,
} from "@earendil-works/pi-durable";
import { afterEach, describe, expect, it } from "vitest";
import { openNodeSqliteStorage } from "../src/storage/sqlite/node.ts";
@@ -281,9 +281,13 @@ describe("Harness root and conversations", () => {
it("binds commits and task creation to the conversation", async () => {
const { harness } = await openHarness(new MemoryStorage());
const conversation = await harness.createConversation({ ownership: { kind: "ownerless" } }, context);
const task: Task<{ n: number }, { phase: "run" }, void, object> = {
definition: { name: "test.work", version: 1, initial: () => ({ phase: "run" }) },
};
const task = defineTask<{ n: number }, { phase: "run" }, null>({
name: "test.work",
version: 1,
initial: () => ({ phase: "run" }),
phases: { run: async () => {} },
abort: async () => {},
});
const taskId = await conversation.commit((tx) => tx.createTask(task, { n: 1 }), context);
const record = await harness.commit((tx) => tx.task(taskId), context);
expect(record?.conversationId).toBe(conversation.id);
+9 -11
View File
@@ -1,6 +1,6 @@
import type { JsonValue } from "@earendil-works/chord";
import { Type } from "@earendil-works/pi-ai";
import { type AnyTask, createRegistry, type ToolRegistration } from "@earendil-works/pi-durable";
import { createRegistry, defineTask, type ToolRegistration } from "@earendil-works/pi-durable";
import { describe, expect, it } from "vitest";
type AppTool = ToolRegistration & { readonly snippet?: string };
@@ -17,16 +17,14 @@ function tool(name: string, extra: Partial<AppTool> = {}): AppTool {
function task(name: string) {
const hooks: { beforeRun(): void } = { beforeRun: () => {} };
return {
definition: {
name,
version: 1,
initial: (_input: undefined) => ({ phase: "run" as const }),
phases: {},
abort: () => {},
hooks,
},
} satisfies AnyTask;
return defineTask<undefined, { phase: "run" }, null, { beforeRun(): void }>({
name,
version: 1,
initial: () => ({ phase: "run" }),
phases: { run: async () => {} },
abort: async () => {},
hooks,
});
}
function names(tools: readonly { readonly name: string }[]): string[] {
@@ -0,0 +1,822 @@
import { mkdtemp, rm } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import type { Context, JsonValue } from "@earendil-works/chord";
import { createModels } from "@earendil-works/pi-ai";
import {
createRegistry,
defineDoc,
defineTask,
Harness,
MemoryStorage,
type Registry,
type Task,
type TaskId,
type TaskRuntime,
} from "@earendil-works/pi-durable";
import { afterEach, describe, expect, it } from "vitest";
import { openNodeSqliteStorage } from "../src/storage/sqlite/node.ts";
import { ControlledStorage, context, flush } from "./session-support.ts";
import { aborted, abortedWith, completed, countingReader, deferred, eventually, openTasks } from "./task-support.ts";
const directories = new Set<string>();
async function sqlitePath(): Promise<string> {
const directory = await mkdtemp(join(tmpdir(), "pi-durable-tasks-"));
directories.add(directory);
return join(directory, "session.sqlite");
}
afterEach(async () => {
for (const directory of directories) await rm(directory, { recursive: true, force: true });
directories.clear();
});
type OpenHarness = Awaited<ReturnType<typeof openTasks>>["harness"];
async function createIn<S extends { phase: string }, R>(
harness: OpenHarness,
task: Task<null, S, R, object>,
): Promise<TaskId<R>> {
const root = await harness.root(context);
return root.commit((tx) => tx.createTask(task, null), context);
}
/** Fake external service whose operations are idempotent by request key. */
class TransferService {
readonly applied = new Map<string, number>();
calls = 0;
async apply(key: string, amount: number): Promise<number> {
this.calls++;
if (!this.applied.has(key)) this.applied.set(key, amount * 10);
return this.applied.get(key)!;
}
}
type TransferState = { phase: "prepare" } | { phase: "apply"; key: string };
/**
* Intent/effect/outcome task: `prepare` commits the intent, `apply` performs the effect and commits the outcome.
* While `interrupt.first` is set, the first `apply` blocks after the effect until the invocation is signalled.
*/
function transferTask(service: TransferService, interrupt: { first: boolean }) {
return defineTask<{ amount: number }, TransferState, { receipt: number }>({
name: "test.transfer",
version: 1,
initial: () => ({ phase: "prepare" }),
phases: {
prepare: async (task, runtime, ctx) => {
await runtime.memo("requested", task.input.amount, ctx);
await runtime.commit(
() => ({ status: "running", checkpoint: { phase: "apply", key: `transfer-${task.id}` } }),
ctx,
);
},
apply: async (task, runtime, ctx) => {
const receipt = await service.apply(task.state.checkpoint.key, task.input.amount);
if (interrupt.first) {
interrupt.first = false;
await aborted(runtime.signal);
}
await runtime.commit(() => completed({ receipt }), ctx);
},
},
abort: async (_task, runtime, ctx) => {
await runtime.commit(() => ({ status: "terminal", outcome: { status: "aborted" } }), ctx);
},
});
}
type Step = { phase: "run" };
/** A versioned one-phase task that completes with `result`, optionally migrating older records. */
function versioned(
version: number,
result: string,
migrate?: (input: JsonValue, checkpoint: JsonValue, fromVersion: number) => { input: null; checkpoint: Step },
) {
return defineTask<null, Step, JsonValue>({
name: "test.versioned",
version,
initial: () => ({ phase: "run" }),
phases: { run: (_task, runtime, ctx) => runtime.commit(() => completed(result), ctx) },
abort: (_task, runtime, ctx) => runtime.commit(() => abortedWith(result), ctx),
...(migrate === undefined ? {} : { migrate }),
});
}
describe("task recovery", () => {
it("resumes an intent/effect/outcome task interrupted after its intent across close and reopen", async () => {
const path = await sqlitePath();
const service = new TransferService();
const Transfer = transferTask(service, { first: true });
const first = await openTasks(await openNodeSqliteStorage(path), [Transfer]);
const root = await first.harness.root(context);
const id = await root.commit((tx) => tx.createTask(Transfer, { amount: 7 }), context);
first.harness.resume();
await eventually(() => service.calls === 1);
await first.harness.close(context);
const second = await openTasks(await openNodeSqliteStorage(path), [Transfer]);
// Open reconciled `running` to `pending` and kept the checkpoint and memos; nothing ran yet.
expect(await second.harness.getTask(id, context)).toMatchObject({
abortRequested: false,
memos: { requested: 7 },
state: { status: "pending", checkpoint: { phase: "apply", key: `transfer-${id}` } },
});
expect(service.calls).toBe(1);
second.harness.resume();
const receipt = await second.harness.waitForTask(id, context);
expect(receipt.state.outcome).toEqual({ status: "completed", result: { receipt: 70 } });
expect(service.calls).toBe(2);
expect(service.applied.size).toBe(1);
await second.harness.close(context);
const third = await openTasks(await openNodeSqliteStorage(path), [Transfer]);
expect(await third.harness.getTask(id, context)).toEqual(receipt);
await third.harness.close(context);
});
it("resumes abort work after close at every direct-task abort stage", async () => {
const path = await sqlitePath();
const log: string[] = [];
const abortGate = { block: true };
const abortReached = deferred();
const runRelease = deferred();
const Abortable = defineTask<null, Step, null>({
name: "test.abortable",
version: 1,
initial: () => ({ phase: "run" }),
phases: {
run: async () => {
log.push("run");
// Ignores the abort signal until released, so the mark is durable while the run is active.
await runRelease.promise;
},
},
abort: async (_task, runtime, ctx) => {
log.push("abort");
if (abortGate.block) {
abortReached.resolve();
await aborted(runtime.signal);
}
await runtime.commit(() => abortedWith("stop"), ctx);
},
});
// Stage 1: the mark is committed while the run invocation is still active.
let opened = await openTasks(await openNodeSqliteStorage(path), [Abortable]);
const id = await createIn(opened.harness, Abortable);
opened.harness.resume();
await eventually(() => log.length === 1);
const aborting = opened.harness.abortTask(id, context);
while (!(await opened.harness.getTask(id, context))?.abortRequested) await flush();
const closing = opened.harness.close(context);
runRelease.resolve();
await closing;
expect(await aborting).toBe("marked");
expect(log).toEqual(["run"]);
// Stage 2: reopen dispatches the abort invocation, never the run; close while it is active.
opened = await openTasks(await openNodeSqliteStorage(path), [Abortable]);
expect(await opened.harness.getTask(id, context)).toMatchObject({
abortRequested: true,
state: { status: "pending" },
});
opened.harness.resume();
await abortReached.promise;
await opened.harness.close(context);
expect(log).toEqual(["run", "abort"]);
// Stage 3: a fresh abort invocation settles the task.
abortGate.block = false;
opened = await openTasks(await openNodeSqliteStorage(path), [Abortable]);
opened.harness.resume();
expect((await opened.harness.waitForTask(id, context)).state.outcome).toEqual({
status: "aborted",
reason: "stop",
});
await opened.harness.close(context);
expect(log).toEqual(["run", "abort", "abort"]);
// Stage 4: the terminal receipt survives reopen and nothing runs again.
opened = await openTasks(await openNodeSqliteStorage(path), [Abortable]);
opened.harness.resume();
await flush();
expect(await opened.harness.abortTask(id, context)).toBe("terminal");
await opened.harness.close(context);
expect(log).toEqual(["run", "abort", "abort"]);
});
});
/**
* Crash simulation: the crashed Harness is abandoned without close, its held storage commit never lands, and its
* blocked handlers never return. A new Harness then opens the same storage.
*/
describe("task crash recovery", () => {
type Log = string[];
/** A task whose run ignores its signal and whose abort handler blocks while `blockAbort` is set. */
function crashTask(log: Log, options: { blockAbort: boolean }) {
return defineTask<null, Step, null>({
name: "test.crash",
version: 1,
initial: () => ({ phase: "run" }),
phases: {
run: async () => {
log.push("run");
await new Promise(() => {});
},
},
abort: async (_task, runtime, ctx) => {
log.push("abort");
if (options.blockAbort) await new Promise(() => {});
await runtime.commit(() => abortedWith("recovered"), ctx);
},
});
}
async function recover(storage: ControlledStorage, log: Log): Promise<{ harness: OpenHarness; id: TaskId }> {
storage.crash();
const { harness } = await openTasks(storage, [crashTask(log, { blockAbort: false })]);
const page = await harness.commit((tx) => tx.scanTasks({ kind: "test.crash" }, 1), context);
return { harness, id: page.items[0]!.id };
}
async function crashedRun(log: Log): Promise<{ storage: ControlledStorage; harness: OpenHarness; id: TaskId }> {
const storage = new ControlledStorage();
const { harness } = await openTasks(storage, [crashTask(log, { blockAbort: true })]);
const id = await createIn(harness, crashTask(log, { blockAbort: true }));
harness.resume();
await eventually(() => log.length === 1);
return { storage, harness, id };
}
it("crash while the mark commit is in storage: the run resumes", async () => {
const log: Log = [];
const { storage, harness, id } = await crashedRun(log);
const held = storage.holdCommits();
void harness.abortTask(id, context);
await held.entered;
const recovered = await recover(storage, log);
expect(await recovered.harness.getTask(id, context)).toMatchObject({
abortRequested: false,
state: { status: "pending" },
});
recovered.harness.resume();
await eventually(() => log.length === 2);
expect(log).toEqual(["run", "run"]);
});
it("crash after the mark before the run joins: only the abort handler runs", async () => {
const log: Log = [];
const { storage, harness, id } = await crashedRun(log);
// The run ignores its signal, so abortTask never finishes joining it.
let joined = false;
void harness.abortTask(id, context).then(() => {
joined = true;
});
while (!(await harness.getTask(id, context))?.abortRequested) await flush();
await flush();
expect(joined).toBe(false);
const recovered = await recover(storage, log);
expect(await recovered.harness.getTask(id, context)).toMatchObject({
abortRequested: true,
state: { status: "pending" },
});
recovered.harness.resume();
expect((await recovered.harness.waitForTask(id, context)).state.outcome).toEqual({
status: "aborted",
reason: "recovered",
});
expect(log).toEqual(["run", "abort"]);
await recovered.harness.close(context);
});
it("crash while the abort handler runs: a fresh abort invocation settles the task", async () => {
const log: Log = [];
const storage = new ControlledStorage();
const Crash = crashTask(log, { blockAbort: true });
const { harness } = await openTasks(storage, [Crash]);
const id = await createIn(harness, Crash);
await harness.abortTask(id, context);
harness.resume();
await eventually(() => log.length === 1);
expect(log).toEqual(["abort"]);
expect(await harness.getTask(id, context)).toMatchObject({ state: { status: "running" } });
const recovered = await recover(storage, log);
recovered.harness.resume();
expect((await recovered.harness.waitForTask(id, context)).state.outcome).toEqual({
status: "aborted",
reason: "recovered",
});
expect(log).toEqual(["abort", "abort"]);
await recovered.harness.close(context);
});
it("crash while the abort outcome is in storage: the abort handler runs again", async () => {
const log: Log = [];
const storage = new ControlledStorage();
const reached = deferred();
const proceed = deferred();
const Crash = defineTask<null, Step, null>({
name: "test.crash",
version: 1,
initial: () => ({ phase: "run" }),
phases: { run: async () => {} },
abort: async (_task, runtime, ctx) => {
log.push("abort");
reached.resolve();
await proceed.promise;
await runtime.commit(() => abortedWith("lost"), ctx);
},
});
const { harness } = await openTasks(storage, [Crash]);
const id = await createIn(harness, Crash);
await harness.abortTask(id, context);
harness.resume();
await reached.promise;
const held = storage.holdCommits();
proceed.resolve();
await held.entered;
const recovered = await recover(storage, log);
recovered.harness.resume();
expect((await recovered.harness.waitForTask(id, context)).state.outcome).toEqual({
status: "aborted",
reason: "recovered",
});
expect(log).toEqual(["abort", "abort"]);
await recovered.harness.close(context);
});
it("crash after the terminal outcome: nothing runs again", async () => {
const log: Log = [];
const storage = new ControlledStorage();
const Crash = crashTask(log, { blockAbort: false });
const { harness } = await openTasks(storage, [Crash]);
const id = await createIn(harness, Crash);
await harness.abortTask(id, context);
harness.resume();
await harness.waitForTask(id, context);
const recovered = await recover(storage, log);
recovered.harness.resume();
await flush();
expect(log).toEqual(["abort"]);
expect((await recovered.harness.getTask(id, context))?.state).toEqual({
status: "terminal",
outcome: { status: "aborted", reason: "recovered" },
});
await recovered.harness.close(context);
});
it("crash while the reservation commit is in storage: the task is still pending", async () => {
const log: Log = [];
const storage = new ControlledStorage();
const Crash = crashTask(log, { blockAbort: false });
const { harness } = await openTasks(storage, [Crash]);
const id = await createIn(harness, Crash);
const held = storage.holdCommits();
harness.resume();
await held.entered;
const recovered = await recover(storage, log);
expect((await recovered.harness.getTask(id, context))?.state.status).toBe("pending");
expect(log).toEqual([]);
await recovered.harness.abortTask(id, context);
recovered.harness.resume();
await recovered.harness.waitForTask(id, context);
expect(log).toEqual(["abort"]);
await recovered.harness.close(context);
});
});
describe("blocked tasks", () => {
it("keeps a task with a missing definition pending, and live for idle waits, until registration", async () => {
const V1 = versioned(1, "v1");
const { harness, registry } = await openTasks(new MemoryStorage(), []);
const id = await createIn(harness, V1);
harness.resume();
await flush();
expect((await harness.getTask(id, context))?.state.status).toBe("pending");
const idle = harness.waitForIdle(context);
await flush();
registry.tasks.add(V1);
await idle;
expect((await harness.getTask(id, context))?.state).toEqual({
status: "terminal",
outcome: { status: "completed", result: "v1" },
});
await harness.close(context);
});
it("keeps a task stored by a newer version pending until a fitting definition is registered", async () => {
const registry = createRegistry();
const old = registry.tasks.add(versioned(1, "old"));
const { harness } = await openTasks(new MemoryStorage(), [], { registry });
const id = await createIn(harness, versioned(2, "new"));
harness.resume();
await flush();
expect((await harness.getTask(id, context))?.state.status).toBe("pending");
registry.batch(() => {
old.dispose();
registry.tasks.add(versioned(2, "new"));
});
expect((await harness.waitForTask(id, context)).state.outcome).toEqual({ status: "completed", result: "new" });
await harness.close(context);
});
it("migrates at reservation, leaves the record unchanged when migration fails, and retries only for a new definition", async () => {
const path = await sqlitePath();
let opened = await openTasks(await openNodeSqliteStorage(path), []);
const id = await createIn(opened.harness, versioned(1, "v1"));
await opened.harness.close(context);
const registry: Registry = createRegistry();
let failures = 0;
const failing = registry.tasks.add(
versioned(2, "v2", () => {
failures++;
throw new Error("cannot migrate");
}),
);
opened = await openTasks(await openNodeSqliteStorage(path), [], { registry });
opened.harness.resume();
await eventually(() => opened.reports.length === 1);
// An unrelated registry change wakes the scheduler without retrying the same failed definition.
registry.tools.add({
name: "unrelated",
description: "unrelated",
parameters: { type: "object", properties: {} } as never,
execute: async () => ({}),
});
await flush();
expect(failures).toBe(1);
expect(opened.reports).toHaveLength(1);
expect(String(opened.reports[0])).toContain("cannot migrate");
expect(await opened.harness.getTask(id, context)).toMatchObject({
version: 1,
state: { status: "pending", checkpoint: { phase: "run" } },
});
const migrations: number[] = [];
registry.batch(() => {
failing.dispose();
registry.tasks.add(
versioned(2, "v2", (_input, checkpoint, fromVersion) => {
migrations.push(fromVersion);
return { input: null, checkpoint: checkpoint as Step };
}),
);
});
const receipt = await opened.harness.waitForTask(id, context);
expect(receipt).toMatchObject({ version: 2, state: { outcome: { status: "completed", result: "v2" } } });
expect(migrations).toEqual([1]);
await opened.harness.close(context);
});
it("blocks an older record whose newer definition has no migration", async () => {
const path = await sqlitePath();
let opened = await openTasks(await openNodeSqliteStorage(path), []);
const id = await createIn(opened.harness, versioned(1, "v1"));
await opened.harness.close(context);
opened = await openTasks(await openNodeSqliteStorage(path), [versioned(2, "v2")]);
opened.harness.resume();
await eventually(() => opened.reports.length === 1);
expect(String(opened.reports[0])).toContain("has no migration from 1");
expect(await opened.harness.getTask(id, context)).toMatchObject({ version: 1, state: { status: "pending" } });
await opened.harness.close(context);
});
it("settles an aborted blocked task as orphaned and retires its documents", async () => {
const Scratch = defineDoc<{ n: number }>({
kind: "test.orphan-scratch",
version: 1,
scope: "task",
initial: () => ({ n: 0 }),
});
const { harness } = await openTasks(new MemoryStorage(), []);
const root = await harness.root(context);
const id = await root.commit(async (tx) => {
const created = await tx.createTask(versioned(1, "x"), null);
(await tx.doc(Scratch, created)).n = 1;
return created;
}, context);
// Before resume: the marking commit settles the blocked task directly.
expect(await harness.abortTask(id, context)).toBe("marked");
expect((await harness.getTask(id, context))?.state).toEqual({
status: "terminal",
outcome: { status: "orphaned", reason: "missing_task" },
});
expect(await harness.snapshot(Scratch, id, context)).toBeUndefined();
await harness.close(context);
});
it("orphans a marked task whose definition disappeared while its run was active", async () => {
const reached = deferred();
const Running = defineTask<null, Step, null>({
name: "test.vanishing",
version: 1,
initial: () => ({ phase: "run" }),
phases: {
run: async (_task, runtime) => {
reached.resolve();
await aborted(runtime.signal);
},
},
abort: async () => {
throw new Error("must not run");
},
});
const registry = createRegistry();
const registration = registry.tasks.add(Running);
const { harness } = await openTasks(new MemoryStorage(), [], { registry });
const id = await createIn(harness, Running);
harness.resume();
await reached.promise;
registration.dispose();
// An active run means the mark does not orphan directly; the scheduler orphans once the run has ended.
expect(await harness.abortTask(id, context)).toBe("marked");
expect((await harness.waitForTask(id, context)).state.outcome).toEqual({
status: "orphaned",
reason: "missing_task",
});
await harness.close(context);
});
it("orphans a reopened abort-marked task without a definition once scheduling resumes", async () => {
const path = await sqlitePath();
let opened = await openTasks(await openNodeSqliteStorage(path), [versioned(1, "x")]);
const id = await createIn(opened.harness, versioned(1, "x"));
expect(await opened.harness.abortTask(id, context)).toBe("marked");
await opened.harness.close(context);
opened = await openTasks(await openNodeSqliteStorage(path), []);
expect(await opened.harness.getTask(id, context)).toMatchObject({
abortRequested: true,
state: { status: "pending" },
});
const idle = opened.harness.waitForIdle(context);
opened.harness.resume();
await idle;
expect((await opened.harness.getTask(id, context))?.state).toEqual({
status: "terminal",
outcome: { status: "orphaned", reason: "missing_task" },
});
await opened.harness.close(context);
});
});
describe("definition handover", () => {
type Handover = { phase: "a" } | { phase: "b" } | { phase: "c" };
type Gates = { readonly a?: Promise<void>; readonly b?: Promise<void>; readonly onEnd?: () => void };
function handoverTask(
label: string,
version: number,
log: string[],
options: { readonly gates?: Gates; readonly migrate?: () => { input: null; checkpoint: Handover } } = {},
) {
const advance =
(phase: "a" | "b", next: Handover) =>
async (_task: unknown, runtime: TaskRuntime<null, Handover, null, object>, ctx: Context) => {
log.push(`${label}:${phase} start`);
await options.gates?.[phase];
await runtime.commit(() => ({ status: "running", checkpoint: next }), ctx);
// Leave room for a wrongly dispatched successor before this invocation ends.
await flush();
options.gates?.onEnd?.();
log.push(`${label}:${phase} end`);
};
return defineTask<null, Handover, null>({
name: "test.handover",
version,
initial: () => ({ phase: "a" }),
phases: {
a: advance("a", { phase: "b" }),
b: advance("b", { phase: "c" }),
c: async (_task, runtime, ctx) => {
log.push(`${label}:c`);
await runtime.commit(() => completed(null), ctx);
},
},
abort: async (_task, runtime, ctx) => {
log.push(`${label}:abort`);
await runtime.commit(() => abortedWith(label), ctx);
},
...(options.migrate === undefined ? {} : { migrate: options.migrate }),
});
}
async function startHandover(log: string[], gates: Gates, storage = new MemoryStorage()) {
const registry = createRegistry();
const old = registry.tasks.add(handoverTask("old", 1, log, { gates }));
const opened = await openTasks(storage, [], { registry });
const id = await createIn(opened.harness, handoverTask("old", 1, log));
opened.harness.resume();
await eventually(() => log.length === 1);
return { ...opened, old, id };
}
it("hands over at the next phase boundary to a same-version replacement without overlap", async () => {
const log: string[] = [];
const gate = deferred();
const { harness, registry, old, id } = await startHandover(log, { a: gate.promise });
registry.batch(() => {
old.dispose();
registry.tasks.add(handoverTask("new", 1, log));
});
gate.resolve();
await harness.waitForTask(id, context);
expect(log).toEqual(["old:a start", "old:a end", "new:b start", "new:b end", "new:c"]);
await harness.close(context);
});
it("hands over to a newer version with a migration", async () => {
const log: string[] = [];
const gate = deferred();
const { harness, registry, old, id } = await startHandover(log, { a: gate.promise });
registry.batch(() => {
old.dispose();
registry.tasks.add(
handoverTask("v2", 2, log, { migrate: () => ({ input: null, checkpoint: { phase: "c" } }) }),
);
});
gate.resolve();
const receipt = await harness.waitForTask(id, context);
expect(receipt.version).toBe(2);
expect(log).toEqual(["old:a start", "old:a end", "v2:c"]);
await harness.close(context);
});
it("hands over to a newer version whose migration fails and leaves the task blocked", async () => {
const log: string[] = [];
const gate = deferred();
const { harness, registry, old, id, reports } = await startHandover(log, { a: gate.promise });
registry.batch(() => {
old.dispose();
registry.tasks.add(
handoverTask("broken", 2, log, {
migrate: () => {
throw new Error("broken migration");
},
}),
);
});
gate.resolve();
await eventually(() => reports.length === 1);
await flush();
expect(await harness.getTask(id, context)).toMatchObject({
version: 1,
state: { status: "pending", checkpoint: { phase: "b" } },
});
expect(log).toEqual(["old:a start", "old:a end"]);
await harness.close(context);
});
it("keeps running under the old definition when the replacement is missing or cannot take the task", async () => {
const log: string[] = [];
const gateA = deferred();
const gateB = deferred();
const { harness, registry, old, id, reports } = await startHandover(log, { a: gateA.promise, b: gateB.promise });
old.dispose();
gateA.resolve();
await eventually(() => log.includes("old:b start"));
// A newer definition without a migration cannot take the task either.
registry.tasks.add(handoverTask("incompatible", 2, log));
gateB.resolve();
await harness.waitForTask(id, context);
expect(log).toEqual(["old:a start", "old:a end", "old:b start", "old:b end", "old:c"]);
expect(reports.map((report) => (report as Error).cause)).toEqual(["missing_task", "incompatible_task"]);
await harness.close(context);
});
it("rejects a runtime commit of the old invocation queued behind its handover commit", async () => {
const log: string[] = [];
const gate = deferred();
const storage = new ControlledStorage();
let held: ReturnType<ControlledStorage["holdCommits"]> | undefined;
let oldRuntime: TaskRuntime<null, Handover, null, object> | undefined;
let harnessRef: OpenHarness | undefined;
const registry = createRegistry();
const Old = defineTask<null, Handover, null>({
name: "test.handover",
version: 1,
initial: () => ({ phase: "a" }),
phases: {
a: async (task, runtime, ctx) => {
oldRuntime = runtime;
await gate.promise;
await runtime.commit(() => ({ status: "running", checkpoint: { phase: "b" } }), ctx);
// Hold the line with an unrelated commit so the handover commit queues behind it.
held = storage.holdCommits();
void harnessRef!.commit(async (tx) => {
await tx.appendEntry(task.conversationId, { kind: "blocker" });
}, ctx);
},
b: async () => {},
c: async () => {},
},
abort: async () => {},
});
const old = registry.tasks.add(Old);
const { harness } = await openTasks(storage, [], { registry });
harnessRef = harness;
const id = await createIn(harness, Old);
harness.resume();
await eventually(() => oldRuntime !== undefined);
registry.batch(() => {
old.dispose();
registry.tasks.add(handoverTask("new", 1, log));
});
gate.resolve();
await eventually(() => held !== undefined);
await held!.entered;
await flush();
// Queued behind the handover commit while the invocation has not ended yet.
const late = oldRuntime!.commit(() => completed(null), context);
held!.release();
await expect(late).rejects.toThrow("invocation has ended");
await harness.waitForTask(id, context);
expect(log).toEqual(["new:b start", "new:b end", "new:c"]);
await harness.close(context);
});
it("preserves an abort mark that races the handover commit; the new definition aborts", async () => {
const log: string[] = [];
const gate = deferred();
const storage = new ControlledStorage();
let held: ReturnType<ControlledStorage["holdCommits"]> | undefined;
// The progress commit landed; hold the next commit, the handover, and queue the abort mark behind it.
const gates: Gates = {
a: gate.promise,
onEnd: () => {
held ??= storage.holdCommits();
},
};
const { harness, registry, old, id } = await startHandover(log, gates, storage);
registry.batch(() => {
old.dispose();
registry.tasks.add(handoverTask("new", 1, log));
});
gate.resolve();
await eventually(() => held !== undefined);
await held!.entered;
const aborting = harness.abortTask(id, context);
held!.release();
expect(await aborting).toBe("marked");
expect((await harness.waitForTask(id, context)).state.outcome).toEqual({ status: "aborted", reason: "new" });
const states = storage.commits.flatMap((writes) =>
writes.flatMap((write) =>
write.type === "task" && write.value.id === id
? [`${write.value.state.status}${write.value.abortRequested ? "+mark" : ""}`]
: [],
),
);
// created, reserved, progress, handover, mark, abort reservation, aborted
expect(states).toEqual([
"pending",
"running",
"running",
"pending",
"pending+mark",
"running+mark",
"terminal+mark",
]);
expect(log).toEqual(["old:a start", "old:a end", "new:abort"]);
await harness.close(context);
});
});
describe("Harness open", () => {
it("releases its registry subscription and closes the Session when open fails", async () => {
const storage = new ControlledStorage();
const log: string[] = [];
const Stuck = defineTask<null, Step, null>({
name: "test.stuck",
version: 1,
initial: () => ({ phase: "run" }),
phases: {
run: async () => {
log.push("run");
await new Promise(() => {});
},
},
abort: async () => {},
});
// Leave a running task behind so open has a reconciliation commit to fail.
const first = await openTasks(storage, [Stuck]);
await createIn(first.harness, Stuck);
first.harness.resume();
await eventually(() => log.length === 1);
const registry = createRegistry();
const reader = countingReader(registry);
storage.failNextCommit(new Error("disk full"));
await expect(Harness.open(storage, { models: createModels(), registry: reader }, context)).rejects.toThrow(
"disk full",
);
expect(reader.subscriptions()).toBe(0);
await expect(storage.task(1 as TaskId, context)).rejects.toThrow();
});
});
File diff suppressed because it is too large Load Diff
+153 -11
View File
@@ -1,6 +1,9 @@
// A tour of the durable Session and Harness APIs in twelve small examples.
// A tour of the durable Session and Harness APIs in fourteen small examples.
// Run from packages/durable:
// node --conditions=source --experimental-strip-types test/scratch.ts
import { mkdtemp, rm } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { BACKGROUND_CONTEXT } from "@earendil-works/chord/context";
import {
type AssistantMessage,
@@ -17,13 +20,14 @@ import {
createSession,
defineDoc,
defineEntry,
defineTask,
type EntryRecord,
Harness,
MemoryStorage,
type PromptInput,
type Task,
type ToolRegistration,
} from "../src/index.ts";
import { openNodeSqliteStorage } from "../src/storage/sqlite/node.ts";
// A Session stores conversations, transcript entries, tasks, and documents.
// MemoryStorage keeps everything in memory; other storage backends keep it on disk.
@@ -109,15 +113,16 @@ console.log("3. parent notes after edit:", await session.snapshot(Notes, chat.id
// conversations. All three are created in one commit, so after a crash either
// all of them exist or none do.
// A task definition needs a name, a version, and the task's starting state.
// This example only creates the task record; nothing runs it yet.
const Supervisor: Task<null, { phase: "ready" }, null, object> = {
definition: {
name: "example.supervisor",
version: 1,
initial: () => ({ phase: "ready" }),
},
};
// A task definition needs a name, a version, the task's starting state, a
// handler for every phase, and an abort handler (example 13 runs a task).
// This example only creates the task record; a plain Session never runs it.
const Supervisor = defineTask<null, { phase: "ready" }, null>({
name: "example.supervisor",
version: 1,
initial: () => ({ phase: "ready" }),
phases: { ready: async () => {} },
abort: async () => {},
});
// "latest" keeps only the current value. "initial" means forks of this
// conversation start without a registry, so a child doesn't inherit its
@@ -478,4 +483,141 @@ console.log(`12. rendered prompt:\n${rendered.join("\n")}`);
prompt.dispose();
// ─── 13. Run a durable task ─────────────────────────────────────────────────
// A task is a small state machine. Its state, the checkpoint, is saved after
// every step, so after a crash the next open continues from the last saved
// step. The usual pattern: save what you are about to do, do it, then save
// the result. A crash between doing and saving reruns that step, so the step
// must be safe to repeat; here the fake payment service ignores a repeated key.
const payments = new Map<string, number>();
type PaymentState = { phase: "prepare" } | { phase: "charge"; key: string };
const Payment = defineTask<{ amount: number }, PaymentState, { receipt: number }>({
name: "example.payment",
version: 1,
initial: () => ({ phase: "prepare" }),
// One handler per phase. Each must save progress through runtime.commit():
// its callback returns the next checkpoint or the final outcome, and that
// state is saved in the same commit as everything else the callback wrote.
phases: {
prepare: async (task, runtime, taskContext) => {
await runtime.commit(
() => ({ status: "running", checkpoint: { phase: "charge", key: `payment-${task.id}` } }),
taskContext,
);
},
charge: async (task, runtime, taskContext) => {
const key = task.state.checkpoint.key;
if (!payments.has(key)) payments.set(key, task.input.amount * 100);
const receipt = payments.get(key)!;
await runtime.commit(
() => ({ status: "terminal", outcome: { status: "completed", result: { receipt } } }),
taskContext,
);
},
},
// Runs instead of the phases after harness.abortTask(); it decides the outcome.
abort: async (_task, runtime, taskContext) => {
await runtime.commit(() => ({ status: "terminal", outcome: { status: "aborted" } }), taskContext);
},
});
// The Harness finds task code by name in the registry. Nothing runs until
// resume(); a host calls it once it is ready for work to start.
registry.tasks.add(Payment);
const paymentId = await root.commit((tx) => tx.createTask(Payment, { amount: 5 }), context);
harness.resume();
// The finished task record is the durable receipt; waitForTask() knows its result type.
const paid = await harness.waitForTask(paymentId, context);
console.log("13. payment outcome:", paid.state.outcome);
await harness.close(context);
// ─── 14. Close, reopen, and continue where the task stopped ─────────────────
// Everything a task needs to continue is in storage, so a new Harness over the
// same storage picks up where the last one stopped. This example keeps its
// storage in a SQLite file so it survives closing.
const directory = await mkdtemp(join(tmpdir(), "pi-durable-scratch-"));
const databasePath = join(directory, "session.sqlite");
let reachedTick = (_n: number): void => {};
const Ticker = defineTask<{ to: number }, { phase: "tick"; n: number }, string>({
name: "example.ticker",
version: 1,
initial: () => ({ phase: "tick", n: 1 }),
phases: {
tick: async (task, runtime, taskContext) => {
const n = task.state.checkpoint.n;
// Save the intent before the effect. A memo keeps the first value
// written under its name, so if the process dies after printing but
// before the next checkpoint is saved, the rerun sees the memo and
// does not print the same tick twice.
if ((await runtime.memo(`printed-${n}`, taskContext)) === undefined) {
await runtime.memo(`printed-${n}`, true, taskContext);
console.log(`14. tick ${n}`);
}
reachedTick(n);
// Save the outcome: the next tick, or the final result.
await runtime.commit(
() =>
n === task.input.to
? { status: "terminal", outcome: { status: "completed", result: `counted to ${n}` } }
: { status: "running", checkpoint: { phase: "tick", n: n + 1 } },
taskContext,
);
// Wait a little between ticks. Closing the Harness cancels this wait;
// the checkpoint saved above is where the next Harness continues.
await runtime.sleep(Date.now() + 50, taskContext);
},
},
abort: async (_task, runtime, taskContext) => {
await runtime.commit(() => ({ status: "terminal", outcome: { status: "aborted" } }), taskContext);
},
});
registry.tasks.add(Ticker);
// First run: start counting to 5, and close the Harness right after tick 2 is
// printed, before its next checkpoint is saved. That is the same situation as
// a crash between the effect and saving its outcome.
const firstRun = await Harness.open(
await openNodeSqliteStorage(databasePath),
{ models: createModels(), registry },
context,
);
const tickerId = await (await firstRun.root(context)).commit((tx) => tx.createTask(Ticker, { to: 5 }), context);
const tickTwo = new Promise<void>((resolve) => {
reachedTick = (n) => {
if (n === 2) resolve();
};
});
firstRun.resume();
await tickTwo;
await firstRun.close(context);
reachedTick = () => {};
const saved = await readTicker();
console.log("14. closed; saved checkpoint:", saved.state, "memos:", saved.memos);
// Second run: nothing to restart by hand. Opening the storage finds the
// unfinished task and resume() continues it. Tick 2 runs again because its
// outcome was never saved, but its memo says it was already printed.
const secondRun = await Harness.open(
await openNodeSqliteStorage(databasePath),
{ models: createModels(), registry },
context,
);
secondRun.resume();
const counted = await secondRun.waitForTask(tickerId, context);
console.log("14. after reopen:", counted.state.outcome);
await secondRun.close(context);
await rm(directory, { recursive: true, force: true });
/** Read the ticker record through a short-lived Harness over the same file. */
async function readTicker() {
const reader = await Harness.open(
await openNodeSqliteStorage(databasePath),
{ models: createModels(), registry },
context,
);
const record = await reader.getTask(tickerId, context);
await reader.close(context);
return record!;
}
@@ -351,8 +351,8 @@ describe("Session document transactions", () => {
it("rejects Tx use after the callback settles", async () => {
const { session, conversationId } = await setupLive();
let captured: Parameters<Parameters<typeof session.commit>[0]>[0] | undefined;
await session.commit((tx) => {
let captured: Parameters<Parameters<typeof session.commitWith>[0]>[0] | undefined;
await session.commitWith((tx) => {
captured = tx;
}, context);
await expect(captured!.doc(LiveDoc, conversationId)).rejects.toThrow("Transaction has settled");
@@ -411,7 +411,7 @@ describe("Session document transactions", () => {
const counter = await session.snapshot(CounterDoc, context);
const commits = storage.commits.length;
await expect(
session.commit(async (tx) => {
session.commitWith(async (tx) => {
(await tx.doc(LiveDoc, conversationId)).message = "lost";
(await tx.doc(CounterDoc)).count = 2;
// Replacing a missing task fails during assembly, after every change was prepared.
+8 -7
View File
@@ -2,6 +2,7 @@ import {
type ConversationId,
defineDoc,
defineDocFamily,
defineTask,
type EntryId,
StorageRejected,
type StorageWrite,
@@ -627,13 +628,13 @@ describe("Session conversation document forks", () => {
scope: "task",
initial: () => ({ value: "task" }),
});
const Work = {
definition: {
name: "fork.scope.work",
version: 1,
initial: () => ({ phase: "start" }),
},
};
const Work = defineTask<null, { phase: "start" }, null>({
name: "fork.scope.work",
version: 1,
initial: () => ({ phase: "start" }),
phases: { start: async () => {} },
abort: async () => {},
});
const { session, storage, publications } = openTestSession();
const parentId = await createConversation(session);
let forkAt!: EntryId;
+5
View File
@@ -57,6 +57,11 @@ export class ControlledStorage extends MemoryStorage {
return { entered: held.entered.promise, release: () => this.#release("find", held) };
}
/** Simulate a crash during the held commit: it never reaches storage, and later commits proceed. */
crash(): void {
this.#commitGate = undefined;
}
failNextCommit(error: Error): void {
this.#commitFailure = error;
}
+18 -14
View File
@@ -3,9 +3,9 @@ import {
type ConversationId,
defineDoc,
defineDocFamily,
defineTask,
type EntryId,
ReadAfterWrite,
type Task,
type TaskId,
type TaskRecord,
} from "@earendil-works/pi-durable";
@@ -14,9 +14,13 @@ import { idFromNumber } from "../src/ids.ts";
import { context, createConversation, documentChanges, flush, openTestSession } from "./session-support.ts";
type Checkpoint = { phase: "start" } | { phase: "next"; step: number };
const WorkTask: Task<{ path: string }, Checkpoint, { ok: boolean }, object> = {
definition: { name: "test.work", version: 1, initial: () => ({ phase: "start" }) },
};
const WorkTask = defineTask<{ path: string }, Checkpoint, { ok: boolean }>({
name: "test.work",
version: 1,
initial: () => ({ phase: "start" }),
phases: { start: async () => {}, next: async () => {} },
abort: async () => {},
});
type Progress = { lines: string[] };
const ProgressDoc = defineDoc<Progress>({
@@ -115,7 +119,7 @@ describe("Session transaction tables", () => {
const { session } = openTestSession();
const conversationId = await createConversation(session);
const taskId = await createTask(session, conversationId);
await session.commit(async (tx) => {
await session.commitWith(async (tx) => {
const task = (await tx.task(taskId))!;
tx.setTask(task);
await expect(tx.task(taskId)).rejects.toBeInstanceOf(ReadAfterWrite);
@@ -193,7 +197,7 @@ describe("Session transaction tables", () => {
expect(await storage.conversation(created.child.id, context)).toEqual(created.child);
await expect(
session.commit(async (tx) => {
session.commitWith(async (tx) => {
const supervisor = (await tx.task(created.supervisorId))!;
tx.setTask({ ...supervisor, conversationId: created.child.id });
}, context),
@@ -206,7 +210,7 @@ describe("Session transaction tables", () => {
const parentId = await createConversation(session);
let movedStagedTaskId: TaskId<{ ok: boolean }> | undefined;
await expect(
session.commit(async (tx) => {
session.commitWith(async (tx) => {
const taskId = await tx.createTask(WorkTask, { path: "move" }, { conversationId: parentId });
movedStagedTaskId = taskId;
tx.setTask({
@@ -232,7 +236,7 @@ describe("Session transaction tables", () => {
let rejectedChildId: ConversationId | undefined;
await expect(
session.commit(async (tx) => {
session.commitWith(async (tx) => {
const supervisorId = await tx.createTask(WorkTask, { path: "aborting" }, { conversationId: parentId });
rejectedChildId = (await tx.createConversation({ ownership: { kind: "task", taskId: supervisorId } })).id;
tx.setTask({
@@ -253,7 +257,7 @@ describe("Session transaction tables", () => {
expect(await storage.conversation(rejectedChildId, context)).toBeUndefined();
await expect(
session.commit(async (tx) => {
session.commitWith(async (tx) => {
const supervisorId = await tx.createTask(WorkTask, { path: "terminal" }, { conversationId: parentId });
await tx.createConversation({ ownership: { kind: "task", taskId: supervisorId } });
tx.setTask({
@@ -271,7 +275,7 @@ describe("Session transaction tables", () => {
).rejects.toThrow("is terminal");
const terminalOwnerId = await createTask(session, parentId);
await session.commit(async (tx) => {
await session.commitWith(async (tx) => {
tx.setTask(terminal((await tx.task(terminalOwnerId))!));
}, context);
await expect(
@@ -313,7 +317,7 @@ describe("Session transaction tables", () => {
const { session, storage } = openTestSession();
const conversationId = await createConversation(session);
const taskId = await createTask(session, conversationId);
await session.commit(async (tx) => {
await session.commitWith(async (tx) => {
const task = (await tx.task(taskId))!;
tx.setTask({
id: task.id,
@@ -367,7 +371,7 @@ describe("Session transaction tables", () => {
const { session } = openTestSession();
const conversationId = await createConversation(session);
const taskId = await createTask(session, conversationId, true);
await session.commit(async (tx) => {
await session.commitWith(async (tx) => {
const task = (await tx.task(taskId))!;
const progress = await tx.doc(ProgressDoc, taskId);
tx.setTask(terminal(task));
@@ -390,7 +394,7 @@ describe("Session transaction tables", () => {
}, context);
await flush();
const published = publications.length;
await session.commit(async (tx) => {
await session.commitWith(async (tx) => {
const task = (await tx.task(taskId))!;
(await tx.doc(StepDoc, taskId, "new", null)).lines.push("created then retired");
tx.setTask(terminal(task));
@@ -421,7 +425,7 @@ describe("Session transaction tables", () => {
);
expect(alive.items).toEqual([]);
await expect(
session.commit(async (tx) => {
session.commitWith(async (tx) => {
tx.setTask(terminal((await tx.task(taskId))!));
}, context),
).rejects.toThrow(`Task ${taskId} is already terminal`);
+110
View File
@@ -0,0 +1,110 @@
import { createModels } from "@earendil-works/pi-ai";
import {
type AnyTask,
createRegistry,
Harness,
type Registry,
type RegistryReader,
type Storage,
} from "@earendil-works/pi-durable";
import { context, flush } from "./session-support.ts";
export type Deferred<T = void> = {
readonly promise: Promise<T>;
resolve(value: T): void;
reject(error: unknown): void;
};
export function deferred<T = void>(): Deferred<T> {
let resolve!: (value: T) => void;
let reject!: (error: unknown) => void;
const promise = new Promise<T>((done, fail) => {
resolve = done;
reject = fail;
});
return { promise, resolve, reject };
}
/** Reject with the signal's reason once it aborts; for handlers that block until cancelled. */
export function aborted(signal: AbortSignal): Promise<never> {
return new Promise((_, reject) => {
if (signal.aborted) reject(signal.reason);
signal.addEventListener("abort", () => reject(signal.reason), { once: true });
});
}
/** Flush macrotask turns until `check` holds. */
export async function eventually(check: () => boolean): Promise<void> {
for (let attempt = 0; attempt < 200; attempt++) {
if (check()) return;
await flush();
}
throw new Error("Condition was not reached");
}
/** Whether `promise` has settled after pending work flushes. */
export async function settled(promise: Promise<unknown>): Promise<boolean> {
let done = false;
promise.then(
() => {
done = true;
},
() => {
done = true;
},
);
await flush();
return done;
}
/** Open a Harness whose registry holds `tasks`; failures passed to `onReport` are collected. */
export async function openTasks(
storage: Storage,
tasks: readonly AnyTask[],
options: { readonly registry?: Registry; readonly now?: () => number } = {},
): Promise<{ readonly harness: Harness; readonly registry: Registry; readonly reports: unknown[] }> {
const registry = options.registry ?? createRegistry();
for (const task of tasks) registry.tasks.add(task);
const reports: unknown[] = [];
const harness = await Harness.open(
storage,
{
models: createModels(),
registry,
onReport: (error) => reports.push(error),
...(options.now === undefined ? {} : { now: options.now }),
},
context,
);
return { harness, registry, reports };
}
/** Next state that completes a task with `result`. */
export function completed<R>(result: R) {
return { status: "terminal", outcome: { status: "completed", result } } as const;
}
/** Next state that aborts a task. */
export function abortedWith(reason: string) {
return { status: "terminal", outcome: { status: "aborted", reason } } as const;
}
/** Registry reader that counts live subscriptions. */
export function countingReader(registry: Registry): RegistryReader & { subscriptions(): number } {
let count = 0;
return {
snapshot: () => registry.snapshot(),
subscribe: (listener) => {
count++;
const unsubscribe = registry.subscribe(listener);
let active = true;
return () => {
if (!active) return;
active = false;
count--;
unsubscribe();
};
},
subscriptions: () => count,
};
}