Skip to content

Build agents

Run & stream

The runtime exposes one scoped Effect execution through run, stream, and start. All three decode input before instructions execute and require native model services. Use runUnknown, streamUnknown, or startUnknown for external values typed as unknown. See Agent definitions and the authoritative runtime model.

Use InMemory.layer from @yielded/agent for in-memory conversations, including attached subagents. Provide it once around the application and reuse a Thread ID for follow-up Runs. It retains complete history updates and shares subagent reservation state for that Scope. IDs are generated automatically, and context preparation is optional. Use PersistentHistory.layer with a store to retain completed runs, or a durable host when execution must recover after process loss.

A valid no-tool answer needs one model call. A designated completion Tool can complete without a follow-up model call. Independent application Tools default to four concurrent handlers, behind the batch’s authorization and approval barrier. Use one when execution must be serial or a fixture deliberately measures a sequential workflow. Continuity fixtures serialize changes to shared notes; Code Mode examples bound generated programs separately. Node host worker concurrency is a separate setting.

OpenAI and xAI Responses requests preserve system instructions in conversation order and place the output contract after the initial system block. Changing instructions in later Runs and appended system context stay after earlier history, preserving its cache prefix. An exact repeated instruction is omitted only when no different system instruction intervenes. Conversation-only history can recover its leading static instructions from those still present in the prepared prompt. Stored history remains intact.

The engine chooses this projection for the actual model selected on each call. The native LanguageModel service’s supportsSystemMessagesInHistory capability takes precedence over the provider default. Without it, only openai and xai bindings use chronological instructions. Other adapters group systems before the conversation, retaining the last equivalent instruction and its native options. Changing that grouped block can invalidate the history cache.

The repository’s Anthropic adapter requires grouping. Chronological Anthropic history needs a release containing the upstream adapter change and a model supporting mid-conversation system messages. That adapter must retain system authority and place later instructions after the corresponding user/tool results, before the next assistant response. A newer model name alone is insufficient.

For xAI, use the native @effect/ai-openai Responses adapter with an OpenAiClient whose apiUrl is https://api.x.ai/v1. Where the host knows the Thread identity, bind the provider and routing key:

Model.make(
"xai",
"grok-4.3",
OpenAiLanguageModel.layer({
model: "grok-4.3",
config: { store: false, prompt_cache_key: threadId },
}),
);

Supply the configured OpenAiClient Layer to this Model, including when returning it from context.prepare as modelCall.model. The routing key encourages server affinity; it does not create a cache entry or guarantee a hit. This path uses Responses; xAI’s Chat Completions header x-grok-conv-id is a separate transport setting. Do not assume OpenAI-specific request options such as promptCacheBreakpoint work on xAI.

For Anthropic, place native options.anthropic.cacheControl on the last retained user/tool message before changing transient context, using context preparation. For example, use { type: "ephemeral", ttl: "5m" }. Request-level automatic caching can instead write after the changing suffix; it alone does not establish reuse of the stable history. Follow Anthropic’s cache placement and model limits. Keep trusted instructions in system messages and untrusted references in user/tool content.

The immutable output contract also retains its message identity across turns, allowing OpenAI’s opt-in native ResponseIdTracker reuse for ordinary append-only prompts. Context preparation, transient references, and appended run status use full requests so provider-held responses cannot replay discarded material. Full requests can still use provider prompt caching.

Keep tools and fixed instructions stable, and append changing context after history. Rewriting or prepending context, compaction, and changing provider settings can invalidate cached prefixes. Provider caching, minimum prompt lengths, and billing depend on the selected provider and configuration; stable ordering does not guarantee a cache hit.

Context preparation is optional. Provide RunContextPreparation to load extra context; without it, Runs use their normal prompt and compaction behavior. See context management for service-based recall and tagged errors.

const program = Effect.gen(function* () {
const result = yield* AgentRuntime.run(agent, input);
return result;
});

run closes run-owned resources before returning decoded output. A self-contained run needs no caller Effect.scoped.

The result contains output, threadId, runId, turns, and finishReason. Budget-limited results also include exhausted, naming "turns", "tool-calls", or "tokens".

runDisposition appears only after ordinary completion when the definition declares one and its selector returns a value. It contains schema-encoded JSON. Decode durable settlement values with the same application schema.

With the default onExhaustion: "final-answer", turn, tool call, or token exhaustion allows one constrained final turn. The result reports finishReason: "budget-exhausted", and turns may exceed maxTurns by one. Duration or cost exhaustion, pending approval, interruption, and output decoding failure remain failures. Set onExhaustion: "fail" to fail before the final turn.

const events = AgentRuntime.stream(agent, input);
const program = events.pipe(
Stream.tap((event) => Effect.log(event._tag)),
Stream.runDrain,
);

Events cover run and turn lifecycle, text and reasoning deltas, tool activity, approval requests, and one terminal classification. Provider SDK chunks do not enter this stable union.

For structured output, treat text deltas as provisional wire data. Show activity or received-character progress until the terminal output passes its Schema; do not display partial JSON as an answer. The demo follows this pattern. Plain-text output can render provisional text directly. Provider parts are copied into bounded, engine-owned data and Schema-validated in order. The engine publishes each part’s semantic events before processing the next part, even within a single provider chunk. A later invalid or interrupted part preserves earlier published progress and reported usage. Transport fragmentation does not create an ownership span per delta.

