Skip to content

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}.


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.
  • Conversation: {id, parent?: {conversationId, at}, owner?: {conversationId, taskId}}. Root is id 1. 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 are pi.user, pi.assistant, pi.system, pi.tool-result, pi.reset, pi.compaction.
    • head makes 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}).
    • scope is session, task or conversation.
    • history is latest or rewindable (snapshotAsOf).
    • fork is initial, current or asOf.
    • 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.
  • Submission: user input or a write handed to a conversation, idempotent by requestId. Its lifecycle is queued → placed → done | unanswered{reason}.
  • A definition is {name, version, initial(input), phases: {[phase]: handler}, abort, migrate?}. The checkpoint carries phase.
  • 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 abort handler runs.
    • A task that finishes while it still owns live work becomes completing and turns terminal later.
    • background: true tasks are a boundary that ordinary aborts and idle checks do not cross.
  • Waiting: waiting{on: TaskId[], policy: allSettled | failFast}. No code runs while a task waits. failFast aborts the remaining siblings.
  • Recovery: on open, running is reset to pending. A reconcile re-derives the abort cascades. Nothing is dispatched until resume() or a progress call.
  • Versioning: a stored version newer than the code blocks the task. An older stored version runs migrate when the task is reserved. A task with a missing definition stays blocked, never failed.
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.live on 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.
  • Inbox placement:
    • Writes are always placed. steer items join after the current tool round, one at a time or all with steeringMode: "all".
    • followUp items start the next run.
    • whenBusy: "reject" throws ConversationBusy.
  • Context derivation:
    1. Start from the newest head marker.
    2. Apply edits, newest wins.
    3. Drop aborted, error and deferred assistant messages.
    4. 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.system entries, 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.
  • 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, afterTools on generation.
      • beforeTool, afterTool on tools.
      • beforeCompact on compaction.
  • Tool: {name, description, parameters (TypeBox), execute(args, api, ctx), replay?, executionMode?}.
    • The api offers output() streaming, details(), commit(), conversation(id) for subagents, env and taskId.
    • A result may carry usage and control: {terminate | handoff}.
  • Observation:
    • viewState() and watch() 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 an ExecutionEnv (fs and exec) per call. The built-in tools read, write, edit and bash use 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.

  1. One serialized commit line. Publication happens only after storage succeeds. StorageRejected rolls back; any other storage failure poisons the session.
  2. Once storage has admitted a commit, cancellation cannot stop it.
  3. Every task transition is decided in one step on the line. A phase without durable progress faults.
  4. Tool intent is durable before execute. beforeTool never reruns. A replay requires safe in both the stored and the current declaration.
  5. A conversation is busy exactly while pi.live.run is set. Every terminal path of a run clears run, generation and tools in the same commit that settles its inputs.
  6. Usage is updated in the same commit as the assistant entry or tool result that produced it.
  7. Abort is bottom-up. Cascade marks are written in a separate reconcile commit and re-derived on open.
  8. New owned work requires a live owner that is not abort-marked and not completing.
  9. IDs are never reused after a committed write. Entries and conversations are immutable.
  10. A watch captures its baseline and registers for later frames atomically.

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 ✅

# 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).

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, completing holds, failFast waits, 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:

  1. The inbox (steers and follow-ups) is a document. It must commit atomically with transcript writes, so it stays inside Storage and is not a separate queue port.
  2. Task dispatch is the TaskDispatcher port.
  3. Ingress, meaning submissions arriving from outside, is a separate adapter. Because Conversation.submit is idempotent by requestId, any at-least-once source can feed it safely by using the message ID as the requestId: effect/persistence PersistedQueue, SQS, Kafka or HTTP. Examples: Ingress.fromQueue(queue) or Ingress.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 teardown
const 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.
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 decodes
yield* 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.

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())),
})

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).

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.

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-style documentState: parity track. Storage already supports parent, key, rewindable history and document.copy, so these add Session code only.

