Skip to content

Test & operate

Operations

A durable host must expose accepted work, explain blocked work, and recover it safely. The same administration contract applies to Node and SQLite class DN and Cloudflare Durable Objects class DC.

DurableAgentRuntime provides five operator functions:

  • explain(submissionId) and explainThread(threadId) return the recovery decision, operator meaning, and expected disposition. They write nothing.
  • verify(threadId) checks the canonical batch digest chain and recomputes Run continuations from their referenced facts. Missing or mismatched progress fails integrity verification.
  • retry(RetryCommand.make({ submissionId, author, reason })) logs the operator and repeats the classifier’s decision. Repairs with a claim annotate their attempt using that claim’s epoch; state-only wakes and marker repairs do not append canonical audit records. Retry refuses settled work and requests awaiting resolveUnknown or resolveApproval.
  • wake(threadId) sends a droppable liveness hint without taking ownership or advancing an epoch.
  • scanObligations(thresholds) reports blocked or aging accepted work.

NodeDurableHost exposes all five. Use vp run admin:durable <explain|verify|retry|wake|obligations> --database <file> on Node. Stop an automatically managed Node host before opening this separate CLI connection; use the live host’s methods for administration while its exclusive database lock is held. Cloudflare Thread Objects expose encoded administration methods through the application’s Worker.

The explain --thread <id> --json payload is an array of recovery explanations, including [] when the lane has no nonterminal work. explain --submission <id> --json returns one explanation object.

evidence.pendingOperations derives unresolved calls from their original committed model responses, including parameters, execution class, kind, and replay hash. unknownCalls[].resolved means a canonical ToolCallSettled exists. Accepted resolution intents appear separately in unknownResolutions; an intent alone is not a recorded result.

The administrative methods above, observation, settlement waits, abort, and unknown or approval resolution consult OperationAuthorizer. Its default allows trusted service holders. Install a real authorizer before exposing these methods outside a trusted host. Denial fails as OperationDenied before protected I/O. The host must authorize admissions before calling submit.

Use explain to inspect unresolved effects and their original declarations. A committed tool call may have executed even when no result was recorded; declaration alone grants no execution authority.

For exact receipts, use ThreadStore.getRecord({ threadId, recordId }) or ThreadStore.getRunInput({ threadId, runId }). The latter returns the original user input, excluding joined inputs, and rejects ambiguous original inputs. Both return an optional canonical envelope and require only ThreadReader. Authorize the owner and locator before reading and verify the returned payload; absence alone does not prove an admission was never accepted.

Native stores provide readIdentity({ threadId }) for the first canonical record and exact worker origin/lineage records with their captured tail and producer epoch. It is one bounded, consistent read, not an authorization grant. Custom ThreadStore adapters must implement this operation; deploy matching Cloudflare client and owner packages for its read-only port protocol.

Use the existing submissionInputRecordId / submissionSettlementRecordId exports from @yielded/agent/submission-ledger, or these @yielded/agent/run-journal locators:

  • workerInputRecordId(messageId) and firstWorkerInputRecordId(worker) for source reservations, including before child admission and after retirement;
  • workerOriginRecordId(threadId) for child origin;
  • workerReportRecordId(destinationIdempotencyKey) and peerMessageRecordId(messageId);
  • agentUpdateRecordId(threadId, runId, updateId), an Effect requiring Crypto.

Worker admission resolves its origin and selected funding Run by exact identity. Capacity comes from live inputs and unresolved effects across the whole Thread; completed reservations do not consume lifetime capacity. A terminal Run does not release an unresolved effect’s charge. Peer admission checks live message capacity and retained message identities; replies use the original envelope’s delivery principal.

MessageDelivery.readPending({ ownerThreadId, limit }) reads current retained deliveries through the existing owner-scoped list with pendingOnly: true, separately from canonical operations. It includes accepted, parked, and future-due entries; processed/refused history is excluded. Cloudflare routes this read to the source owner. Bounds, missing capability, malformed data, or a mismatched owner fail closed. A no-receipt action delivery remains uncertain even before a canonical worker reservation exists. A native receipt proves accepted admission, not destination materialization; worker admission reserves its source input before returning a receipt. Applications can leave accepted worker discovery to canonical worker inputs and recognize native completion/update reports without reconstructing transport validation. These reads grant no execution authority.

Aborting a submission retains its unknown outcomes. Execution decisions such as AbortSubmission and SafeToRetry do not establish whether an effect happened. Later supplier reconciliation, CompletedWithResult, or NeverHappened can close that original effect while preserving the terminal settlement. Same-format transfer imports into an empty Thread; see adopting these contracts for the supported baseline.

scanObligations scans current ledger state. It returns submissionId, threadId, state, age, severity, and one of these blockers: unknown, approval, waitingForChild, ready-aged, or running-aged.

Run the scan from cron, an alarm, or your monitoring service. Alert on every unknown outcome, every overdue row, and a growing approval backlog. The framework starts no monitoring daemon and sends no alerts.