Once stream consumption starts, a scoped producer advances until its bounded event buffer fills. Slow consumption backpressures publication, but individual pulls do not pace tool execution. The producer captures its execution Context at startup: provide services around the whole stream, rather than changing them around individual pulls. Completion, failure, and interruption close its resources. Interrupting the only ephemeral consumer interrupts the run.

Published ToolProgress results are owned JSON snapshots. Their cumulative UTF-8 JSON size is limited to 8 MiB per run, shared by application and provider progress. This also bounds progress payloads retained for detached replay and applies even when no progress is observed. Oversized progress fails with ModelProtocolError without truncation. Application progress must contain plain JSON data; accessors, custom serialization, and non-finite numbers fail with the same error. Terminal tool results use toolResultBounds separately.

Lower the progress allowance with bufferLimits on run, stream, or start. Larger values cannot raise the engine’s ceiling:

import { type
(alias) interface RunBufferLimits
import RunBufferLimits

Tightening-only memory limits for one Run. The engine supplies finite ceilings for every field; callers may lower them for a deployment or test but cannot widen the engine defaults.

RunBufferLimits
} from "@yielded/agent/run-options";
export const
const progressBufferLimits: RunBufferLimits
progressBufferLimits
:
(alias) interface RunBufferLimits
import RunBufferLimits

Tightening-only memory limits for one Run. The engine supplies finite ceilings for every field; callers may lower them for a deployment or test but cannot widen the engine defaults.

RunBufferLimits
= {
RunBufferLimits.maxToolProgressBytes?: number | undefined

Maximum cumulative UTF-8 JSON bytes in Tool progress, including unobserved and provider progress. Defaults to 8 MiB. Invalid or oversized application progress fails with ModelProtocolError; terminal Tool results use the Agent's toolResultBounds.

maxToolProgressBytes
: 1024 * 1024,
};

maxRunEvents bounds progress produced by stream and start. Headless run and durable execution do not produce a public event sequence. Use Agent policy limits to bound their execution; native model-response and tool-result validation and bounds still apply. maxBufferedEvents lowers the public stream queue capacity from its default of 1,024 events.

Declare updates on an Agent to give it an emit_update tool for structured intermediate findings:

observe-updates.ts
import {
import AgentRuntime
AgentRuntime
,
import AgentUpdates
AgentUpdates
} from "@yielded/agent";
import {
import Effect
Effect
,
import Stream
Stream
} from "effect";
import {
const HotelResearcher: Definition<Struct<{
readonly city: String;
readonly area: String;
readonly sources: $Array<Struct<{
readonly url: String;
readonly notes: String;
}>>;
}>, Struct<{
readonly hotels: $Array<String>;
readonly summary: String;
}>, string, Toolkit<{
readonly emit_update: Tool<"emit_update", {
readonly parameters: Struct<{
readonly value: NoInfer<TaggedStruct<"AreaConcern", {
readonly area: String;
readonly finding: String;
readonly sources: $Array<String>;
}>>;
}>;
readonly success: Struct<...>;
readonly failure: typeof AgentUpdates.UpdateError;
readonly failureMode: "return";
}, AgentUpdates.Emitter>;
}>, undefined, undefined, NoInfer<TaggedStruct<...>>> & {
...;
}
HotelResearcher
} from "./background-updates.ts";
const
const updates: Stream.Stream<AgentUpdates.Update, AgentRuntime.AgentRuntimeFailure<Definition<Struct<{
readonly city: String;
readonly area: String;
readonly sources: $Array<Struct<{
readonly url: String;
readonly notes: String;
}>>;
}>, Struct<{
readonly hotels: $Array<String>;
readonly summary: String;
}>, string, Toolkit<{
readonly emit_update: Tool<"emit_update", {
readonly parameters: Struct<{
readonly value: NoInfer<TaggedStruct<"AreaConcern", {
readonly area: String;
readonly finding: String;
readonly sources: $Array<String>;
}>>;
}>;
readonly success: Struct<...>;
readonly failure: typeof AgentUpdates.UpdateError;
readonly failureMode: "return";
}, AgentUpdates.Emitter>;
}>, undefined, undefined, NoInfer<TaggedStruct<...>>> & {
...;
}, never, never>, AgentRuntime.AgentRuntimeRequirements<...>>
updates
=
import AgentRuntime
AgentRuntime
.
stream<Definition<Struct<{
readonly city: String;
readonly area: String;
readonly sources: $Array<Struct<{
readonly url: String;
readonly notes: String;
}>>;
}>, Struct<{
readonly hotels: $Array<String>;
readonly summary: String;
}>, string, Toolkit<{
readonly emit_update: Tool<"emit_update", {
readonly parameters: Struct<{
readonly value: NoInfer<TaggedStruct<"AreaConcern", {
readonly area: String;
readonly finding: String;
readonly sources: $Array<String>;
}>>;
}>;
readonly success: Struct<...>;
readonly failure: typeof AgentUpdates.UpdateError;
readonly failureMode: "return";
}, AgentUpdates.Emitter>;
}>, undefined, undefined, NoInfer<TaggedStruct<...>>> & {
...;
}, never, never>(agent: Definition<...> & {
...;
}, input: NoInfer<{
...;
}>, options?: RunOptions<...> | undefined): Stream.Stream<...>
export stream

Accept schema-encoded input, retaining runtime validation. Use streamUnknown for external data.

