Durable agents in Effect: analysis and plan
Goal: reproduce the capabilities and durability guarantees of @earendil-works/pi-durable
(repos/pi/packages/durable) with Effect v4 (repos/effect), using effect/ai for model access.
Sources: repos/pi/packages/durable/README.md, docs/spec.md (normative), src/; Effect LLMS.md,
MIGRATION.md, packages/effect/src/{ai,workflow,persistence,sql,eventlog,reactivity}.
Part 1 — What pi-durable is
Section titled “Part 1 — What pi-durable is”1.1 Core claim
Section titled “1.1 Core claim”Conversations, model turns, tool calls, and your own state are committed to storage before anything is shown. If the process dies mid-turn, reopening the storage picks the work up where it stopped.
It provides this with three layers, from the bottom up:
| Layer | Responsibility |
|---|---|
| Storage | A backend interface (Memory, SQLite, JSONL) that applies a batch of StorageWrites atomically and returns a monotonically increasing Seq. |
| Session | A single serialized commit line. Transactions group entry appends, typed document edits, task changes and submission changes. Publication to observers happens only after storage succeeds. |
| Harness | A durable task scheduler plus the built-in agent tasks (pi.generation, pi.tool, pi.compaction), submissions/inbox, the extension registry, the system prompt, views and events. |
1.2 Data model
Section titled “1.2 Data model”- Conversation:
{id, parent?: {conversationId, at}, owner?: {conversationId, taskId}}. Root is id1. It is immutable; a fork points at a parent entry. - Entry: an immutable transcript record
{id, conversationId, kind, model?: Message[], data?, head?, edits?, byTaskId?}. Built-in kinds arepi.user,pi.assistant,pi.system,pi.tool-result,pi.reset,pi.compaction.headmakes an entry a context head (resets and compactions).edits(omit/replace) rewrite earlier entries in the model’s view without mutating them.
- Document: typed JSON next to the transcript, defined with
defineDoc({kind, version, scope, history, fork, initial, migrate, checkpointWhen}).scopeissession,taskorconversation.historyislatestorrewindable(snapshotAsOf).forkisinitial,currentorasOf.- Documents are stored as base and delta revisions. A document has incarnations with a
[createdAt, retiredAt)lifetime. - Built-ins:
pi.agent: model, thinking level, extensions, tools, instructions, cwd.pi.live: the running run, generation partial, tool slots and compactions.pi.inbox: queued submissions.pi.usage: tokens and cost.
- Task: a durable state machine
{id, conversationId, kind, version, input, owner?, background, abortRequested, state}.- State is one of
pending | running | waiting{on, policy} | completing{outcome} | terminal{outcome}. - Outcome is one of
completed | failed | aborted | orphaned | faulted.
- State is one of
- Submission: user input or a write handed to a conversation, idempotent by
requestId. Its lifecycle isqueued → placed → done | unanswered{reason}.
1.3 Task semantics (the durable heart)
Section titled “1.3 Task semantics (the durable heart)”- A definition is
{name, version, initial(input), phases: {[phase]: handler}, abort, migrate?}. The checkpoint carriesphase. - A phase handler advances only by calling
runtime.commit(tx => newState). The new state is committed atomically with the entries and documents the handler wrote. - A phase that returns without changing the checkpoint is
faulted(“no durable progress”). Each phase is therefore a step between two durable checkpoints. There is no replay of a function body: on a crash the task re-enters at its last committed phase. - Ownership tree: tasks and conversations can be owned by a task.
- Abort is bottom-up: owned work drains first, then the owner’s
aborthandler runs. - A task that finishes while it still owns live work becomes
completingand turns terminal later. background: truetasks are a boundary that ordinary aborts and idle checks do not cross.
- Abort is bottom-up: owned work drains first, then the owner’s
- Waiting:
waiting{on: TaskId[], policy: allSettled | failFast}. No code runs while a task waits.failFastaborts the remaining siblings. - Recovery: on open,
runningis reset topending. A reconcile re-derives the abort cascades. Nothing is dispatched untilresume()or a progress call. - Versioning: a stored version newer than the code blocks the task. An older stored version runs
migratewhen the task is reserved. A task with a missing definition stays blocked, never failed.
1.4 Built-in agent tasks
Section titled “1.4 Built-in agent tasks”submit(input) → pi.user, live.run set (busy) pi.generation prepare → request → (answer | toolRound | retry | poll | overflow→compaction) answer: pi.assistant, onYield hook, final inbox boundary, endRun → submission done toolRound: pi.tool × n owned by generation; waiting{allSettled} pi.tool call (resolve, validate, beforeTool) → intent commit → execute → afterTool → pi.tool-result on crash: rerun if replay:"safe" (stored AND current), else an "interrupted" result with partial output pi.generation tools phase → afterTools, addTools, terminate/handoff, postTools boundary → next generation pi.compaction select cut → summarize (own retry) → place pi.compaction entry (direct or via admission)- Streaming persistence:
- Generation partials are committed to
pi.liveon a trailing 100 ms timer, with one commit in flight. - Tool output uses an adaptive throttle of
max(100ms, bytes/100KiB/s). - A crash therefore loses at most about 100 ms of output.
- Generation partials are committed to
- Inbox placement:
- Writes are always placed.
steeritems join after the current tool round, one at a time or all withsteeringMode: "all". followUpitems start the next run.whenBusy: "reject"throwsConversationBusy.
- Writes are always placed.
- Context derivation:
- Start from the newest head marker.
- Apply edits, newest wins.
- Drop aborted, error and deferred assistant messages.
- Move each tool result directly after its assistant message, synthesizing missing results.
- System prompt: extension sections are rendered on each request. Only changes are stored as positional
pi.systementries, which keeps provider prompt caches warm. - Compaction: automatic (blocking above
contextWindow - reserveTokens, background earlier) or manual. After a context-overflow error it compacts and retries once. Stale summaries are dropped.
1.5 Extensibility and observation
Section titled “1.5 Extensibility and observation”- Extension:
{name, tools, sections, hooks, wraps, tasks}.- Extensions are installed into a live Registry. Reinstalling under the same name hot-swaps the extension.
- Each conversation stores extension names, which are resolved again in every task phase.
- Hooks:
beforeRequest,afterResponse,onYield,afterToolson generation.beforeTool,afterToolon tools.beforeCompacton compaction.
- Tool:
{name, description, parameters (TypeBox), execute(args, api, ctx), replay?, executionMode?}.- The
apioffersoutput()streaming,details(),commit(),conversation(id)for subagents,envandtaskId. - A result may carry
usageandcontrol: {terminate | handoff}.
- The
- Observation:
viewState()andwatch()deliver a structural view (entries plus docs) with per-commit ops. A watch keeps at most 100 pending frames and collapses them into a snapshot on overflow.watchEvents()gives coding-agent style events (message_update,tool_execution_*, …) derived from commits.taskGraph()shows the live task tree.
- Environment:
env({conversationId, cwd, read})builds anExecutionEnv(fs and exec) per call. The built-in toolsread,write,editandbashuse it. - Settings: live getters (retry, stream timeout, compaction, tool execution, queue modes). They are never persisted.
1.6 Invariants a reimplementation must keep
Section titled “1.6 Invariants a reimplementation must keep”These are condensed from the spec; full list in the agent analysis.
- One serialized commit line. Publication happens only after storage succeeds.
StorageRejectedrolls back; any other storage failure poisons the session. - Once storage has admitted a commit, cancellation cannot stop it.
- Every task transition is decided in one step on the line. A phase without durable progress faults.
- Tool intent is durable before
execute.beforeToolnever reruns. A replay requiressafein both the stored and the current declaration. - A conversation is busy exactly while
pi.live.runis set. Every terminal path of a run clearsrun,generationandtoolsin the same commit that settles its inputs. - Usage is updated in the same commit as the assistant entry or tool result that produced it.
- Abort is bottom-up. Cascade marks are written in a separate reconcile commit and re-derived on open.
- New owned work requires a live owner that is not abort-marked and not completing.
- IDs are never reused after a committed write. Entries and conversations are immutable.
- A watch captures its baseline and registers for later frames atomically.
Part 2 — What Effect v4 gives us
Section titled “Part 2 — What Effect v4 gives us”| Need | Effect building block | Fit |
|---|---|---|
| Model calls, streaming | effect/ai LanguageModel.streamText, Response.StreamPart (text, reasoning and tool-params deltas, tool-call, finish{reason, usage}), AiError.isRetryable / retryAfter |
✅ direct |
| Messages / transcript content | Prompt (Message, toolResultPart, fromResponseParts), Response.*Encoded schemas |
✅ use as the model payload of entries |
| Tool declarations | Tool.make({parameters, success, failure}), Toolkit, disableToolCallResolution: true so tool execution stays with us |
✅ |
| Providers | @effect/ai-anthropic, -openai, -openai-compat, -openrouter; LanguageModel.make for a faux test model |
✅ |
| Durable workflows | effect/workflow (Workflow, Activity, DurableDeferred, DurableClock). The only durable engine is ClusterWorkflowEngine + SingleRunner (SQL) |
⚠️ semantic mismatch, see 3.1 |
| SQL storage | effect/sql SqlClient.withTransaction, Migrator, @effect/sql-sqlite-node / -bun / libsql / sqlite-do |
✅ |
| File storage | FileSystem (open, writeAll, sync, rename) |
✅ for JSONL |
| Commit line | Semaphore.make(1) or a single consumer fiber over a Queue |
✅ |
| Publication / views | PubSub, SubscriptionRef.changes, Stream (groupedWithin, throttle, sliding buffers), Reactivity |
✅ |
| Cancellation | Fiber interruption, uninterruptibleMask, Scope, FiberMap keyed by task id |
✅ replaces Chord Context |
| Throttled persistence | Stream.groupedWithin / aggregateWithin, or a small Ref + Clock throttle |
✅ |
| Retry | Schedule + AiError.retryAfter; ExecutionPlan for provider fallback |
✅ but must be durable (see 3.3) |
| Typed state + migration | Schema (TaggedUnion, Class, toCodecJson, decodeTo for version upgrades), JsonPatch for deltas |
✅ |
| Bash / fs tools | effect/process ChildProcessSpawner, FileSystem, Path, NodeServices.layer |
✅ |
| Remote clients | effect/rpc with streaming RPCs; Atom for UI |
✅ later |
| Tests | @effect/vitest, TestClock |
✅ |
Part 3 — Design
Section titled “Part 3 — Design”3.0 Decisions
Section titled “3.0 Decisions”| # | Decision |
|---|---|
| D1 | Service-function API. No stateful handle objects. Functions live in modules (Conversation.submit(id, …)), take branded IDs, and require services in R. Same intent as pi, written Effect-natively. |
| D2 | Anthropic first, via @effect/ai-anthropic. A scripted faux LanguageModel is used for tests. |
| D7 | effect/ai is the model layer; we borrow, not wrap. LanguageModel, Model, Prompt, Response, Tool, Toolkit, AiError, Tokenizer, IdGenerator, ExecutionPlan and GenAI telemetry are used as-is. We add only what durability needs that effect/ai doesn’t have (see 3.4). |
| D3 | Published Effect v4 from npm (effect@^4.0.0, @effect/ai-anthropic, @effect/sql-sqlite-node, @effect/platform-node, @effect/vitest). repos/effect is reference source only. |
| D4 | Core first. Feature parity is tracked in the parity checklist and checked off as it lands. |
| D5 | Ports and adapters through Layers. Every piece whose technology could reasonably vary is a Context.Service port with one or more adapter layers. The harness core depends only on ports. |
| D6 | Own durable kernel, not effect/workflow (rationale below). |
3.1 Why not effect/workflow at the core
Section titled “3.1 Why not effect/workflow at the core”effect/workflow is replay-based. A workflow body re-runs from the top and memoised Activity results short-circuit
the completed steps. pi’s model is different, and its guarantees depend on that difference:
- Atomic checkpoint plus domain writes. pi commits a task’s new phase in the same transaction as the entries,
document deltas and child tasks it produced. Workflow stores Activity results in its own
MessageStorage, separate from our transcript. Bridging the two allows “entry appended but step not recorded”. - Ownership tree semantics have no Workflow equivalent: bottom-up abort,
completingholds,failFastwaits, background boundaries, and orphaning of blocked tasks. - Durability needs cluster. A durable Workflow needs
ClusterWorkflowEngine+SingleRunner+ SQL. That rules out the memory and file backends, and it brings in sharding. - Hot reload and versioning. pi resolves definitions per phase from a live registry and migrates checkpoints explicitly. Replay-based engines are fragile when code changes under them.
So we build Storage → Session → Scheduler as Effect services and keep pi’s explicit checkpoint/phase model, with
Schema-typed checkpoints. A WorkflowEngine adapter on our storage stays possible later.
3.2 Layer architecture: ports and adapters
Section titled “3.2 Layer architecture: ports and adapters” ┌───────────────────────── core (fixed) ─────────────────────────┐ Conversation / Submission / Task / State functions ──▶ Harness ──▶ Scheduler ──▶ Session (commit line) └───────┬──────────┬───────────┬───────────┬──────────┬───────────┘ ▼ ▼ ▼ ▼ ▼ ports (swappable): Storage TaskDispatcher ModelCatalog ExecEnv Settings StorageLease Registry ViewBuffer effect/ai services used as-is: LanguageModel (via Model layers) Tokenizer IdGenerator (+ Clock, Tracer)Harness.layer has every port in its requirements, so forgetting an adapter is a compile error, not a runtime
error. Convenience bundles (Harness.layerInProcess, Testing.layer) provide the common defaults.
| Port | Responsibility | Core adapters | Later / possible adapters |
|---|---|---|---|
Storage |
Atomic batch commit(writes) → Seq, ID minting, reads/scans (spec §storage). The only source of durability. |
MemoryStorage.layer (two-phase prepare/apply). SqlStorage.layer, which requires any SqlClient, so the driver is itself swapped by layer: @effect/sql-sqlite-node, -sqlite-bun, -libsql, -sqlite-do, -d1, -pglite/-pg. |
JsonlStorage.layer (on FileSystem), remote/KV storage |
StorageLease |
Exclusive ownership of a storage by one process (pi assumes it and does not enforce it). | StorageLease.unchecked (noop) |
file lock, SQL lease row with heartbeat |
TaskDispatcher |
Where and how reserved task invocations run: concurrency limits, fiber placement, per-kind pools. The scheduler decides what runs; the dispatcher decides how. | TaskDispatcher.inProcess({ concurrency }) (FiberMap<TaskId> in the Harness scope) |
per-kind limits/priorities, worker threads (effect/workers), rate-limited model pool (RateLimiter) |
ModelCatalog |
The one model-side piece effect/ai lacks: a stored ModelRef {provider, modelId} must resolve to an effect/ai Model (which is already a Layer<LanguageModel | ProviderName | ModelName>) after a restart. It also holds the metadata effect/ai doesn’t model: context window, pricing, a thinking-level mapping and an overflow predicate. |
ModelCatalog.make([AnthropicLanguageModel.model("claude-sonnet-5-5"), …]) plus the anthropicProfile metadata; FauxLanguageModel.model(script) (built with LanguageModel.make) |
any @effect/ai-* model(...): OpenAI, OpenRouter, … with no harness change |
ExecEnv |
Per-call execution environment built from {conversationId, cwd}: files and processes for tools and sections. |
NodeExecEnv.layer (on FileSystem + ChildProcessSpawner, so those swap too), MemoryExecEnv.layer for tests |
container/sandbox per conversation, remote exec |
Settings |
Run policy read at each use (retry, stream timeout, tool execution, queue modes, compaction). Never stored. | Settings.layer(static), Settings.layerConfig (from Config) |
file-watched live settings (SubscriptionRef) |
Registry |
Installed extensions; hot reload; a snapshot per task phase. | Registry.layer, plus Extension.layer(ext), which installs on build and uninstalls on scope close (rebuilding the layer is a reload) |
— |
ViewBuffer |
Backpressure policy for watchers (pi: at most 100 frames, then collapse into a snapshot). | ViewBuffer.collapsing(100) |
sliding, unbounded (tests) |
Token counting uses effect/ai’s Tokenizer service. We ship a heuristic layer, Tokenizer.make with chars/4 as
pi does, and any exact tokenizer layer can replace it.
Not swappable on purpose: Session (the commit line, publication and poisoning) and the Scheduler state
machine. They are the semantics; swapping them would change the guarantees. They are internal services built
inside Harness.layer.
About queues. pi has three queue-like things, and each lands in a different place:
- The inbox (steers and follow-ups) is a document. It must commit atomically with transcript writes, so it
stays inside
Storageand is not a separate queue port. - Task dispatch is the
TaskDispatcherport. - Ingress, meaning submissions arriving from outside, is a separate adapter. Because
Conversation.submitis idempotent byrequestId, any at-least-once source can feed it safely by using the message ID as therequestId:effect/persistencePersistedQueue, SQS, Kafka or HTTP. Examples:Ingress.fromQueue(queue)orIngress.fromStream(stream), later.
Composition (production)
const StorageLive = SqlStorage.layer.pipe( Layer.provide(SqliteClient.layer({ filename: "./session.sqlite" })), // swap driver here)// Plain effect/ai: the client layer, plus effect/ai Model values registered in the catalog.const ModelsLive = ModelCatalog.layer({ default: AnthropicLanguageModel.model("claude-sonnet-5-5"), models: [ AnthropicLanguageModel.model("claude-sonnet-5-5"), AnthropicLanguageModel.model("claude-haiku-4-5-20251001"), ], profiles: [anthropicProfile], // context windows, pricing, thinking, overflow predicate}).pipe( Layer.provide(AnthropicClient.layerConfig({ apiKey: Config.redacted("ANTHROPIC_API_KEY") })), Layer.provide(FetchHttpClient.layer),)const HarnessLive = Harness.layer.pipe( Layer.provide([ StorageLive, ModelsLive, NodeExecEnv.layer, Settings.layerConfig, Registry.layer, TaskDispatcher.inProcess({ concurrency: 16 }), HeuristicTokenizer.layer, // effect/ai Tokenizer ViewBuffer.collapsing(100), StorageLease.unchecked, ]), Layer.provideMerge(Extension.layer(CodingTools)), Layer.provide(NodeServices.layer),)Composition (tests and crash simulation)
const store = MemoryStorage.makeStore() // survives Harness scope teardownconst TestHarness = (script: FauxScript) => Testing.layer({ storage: MemoryStorage.layerFromStore(store), models: ModelCatalog.layer({ default: FauxLanguageModel.model(script) }) })
// "Crash": build a Harness, start a run, then close its scope mid-turn (all fibers interrupted, no outcomes// written), then build a fresh TestHarness on the same store and assert recovery.3.3 Mapping pi concepts to Effect
Section titled “3.3 Mapping pi concepts to Effect”| pi | Effect-native form |
|---|---|
Chord Context on every call |
Ambient fiber. Work runs on scheduler fibers in the Harness scope, so interrupting a waiter (Submission.await) only stops the wait. |
withoutAbortSignal around commits |
Effect.uninterruptible around storage commit, adopt and publish |
| Promise-chain commit line | Semaphore(1) held across the commit body, the storage commit and synchronous publication; a SessionPoisoned latch |
tx object passed to callbacks |
Tx service provided inside Conversation.commit(id, body) / Session.commit(body). Helpers are Tx.append, Tx.update(state, …) and Tx.createTask, so they compose without threading tx. |
defineDoc with Chord drafts |
State.conversation({kind, version, schema, history, fork, initial, migrate?, checkpointWhen?}) (also State.session, State.task). Updates are pure functions (a) => a; deltas are stored as JsonPatch ops. |
| Chord view state and ops | Stream<ViewFrame>, where ViewFrame = {seq, ops: JsonPatch[]} | {seq, snapshot} |
defineTask({phases, abort}) |
Task.make(name, {version, input: Schema, checkpoint: Schema.TaggedUnion, result?, initial, phases, abort, migrate?}). Phases are keyed by checkpoint _tag, and missing phases are a type error. |
runtime.commit(tx => state) |
TaskRuntime.commit(body), where body returns Next.continue(cp) | Next.wait(cp, on, policy) | Next.complete(r) | Next.fail(e) | Next.aborted(reason?) |
defineTool (TypeBox) |
A plain effect/ai Tool.make(name, {description, parameters, success, failure, dependencies: [ToolCall]}) with handlers from Toolkit.toLayer. Durable policy is expressed with effect/ai annotations: Tool.Idempotent ⇔ pi’s replay: "safe". Tool.Readonly and Tool.Destructive are available to hooks such as plan mode. The only new piece is the ToolCall service (taskId, details, commit, conversation, env), provided per call. |
defineExtension |
Extension.make(name, {toolkit, sections, hooks, wraps, tasks}), where toolkit is an effect/ai Toolkit. Handlers may need services; Extension.layer(ext) resolves them at install, so the registry only holds R = never values. |
hook(ToolTask, {beforeTool}) |
Hook.tool({ beforeTool: (call) => Effect<Option<Block | Rewrite>> }), typed per built-in task |
| Errors | Schema.TaggedError classes: ConversationBusy, StorageRejected, UnknownConversation, … Poisoning is a defect. |
| Partial throttle (100 ms trailing, one commit in flight) | A flusher fiber per request drains a Ref<Partial> every 100 ms, with a final flush before classification |
| Durable retry/backoff | Classify the AiError (isRetryable, retryAfter, plus the profile’s overflow predicate), then commit Retry{until}. Sleep with Clock until until, so a restart continues the remaining backoff. It is not Effect.retry. |
3.4 Model and tools: what we borrow from effect/ai and what we add
Section titled “3.4 Model and tools: what we borrow from effect/ai and what we add”Borrowed as-is
| effect/ai piece | Role in the harness |
|---|---|
LanguageModel service + LanguageModel.streamText |
The only way generation talks to a model. Called with disableToolCallResolution: true, so tool calls come back as parts and we run them as durable tasks. |
Model (AnthropicLanguageModel.model(id, config)) |
The unit of model selection. It is a Layer; generation does Effect.provide(model) per request. ProviderName / ModelName give us the provider/model key for pi.usage. |
AnthropicLanguageModel.withConfigOverride / Config |
Per-request provider config: thinking budget, max tokens, cache control. Thinking level maps to it through the profile. |
Prompt (Message, fromResponseParts, toolMessage) |
Entry payloads. A pi.user, pi.assistant or pi.tool-result entry stores encoded Prompt.Messages, so context derivation produces a Prompt directly. |
Response.StreamPart / Response.Usage / FinishReason |
Stream accumulation into the durable partial (pi.live), usage accounting, and stop classification (stop, length, tool-calls, error, pause). |
Tool / Toolkit (make, merge, toLayer, WithHandler.handle) |
Tool definition, provider declaration, parameter decoding and validation, and execution. The pi.tool task calls handle(name, params, callId); preliminary results are the streamed output that we commit with the adaptive throttle. |
Tool.Idempotent, Readonly, Destructive annotations |
Replay policy (Idempotent ⇒ safe to rerun after a crash) and hook inputs |
AiError + reasons (isRetryable, retryAfter) |
Typed error channel and retry classification. AiErrorReason is a Schema, so failed attempts are stored as-is. |
Tokenizer |
Token estimates for compaction thresholds |
IdGenerator |
Tool-call IDs in the faux model; any IDs we mint for parts |
LanguageModel.make |
FauxLanguageModel: the scripted test model is just another Model |
ExecutionPlan / Effect.withExecutionPlan |
Provider fallback, applied around a single request attempt |
GenAI telemetry (Telemetry, provider spans) |
Request tracing for free; we add harness spans (Effect.fn) for tasks and commits |
Added by us (durability only)
| Addition | Why effect/ai can’t cover it |
|---|---|
ModelCatalog (a stored ModelRef resolves to Model) |
Model choice is persisted per conversation and must resolve after a restart; effect/ai models are code values. |
ModelProfile (context window, pricing, thinking mapping, overflow predicate) |
effect/ai has no model metadata, no cost, and no “context too long” reason (Anthropic reports it as InvalidRequestError). |
The pi.generation task around streamText |
The checkpoint, the 100 ms partial commits, durable retry until, and turning tool-call parts into child tasks |
The pi.tool task around Toolkit.handle |
Intent commit before execution, replay vs interrupted result, durable output, truncation |
The ToolCall service (handler dependency) |
taskId, details, commit, conversation (subagents) and env: harness capabilities a handler can yield* |
We don’t use Chat: it persists the whole history as one blob after each call, while we need an append-only,
atomically committed transcript. A future Chat-compatible facade over a conversation is possible.
Generation request (sketch)
const request = Effect.fn("Generation.request")(function* (agent: ResolvedAgent, prompt: Prompt.Prompt) { const model = yield* ModelCatalog.resolve(agent.model) // effect/ai Model (a Layer) const profile = yield* ModelCatalog.profile(agent.model) return yield* LanguageModel.streamText({ prompt, toolkit: agent.toolkit, // Toolkit.merge of selected extensions disableToolCallResolution: true, }).pipe( Stream.runFoldEffect(Partial.empty, Partial.accumulate), // deltas → Ref, flushed to pi.live every 100 ms profile.withThinking(agent.thinkingLevel), // e.g. AnthropicLanguageModel.withConfigOverride Effect.provide(model), )})Tool execution inside pi.tool (sketch)
const results = yield* toolkit.handle(call.name, call.params, call.id) // effect/ai validates and decodesyield* results.pipe( Stream.runForEach((r) => r.preliminary ? DurableOutput.offer(r.encodedResult) : Final.set(r)), Effect.provideService(ToolCall, makeToolCall(task)), // per-call harness capabilities)Not yet covered: deferred/poll responses. The checkpoint slot is reserved.
3.5 Public API (service functions)
Section titled “3.5 Public API (service functions)”import { Conversation, Submission, Tx, State, Entry } from "@effective-harness/core"
const program = Effect.gen(function* () { const root = yield* Conversation.root({ agent: { model: ModelRef.make("anthropic", "claude-sonnet-5-5") } })
const sub = yield* Conversation.submit(root, Submission.input("What is the capital of France?", { requestId: "q1" })) const settled = yield* Submission.await(sub) // Settled.Done | Settled.Unanswered if (settled._tag === "Done") yield* Console.log(Entry.text(settled.answer))
yield* Conversation.events(root).pipe(Stream.runForEach(render), Effect.forkScoped) yield* Conversation.submit(root, Submission.input("Use pnpm", { whenBusy: "steer" })) // may fail ConversationBusy yield* Conversation.configure(root, { thinkingLevel: "high", instructions: "Only read files." }) yield* Conversation.commit(root, Effect.gen(function* () { yield* Tx.append(Entry.make("app.note", { text: "user opened a file" })) yield* Tx.update(Todos, (t) => ({ items: [...t.items, "write docs"] })) })) yield* Conversation.abort(root)}).pipe(Effect.provide(HarnessLive))Each module function has the shape (id, …) => Effect<A, DomainError, Harness>. IDs (ConversationId, EntryId,
TaskId, SubmissionId) are branded Schema types. Recovery runs when Harness.layer is built
(resume: true by default) or later through Harness.resume.
Defining your own tool and task:
// A plain effect/ai tool. `Idempotent` = safe to rerun after a crash; ToolCall gives durable capabilities.const Count = Tool.make("count", { description: "Count from 1 to n", parameters: Schema.Struct({ n: Schema.Int }), success: Schema.String, dependencies: [ToolCall],}).annotate(Tool.Idempotent, true)
const CountToolkit = Toolkit.make(Count)const CountExtension = Extension.make("count", { toolkit: CountToolkit, handlers: CountToolkit.toLayer({ count: Effect.fn(function* ({ n }, { preliminary }) { for (let i = 1; i <= n; i++) yield* preliminary(`${i}`) // streamed, durably committed output return `counted to ${n}` }), }),})
const Checkout = Task.make("app.checkout", { version: 1, input: Schema.Struct({ cards: Schema.Array(Card) }), checkpoint: Schema.TaggedUnion({ Pay: {}, Decide: { payments: Schema.Array(TaskId) } }), result: Schema.String, initial: () => ({ _tag: "Pay" }), phases: { Pay: (task) => TaskRuntime.commit(Effect.gen(function* () { const payments = yield* Effect.forEach(task.input.cards, (card) => Tx.createTask(Payment, { card }, { ownership: { kind: "task", taskId: task.id } })) return Next.wait({ _tag: "Decide", payments }, payments, "failFast") })), Decide: (task) => Effect.gen(function* () { const outcomes = yield* TaskRuntime.outcomes(task.checkpoint.payments) const paid = outcomes.every((outcome) => outcome.status === "completed") yield* TaskRuntime.commit(Effect.succeed(paid ? Next.complete("order placed") : Next.fail("payment failed"))) }), }, abort: () => TaskRuntime.commit(Effect.succeed(Next.aborted())),})Part 4 — Implementation plan
Section titled “Part 4 — Implementation plan”4.1 Repository layout (pnpm workspace)
Section titled “4.1 Repository layout (pnpm workspace)”| Package | Contents | Depends on |
|---|---|---|
packages/core (@effective-harness/core) |
IDs and schemas; port definitions; MemoryStorage; Session; State, Entry, Task, Scheduler; Harness and built-in tasks; Conversation / Submission APIs; TaskDispatcher.inProcess; ModelCatalog; HeuristicTokenizer; ToolCall; Settings; Registry |
effect (incl. effect/ai) |
packages/core/testing (subpath) |
FauxLanguageModel (via LanguageModel.make), MemoryExecEnv, Testing.layer, storage conformance suite |
@effect/vitest |
packages/storage-sql |
M0 smoke test of the SQLite driver; SqlStorage and NodeSqliteStorage landed in core (M4) |
@effect/sql-sqlite-node |
packages/anthropic |
anthropicProfile only: context windows, pricing, thinking → withConfigOverride, overflow predicate. The models themselves are plain AnthropicLanguageModel.model(...). |
@effect/ai-anthropic |
packages/node |
NodeExecEnv, coding tools (read, write, edit, bash) |
@effect/platform-node |
examples/ |
Ports of pi’s test/examples (chat, print, coding agent, inbox, crash recovery) |
all |
Later: packages/storage-jsonl, packages/rpc (remote views and control over effect/rpc).
4.2 Milestones
Section titled “4.2 Milestones”Each milestone ports the matching pi tests from repos/pi/packages/durable/test/ and ticks items in
the parity checklist. Testing uses a faux model and TestClock. Crashes are simulated by closing the
Harness scope and reopening the same store.
Core track:
| # | Milestone | Contents | Pi tests to port |
|---|---|---|---|
| M0 ✅ | Scaffold | pnpm workspace, TypeScript 7, vitest + @effect/vitest, oxlint + dprint, GitHub Actions CI. Smoke tests prove the published v4 packages work, including the D7 assumption: streamText with disableToolCallResolution passes tool calls through unexecuted and Toolkit.handle runs them later. |
— |
| M1 ✅ | Kernel | See 4.4. IDs, records, Storage port + MemoryStorage, conformance suite, Session commit line, Tx, conversations, entries, typed documents (versions, migration, scopes, base/delta), publication, poisoning, document watches |
storage conformance, session-tables, session-documents, session-checkpoints-migrations, session-watches (forks and Chord-state tests excluded) |
| M2 ✅ | Scheduler | See 4.6. Task.make, states and outcomes, reservation, phase stepping and no-progress fault, waits (allSettled / failFast), ownership, completing, bottom-up abort, background boundary, idle, recovery, versioning and blocking, TaskDispatcher.inProcess, TaskRegistry |
harness-tasks(-recovery), harness-lifecycle, harness-ownership, 24-child-tasks |
| M3 | Agent core | built-in entries and docs, context derivation, ModelCatalog + FauxLanguageModel + anthropicProfile, pi.generation on LanguageModel.streamText (prepare/request/retry, throttled partials), pi.tool on Toolkit.handle (intent, Idempotent replay, preliminary → output throttle, truncation), submissions/inbox/idempotency, conversation abort, usage and cost, basic configure |
harness-generation(-recovery), harness-tools(-recovery), harness-submissions, harness-inbox, harness-context, harness-output, harness-live-deltas |
| M4 | SQL storage | SqlStorage on SqlClient, migrations, conformance suite on sqlite-node, crash-recovery tests on disk |
sqlite-*, storage-runtime-boundary |
| M5 | Usable agent | Registry + Extension.layer, tools and sections from extensions, system prompt with pi.system diffs, ExecEnv + NodeExecEnv, coding tools, conversation view stream, AgentEvent stream, CLI example (print/chat) against Anthropic |
harness-registry, harness-prompt, harness-view, harness-events, tools, env-node |
Examples track. The goal is to run every pi example (repos/pi/packages/durable/test/examples) on memory storage.
Each lands in examples/src with the same number as an Effect program, and examples/test/examples.test.ts snapshots
its output. Examples that pi runs on SQLite or JSONL reopen the same MemoryStore instead.
| # | Milestone | Contents | Examples |
|---|---|---|---|
| M3a ✅ | History and forks | State.snapshotAsOf, Tx.forkConversation with asOf / current / initial copies, State.live, Conversation.page, the examples package |
00–05 |
| M3b ✅ | Harness core | Harness layer, entry kinds, extensions and the registry, the pi.agent state and configure, settings, context derivation, conversation forks with agents |
06–10, 12, 13 |
| M3c ✅ | Agent runs | ModelCatalog and a faux model, generation and tool tasks, submissions and the inbox, the system prompt, usage |
11, 14–16, 18, 20, 31 |
| M3d ✅ | Tools, views, events | ExecEnv and coding tools, hooks and wraps, conversation views, agent events, tool output |
17, 19, 21, 26–30 |
| M3e | Compaction and subagents | compaction (automatic, manual, overflow), foreground and background subagents, the task graph | 22, 23, 24, 25 |
Parity track (after core): hooks and wraps, compaction, reset/handoff, forks, rewindable docs and
asOf, JSONL storage, task graph and inspect, subagents (foreground and background), tool control
(terminate/handoff, addTools), deferred/poll, RPC remote views, StorageLease adapters, ingress adapters.
4.4 M1 spec (kernel)
Section titled “4.4 M1 spec (kernel)”In scope. Everything below the Harness that M2–M4 build on:
| Module | Contents |
|---|---|
Ids |
Branded IDs (ConversationId, EntryId, TaskId, SubmissionId, DocumentId, Seq), ROOT_CONVERSATION_ID = 1 |
Records |
Plain JSON record types: conversations, entries, tasks, submissions, document records and content, the StorageWrite union, queries, pages |
Errors |
Schema.TaggedErrors: StorageRejected, StorageFailure, SessionPoisoned, ReadAfterWrite, TransactionSettled, ConversationNotFound, TaskNotFound, OwnerNotLive, StateError |
Storage |
The port. Every method returns an Effect; reads return Option; commit(writes) fails with StorageRejected (nothing applied) or StorageFailure (uncertain) |
MemoryStorage |
Reference adapter. MemoryStore outlives a layerFromStore(store) handle, which is how tests simulate a crash and reopen |
testing/StorageConformance |
registerStorageConformance(name, makeLayer): 21 cases ported from pi’s conformance suite; M4’s SqlStorage must pass the same suite |
Session |
commit(body) with Tx provided; snapshot, watch, subscribeCommits, readOnLine; Session.layer(options) with a conversationCreated hook for M3 |
Tx |
Table reads (before the first table write), createConversation, appendEntry, and documents: doc, update, set, retire |
State |
Durable typed state (pi’s “documents”, renamed because the term only makes sense with Chord): State.session / State.conversation / State.task definitions; Todos.of(id) gives a typed Ref; State.snapshot, State.watch |
Conversation |
root, create, get, entries |
Done criteria (met): pnpm verify is green with 53 tests: the conformance suite against MemoryStorage, atomicity and rollback,
ReadAfterWrite, owner liveness, publication only after success, StorageRejected without poisoning, poisoning after an uncertain
failure, abandonment when interrupted before admission, completion when interrupted after admission, document creation, patches,
empty-change suppression, retirement and recreation, semantics and version checks, checkpointWhen, read-only migration then base on
first write, and watches: baseline, ordered frames, collapse, retirement, version replacement, and closing with the Session.
Out of scope, with where it lands:
- Task staging in
Tx(createTask, task replacement), task-document retirement at terminal settlement, owner validation of new tasks: M2 (landed; see 4.6). - Submission staging (
createSubmission,settle,place): M3. - Forks and
snapshotAsOf, document families, Chord-styledocumentState: parity track. Storage already supportsparent,key, rewindable history anddocument.copy, so these add Session code only.
Deliberate deviations from pi:
- No drafts. Documents are immutable Schema values updated with pure functions. The delta is
JsonPatch.get(committed, staged)at assembly. So a change that leaves the value equal writes and publishes nothing; pi publishes “structural no-op” batches. - Typed values. Stored JSON is the definition’s
Schema.toCodecJsonencoding; reads decode and fail withStateError. - Watches are Streams.
State.watch(ref)gives{ initial, changes: Stream<Frame> }scoped to the caller. Frames carry the decoded value and the JSON patch. The collapse limit is theSession.MaxPendingWatchFramesreference (default 100). - Cancellation is interruption. No
Contextparameter. A storage commit runs in an uninterruptible region. Closing the Session layer’s scope plays the role ofclose(). - Memory storage rejects invalid batches with
StorageRejected(validation happens before any state change), where pi throws a plain error that poisons the Session.
4.6 M2 spec (scheduler)
Section titled “4.6 M2 spec (scheduler)”In scope. The durable task engine that M3 runs model turns and tool calls on. pi’s harness/scheduler.ts (~1,300 lines) and spec
§5 are the oracle; every rule below is pi’s unless it is listed under deviations.
| Module | Contents |
|---|---|
Task |
Task.make(name, { version, input, checkpoint, result?, initial, phases, abort, migrate? }); Task.await(id) |
Next |
Next.continue(cp), Next.wait(cp, on, policy), Next.complete(result), Next.fail(error, result?), Next.aborted(reason?, result?) |
TaskRuntime |
The service a running phase or abort handler sees: commit, memo, outcomes, getTask, waitForTask, sleep, now, taskId, conversationId |
Tx (additions) |
Tx.createTask(task, input, { ownership, background? }), with owner checks; an internal task-replacement call used only by the scheduler. Task-scoped state is retired in the commit that makes its task terminal |
Scheduler |
Scheduler.layer(options?); resume, abortTask, abortConversation, waitForIdle, inspect. Recovery runs when the layer is built |
TaskDispatcher |
Port: how reserved invocations run. TaskDispatcher.inProcess({ concurrency }) keeps them in a FiberMap keyed by task ID |
TaskRegistry |
Port with a simple adapter until M5: TaskRegistry.layer(tasks), install / uninstall (scoped), change notification |
Definition.
checkpointis aSchema.TaggedUnion.phaseshas one handler per tag:{ [K in Tag]: (task: RunningTask<I, Case<K>>) => Effect<void, unknown, TaskRuntime> }, so a missing handler is a type error.abortis one handler for the whole task.input,checkpointandresultare stored as theirSchema.toCodecJsonencoding. A stored checkpoint that no longer decodes faults the task.RunningTaskis the committed record with decodedinputandcheckpoint:id,conversationId,kind,version,owner?,background,abortRequested,memos?.- A handler fails with any error or defect; the scheduler turns it into a
faultedoutcome. Interruption (an abort mark, or closing the layer) is not a fault.
Commit. TaskRuntime.commit(body) runs body as one Session.commit, with entries attributed to the task (byTaskId).
- Gate, evaluated on the line before
body: the invocation has not ended, the scheduler is not closing, the task isrunning, and a run invocation’s task has no abort mark. bodyreturns aNextor nothing. Nothing leaves the state unchanged. The new state is written in the same transaction as the entries, state and child tasksbodystaged.Next.continuereplaces the checkpoint.Next.waitparks the task.Next.complete,Next.failandNext.aborteddecide the outcome. A decided outcome is stored ascompletinginstead ofterminalwhile the task’s ordinary owned work is live, judged on the commit’s candidate records (rule 4).- A wait is validated: every
ontask exists and is not the waiter or on its owner chain;failFastonly names tasks the waiter owns; an abort handler cannot wait. - The checkpoint and result are validated against the definition’s codecs at runtime.
commitis not generic over the checkpoint, so a wrong shape fails the commit and faults the phase; a typedcommitis a possible later addition. - A terminal or
completingstate drops the memos.memo(name, candidate)is first-writer-wins: one gated commit returns the durable winner.
Invocations and steps. Reservation commits pending (or a waiting task whose on are all terminal) to running and registers an invocation.
One invocation is one fiber started through the TaskDispatcher: a run invocation executes phases in sequence, and an abort invocation runs the abort handler once.
Before every phase, including the first, a step runs as one commit on the line. It reads the committed task and applies the first matching rule:
- Terminal,
completingorwaiting: end the invocation. - Closing: end it; the checkpoint and any abort mark stay for reopen.
- Run invocation with a durable abort mark: end it; a fresh abort invocation starts once the task’s ordinary owned work is no longer live.
- The previous phase failed: write terminal
faulted. - Checkpoint changed: continue with the handler for the new tag.
- Checkpoint unchanged (compared by JSON value): write
faulted(“returned without durable progress”).
Before the first phase only rules 1–3 apply. After an abort handler only rules 1, 2 and 4 apply, and a handler that returns without a terminal outcome faults. If storage rejects a step’s fault write, the invocation still ends and the task stays running, so the next reservation runs it again.
An invocation ends inside the step that decides it, so a commit it issues later (a leaked fiber) fails with InvocationEnded.
Ownership. Tasks and conversations form one tree: a task’s parent is its owner task, else its conversation; a conversation’s parent is its owner task, if any.
Tx.createTasknames its owner:{ kind: "task", taskId }(a child, which lives in the owner’s conversation and cannot be background) or{ kind: "conversation", conversationId }. New owned work (tasks and, from M1, conversations) needs a live owner in the commit’s final candidate: notcompleting, terminal or abort-marked, elseOwnerNotLive.- The ordinary owned work of a task is every live task below it that is reachable without crossing a background task. Background tasks are boundaries.
- Rule 4, the hold. An outcome a task commits, or one the scheduler writes (
faulted,orphaned), is storedcompletingwhile its ordinary owned work is live. A reconcile commit writesterminalonce none is live, evaluated after every commit and at open. A held outcome is final: no handler runs again and no definition is needed.abortTaskon it only marks it. Task-scoped state is retired in the final commit, and so are task waiters. - Abort flows down. A live owner’s cancellation intent (an abort mark, or a held outcome other than
completed) marks its ordinary owned work in a separate reconcile commit. Terminal owners never cascade. The marks are re-derived at open, so a crash between an owner’s commit and its cascade loses nothing. - Abort runs bottom-up. An abort-marked task is not reserved while its ordinary owned work is live, so the abort handler sees final outcomes below it.
A waiting task with a mark leaves its wait early for its abort handler.
abortTaskcommits the mark, interrupts and joins an active run invocation, and returns"marked"(or"terminal"). - Waiting. A waiting task runs no code.
allSettledresumes when every task inonis terminal. WithfailFast, the first task inonthat holds or ends non-completedgets every other live task inonmarked, in the next reconcile commit; the waiter itself is not marked and resumes once all ofonis terminal. - Idle.
waitForIdle(conversationId?)waits until ordinary traversal from that conversation, or from every ownerless conversation, reaches no live non-background task.completingand blocked tasks count as live. abortConversation(id, { background? })marks the live tasks that traversal reaches (all of them withbackground: true) and waits until they are terminal and the conversation is idle. Input withdrawal joins it in M3.
Versioning and blocking. Resolution happens at reservation, and again for an abort invocation:
- No definition: stays
pending, blockedmissing_task. Stored version newer than the code: blockedtask_too_old. Older: runsmigrate(input, checkpoint, fromVersion), then commits the migrated record andrunningatomically; a failure (or a missingmigrate) blocks itmigration_failed, is reported once, and is retried only for a different definition object. - The scheduler never terminalizes a task because its code is missing. A blocked task stays live, still counts for idle, and is reconsidered whenever the registry changes. Blocked reasons are derived, never stored;
inspectshows them. - Only an abort orphans: an abort-marked task no definition can take is settled
orphaned(reason = the blocked reason), directly byabortTaskwhen nothing it owns is live, otherwise when its abort invocation would be reserved. Orphaning runs thesettleOutcomehook that M3 uses for run cleanup.
Recovery. Building Scheduler.layer scans live tasks in one commit, resets running to pending (checkpoint and mark kept), loads the failFast waiters, re-derives cascades and finalizes held outcomes. With resume: true (default) it then schedules; with false nothing is dispatched until Scheduler.resume.
Closing the layer’s scope is pi’s close(): it seals the scheduler, rejects waiters with SchedulerClosed, interrupts every invocation, and returns once they are gone and admitted commits settled. It writes no outcomes and no marks.
Effect mapping.
- Each invocation is a fiber in the layer scope (via
TaskDispatcher); one driver fiber serializes reconcile and reservation passes and is woken by commit publications (subscribeCommitsupdates a mirror of live tasks synchronously on the line). - Abort is pi’s durable mark plus fiber interruption. A committed
Nextstays committed: commits are uninterruptible once admitted. sleep(until)andnowreadClock(TestClockdrives them);untilis epoch milliseconds, so a restart continues the remaining time. Sleeps longer than the timer limit are split.- Handlers see only
TaskRuntime. Services a task needs are closed over when it is built (a Layer that builds theTask); M5’sExtension.layerformalizes this.
Changes to M1 that M2 needed.
Tx.createTask, the internalsetTask(replace one task record; only the scheduler uses it), candidate-aware owner validation (a new owned conversation or child task is judged against the owner as the commit leaves it), and task-scoped state retired in the commit that makes its task terminal. State staged by that same commit retires with it, and state access after a terminal candidate failsOwnerNotLive.Session.commit(body, { byTask })attributes appended entries to the committing task (byTaskId).- The commit line is now strictly first come, first served (
internal/fifoLock.ts). Effect’sSemaphoreis not fair: a fiber arriving as a permit is released can overtake queued ones. pi’s Promise-chain line never does, and several rules lean on it (an abort mark queued behind a reservation must land before that task’s first step). Interrupting a waiter removes it from the queue; a lock granted to a waiter that is then interrupted passes on. - New errors:
TaskRejected,InvocationEnded,TaskNotTerminal,SchedulerClosed.
Deviations from pi.
- No handover. pi hands an invocation over to a replacement definition at a phase boundary. Hot reload belongs to M5’s registry, so M2 resolves definitions only at reservation.
- No agent, hooks, settings, models or env on the runtime, and no document watches through it: those arrive with M3 and M5.
- Interruption instead of
AbortSignal. There are no handler contexts or runtime-bound signals. A handler that must finish its work after an abort mark or a close wraps it inEffect.uninterruptible; interruption lands at the next interruptible point, so code after an uninterruptible region does not run. A leaked fiber’s runtime calls fail withInvocationEnded. Next.abortedis added to the plan’s four constructors: an abort handler must be able to endaborted.- One fiber per invocation, not per phase. An invocation is a fiber started through the
TaskDispatcherand runs its phases in sequence, with the step between them. A mark interrupts the fiber instead of signalling a controller. - Scheduling starts only at
resume. pi also starts it from progress calls on a paused Harness (submit, wait, abort). HereScheduler.layerresumes by default andresume: falsestays paused untilScheduler.resume; the Harness layer in M3 decides how its own calls behave. commitis not generic over the checkpoint. Checkpoints and results are validated against the definition’s codecs when committed, so a wrong shape fails the commit and faults the phase. Thephasesmap and each handler’stask.checkpointare fully typed, andinitialmust produce a checkpoint of the schema (checked withconst CandNoInfer).- Wait validation fails the commit with
TaskRejected/TaskNotFound(pi throws); the handler decides what to do with it.
Behaviours kept from pi that are easy to trip over.
- A reservation or reconcile commit that storage rejects is reported and retried on the next wakeup, not on a timer: any commit that changes a task, a registry change, or
Scheduler.resume(which is safe to call repeatedly) wakes the driver. - If the step cannot write a fault (storage keeps rejecting it), the invocation ends, the task stays
running, and the next reservation runs the phase again with no backoff. - With a concurrency limit, an invocation that waits for another task inside its handler (
TaskRuntime.waitForTask) keeps its slot;Next.waitfrees it. TaskDispatcher.startreturns a handle that interrupts exactly that invocation, including one still waiting for a slot; the scheduler ends the invocation itself after interrupting it, so a fiber that never began cannot leak.
Tests ported from pi. Rewritten with @effect/vitest, TestClock and MemoryStorage.makeStore(). Crashes close the layer scope (or, for a commit that never lands, make storage fail it first) and reopen the same store. Handlers that “ignore their signal” in pi use
Effect.uninterruptible. Each pi test below that is not listed as ported has a reason.
| pi file | Ported | Deferred, or not applicable (to) |
|---|---|---|
harness-tasks (scheduler-tasks.test.ts) |
phases; faults (no progress, fail, die, throw); terminal kept after a later failure; checkpoint compared by value; atomic results, memos, task state; waits; ended runtime; clock; reservation retry (×2); fault-write retry; idle; unknown and terminal tasks; waits and close; abort (×8); close (×5) | registry snapshot per phase and agent() (M3, M5); document watches through the runtime (M3); “orders runtime commits against the step” and “close seals during a step that decided to continue” (no queued detached commits; no snapshot hook), covered by the leaked-fiber and close tests |
harness-tasks-recovery (scheduler-recovery.test.ts) |
intent/effect/outcome resume; every abort stage across reopen; six crash cases; blocked tasks (missing, too old, migration failing and retried, no migration, orphaned ×3) | definition handover (×7, M5); SQLite variants (M4) |
harness-lifecycle (scheduler-recovery.test.ts) |
close joins an uninterruptible handler; reopen on the same storage runs no old code; paused scheduling; registry change before resume; open failure | Harness open/view/inspect-after-close tests (M3); “progress calls schedule” (deviation 6); commit-cancelled-in-storage (M1 covers it) |
harness-ownership (scheduler-ownership.test.ts) |
holds; idle; every cascade case; background boundaries; conversation abort; reopen derivations (×3); orphan cascade; cancelled waits; waiting child; rejected-cascade retries (×2) | queued-input withdrawal, owned-conversation handles, subagent tools (M3, P) |
examples/24-child-tasks (checkout.test.ts) |
four payments with a declined card (failFast), a cancelled checkout, a restart mid-flight |
the task-graph print (parity); Scheduler.inspect prints the same tree |
Also covered, beyond pi’s tests: bottom-up abort order, failFast and allSettled, invalid waits, finishing commits and holds, the settleOutcome hook, the dispatcher concurrency limit, the FIFO commit line, type-level checks (task-types.test.ts), and the Tx additions (task-tx.test.ts).
Out of scope, with where it lands: submissions and input withdrawal (M3); the Harness layer and built-in tasks (M3); definition handover and hot reload (M5); task graph and watch (parity).
Done criteria (met): pnpm verify is green with 157 tests, 104 of them new in M2. Every ported test above passes, each rule above has a test, and a mutation pass over the scheduler’s key rules (no-progress fault, holds, failFast marking, bottom-up abort wait, cascades, finalization, abort-mark gating, step rule 3, background boundary, orphaning, migration memo, abort-handler fault, interruption on mark, waiters, reset at open) is killed by the suite, except one equivalent mutant: terminal owners cannot cascade because they are never in the live set.
4.7 M3 spec (agent harness)
Section titled “4.7 M3 spec (agent harness)”M3 is split into the examples track’s milestones (4.2). The decisions that shape all of them:
- Service functions, no handles. pi’s
HarnessandConversationobjects become modules of functions over services:Conversation.root({ agent, init }),create,fork,agent,configure,context,page,abort,waitForIdle;Agent.get,Agent.configure(inside any commit);Registry.install/uninstall;Task.await/Task.get. Harness.layer(options)builds the Session with the built-in creation hook, the scheduler,Settingsand an in-process dispatcher over theStorageandRegistryin context. The registry is application-owned and also serves the scheduler’sTaskRegistry, with built-in tasks first; installing an extension unblocks the tasks it brings.- Transcript messages are the harness’s own
Messageschema, shaped like pi-ai’s: assistant messages carry their stop reason and usage, and positional system messages patch named prompt sections and the offered tools. effect/ai’sPromptcannot express either, and context derivation depends on both, so generation converts a derived context to aPromptat request time. This replaces 3.4’s plan to store encodedPrompt.Messages. - Tools are the harness’s own
Tool.make(name, { description, parameters, execute, replay?, executionMode? }), withexecute: (args) => Effect<Tool.Result, unknown, ToolCall>. Wrappers (Extension.wrapTool) must replaceexecuteand per-conversation selection works by name, which effect/ai’s handler layers do not allow.Tool.declarationderives the JSON Schema with effect/ai’sTool.getJsonSchemaFromSchema, and generation offers effect/aiTool.dynamictools. - Sections render a string,
undefined, or anEffectthat may read committed state (R = Session). ExecEnvis a port (ExecEnvBuilder) whose builder may read committed state, so a conversation’s sandbox can live in its own state.Task.make’sabortis optional and defaults to committingNext.aborted().
M3b (done) covers 06–10, 12 and 13: the registry with validation and in-place replacement, agent resolution (selection,
same-name replacement, wrappers, filters, instructions), configure, the creation hook (an ownerless conversation starts
{}, a task-owned one copies its owner’s agent, a fork keeps its asOf copy, then conversationCreated, agent and
init), and context derivation. 21 tests are ported from pi’s harness-registry, harness-conversations and
harness-context.
M3c (done) covers 11, 14, 15, 20 and 31 (16 and 18 need a real provider and the coding tools, so they move to M3d):
pi.generation (prepare with positional pi.system diffs, request over LanguageModel.streamText with throttled partials
in pi.live, durable retry, classification, hooks, abort) and pi.tool (intent before execute, replay policy, adaptive
progress, output bounds, hooks, abort), submissions and the inbox (Conversation.submit, Submission.wait / status /
abort, boundaries, queue modes, withdrawal on conversation abort through the scheduler’s withdrawInputs hook), the
pi.usage ledger, ModelCatalog and FauxModel. 144 tests are ported from pi’s harness-generation(-recovery),
harness-prompt, harness-tools(-recovery), harness-output, harness-inbox and harness-submissions. Deviations:
- Calls to tools the request did not offer. effect/ai decodes response parts against the request’s toolkit, so such a
call is invalid output for the whole attempt (retryable, like any invalid output) instead of a
tool_unavailableresult. A tool deactivated after the request was sent still getstool_unavailable. - No deferred responses (pi’s poll/cancel), no
prepareArguments,addTools,control.handoff, per-tool output limits or the Session-wide usage total yet (parity track). - Scheduling starts with the layer (
resume: falsedefers it toScheduler.resume), not lazily on the first submit.
M3d (done) covers 16–19, 21 and 26–30. ExecEnv mirrors pi’s FileSystem + Shell with typed FileError / ExecError;
NodeExecEnv implements it on node built-ins (bash in its own process group, killed on interruption or timeout; spill
files), and NodeExecEnv.layer follows each conversation’s cwd. CodingTools brings read, write, edit (fuzzy
matching, own line diff instead of the diff package) and bash (tail-retained output, per-tool outputLimits,
ToolCall.diagnostic), with per-file serialization by env ID and canonical path. View.watch / Conversation.watch stream
a conversation’s view as JSON Patch frames, and AgentEvent.watch derives pi’s agent events from them; both collapse to a
snapshot past MaxPendingWatchFrames. 117 tests are ported from pi’s env-node(-spill), env-truncate, tools,
harness-view, harness-events and harness-live-deltas. Deviations: JSON Patch has no string append, so text and output
growth arrive as replace and events recover deltas by prefix; interruption replaces abort signals (no aborted exec code,
no text line readers, truncateFile or flushFile); edit’s argument repair is exported but not applied, since tools have
no prepareArguments hook yet; examples 16 and 18 use @effect/ai-openai when OPENAI_API_KEY is set.
M3e, subagents and the task graph (22–24). ConversationHandle is pi’s invocation-bound handle (id, submit →
SubmissionHandle with status / wait / abort, abort, waitForIdle), from TaskRuntime.conversation(id) and
ToolCall.conversation(id) (Option.none() for a missing conversation). Every operation fails with InvocationEnded once
the invocation ended (a tool’s handle also once its call returned), and one in flight is interrupted with it. submit
admits through the task’s commit gate without attributing the user entry to the task, so a submit queued on the line
behind an abort mark is rejected; waits run off the line. The scheduler gets the admission from the Harness as
Scheduler.Options.submitter. TaskGraph.get / TaskGraph.watch give pi’s task graph (live tasks with kind, owner,
owned conversations and pending / running / waiting with on / completing with the held outcome) built on the line
from committed records and advanced by publications as JSON Patch frames. 13 tests: 4 ported from pi’s
harness-task-graph, the 7 handle and subagent tests of harness-ownership (harness-handles.test.ts), and 2 new ones
(a task phase’s handle; a watch ends with the Session). Deviations:
no shared graph mount (each watch builds its own value, so pi’s mount-sharing and cancelled-acquisition tests do not
apply); phase is the checkpoint’s _tag; handle submissions take input drafts only, as pi’s. The examples replace
timing with deferreds: 22’s scripted model answers the child once the UI attached to it and reports once the UI printed
the child’s events; 23’s chapter walk-through and whale answer wait (until stopped, until the restart) instead of
streaming at 50 tokens per second, the host waits for the main agent’s answer before the restart, and the UI is uncoloured;
24 waits until four cards are charged instead of 20 ms and reopens the memory store instead of SQLite.
M3e, compaction (25). pi.compaction (internal/compaction.ts) selects the cut (keepRecentTokens, candidates as pi),
asks beforeCompact (Hooks.compaction: decline, or supply the summary), and summarizes the serialized prefix with pi’s
prompts through the conversation’s model, pinning model, thinking level, stream timeout and maxTokens; retries follow the
retry policy and show in pi.live.compactions. A blocking compaction appends its Entry.Compaction summary (head = first
kept entry, data: { reason }); a conversation-owned one admits it as a write submission (compaction:<taskId>), placed
at once, at the next boundary, or settled stale. Generation’s Prepare starts a blocking compaction above
contextWindow - reserveTokens and waits for it, or a background one past backgroundTokens when none is listed; an
overflow error with a cut appends the error, waits for an overflow compaction and prepares again with the same attempt,
failing with the overflow text when that compaction placed nothing. Conversation.compact(id, instructions?) returns the
task ID; its result (Live.CompactionResult) carries the summary’s submissionId. Fault and orphan cleanup removes the
status. 98 tests are ported from pi’s harness-compaction (recovery reopens the memory store). Deviations: estimates use
pi’s chars/4 heuristic behind internal/tokens.ts, not yet an effect/ai Tokenizer; effect/ai has no provider-neutral
answer cap, so maxTokens applies through the new optional CatalogModel.withMaxTokens, and there is no
cacheRetention or deferred option to clear; a tool call in a summary is invalid output for the empty toolkit (retried,
then Summarization failed: …) instead of attempted to call a tool; a failed attempt reports no usage, since effect/ai
errors carry none; the overflow detail is the AiError message, which contains the provider text. Not ported: the storage
commit-rejection test (no injection hook in the test support) and “enables scheduling right after open” (scheduling starts
with the layer). Example 25 holds the busy answer with deferreds and compacts once that answer’s request arrived, so the
manual compaction never races the preparation’s background check.
Storage adapters, JSONL (parity track). JsonlStorage.layer({ directory, fsync?, fileSystem? }) ports pi’s JSONL
storage: main.jsonl holds one commit marker per commit with the small writes inline, live task records and document
content go to task-<id>.jsonl / doc-<id>.jsonl sidecars appended (and with fsync, flushed) before the marker,
and sidecars of terminal tasks and retired or rebased current-only documents are reclaimed after it (temp file plus
rename). Opening replays the log into a MemoryStore index through the new internal MemoryStorage.openIndex
(prepare at an explicit sequence, then apply, like pi’s prepareCommit), drops torn last lines, truncates
unconfirmed sidecar tails and finishes interrupted reclamation. 83 tests: the conformance suite twice (plain, and
reopening after every commit) and the 41 cases of pi’s jsonl-storage, plus a closed-handle/copy test and a Harness
that resumes a turn interrupted mid-stream from the same directory. Deviations: the file system is a small Effect port
(JsonlStorage.FileSystem, default nodeFileSystem over node:fs/promises) instead of effect’s FileSystem, which
would need @effect/platform-node; failures are StorageFailure (open, corruption, poisoning) and StorageRejected
(invalid or unserializable batch), not pi’s JsonlCorruptionError / JsonlStoragePoisonedError classes; commits take
a one-permit lock and run uninterruptibly; nothing locks the directory, so one process must own it, as in pi
(StorageLease stays open).
Storage adapters, SQL (M4). SqlStorage.layer(options?) (Layer<Storage, StorageFailure, SqlClient>, in core:
it needs only effect/sql) ports pi’s portable SQLite storage: one JSON record column per row plus the indexed
columns pi has, one SQL transaction per commit (metadata row read, validation against stored state, writes, metadata
advance), document revisions as base/delta rows keyed by commit sequence (latest-only documents drop older revisions on
a new base or retirement), document.copy materialized inside the commit’s transaction, and document(id, at) read in
a transaction so a concurrent base replacement cannot split it. mintId hands out candidates from memory above the
persisted floor, and a commit raises the floor only when it succeeds. Migrations are pi’s (SqlStorage.migrations,
currentSchemaVersion, migrate(steps?)): a durable_schema version row, all pending steps in one transaction,
refusal of a newer database. NodeSqliteStorage.layer({ filename, busyTimeout?, walAutoCheckpointPages?, migrations? }) puts it on @effect/sql-sqlite-node (built on node:sqlite, so no native module and no new package)
with pi’s pragmas (WAL, synchronous = NORMAL, auto-checkpoint, 5 s busy timeout, directory creation, a truncating WAL
checkpoint on close); NodeSqliteStorage.layerSqlClient(options) is the bare client. 80 tests: the conformance suite
three times (:memory:, a file, and a file reopened after every commit), pi’s sqlite-storage and
sqlite-migrations cases, the applicable sqlite-facade cases (reads queue behind an open transaction, IDs adopted
only on success, consistent document reads), and a Harness that resumes a turn interrupted mid-request after the file
is closed and reopened. Deviations: effect’s Migrator is not used (it creates its table outside the transaction,
cannot refuse a newer schema and turns failures into defects); deltas are stored as JSON Patch, so a pi database is
not readable; the schema is SQLite-only (STRICT tables, json_valid) and the layer refuses other dialects, while the
queries themselves are portable (ON CONFLICT … DO NOTHING/UPDATE); validation failures are StorageRejected, while SQL
errors and failed COMMIT/ROLLBACK (defects of withTransaction) are StorageFailure, even when the rollback
succeeded; pi’s own SqliteDatabase facade (serial queue, transaction handles, prepare cache, read draining on close)
is replaced by SqliteClient, so its facade tests are not ported, and reads still running when the layer closes are
not drained (the Harness scope closes first). packages/storage-sql (the M0 driver smoke test) is now redundant.
4.3 Risks and gaps
Section titled “4.3 Risks and gaps”- No Chord wire compatibility. Views emit JSON Patch ops, so we have API parity but not wire compatibility with pi clients. This is intentional.
- Overflow classification and pricing are not in
effect/ai.anthropicProfileowns them. Overflow detection matches on theInvalidRequestErrordescription, which is brittle; it is covered by a recorded fixture. - Dynamic toolkits. Per-conversation tool selection means merging toolkits at request time, so the merged
Toolkittypes are lost (Tool.Any). That is acceptable inside the harness; user-facing tools stay fully typed. - Effect v4 is newly stable (4.0.0). Unstable modules (
effect/ai,effect/sql) may still change in minor versions, so pin exact versions. - The spec is large (about 4.5k lines). It is the acceptance oracle; port its tests, not just the README examples.
- In-process commit line means one process per storage, the same as pi.
StorageLeasemakes this explicit; multi-process would need a differentSessionand is out of scope.
4.5 Open decisions
Section titled “4.5 Open decisions”- Local model end-to-end tests (discussed after M0, deferred): run a small tool-calling model such as Qwen3 through Ollama with
@effect/ai-openai-compat, as an opt-in e2e tier next to the scripted faux model. This would be decided before M3.