After authorization, abort a submission with:

import {
class AbortCommand

A durable abort command (durability §13). Field bounds are exactly those of the canonical AbortRequested payload so the intent can become canonical without re-validation.

AbortCommand
} from "@yielded/agent/submission-ledger";
import {
class DurableAgentRuntime

Durable Agent Runtime coordinator (deployment class DN; D1/D2). It coordinates the SubmissionLedger, ThreadStore, and WakeScheduler ports so that once submit returns a Receipt, the Submission settles exactly once (DUR-001/DUR-002) while every external effect remains at-least-once (DUR-003) — this runtime never claims exactly-once side effects.

  • submit(agent, input, options) — durable admission → Thread materialization → ThreadCreated → readiness → Receipt, with failpoints between the steps; a retry with the same (thread, principal, idempotencyKey) resumes and returns the same Receipt.
  • awaitSettlement(receipt) — authorized once before the first ledger read, for the lifetime of this wait; interrupting it detaches the caller only and never cancels accepted work.
  • observe(receipt, {after}) — canonical record observation from a stored offset.
  • abort(command) — authorized before ledger access; durable idempotent abort intent; inactive work settles aborted through recovery, an active worker makes the command canonical before interrupting its Run (§13), settled work fails with SettlementConflict (DUR-012). A joined Submission fails with a typed JoinedToHost conflict carrying the host identity — it settles with its host, so the abort target is the host; aborting a joining Submission records the intent, honored only if the host has not consumed the input (revert-then-abort, plan §2.5).
  • resolveUnknown(command) — the authorized DUR-017 resolution path for Unknown Outcomes: the durable intent is idempotent per (submission, tool call) and conflicts typed on divergence; the canonical resolution records are applied by recovery or the next Attempt, and the lane wakes once every marked call is covered.
  • resolveApproval(command) — the durable approval decision path (plan §2.6): the intent is idempotent per (submission, tool call) with a typed ApprovalConflict on divergence; once every pending call of the suspension reason is decided the lane wakes (suspended → input-applied) and the next Attempt resumes the declared batch without model re-invocation, appending the canonical ToolApprovalDecided before honoring the decision.
  • processThread(agent, threadId) — drain one lane: fenced FIFO-head claims, canonical input apply, split response/results Turn commits (plan §2.1), reconcile-then-mark for open ordinary Tool Calls (DUR-009, never an automatic replay), declared-batch resume without model re-invocation (§15), and terminalization. An active host Run claims the contiguous ready prefix of later queued Submissions at every safe Turn seam (Joining/Joined, plan §2.5): the queued input becomes canonical (input:{sid}) before the next model request, reattaches through the prompt-coverage rule after a crash, and the joined Submissions settle with the host outcome (DUR-002/DUR-016).
  • processThreadResolved(threadId) / runResolvedWorker — drain or continuously process lanes using registrations owned by the runtime Layer: each stable agentId selects one current Binding. Unfinished operation contracts gate handler execution independently of immutable admission evidence; missing bindings release the claim with a typed refusal.
  • discoverWork() / discoverWorkThreads() — enumerate native metadata without granting execution authority. Missing or incomplete indexes fail explicitly; rebuildWorkIndex() reconstructs them in bounded, resumable pages independently of ordinary recovery.
  • runRecovery() — repair at most 32 native owners per pass, returning a cursor for remaining work. Follow that cursor to finish the sweep. Selected admissions use classifyRecovery; ready input and model work remain deferred for a fenced worker claim. Terminal Runs can retain factual tool closures, frozen deliveries, child accounting and worker acknowledgements. These repairs use original evidence and never invoke a Tool handler or reopen a settled Run.
  • explain/explainThread, verify, retry, wake, scanObligations — the P7 administrative operations (plan §3) over the same two ports, identical on DN and DC. explain and verify are strictly read-only; retry re-drives exactly one classified repair with mandatory author/reason audit and typed refusals; scanObligations is the scan-based DUR-017/OPS-001 obligation surface. Every one of them (plus observe, resolveUnknown, resolveApproval) consults the OperationAuthorizer reference fail-closed — the default Layer preserves the service-possession behavior, and a host-supplied authorizer turns denials into the typed OperationDenied.

DurableAgentRuntime
} from "@yielded/agent/durable-agent-runtime";
import {
import Effect
Effect
} from "effect";
const
const abortSubmission: (command: AbortCommand) => Effect.Effect<AbortIntent, DurableAbortFailure, DurableAgentRuntime>
abortSubmission
=
import Effect
Effect
.
const fn: (name: string, options?: SpanOptionsNoTrace) => Effect.fn.Traced (+63 overloads)
fn
("abortSubmission")(function* (
command: AbortCommand
command
:
class AbortCommand

A durable abort command (durability §13). Field bounds are exactly those of the canonical AbortRequested payload so the intent can become canonical without re-validation.

AbortCommand
) {
const
const runtime: {
readonly bindingRegistryKey: string;
readonly workerHost: (request: {
readonly sourceThreadId: ThreadId;
readonly principal: Principal;
readonly sourceSubmissionId?: SubmissionId;
}) => Effect.Effect<SubagentHost["Service"], WorkerError>;
readonly messagingHost: (request: {
readonly sourceThreadId: ThreadId;
readonly principal: Principal;
}) => Effect.Effect<MessagingHost["Service"], MessagingError>;
... 27 more ...;
readonly recoverWork: (request: WorkRecoveryRequest) => Effect.Effect<WorkRecoveryReport, DurableWorkerFailure>;
}
runtime
= yield*
class DurableAgentRuntime

Durable Agent Runtime coordinator (deployment class DN; D1/D2). It coordinates the SubmissionLedger, ThreadStore, and WakeScheduler ports so that once submit returns a Receipt, the Submission settles exactly once (DUR-001/DUR-002) while every external effect remains at-least-once (DUR-003) — this runtime never claims exactly-once side effects.

  • submit(agent, input, options) — durable admission → Thread materialization → ThreadCreated → readiness → Receipt, with failpoints between the steps; a retry with the same (thread, principal, idempotencyKey) resumes and returns the same Receipt.
  • awaitSettlement(receipt) — authorized once before the first ledger read, for the lifetime of this wait; interrupting it detaches the caller only and never cancels accepted work.
  • observe(receipt, {after}) — canonical record observation from a stored offset.
  • abort(command) — authorized before ledger access; durable idempotent abort intent; inactive work settles aborted through recovery, an active worker makes the command canonical before interrupting its Run (§13), settled work fails with SettlementConflict (DUR-012). A joined Submission fails with a typed JoinedToHost conflict carrying the host identity — it settles with its host, so the abort target is the host; aborting a joining Submission records the intent, honored only if the host has not consumed the input (revert-then-abort, plan §2.5).
  • resolveUnknown(command) — the authorized DUR-017 resolution path for Unknown Outcomes: the durable intent is idempotent per (submission, tool call) and conflicts typed on divergence; the canonical resolution records are applied by recovery or the next Attempt, and the lane wakes once every marked call is covered.
  • resolveApproval(command) — the durable approval decision path (plan §2.6): the intent is idempotent per (submission, tool call) with a typed ApprovalConflict on divergence; once every pending call of the suspension reason is decided the lane wakes (suspended → input-applied) and the next Attempt resumes the declared batch without model re-invocation, appending the canonical ToolApprovalDecided before honoring the decision.
  • processThread(agent, threadId) — drain one lane: fenced FIFO-head claims, canonical input apply, split response/results Turn commits (plan §2.1), reconcile-then-mark for open ordinary Tool Calls (DUR-009, never an automatic replay), declared-batch resume without model re-invocation (§15), and terminalization. An active host Run claims the contiguous ready prefix of later queued Submissions at every safe Turn seam (Joining/Joined, plan §2.5): the queued input becomes canonical (input:{sid}) before the next model request, reattaches through the prompt-coverage rule after a crash, and the joined Submissions settle with the host outcome (DUR-002/DUR-016).
  • processThreadResolved(threadId) / runResolvedWorker — drain or continuously process lanes using registrations owned by the runtime Layer: each stable agentId selects one current Binding. Unfinished operation contracts gate handler execution independently of immutable admission evidence; missing bindings release the claim with a typed refusal.
  • discoverWork() / discoverWorkThreads() — enumerate native metadata without granting execution authority. Missing or incomplete indexes fail explicitly; rebuildWorkIndex() reconstructs them in bounded, resumable pages independently of ordinary recovery.
  • runRecovery() — repair at most 32 native owners per pass, returning a cursor for remaining work. Follow that cursor to finish the sweep. Selected admissions use classifyRecovery; ready input and model work remain deferred for a fenced worker claim. Terminal Runs can retain factual tool closures, frozen deliveries, child accounting and worker acknowledgements. These repairs use original evidence and never invoke a Tool handler or reopen a settled Run.
  • explain/explainThread, verify, retry, wake, scanObligations — the P7 administrative operations (plan §3) over the same two ports, identical on DN and DC. explain and verify are strictly read-only; retry re-drives exactly one classified repair with mandatory author/reason audit and typed refusals; scanObligations is the scan-based DUR-017/OPS-001 obligation surface. Every one of them (plus observe, resolveUnknown, resolveApproval) consults the OperationAuthorizer reference fail-closed — the default Layer preserves the service-possession behavior, and a host-supplied authorizer turns denials into the typed OperationDenied.

DurableAgentRuntime
;
return yield*
const runtime: {
readonly bindingRegistryKey: string;
readonly workerHost: (request: {
readonly sourceThreadId: ThreadId;
readonly principal: Principal;
readonly sourceSubmissionId?: SubmissionId;
}) => Effect.Effect<SubagentHost["Service"], WorkerError>;
readonly messagingHost: (request: {
readonly sourceThreadId: ThreadId;
readonly principal: Principal;
}) => Effect.Effect<MessagingHost["Service"], MessagingError>;
... 27 more ...;
readonly recoverWork: (request: WorkRecoveryRequest) => Effect.Effect<WorkRecoveryReport, DurableWorkerFailure>;
}
runtime
.
abort: (command: AbortCommand) => Effect.Effect<AbortIntent, DurableAbortFailure>
abort
(
command: AbortCommand
command
);
});