stream
(
const HotelResearcher: Definition<Struct<{
readonly city: String;
readonly area: String;
readonly sources: $Array<Struct<{
readonly url: String;
readonly notes: String;
}>>;
}>, Struct<{
readonly hotels: $Array<String>;
readonly summary: String;
}>, string, Toolkit<{
readonly emit_update: Tool<"emit_update", {
readonly parameters: Struct<{
readonly value: NoInfer<TaggedStruct<"AreaConcern", {
readonly area: String;
readonly finding: String;
readonly sources: $Array<String>;
}>>;
}>;
readonly success: Struct<...>;
readonly failure: typeof AgentUpdates.UpdateError;
readonly failureMode: "return";
}, AgentUpdates.Emitter>;
}>, undefined, undefined, NoInfer<TaggedStruct<...>>> & {
...;
}
HotelResearcher
, {
city: string
city
: "Johannesburg",
area: string
area
: "Rosebank",
sources: readonly {
readonly url: string;
readonly notes: string;
}[]
sources
: [],
}).
Pipeable.pipe<Stream.Stream<AgentUpdateEmitted | RunStarted | TurnStarted | ModelStarted | ModelRestarted | TextDelta | ReasoningDelta | ToolCallDeclared | ToolCallStarted | ToolProgress | ToolCallSucceeded | ToolCallFailed | ApprovalRequested | TurnCompleted | BudgetWarning | CompactionPerformed | ... 10 more ... | SubagentJoined, AgentRuntime.AgentRuntimeFailure<...>, AgentRuntime.AgentRuntimeRequirements<...>>, Stream.Stream<...>, Stream.Stream<...>>(this: Stream.Stream<...>, ab: (_: Stream.Stream<...>) => Stream.Stream<...>, bc: (_: Stream.Stream<...>) => Stream.Stream<...>): Stream.Stream<...> (+21 overloads)
pipe
(
import Stream
Stream
.
const filter: <AgentUpdateEmitted | RunStarted | TurnStarted | ModelStarted | ModelRestarted | TextDelta | ReasoningDelta | ToolCallDeclared | ToolCallStarted | ToolProgress | ToolCallSucceeded | ToolCallFailed | ApprovalRequested | TurnCompleted | BudgetWarning | CompactionPerformed | ... 10 more ... | SubagentJoined, AgentUpdateEmitted>(refinement: Refinement<...>) => <E, R>(self: Stream.Stream<...>) => Stream.Stream<...> (+3 overloads)

Filters a stream to the elements that satisfy a predicate.

Example (Filtering stream values)

import { Effect, Stream } from "effect"
const program = Effect.gen(function*() {
const stream = Stream.make(1, 2, 3, 4).pipe(
Stream.filter((n) => n % 2 === 0)
)
const values = yield* Stream.runCollect(stream)
values // => [ 2, 4 ]
})
await Effect.runPromise(program)

@category ― filtering

@since ― 2.0.0

filter
((
event: AgentUpdateEmitted | RunStarted | TurnStarted | ModelStarted | ModelRestarted | TextDelta | ReasoningDelta | ToolCallDeclared | ToolCallStarted | ToolProgress | ToolCallSucceeded | ToolCallFailed | ApprovalRequested | TurnCompleted | BudgetWarning | CompactionPerformed | ... 10 more ... | SubagentJoined
event
) =>
event: AgentUpdateEmitted | RunStarted | TurnStarted | ModelStarted | ModelRestarted | TextDelta | ReasoningDelta | ToolCallDeclared | ToolCallStarted | ToolProgress | ToolCallSucceeded | ToolCallFailed | ApprovalRequested | TurnCompleted | BudgetWarning | CompactionPerformed | ... 10 more ... | SubagentJoined
event
.
_tag: "AgentUpdateEmitted" | "RunStarted" | "TurnStarted" | "ModelStarted" | "ModelRestarted" | "TextDelta" | "ReasoningDelta" | "ToolCallDeclared" | "ToolCallStarted" | "ToolProgress" | "ToolCallSucceeded" | "ToolCallFailed" | "ApprovalRequested" | "TurnCompleted" | "BudgetWarning" | "CompactionPerformed" | "RunCompleted" | "RunFailed" | "RunInterrupted" | "RunSuspended" | "SubagentRequested" | "SubagentStarted" | "SubagentProgress" | "SubagentCompleted" | "SubagentFailed" | "SubagentInterrupted" | "SubagentJoined"
_tag
=== "AgentUpdateEmitted"),
import Stream
Stream
.
const map: <AgentUpdateEmitted, AgentUpdates.Update>(f: (a: AgentUpdateEmitted, i: number) => AgentUpdates.Update) => <E, R>(self: Stream.Stream<AgentUpdateEmitted, E, R>) => Stream.Stream<AgentUpdates.Update, E, R> (+1 overload)

Transforms the elements of this stream using the supplied function.

Example (Mapping stream values)

import { Effect, Option, Stream } from "effect"
const stream = Stream.fromArray([1, 2, 3]).pipe(Stream.map((n, i) => n + i))
await Effect.runPromise(Stream.runCollect(stream)) // => [1, 3, 5]

@category ― mapping

@since ― 2.0.0

map
((
event: AgentUpdateEmitted
event
) =>
event: AgentUpdateEmitted
event
.
update: AgentUpdates.Update
update
),
);
// Provide the model and runtime services from the getting-started guide.
export const
const observe: Effect.Effect<void, AgentRuntime.AgentRuntimeFailure<Definition<Struct<{
readonly city: String;
readonly area: String;
readonly sources: $Array<Struct<{
readonly url: String;
readonly notes: String;
}>>;
}>, Struct<{
readonly hotels: $Array<String>;
readonly summary: String;
}>, string, Toolkit<{
readonly emit_update: Tool<"emit_update", {
readonly parameters: Struct<{
readonly value: NoInfer<TaggedStruct<"AreaConcern", {
readonly area: String;
readonly finding: String;
readonly sources: $Array<String>;
}>>;
}>;
readonly success: Struct<...>;
readonly failure: typeof AgentUpdates.UpdateError;
readonly failureMode: "return";
}, AgentUpdates.Emitter>;
}>, undefined, undefined, NoInfer<TaggedStruct<...>>> & {
...;
}, never, never>, AgentRuntime.AgentRuntimeRequirements<...>>
observe
=
import AgentUpdates
AgentUpdates
.
const observe: <NoInfer<TaggedStruct<"AreaConcern", {
readonly area: String;
readonly finding: String;
readonly sources: $Array<String>;
}>>, AgentRuntime.AgentRuntimeFailure<Definition<Struct<{
readonly city: String;
readonly area: String;
readonly sources: $Array<Struct<{
readonly url: String;
readonly notes: String;
}>>;
}>, Struct<{
readonly hotels: $Array<String>;
readonly summary: String;
}>, string, Toolkit<{
readonly emit_update: Tool<"emit_update", {
readonly parameters: Struct<...>;
readonly success: Struct<...>;
readonly failure: typeof AgentUpdates.UpdateError;
readonly failureMode: "return";
}, AgentUpdates.Emitter>;
}>, undefined, undefined, NoInfer<TaggedStruct<...>>> & {
...;
}, never, never>, AgentRuntime.AgentRuntimeRequirements<...>>(agent: BoundSource<...>, updates: Stream.Stream<...>) => Stream.Stream<...>