Deliberate deviations from pi:

  1. 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.
  2. Typed values. Stored JSON is the definition’s Schema.toCodecJson encoding; reads decode and fail with StateError.
  3. 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 the Session.MaxPendingWatchFrames reference (default 100).
  4. Cancellation is interruption. No Context parameter. A storage commit runs in an uninterruptible region. Closing the Session layer’s scope plays the role of close().
  5. Memory storage rejects invalid batches with StorageRejected (validation happens before any state change), where pi throws a plain error that poisons the Session.

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.

  • checkpoint is a Schema.TaggedUnion. phases has one handler per tag: { [K in Tag]: (task: RunningTask<I, Case<K>>) => Effect<void, unknown, TaskRuntime> }, so a missing handler is a type error. abort is one handler for the whole task.
  • input, checkpoint and result are stored as their Schema.toCodecJson encoding. A stored checkpoint that no longer decodes faults the task.
  • RunningTask is the committed record with decoded input and checkpoint: id, conversationId, kind, version, owner?, background, abortRequested, memos?.
  • A handler fails with any error or defect; the scheduler turns it into a faulted outcome. 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 is running, and a run invocation’s task has no abort mark.
  • body returns a Next or nothing. Nothing leaves the state unchanged. The new state is written in the same transaction as the entries, state and child tasks body staged.
  • Next.continue replaces the checkpoint. Next.wait parks the task. Next.complete, Next.fail and Next.aborted decide the outcome. A decided outcome is stored as completing instead of terminal while the task’s ordinary owned work is live, judged on the commit’s candidate records (rule 4).
  • A wait is validated: every on task exists and is not the waiter or on its owner chain; failFast only names tasks the waiter owns; an abort handler cannot wait.
  • The checkpoint and result are validated against the definition’s codecs at runtime. commit is not generic over the checkpoint, so a wrong shape fails the commit and faults the phase; a typed commit is a possible later addition.
  • A terminal or completing state 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:

  1. Terminal, completing or waiting: end the invocation.
  2. Closing: end it; the checkpoint and any abort mark stay for reopen.
  3. 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.
  4. The previous phase failed: write terminal faulted.
  5. Checkpoint changed: continue with the handler for the new tag.
  6. 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.createTask names 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: not completing, terminal or abort-marked, else OwnerNotLive.
  • 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 stored completing while its ordinary owned work is live. A reconcile commit writes terminal once none is live, evaluated after every commit and at open. A held outcome is final: no handler runs again and no definition is needed. abortTask on 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. abortTask commits the mark, interrupts and joins an active run invocation, and returns "marked" (or "terminal").
  • Waiting. A waiting task runs no code. allSettled resumes when every task in on is terminal. With failFast, the first task in on that holds or ends non-completed gets every other live task in on marked, in the next reconcile commit; the waiter itself is not marked and resumes once all of on is terminal.
  • Idle. waitForIdle(conversationId?) waits until ordinary traversal from that conversation, or from every ownerless conversation, reaches no live non-background task. completing and blocked tasks count as live.
  • abortConversation(id, { background? }) marks the live tasks that traversal reaches (all of them with background: 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, blocked missing_task. Stored version newer than the code: blocked task_too_old. Older: runs migrate(input, checkpoint, fromVersion), then commits the migrated record and running atomically; a failure (or a missing migrate) blocks it migration_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; inspect shows them.
  • Only an abort orphans: an abort-marked task no definition can take is settled orphaned (reason = the blocked reason), directly by abortTask when nothing it owns is live, otherwise when its abort invocation would be reserved. Orphaning runs the settleOutcome hook 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 (subscribeCommits updates a mirror of live tasks synchronously on the line).
  • Abort is pi’s durable mark plus fiber interruption. A committed Next stays committed: commits are uninterruptible once admitted.
  • sleep(until) and now read Clock (TestClock drives them); until is 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 the Task); M5’s Extension.layer formalizes this.

Changes to M1 that M2 needed.

  • Tx.createTask, the internal setTask (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 fails OwnerNotLive.
  • 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’s Semaphore is 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.

  1. 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.
  2. No agent, hooks, settings, models or env on the runtime, and no document watches through it: those arrive with M3 and M5.
  3. 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 in Effect.uninterruptible; interruption lands at the next interruptible point, so code after an uninterruptible region does not run. A leaked fiber’s runtime calls fail with InvocationEnded.
  4. Next.aborted is added to the plan’s four constructors: an abort handler must be able to end aborted.
  5. One fiber per invocation, not per phase. An invocation is a fiber started through the TaskDispatcher and runs its phases in sequence, with the step between them. A mark interrupts the fiber instead of signalling a controller.
  6. Scheduling starts only at resume. pi also starts it from progress calls on a paused Harness (submit, wait, abort). Here Scheduler.layer resumes by default and resume: false stays paused until Scheduler.resume; the Harness layer in M3 decides how its own calls behave.
  7. commit is 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. The phases map and each handler’s task.checkpoint are fully typed, and initial must produce a checkpoint of the schema (checked with const C and NoInfer).
  8. 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.wait frees it.
  • TaskDispatcher.start returns 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.

M3 is split into the examples track’s milestones (4.2). The decisions that shape all of them:

  • Service functions, no handles. pi’s Harness and Conversation objects 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, Settings and an in-process dispatcher over the Storage and Registry in context. The registry is application-owned and also serves the scheduler’s TaskRegistry, with built-in tasks first; installing an extension unblocks the tasks it brings.
  • Transcript messages are the harness’s own Message schema, 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’s Prompt cannot express either, and context derivation depends on both, so generation converts a derived context to a Prompt at request time. This replaces 3.4’s plan to store encoded Prompt.Messages.
  • Tools are the harness’s own Tool.make(name, { description, parameters, execute, replay?, executionMode? }), with execute: (args) => Effect<Tool.Result, unknown, ToolCall>. Wrappers (Extension.wrapTool) must replace execute and per-conversation selection works by name, which effect/ai’s handler layers do not allow. Tool.declaration derives the JSON Schema with effect/ai’s Tool.getJsonSchemaFromSchema, and generation offers effect/ai Tool.dynamic tools.
  • Sections render a string, undefined, or an Effect that may read committed state (R = Session).
  • ExecEnv is a port (ExecEnvBuilder) whose builder may read committed state, so a conversation’s sandbox can live in its own state.
  • Task.make’s abort is optional and defaults to committing Next.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_unavailable result. A tool deactivated after the request was sent still gets tool_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: false defers it to Scheduler.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.

  • 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. anthropicProfile owns them. Overflow detection matches on the InvalidRequestError description, 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 Toolkit types 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. StorageLease makes this explicit; multi-process would need a different Session and is out of scope.
  • 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.