Create the command with AbortCommand.make({ submissionId, author, reason }). On Cloudflare, obtain client from yield* CloudflareThreadClient, then call client.abort(receipt.threadId, command). The runtime checks authorization before reading or mutating the target.

Recovery claims the unknown Submission with its abort intent, handles attached children, and records an aborted settlement. It does not replay uncertain ordinary tools. The unknown evidence and first abort audit remain. Abort cannot roll back an external effect. A SettlementConflict reports a terminal result that won before the abort.

Do not edit ledger state, wake the lane in a loop, resolve every open call separately, or clean up children by hand. The runtime owns those steps after it durably accepts the parent abort. It never chooses parent abort merely because a child is unknown. The host makes that decision.

There is no automatic inactivity timeout. A parked unknown Submission stays quiet while later input in the same Thread can run. Approval, joining, and joined states retain their ordering barriers. Authorized resolution or abort restores maintenance for the parked work; a live owner still prevents another claim in that Thread.

Persist the admission Receipt. Its submission, receipt, thread, and queue identifiers remain stable across retries and replacement attempts.

awaitSettlement authorizes the receipt’s thread and submission before lookup. A receipt that mixes identifiers fails as OperationDenied. Authorization lasts for one wait. To enforce revocation, interrupt the wait and start another. Interrupting a wait does not abort work.