Decode an encoded update stream without hiding its failures or requirements.

observe
(
const HotelResearcher: Definition<Struct<{
readonly city: String;
readonly area: String;
readonly sources: $Array<Struct<{
readonly url: String;
readonly notes: String;
}>>;
}>, Struct<{
readonly hotels: $Array<String>;
readonly summary: String;
}>, string, Toolkit<{
readonly emit_update: Tool<"emit_update", {
readonly parameters: Struct<{
readonly value: NoInfer<TaggedStruct<"AreaConcern", {
readonly area: String;
readonly finding: String;
readonly sources: $Array<String>;
}>>;
}>;
readonly success: Struct<...>;
readonly failure: typeof AgentUpdates.UpdateError;
readonly failureMode: "return";
}, AgentUpdates.Emitter>;
}>, undefined, undefined, NoInfer<TaggedStruct<...>>> & {
...;
}
HotelResearcher
,
const updates: Stream.Stream<AgentUpdates.Update, AgentRuntime.AgentRuntimeFailure<Definition<Struct<{
readonly city: String;
readonly area: String;
readonly sources: $Array<Struct<{
readonly url: String;
readonly notes: String;
}>>;
}>, Struct<{
readonly hotels: $Array<String>;
readonly summary: String;
}>, string, Toolkit<{
readonly emit_update: Tool<"emit_update", {
readonly parameters: Struct<{
readonly value: NoInfer<TaggedStruct<"AreaConcern", {
readonly area: String;
readonly finding: String;
readonly sources: $Array<String>;
}>>;
}>;
readonly success: Struct<...>;
readonly failure: typeof AgentUpdates.UpdateError;
readonly failureMode: "return";
}, AgentUpdates.Emitter>;
}>, undefined, undefined, NoInfer<TaggedStruct<...>>> & {
...;
}, never, never>, AgentRuntime.AgentRuntimeRequirements<...>>
updates
).
Pipeable.pipe<Stream.Stream<{
readonly _tag: "AreaConcern";
readonly area: string;
readonly finding: string;
readonly sources: readonly string[];
}, AgentRuntime.AgentRuntimeFailure<Definition<Struct<{
readonly city: String;
readonly area: String;
readonly sources: $Array<Struct<{
readonly url: String;
readonly notes: String;
}>>;
}>, Struct<{
readonly hotels: $Array<String>;
readonly summary: String;
}>, string, Toolkit<{
readonly emit_update: Tool<"emit_update", {
readonly parameters: Struct<{
readonly value: NoInfer<TaggedStruct<"AreaConcern", {
readonly area: String;
readonly finding: String;
readonly sources: $Array<...>;
}>>;
}>;
readonly success: Struct<...>;
readonly failure: typeof AgentUpdates.UpdateError;
readonly failureMode: "return";
}, AgentUpdates.Emitter>;
}>, undefined, undefined, NoInfer<TaggedStruct<...>>> & {
...;
}, never, never>, AgentRuntime.AgentRuntimeRequirements<...>>, Effect.Effect<...>>(this: Stream.Stream<...>, ab: (_: Stream.Stream<...>) => Effect.Effect<...>): Effect.Effect<...> (+21 overloads)
pipe
(
import Stream
Stream
.
const runForEach: <{
readonly _tag: "AreaConcern";
readonly area: string;
readonly finding: string;
readonly sources: readonly string[];
}, void, never, never>(f: (a: {
readonly _tag: "AreaConcern";
readonly area: string;
readonly finding: string;
readonly sources: readonly string[];
}) => Effect.Effect<void, never, never>) => <E, R>(self: Stream.Stream<{
readonly _tag: "AreaConcern";
readonly area: string;
readonly finding: string;
readonly sources: readonly string[];
}, E, R>) => Effect.Effect<void, E, R> (+1 overload)

Runs the provided effectful callback for each element of the stream.

Example (Running an effect for each value)

import { Effect, Stream } from "effect"
const stream = Stream.make(1, 2, 3)
const values: Array<string> = []
const program = Effect.gen(function*() {
yield* Stream.runForEach(stream, (n) => Effect.sync(() => values.push(`Processing: ${n}`)))
})
await Effect.runPromise(program)
values // => ["Processing: 1", "Processing: 2", "Processing: 3"]

@category ― destructors

@since ― 2.0.0

runForEach
((
finding: {
readonly _tag: "AreaConcern";
readonly area: string;
readonly finding: string;
readonly sources: readonly string[];
}
finding
) =>
import Effect
Effect
.
const log: (...message: ReadonlyArray<any>) => Effect.Effect<void>

Logs one or more messages using the default log level.

Example (Logging at the default level)

import { Effect, Logger } from "effect"
const output: Array<unknown> = []
const program = Effect.gen(function*() {
const result = 2 + 2
yield* Effect.log("Result:", result)
return result
})
const logger = Logger.make<unknown, void>(({ logLevel, message }) => {
void output.push(`${logLevel}: ${Array.isArray(message) ? message.map(String).join(" ") : String(message)}`)
})
const runnable = Effect.provide(program, Logger.layer([logger]))
void output.push(Effect.runSync(runnable))
output // => ["Info: Result: 4", 4]

@category ― logging

@since ― 2.0.0

log
(
finding: {
readonly _tag: "AreaConcern";
readonly area: string;
readonly finding: string;
readonly sources: readonly string[];
}
finding
.
area: string
area
,
finding: {
readonly _tag: "AreaConcern";
readonly area: string;
readonly finding: string;
readonly sources: readonly string[];
}
finding
.
finding: string
finding
)),
);