Need API or record
admitted or queued SubmissionLedger.lookup returns admitted or ready
execution stage lookup returns running or input-applied
suspension or joined host loadRecoverySnapshot returns suspension and host linkage
unknown outcome lookup returns unknown; explain adds calls, audits, resolutions, and abort intent
terminal outcome awaitSettlement(receipt) returns the durable settlement
budget or stop detail SubmissionSettled records finishReason, exhausted, and policyLimit

Use runtime.observe(receipt, { after }) for live progress. Filter the thread stream by submission or run identifier. Joined input shares its host run, while its own terminal record keeps the original submission identifier.

import { DurableAgentRuntime, type Receipt } from "@yielded/agent/durable-agent-runtime";
import { type ObservationOffset } from "@yielded/agent/records";
import { Effect, Stream } from "effect";
const observeOutcome = (receipt: Receipt, after?: ObservationOffset) =>
Stream.unwrap(
Effect.gen(function* () {
const runtime = yield* DurableAgentRuntime;
return runtime.observe(receipt, { after }).pipe(
Stream.filter(
({ record }) =>
record.payload._tag === "SubmissionSettled" &&
record.payload.submissionId === receipt.submissionId,
),
Stream.take(1),
);
}),
);

On Cloudflare, read bounded pages with readPage, save the last sequence, and call awaitProgress after an empty page. Avoid readAll and repeated explain calls in polling code. Keep reads incremental and retain only the application’s projection.

The public Cloudflare client cannot read the complete local recovery snapshot or scan all nonterminal rows. If an external Worker needs suspension, FIFO blocker, input marker, or child references, add an authorized Schema-backed read-only RPC to the owning Thread Object. Use its local ledger and recovery snapshot. Do not copy the recovery classifier into the Worker.

Persist observation cursors and make downstream projection or delivery idempotent. A crash after an external side effect can redeliver the record. Streams, notifications, callbacks, and process finalizers provide no exactly-once delivery guarantee.

Use operation spans to diagnose storage and recovery. Adapter storage spans cover append, claim, renew, release, publish, and finalize boundaries; recovery spans cover recovery operations. Hydration and decoding stay within the enclosing operation span, without separate spans for individual records or decoded rows.

Use a file-consistent SQLite snapshot. Copy the database, WAL, and SHM files while no process owns the database, or use VACUUM INTO or SQLite’s online backup API. The supported DN shape has one process owner per database file. With an automatically managed host, run online backup through its shared connection; another connection cannot read the live database. Stop the host before using an external backup tool.

A restore has four rules:

  1. Pre-backup history must pass the same verify checks as the original.
  2. Fence or terminate every producer from the original timeline before serving restored data. Restored storage rejects post-backup ownership tokens.
  3. Treat external effects recorded only after the backup as unknown. Never replay them automatically. Resolve them through resolveUnknown using supplier records.
  4. Receipts issued after the backup are gone. Clients must resubmit with idempotency keys. Reconcile any external effects from those lost submissions.

SQLite-backed Durable Objects provide 30-day point-in-time recovery through getCurrentBookmark, getBookmarkForTime, and onNextSessionRestoreBookmark. Miniflare does not implement these APIs, so follow this hosted runbook:

  1. stop admission for affected threads and let current alarms drain;
  2. obtain the desired bookmark;
  3. call onNextSessionRestoreBookmark(bookmark) inside the Object, then ctx.abort();
  4. apply the four restore rules above, including unknown external effects and lost receipts;
  5. run verifyEncoded and obligationsEncoded before reopening admission.

Cross-Object children created after the bookmark need parent recovery. A restored Object replaces old local producer epochs automatically.

Put each runtime storage domain behind authenticated ingress. Yielded Agent provides no tenant authentication or row-level tenant isolation. Receipts, IDs, principal strings, and audit authors identify records. They grant no access.

Schedule and subscription owners are tenant-qualified. Ordinary thread requests and records have no tenant field. A different principal does not create a separate thread log.

  • On Node, bind each tenant to its own SQLite database or enforce application isolation on every read, write, scan, worker, and administration path.
  • On Cloudflare, derive tenant-qualified Thread Object addresses in trusted Worker code. CloudflareThreadClient does not infer or authorize a tenant from a thread ID.

Keep tenant addressing stable across admission, history, children, schedules, subscriptions, exports, backups, and operator access. Always check that a submission belongs to its selected thread.

Operation Host responsibility Framework boundary
admission and prepared delivery authenticate ingress, derive principal, authorize agent and destination Schema validation and idempotency grant no access; authorizer is not an admission hook
reads and observation authorize the thread and submission before returning data runtime observation and waits consult the authorizer; direct store reads do not
abort and resolution authorize the decision and supply trusted audit fields runtime checks authorization before target I/O
administration and scans restrict operators to the selected storage domain default authorizer allows all; scans have no automatic tenant filter

possessionOperationAuthorizer allows every request. Use it only behind a trusted host. Install operationAuthorizerLayer when constructing the durable runtime. The runtime captures the policy at Layer acquisition, so a Layer around a later method call cannot replace it.

Authenticate callers before invoking runtime methods. Some operation requests carry no principal, and audit fields provide no authorization. Keep raw stores, ledgers, runtime services, and Durable Object RPCs away from untrusted callers.

Use RunToolAuthorization to check proposed tool calls. OperationAuthorizer protects the runtime operations listed above; the host authorizes admission. Recheck policy after durable suspension. Approval applies to one exact action and expires. Parent approval does not authorize child actions.

Delegation grants are immutable ceilings. Each child action must satisfy current policy, its grant, target requirements, and resource scope. Model parameters cannot supply bindings, secrets, identity, or policy. Validate and bound child output before use.

Keep secrets as handles. Redact diagnostics before storage or export. Generated code needs an isolated executor with host tools behind the validated broker. Local sandbox processes have no isolation. Exact-host browser checks do not provide connection-time network isolation. Enforce read-only SQL through database permissions and host tenant scope.

Scheduling delivers encoded agent input through durable admission. A due occurrence remains its obligation until it records a Receipt or proves permanent refusal. The host provides ScheduleAuthorizer and owner-scoped management. Keep ScheduleDriver and ScheduleStore inside the privileged host.

Preparation freezes the authorized envelope and advances the cursor atomically. Recover pending delivery before preparing another occurrence. A lost admission reply retries the same envelope and idempotency key. Transient or ambiguous failure stays pending. Automatic delivery retries stop after maxAutomaticAttempts (default 8); parked work retains its envelope and receipt uncertainty. Scheduling.recover re-arms that same identity, including after cancellation. Recovery advances an attempt generation (pass the observed pending.retry.generation to recover), so a stale failed attempt cannot spend the new retry allowance. A late valid Receipt can still complete the original obligation. Permanent refusal requires proof that admission did not occur and the unchanged request cannot succeed.

Recurring downtime coalesces to the latest due firing. Pause and cancel stop new preparation while pending delivery continues. Cancel is irreversible. Revocation blocks future preparation but must finish already authorized envelopes. Resuming a paused schedule skips missed recurring times.

Use { _tag: "Interval", everyMillis: 60_000 } for an interval, or { _tag: "Cron", expression: "0 9 * * *", timeZone: "America/New_York" } for an explicit IANA cron zone. Cron without a zone keeps its existing UTC default; no host-local zone is inferred.

Quotas count pending delivery plus active or paused cursors. Terminal records retain replay evidence without using capacity when no delivery remains. Schedule IDs and creation evidence are never recycled. One corrupt due record cannot block later records; retry a failed sweep after the recovery poll.

Node runs one Scope-owned indexed polling driver. Cloudflare commits schedule changes and alarms together, pre-arms recovery before admission, and fences alarm acknowledgements by generation. Storage failure that prevents alarm repair needs a later wake or operator action after storage recovers.

See the compiling Node and Cloudflare examples.

Subscriptions retain normalized events, select matching registrations, prepare agent input, and deliver through durable admission. EventAcknowledgement confirms retained intake. A Receipt confirms admission. No run or waiter stays open to watch the source.