The researcher definition declares the finding Schema separately from its final output. AgentUpdateEmitted contains an accepted, encoded update; AgentUpdates.observe decodes it through that Schema, including any required decoding services. Emitting a finding lets the Agent continue working and does not complete the run.

Application tools can also emit findings with AgentUpdates.emit. See update delivery guarantees for stable keys, limits, and automatic delivery from a background worker to its parent.

A voice adapter can delegate to the same agent and Thread as a text interface. Keep its media Scope separate from accepted durable work: closing a call or stopping playback closes media and observation, while the durable runtime retains its accepted-work obligation. Use the original idempotency key and frozen input to reconcile uncertain admission. A transcript delta is context, not an instruction to admit another Run. Corrections use ordinary queued input and steering.

The travel planner demonstrates GPT-Live client delegation. Its adapter constructs schema-validated planner requests from attributed transcripts, uses the existing planner admission path, and reconciles receipts against canonical settlement. It tracks corrections by work identity rather than treating every caption as a new task. Spoken and typed input share the demo conversation; attributed speech context is separate from the user’s work request and visible message. Typed results and later research answers return to the active voice exchange. Reconnect restores relevant conversation history without restarting accepted work. The demo retains undelegated speech in the current tab; it does not add a durable partial-transcript journal.

For application-selected previews, decorate the native Effect AI LanguageModel service in the model Layer. Its streamText exposes ordered text-delta, tool-params-start, and tool-params-delta parts. The demo selects native text and only the designated deliver_response.message field; it retains a bounded provisional preview and fences writes by Submission and Attempt. Parsing and presentation belong to the adapter. Never forward reasoning, arbitrary tool arguments, or diagnostics to a voice provider. Provider protocol interception is unnecessary for this public-output path, and no additional SDK output hook is required.

Keep generation, schema validation, durable settlement, provider acknowledgment, and actual audio playback distinct. A partial completion-tool argument is provisional even when it resembles a complete sentence. A provider acknowledgment does not prove that the user heard the result. Use canonical, schema-decoded output for final answers, including after reconnect. The official Live delegation guide describes the provider-specific half of this integration.

const program = Effect.gen(function* () {
const detached = yield* AgentRuntime.start(agent, input);
const result = yield* detached.await;
const completeTrace = yield* detached.events;
return { result, completeTrace };
}).pipe(Effect.scoped);

start requires a caller Scope. observe replays prior events, follows new events, and ends when the run settles. events returns the complete replay after settlement. Execution resources close before await returns, while replay remains available until the owner closes.

Observers cannot backpressure execution. Closing the owner interrupts active work and observers. The handle remains process-local and never creates a daemon fiber.

Use a Cloudflare thread object to accept work that survives eviction. For a process with SQLite, use the Node host.

Platform hosts assemble storage and runtime services for you. When building your own host, DurableAgentRuntime.layer supplies default prompt preparation and tool authorization. Use layerWithServices to supply your own service layers. It requires RunToolAuthorization and captures RunContextPreparation when provided.

Here is the default authorization policy; replace it with your application’s implementation:

import {
class RunToolAuthorization

Host action-time authority for native and programmatic application Tools. Implementations close over their dependencies at Layer construction and return a denial when execution is not authorized. A dependency or validation failure instead fails with AgentToolAuthorizationCheckError and retains its original Cause. Defects and interruption remain in the Effect Cause channel. Durable coordinators capture this service once and retain it across replacement Attempts. Ephemeral Runs also resolve this service at their Run boundary. A typed per-run RunOptions.toolAuthorization overrides it while retaining its own error and requirement channel.

RunToolAuthorization
} from "@yielded/agent/run-options";
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 Layer
Layer
} from "effect";
export const
const RuntimeLive: Layer.Layer<DurableAgentRuntime, never, SubmissionLedger | ThreadStore | RunStorage | WakeScheduler | DurableRuntimeFailpoint | DurableRuntimeConfig | ToolReconciler | Crypto>
RuntimeLive
=
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
.
DurableAgentRuntime.layerWithServices: Layer.Layer<DurableAgentRuntime, never, RunToolAuthorization | SubmissionLedger | ThreadStore | RunStorage | WakeScheduler | DurableRuntimeFailpoint | DurableRuntimeConfig | ToolReconciler | Crypto>

Captures optional context preparation; Tool authorization remains required in R.

layerWithServices
.
Pipeable.pipe<Layer.Layer<DurableAgentRuntime, never, RunToolAuthorization | SubmissionLedger | ThreadStore | RunStorage | WakeScheduler | DurableRuntimeFailpoint | DurableRuntimeConfig | ToolReconciler | Crypto>, Layer.Layer<DurableAgentRuntime, never, SubmissionLedger | ThreadStore | RunStorage | ... 4 more ... | Crypto>>(this: Layer.Layer<...>, ab: (_: Layer.Layer<...>) => Layer.Layer<...>): Layer.Layer<...> (+21 overloads)
pipe
(
import Layer
Layer
.
const provide: <never, never, RunToolAuthorization>(that: Layer.Layer<RunToolAuthorization, never, never>) => <RIn2, E2, ROut2>(self: Layer.Layer<ROut2, E2, RIn2>) => Layer.Layer<ROut2, E2, Exclude<RIn2, RunToolAuthorization>> (+3 overloads)

Feeds the output services of the dependency layer into the requirements of this layer, returning a layer that only provides the services from this layer.

When to use

Use when you need to hide an implementation dependency layer from callers.

Details

In serviceLayer.pipe(Layer.provide(dependencyLayer)), the dependency layer is built first and is used to satisfy the requirements of serviceLayer.

Example (Providing layer dependencies)

import { Context, Effect, Layer } from "effect"
class Database extends Context.Service<Database, {
readonly query: (sql: string) => Effect.Effect<string>
}>()("Database") {}
class UserService extends Context.Service<UserService, {
readonly getUser: (id: string) => Effect.Effect<{
id: string
name: string
}>
}>()("UserService") {}
class Logger extends Context.Service<Logger, {
readonly log: (msg: string) => Effect.Effect<void>
}>()("Logger") {}
// Create dependency layers
const databaseLayer = Layer.succeed(Database, {
query: Effect.fn("Database.query")((sql: string) => Effect.succeed(`DB: ${sql}`))
})
const logs: Array<string> = []
const loggerLayer = Layer.succeed(Logger, {
log: Effect.fn("Logger.log")((msg: string) => Effect.sync(() => logs.push(`[LOG] ${msg}`)))
})
// UserService depends on Database and Logger
const userServiceLayer = Layer.effect(UserService, Effect.gen(function*() {
const database = yield* Database
const logger = yield* Logger
return {
getUser: Effect.fn("UserService.getUser")(function*(id: string) {
yield* logger.log(`Looking up user ${id}`)
const result = yield* database.query(
`SELECT * FROM users WHERE id = ${id}`
)
return { id, name: result }
})
}
}))
// Provide dependencies to UserService layer
const userServiceWithDependencies = userServiceLayer.pipe(
Layer.provide(Layer.mergeAll(databaseLayer, loggerLayer))
)
// Now UserService layer has no dependencies
const program = Effect.gen(function*() {
const userService = yield* UserService
return yield* userService.getUser("123")
}).pipe(
Effect.provide(userServiceWithDependencies)
)
Effect.runSync(program) // => { id: "123", name: "DB: SELECT * FROM users WHERE id = 123" }
logs // => ["[LOG] Looking up user 123"]

@see ― provideMerge for retaining the dependency services

@category ― providing services

@since ― 2.0.0

provide
(
class RunToolAuthorization

Host action-time authority for native and programmatic application Tools. Implementations close over their dependencies at Layer construction and return a denial when execution is not authorized. A dependency or validation failure instead fails with AgentToolAuthorizationCheckError and retains its original Cause. Defects and interruption remain in the Effect Cause channel. Durable coordinators capture this service once and retain it across replacement Attempts. Ephemeral Runs also resolve this service at their Run boundary. A typed per-run RunOptions.toolAuthorization overrides it while retaining its own error and requirement channel.

RunToolAuthorization
.
RunToolAuthorization.allowAll: Layer.Layer<RunToolAuthorization, never, never>

Explicit compatibility policy: Tool execution requires no additional host authorization.

allowAll
),
);

This layer still requires SubmissionLedger, ThreadStore, RunStorage, WakeScheduler, DurableRuntimeFailpoint, DurableRuntimeConfig, ToolReconciler, and Crypto.Crypto. Build RunStorage.layer() from the same adapter’s ThreadStore, SubmissionLedger, and SettlementPublisher. The publisher checks authority and appends the canonical settlement in one storage transaction; there is no fallback for independently supplied stores. Provide these services before acquiring the runtime.