Set expiresAtMillis: null for no time-based expiry; a once subscription is still consumed by selection, and cancellation remains explicit. Finite deadlines keep their existing maximum-lifetime validation. Pausing does not erase selected delivery evidence.

Without an explicit retention policy, records remain for the partition lifetime. Set SubscriptionLimits.retention to bound completed event payloads and deliveries. Intake then requires a stable authenticated occurredAtMillis from the source adapter and rejects fresh identities outside replayHorizonMillis. This horizon is fixed for the partition once used. Completed work becomes a compact deduplication tombstone until that horizon ends; duplicate intake cannot reopen routing. completedRetentionMillis controls completed payload retention and maxTombstones bounds deduplication storage. Backpressure remains explicit if live evidence or in-horizon tombstones fill capacity. Idle native maintenance expires old tombstones.

Selected, prepared, parked, recovery-related and admitted-but-unsettled evidence is protected. Node and Cloudflare admission adapters use the runtime’s canonical submissionStatus. Enabling retention with a custom PreparedInputAdmission requires that observation capability. An unavailable probe is logged and retried without releasing evidence. Reconciliation of a payload already reclaimed reports event-reclaimed rather than inventing a new delivery.

Each maintenance pass examines at most batchSize event candidates and batchSize delivery candidates, plus one referenced event per delivery. Indexed relationship checks protect recovery and unsettled work. Independent durable cursors advance past corrupt or protected candidates; corrupt evidence is preserved and reported. Idle maintenance rearms while retained events remain. The in-memory adapter copies its bounded maps/indexes, but does not decode the whole partition.

EventSource.occurredAtMillis must derive a stable timestamp from authenticated source facts. Do not substitute intake time or generate a new timestamp on replay. Exact retained identities replay before current horizon/capacity checks, including their timestamp and payload digest. Expired pruned identities are rejected by the fixed horizon. Omit retention when the source cannot supply trusted occurrence time; there is no undated-event pruning exception.

Each event and registration belongs to one stable SourcePartition with a tenant ID and source address. Keep the address unchanged across deployments and source versions. There is no cross-partition transaction or global subscription directory.

Use makeEventSource and EventSources for versioned event schemas, identity, matching, and optional reconciliation. Use makeSubscriptionInputBinding and SubscriptionInputBindings for destination preparation. Their callbacks and Schema codecs receive a fresh Scope per operation, so acquired resources finalize when that operation completes, fails, or is interrupted. Other service dependencies are captured at host assembly. Keep old source versions and bindings installed while retained work needs them. Missing or ambiguous bindings leave selected delivery pending as unsupported-binding. Persisted records contain no callbacks, Schemas, Effects, credentials, or captured services.

SubscriptionAuthorizer covers management, intake, reconciliation, and preparation. Recovery has its own reconcile decision. Keep stores, intake, and drivers out of model tool environments. Restricted tools must bind owner, agent, principal, source catalog, and thread in the host. A host may also permit a deterministic fresh thread for each selected event.

Use getSubscription to inspect the current revision, state, and configuration fingerprint. updateSubscription, pauseSubscription, resumeSubscription, and recoverSubscription require the expected configurationRevision. cancelSubscription also accepts an optional expected revision. Revision conflicts expose the current revision/state; inspect the fingerprint after a lost management reply to determine whether the intended configuration won. Creation replay always compares the original request fingerprint, even after edits.

Every configuration/control revision advances the registration’s eligibility ordinal. It can select only events accepted after that revision. Events accepted earlier but not yet selected do not gain eligibility under the new configuration, including changed matching keys. Selected deliveries retain their captured configuration; pausing stops new selection while selected delivery continues. Cancellation and captured expiry can refuse unprepared work. Prepared envelopes remain unchanged. An uncertain ordinary tool is never replayed automatically.

Source recovery completions carry the captured registration revision. Pause preserves its recovery intent, resume restores polling, and callbacks from an older revision cannot overwrite a new one. Recovery reads share the per-item timeout and failure boundary. An unreadable record is reported without guessing its revision or blocking routing, delivery, and maintenance for other records. recoverSubscription re-arms source reconciliation where the configured source supports it; recoverDelivery(scope, key, expectedGeneration) re-arms an individual parked delivery without changing its identity or envelope. Repeating the same recovery generation is an idempotent no-op.

Intake deduplicates by tenant, source address, and logical event ID. Conflicting payload or source version fails. Each event records a registration cutoff. Duplicate intake cannot move that cutoff or reopen routing. Selection atomically advances the cursor, creates delivery obligations, and consumes once registrations. Continuous subscriptions retain each event separately. Thread admission order decides execution order.

Preparation rechecks authority, validates input, and freezes destination, principal, digests, authorization metadata, and admission key. Cancellation, expiry, or revocation blocks new preparation. Already prepared envelopes continue unchanged. Expiry never sends agent input.

Lost admission replies retry the exact envelope and key. After maxAutomaticAttempts (default 8), delivery parks until recoverDelivery explicitly re-arms it. Cancellation never erases a prepared uncertain envelope. Capacity waits spend the same automatic retry allowance; hosts should size it for expected job duration and explicitly recover parked work. Event routing and source recovery use their existing bounded sweeps and retry deadlines; settlement probes retry conservatively. The delivery attempt cap does not cap all background maintenance.

Optional admissionGroup permits one actually unsettled submission per group in a destination thread. Admission, suspension, unknown outcome and canonical settlement publication retain occupancy until canonical settlement finalization. A lagging finalization conservatively holds capacity until repair. FreshThread per event does not provide exclusion across threads. Schedules retain one frozen pending occurrence and coalesce missed recurring times; distinct events keep separate durable obligations and produce explicit backpressure when backlog capacity is exhausted.

admissionFence captures bounded { policyId, key, revision } coordinates. Install SubmissionAdmissionFence when acquiring the destination ledger. Exact retained requests replay before policy and occupancy checks; changed input, group or fence conflicts. Fresh admission checks policy in the same local transaction as the ledger insert. SQL hosts must read policy through that transaction’s SqlClient, and all policy writers must share its authority. Remote policy reads cannot fence remote mutations. Memory executes the bounded callback synchronously inside its atomic mutation; an asynchronous callback fails closed as unavailable and is interrupted. Do not fork, reenter admission, or retain transaction resources in the callback.

AdmissionPolicyError.reason distinguishes refused (conclusive stale/unsupported policy), unavailable (retry without assuming nonadmission), and occupied (group capacity). Uninterpreted fences fail closed. Ingress authorization is still required for replays. Put no secrets in fence coordinates or error codes. Public status excludes payloads, parameters, context and credentials. Execution remains at least once; provider failures never invent event completion.

Cloudflare hosts can install SubscriptionPartitionAlarmExtension with handlers built by makeSubscriptionPartitionAlarmHandler. Each handler owns one non-framework tag, a payload Schema and a bounded timeout (at most 30 seconds). The factory captures host services but defers Subscriptions, SubscriptionIntake and SubscriptionDriver requirements, including those in payload decoders, until invocation. makeSubscriptionPartitionObjectClass supplies these native services from the addressed partition after building the host Layer. Keep their lookup inside the callback or decoder; yielding them while building the host Layer still requires them at assembly. Handlers retain only their required native services in R; host requirements remain on the factory Effect. Even if native services were present at assembly, the invocation uses its own instances.

An existing handle: () => runtime.process can therefore keep SubscriptionIntake in the process Effect’s requirements. Map its expected failures to SubscriptionAlarmExtensionError; no extra handler argument or client Layer is needed. See the compiling Cloudflare example below for a handler that accepts an event through native intake.

Each codec/callback invocation owns a fresh Scope. Invocation cleanup runs on success, typed failure, defect, timeout and interruption; captured host services keep their host lifetime.

The native multiplexer processes at most 16 alarms per invocation and durably retries failed rows independently, so a failed or unknown extension cannot block the subscription driver. Installing handlers reserves eight minutes of the twelve-minute invocation budget for ancillary work; native driver limits must fit the remaining four minutes. Unknown, ambiguous, malformed and reserved @yielded/agent/ tags fail closed and are reported. A replacement alarm survives acknowledgement of its earlier version. The host owns external-effect idempotency, uncertainty, payloads and transactional prearming; this extension defines no provider scheduler.

Table layout and record meaning have independent versions. The unreleased protocol accepts fresh layout-21 stores and effect-agent/thread@3 records on SQLite, PostgreSQL, and Cloudflare. Opening predecessor or ambiguous stores fails before DDL or payload mutation. Keep them with their matching release; no historical decoder, layout upgrade, converter, or mixed-format runtime is included.

Quiesce the source and retain a backup before transferring it. streamExport({ threadId }) yields bounded pages through ThreadExportSource; provide its Layer from the source ThreadStore separately from the destination import Layer. ThreadImport.import(pages) consumes an archive Stream within one destination transaction; reencodeThread(source) composes export pages with that import. Check the restored tail before resuming accepted work under compatible Bindings. SQLite’s admin CLI saves pages as NDJSON:

vp run admin:durable export --database source.sqlite --thread THREAD --output thread.ndjson
vp run admin:durable import --database destination.sqlite --input thread.ndjson
vp run admin:durable reencode --database source.sqlite --thread THREAD --target destination.sqlite

Every page binds the canonical tail and a revision of independent admissions, accepted commands, and delivery facts. Changes during export require a new export. Import preserves exact record wire, batch producers, digest anchors, queue order, admission time, Receipts, principals, keys, and opaque admission fences and groups. Destination admission policy and active-group constraints still apply. Truncated pages, missing cross-range evidence, contradictory identities, or unavailable policy checks fail without publishing staged work. Keep the source when validation fails.