The runtime captures its services at acquisition. Supplying a different layer around a later worker call does not replace them. Acquire service dependencies in their layers and keep them alive for the runtime’s Scope. Durable service hooks must have no unresolved dependencies. Preparation failures retain their AgentInputError, MemoryRecallError, or CompactionError tags; RunContextPreparationError is their type union, not a wrapper. Durable execution records failed Runs in Settlements with bounded, structured causal diagnostics and execution correlation; it does not reconstitute the original error object from storage. Programmatic Worker inspection and awaiting expose this evidence as diagnostic, while generated inspect tools and completion reports omit it. Authorization returns an allowed or denied decision, or fails with AgentToolAuthorizationCheckError when its check could not finish. Configure prompt preparation and tool authorization in their respective services.

Each turn follows this sequence:

flowchart TD
  accTitle: Agent turn sequence
  accDescr: A turn prepares context, streams and reduces one model response, decodes the complete tool batch, executes bounded handlers, commits results in declaration order, drains steering, evaluates the stop policy, and drains follow-up only if otherwise complete.
  context["prepare context"] --> model["stream and reduce one model response"]
  model --> batch["decode the complete tool batch"]
  batch --> handlers["execute bounded tool handlers"]
  handlers --> commit["commit results in declaration order"]
  commit --> steering["drain steering"]
  steering --> stop["evaluate stop policy"]
  stop --> followup["drain follow-up only if otherwise complete"]

run and stream use the same loop. Completion tools also drain steering after their results commit: new input continues the active run at its next model turn. When the completing run has exhausted its budget, follow-ups stay queued for a new run. Durable runs drain the ready input prefix together, subject to the host ledger’s joining policy and the runtime batch bound. Each joined input retains its receipt; only a later host response covers it. Inputs rejected by the prompt callback remain queued for their own run without cancelling the host’s model call. Recovery restores all previously consumed, uncovered joins before the next call; the batch bound applies to newly ready inputs.

Set policy: { restartOnJoinedInput: true } on an agent to let eligible joined input replace a running model call. The first call starts immediately. The runtime cancels only the disposable model stream, before any response commits or application tool starts, and restarts with the combined input. It permits two replacements per run, including across durable recovery; later joins use the ordinary seams. Calls exposing hosted tools without Tool.Readonly keep seam steering because remote execution may precede streamed evidence. Hosted web search and file search are read-only; they retain restart. Joining authority, receipts, and settlement are unchanged.

Streaming clients must clear text and reasoning drafts for the turnId in a ModelRestarted event. Its reason is joined-input; the replacement has a new turn ID. Durable hosts retain the replacement count and any reported usage before the next call. Cancelled chat <model> spans end interrupted with effect_agent.model.outcome=aborted and effect_agent.model.abort_reason=joined-input. Unreported usage remains unknown. Custom input adapters supply RunInputHook.awaitJoin to enable the same behavior; durable adapters also implement RunDurabilityHook.commitModelRestart.

RunOptions accepts per-run capability hooks. This process-local example uses in-memory history. history provides an initial Prompt, and onHistory receives incremental updates.

const options: RunOptions<AppError, AppRequirements> = {
threadId,
history,
input: toRunInputHook(commands),
approval: toRunApprovalHook(approvalPolicy),
budget: toRunBudgetHook(budget),
context: toRunContextHook(contextTransform),
scheduling: toRunSchedulingHook({ mode: "bounded", concurrency: 2 }),
onHistory,
};

Hook errors join the run error channel, and their services join R. onHistory runs inline. Writes completed before a later failure or interruption remain caller-owned. Persistent history rejects competing history and input queue hooks before model or tool execution.

Approval predicates, parameter decoding, and approval.request share the run’s duration deadline. Hooks that retain canonical approval facts use prepareBatch for requests and commit for the prepared decision; these storage commits finish outside the preparation timer before execution.

Pass prompt preparation as context and tool authorization as toolAuthorization when needed. Ephemeral runs read these options; providing the durable service layers alone does not install per-run hooks.

A tool may fail and the model may still complete the run. Install toolFailureObserverLayer from @yielded/agent to report such failures.

This observer covers failures contained as results, including programmatic broker outcomes. It does not duplicate model-declared failures that propagate through the run’s Effect error channel, or defects and interruptions. Native model-call argument rejections do not invoke it because no handler started; their failed results and warning telemetry remain available. Use ToolCallFailed.failureHandling and tool telemetry to distinguish returned failures from propagated ones, and handle the run’s Effect exit separately.

import { toolFailureObserverLayer } from "@yielded/agent/run-options";
import { Effect, ErrorReporter } from "effect";
const failureReporting = toolFailureObserverLayer({
observe: (observation) =>
observation.cause === undefined ? Effect.void : ErrorReporter.report(observation.cause),
});
const program = AgentRuntime.run(agent, input).pipe(Effect.provide(failureReporting));

The engine does not forward observations to ErrorReporter by itself. Choose what to record and redact. The observer runs inline at most once per in-memory Attempt. Replacement Attempts may repeat an observation. Nothing here is serialized into thread history.

Observer defects cannot change the tool result, though a slow observer holds a tool permit. Avoid calling the broker, running another agent, or interrupting the observer itself. Durable hosts accept the same observer through their platform options. Call-local telemetry and any applicable failure observer finish before the terminal Tool event is published, so stopping observation at that event does not skip them.

A run Scope owns its model stream, tool fibers, and run-local resources. Closing it interrupts children and runs finalizers. run completes cleanup before returning; stream closes its resources when consumption completes, fails, or is interrupted. Services from an enclosing application layer remain available to other runs until the application Scope closes. Retrying a whole run can repeat external effects; use durable recovery when work must survive interruption without automatically replaying uncertain ordinary tools.

Wrap several runs with one Effect.provide(AppLive) to reuse shared services. Keep caller scoping for start, explicit resource acquisition, and any operation that requires Scope.

History retention waits for run-local cleanup, result validation, and commit before publishing RunCompleted. Interrupting a waiter for durable accepted work only detaches that waiter. Abort a durable Submission with an explicit persisted command. See Persistence & durability.

With an Effect tracer installed, filter gen_ai.operation.name to find agent work:

Operation Span name Identity
invoke_agent invoke_agent <agent definition ID> Agent, Thread, Run
chat chat <configured model name> Agent, Thread, Run, and Turn for ordinary model calls
execute_tool execute_tool <tool name> Agent, Thread, Run, Turn, Tool Call

All three carry gen_ai.agent.name (the definition ID), gen_ai.agent.id (the Thread-backed instance), and gen_ai.conversation.id (the Thread ID). Existing agentId, threadId, runId, and applicable turnId attributes remain available. The invoke_agent, chat, and execute_tool spans retain their identity and outcome attributes. Agent span names replace AgentRuntime.run; update filters using that old name.

Successful tool executions log at Debug; failures log at Warning. The default logger omits successful tool logs. Tool spans retain their identity and outcome attributes at either log level.

Model calls label the existing Effect AI LanguageModel.streamText span rather than creating a second model-call span. The configured model and provider are recorded as gen_ai.request.model and gen_ai.provider.name; native providers retain their response and token-usage annotations. Each retry and compaction summary has its own model span. These labels add identifiers, not prompts, instructions, or tool payloads.

Library generators use Effect.fnUntraced; direct Effect-returning helpers need no wrapper. Spans are reserved for agent run/turn, model and tool calls, browser and transport operations, adapter storage operations, and recovery. The export check enforces their owning modules and names across Effect.fn, Effect and Stream span primitives, and the storage span wrapper. Helpers for stream parts, individual records, per-tool batch orchestration, and decoded rows are always untraced; actual tool calls retain their execute_tool spans. Use the enclosing operation spans for timing and failure diagnostics, and adjust filters that relied on private-helper span names.

An agent span covers one active execution Scope. A durable Run resumed by another Attempt can produce another span with the same Run ID; the span is not the entire wall-clock lifetime of a suspended Run. See Cloudflare tracing for the dashboard setup and alarm-root behavior.

AgentRuntime results and RunCompleted expose optional usage (own calls, including compaction) and delegatedUsage (all attached descendants). Combine them once with Usage.sumRunTotals; background workers are excluded. Child events repeat cumulative totals, so deduplicate by Run ID.

For failed or interrupted Runs, read handle.usageReport after handle.await settles or its owning Scope closes. Earlier reads are live snapshots and may omit in-flight usage. Durable settlements retain their richer per-model usageSummary for the Run’s own calls.

A cost estimator receives the configured binding name in request.model and the actual provider-reported identity in request.response. Use the latter for response-sensitive pricing. request.webSearchCalls counts observed hosted web-search calls for adding the provider search fee; Run totals and durable per-model summaries retain webSearchCalls. It excludes OpenAI page/find actions and unobserved work, and legacy records may omit it. request.finishMetadata carries native Effect AI finish metadata only during estimation; the engine never persists provider HTTP details or raw metadata in accounting records. A summarizer uses the same estimator with purpose: "summary".

Calls and Run totals retain usageStatus and pricingStatus. Missing legacy status is unknown, and a numeric zero without an estimate is not evidence of free execution. Run summaries distinguish complete, partial, and unknown coverage; unobservedModelCalls counts observed calls without retained accounting, which are excluded from numeric token and call totals. A canonical ModelResponseInterrupted can indicate additional unquantified provider work beyond this count.

An explicitly configured costBudgetMicrousd fails with a typed cost-policy error when the estimator reports unknown pricing, after retaining the call’s usage. Uncapped Runs may continue; legacy numeric estimates and estimates without a status remain trusted host estimates.

Response records also retain each Turn’s missing-call count, so approval and child suspension preserve incomplete accounting when a fresh runtime resumes the Run.

Headless runtime hosts use AgentRuntime.executeWithUsageAccountingUnknown with the inward ModelUsageAccounting and AgentUpdateAcceptance services from RunOptions. These dependencies remain visible in R. The native durable runtime supplies its canonical Turn accumulator and Attempt-bound update acceptance at composition, retaining updates before acknowledging them. Hosts that need public progress can use streamWithUsageAccountingUnknown with the same services. Ordinary stream, run, and start calls provide ephemeral accounting and update acceptance.

Canonical response records own committed per-call usage. Terminal settlement uncommittedModelUsage retains only staged calls not already present in a response record; its charges are already included in usageSummary, so do not add them a second time. An isolate loss before either response or settlement commit cannot prove the lost call’s usage or cost.