Import rebuilds ledger state and native indexes from facts. Claims, leases, execution ownership, application checkpoints, and recovery caches start empty. Unresolved ordinary tools remain Unknown; settlement does not close surviving effects or deliveries. Closed historical workers and children can transfer; live obligations owned by another store must finish through their owning workflow before a single-Thread import. Retained deliveries transfer without their operational lease.

Storage operation Bound
Atomic canonical batch 256 records, 16 MiB encoded batch
Physical range Complete batches; 1,024 records and 32 MiB including batch and record payload copies
Export/import page 256 records or independent facts, 32 MiB encoded page
Work-index rebuild 256 records by default (lower with limit), 32 MiB per pass
Archive range listing 32 range descriptors

These are working-set bounds, independent of Thread age. The host owns total storage, archive maintenance, live input/delivery capacity, and per-Run/model limits. A long import holds its transaction until validation completes. Memory storage retains all facts in volatile process memory and needs an explicit host memory allowance.

Use ThreadStore.archives for storage-owner maintenance: seal the current complete-batch range, then archive its firstSequence under the current producer epoch. SQL adapters copy and verify bounded contents, publish a stable locator, and only then remove hot payload copies in the same transaction. Interruption leaves a reachable copy. Ranges share the Thread’s fence and continuous sequence; they do not reset grants, Run budgets, or live obligations. Exact old-record reads and batch retries follow the native identity locator, without traversing intervening ranges.

ThreadExportRecord carries original record wire beside its typed view. Custom exporters must preserve that wire through the codec; reconstructing an envelope through the ordinary record Schema discards additive fields. Import never replaces a nonempty Thread or an existing Submission, Receipt, or scoped Thread/principal/key. The same principal/key in another Thread remains valid.

Cloudflare’s separate Schedule and Subscription stores use their current fresh layouts and reject predecessor schemas without mutation. Their retained input and destination ownership remain independent of canonical Run progress.

Application checkpoint consumers validate their own Schema. verifyOnOpen audits canonical history and application projections. Canonical continuation references, revisions, and original context are validated before import and recovery; missing or corrupt progress leaves work owed. See Run continuations.

Keep the source versions and input bindings needed to finish retained deliveries. Register the current binding for each stable Agent ID; unfinished operations retain their original replay contracts. See deployment continuity. Custom stores must implement revision and retry-generation fencing, bounded retention cursors, and canonical observation.

Node uses a Scope-owned polling driver. Cloudflare commits work and required alarms together and re-arms after failed passes. If storage prevents mutation and alarm repair, restore storage and send a new wake or intervene as an operator.

See the compiling Node and Cloudflare examples.

Import GitHub integration from @yielded/agent/git-hub-workflow-source. makeGitHubWorkflowRunSource watches one repository, run ID, attempt, and expected head SHA. It reports successful and unsuccessful completion. It does not aggregate every check for a commit.

acceptVerifiedGitHubWorkflowRunWebhook verifies raw bytes before parsing a completed workflow_run event. Webhook and exact-attempt API observations normalize to one completion identity. A webhook delivery ID alone cannot deduplicate them.

Registration arms reconciliation before provider reads. An already completed attempt can notify a new watch without a check-then-subscribe race. Cancellation, expiry, and once selection stop provider polling. GitHub does not automatically redeliver failed webhooks, so reconciliation checks the registered attempt while it remains retained, readable, and authorized. It does not provide general historical replay.

When composing additional SQLite adapters with a Node host, provide the host Layer to the adapter. The host exposes its SqlClient, SqliteStorageConfig, and SqliteStorageFailpoint so subscriptions and runtime operations share one serialized connection, including during startup recovery.

MessageDeliveryStore retains a host-prepared input independently of either Thread’s current Run. pending means the obligation is saved; an accepted Receipt parks the delivery as awaiting-settlement, with no status-poll deadline. processed means that exact Receipt has a terminal Settlement. Native workers and peer inputs acknowledge their source before finalizing the destination ledger. A lost acknowledgement is replayed from canonical settlement during recovery. A lost admission reply reuses the frozen envelope and admission key. It cannot create a replacement input.

Node hosts run a scoped, bounded polling loop over the delivery deadline index. Cloudflare Thread Objects prearm their maintenance generation before message writes, retain the earliest alarm, and run one bounded delivery wave beside source work. Both paths recover without wake hints, including after the source and destination Runs settle. Interrupting a host releases live resources while retaining the obligation for recovery.

Automatic retry has finite attempt and deadline bounds. refused and parked remain inspectable; an explicit driver retry renews a parked obligation’s deadline without changing its envelope. Generic host-prepared inputs without native source provenance remain dormant until the host supplies an exact Complete acknowledgement or explicitly retries the retained Receipt once. Completed rows retain exact deduplication evidence and do not consume pending-delivery capacity. The host owns total storage quotas.

These are trusted host ports. Authenticate the sender, authorize routing and encode the destination input before preparing an envelope. Possession of a message key or Receipt is not an authorization decision. Reads require an owner Thread; Cloudflare also enforces that each Object’s delivery store belongs to its own Thread.