@lunora/agent is Lunora's durable-agent primitive. defineAgent declares a
model + tools + memory; the definition compiles onto a Cloudflare Workflow
where each LLM turn and each tool call is a named durable step, and every
message persists idempotently into DO SQLite thread tables. The client
subscribes to the thread query and watches the conversation stream in live, so
there is no new transport and no HTTP streaming infrastructure to run.
pnpm add @lunora/agentDeclaring an agent
// lunora/agents.ts
import { defineAgent, defineAgentTool } from "@lunora/agent";
import { jsonSchema } from "@lunora/ai";
export const support = defineAgent({
instructions: "You are a helpful support agent.",
maxTurns: 8,
model: "@cf/meta/llama-3.3-70b-instruct-fp8-fast", // Workers AI id, AI SDK model, or (env) => model
tools: {
getWeather: defineAgentTool({
description: "Look up the current weather for a city.",
execute: async ({ city }, { idempotencyKey, run, threadKey }) => {
// `run` dispatches any Lunora function; `idempotencyKey` is the
// durable step name — dedupe side effects on it.
return run({ __lunoraRef: "weather:lookup" }, { city });
},
inputSchema: jsonSchema({ properties: { city: { type: "string" } }, required: ["city"], type: "object" }),
}),
},
});Declaring an agent is enough: codegen auto-registers the runtime functions the
durable loop persists through (agents:agentAppendMessage, agents:agentEnsureThread,
agents:agentPatchThread, agents:agentResolveApproval) plus the public thread
queries (agents:agentMessages, agents:agentThread), and wires the typed
ctx.agents.<name> producer.
Merge the thread tables into your schema. They auto-prefix to
agent_threads / agent_messages and can never collide with app tables:
// lunora/schema.ts
import { agentExtension } from "@lunora/agent";
export default defineSchema({/* your tables */}).extend(agentExtension);Scaffold a new agent with the generator (appends to lunora/agents.ts,
creating it if missing, and wires the worker entry):
vis generate lunora-agent --name=supportRunning + observing
// inside a mutation or action:
const { id } = await ctx.agents.support.run({ input: message, owner: ctx.auth.userId, threadKey, title: "Support chat" });
// client: subscribe to the thread — user turns, tool calls, tool results and
// assistant replies stream in as they persist (reactive subscription, so it
// survives reconnects and multiple observers for free):
useSubscription(api.agents.agentMessages, { key: threadKey });Reuse the same threadKey to continue a conversation; the loop reads the
persisted history each turn.
ctx.agents.<name> is the typed producer surface on Mutation/Action ctx:
| Method | Does |
|---|---|
run(params) | Start (or continue) a run; returns { id } (the Workflow instance id). |
cancel(id) | Terminate the in-flight instance and mark the thread "cancelled". |
status(id) | Read the instance's current status. |
sendEvent(id, evt) | Deliver an event to a hibernated run (the primitive HITL approvals build on). |
Loop control
defineAgent exposes the AI SDK's loop-shaping seams, evaluated per turn inside
the durable loop:
maxTurns: hard cap on LLM turns. Hitting it stops the run withresult.stopped === "maxTurns".stopWhen: a stop condition (e.g.stopWhen: hasToolCall("finalize")); stopping this way reportsresult.stopped === "stopCondition". A clean final answer (no pending tool calls) reportsresult.stopped === "final".prepareStep: mutate the request just before a turn (swap the model, trim messages, force a tool). Must be deterministic, because it runs on replay.repairToolCall: repair a malformed tool call the model emits (AI SDKexperimental_repairToolCall): given{ toolCall, error, tools, inputSchema }, return a corrected call ornullto give up. Runs inside the turn, so keep it deterministic.compaction: automatic history compaction. Setcompaction: { maxMessages }and once the thread history exceedsmaxMessages, the loop summarizes the older messages (all but the most recentkeepRecent, defaultceil(maxMessages / 2)) into one system-message brief and prompts the model with that brief plus the recent tail, so context stays bounded as a conversation grows. The summary is produced inside the turn's memoized durable step (replay-safe) by a cheapercompaction.modelif set, else the agent's model;prepareStepstill runs after and can override further. Absent, the full history is sent every turn.output: aFlexibleSchema(zod orjsonSchema(...)) for a structured final answer; the typed value is surfaced onresult.output.onStepFinish: a callback after each turn (observability/side channels).usage: token usage is accumulated across turns and surfaced on the run result and persisted onto the thread.telemetry: an ai@7TelemetryOptionshanded to every turn; settelemetry.integrationsto trace turns and tool calls. The@lunora/agent/telemetrysubpath ships ready-made ones:consoleTelemetry(zero-dependency),combineTelemetry,otlpTelemetry(shipsgen_ai.*spans over OTLP to any collector), and the dependency-injectedsentryTelemetry/braintrustTelemetrybridges (privacy-safe by default, with no prompt/output recorded without opt-in).
import { hasToolCall } from "@lunora/ai";
export const support = defineAgent({
maxTurns: 6,
model: "@cf/meta/llama-3.3-70b-instruct-fp8-fast",
output: jsonSchema({ properties: { answer: { type: "string" }, confidence: { type: "number" } }, required: ["answer"], type: "object" }),
prepareStep: async ({ stepNumber }) => (stepNumber === 0 ? { toolChoice: "required" } : {}),
stopWhen: hasToolCall("finalize"),
});Telemetry integrations
The @lunora/agent/telemetry subpath ships ready-made Telemetry integrations
for the telemetry.integrations array; each traces every LLM turn and tool call
in the durable loop. They are privacy-safe by default: recordInputs and
recordOutputs both default to false, so no prompt, message, tool argument,
generated text, or tool result is recorded without an explicit opt-in. Only
structural metadata is recorded (model id, finish reason, token counts, tool
name, timing, success/failure).
consoleTelemetry: a zero-dependency structured console tracer.combineTelemetry: fan the lifecycle out to several integrations and nest their execution wrappers (executeLanguageModelCall,executeTool) right-to-left, so the first integration is outermost.sentryTelemetry/braintrustTelemetry: dependency-injected bridges. You pass your own initialized Sentry namespace / Braintrust logger, so the heavy vendor SDK is never imported by Lunora. Sentry turns model calls and tool executions into spans (op: "gen_ai.generate"/"gen_ai.execute_tool") and routes errors toSentry.captureException; Braintrust logstype: "llm"/type: "tool"logger.tracedspans.otlpTelemetry: ships agen_ai.*generation span per model turn (and a span per tool call) over OTLP-over-HTTP to any collector, the Lunora cloud or your own, so agent generations land in the same trace store as the rest of your app's telemetry. Zero SDK dependency (it justfetches). Pass{ endpoint, token }(read from the injectedLUNORA_OTLP_ENDPOINT/LUNORA_OTLP_TOKENon the platform) and, optionally, atraceIdto group one run's spans into a single trace. Two deliberate differences from the SDK bridges: noonError(a failed span already carriesstatus.code === 2) and flat spans (OTLP has no ambient parent to nest under).
import { defineAgent } from "@lunora/agent";
import { combineTelemetry, consoleTelemetry, otlpTelemetry, sentryTelemetry } from "@lunora/agent/telemetry";
import * as Sentry from "@sentry/cloudflare";
export const support = defineAgent({
name: "support",
telemetry: {
isEnabled: true,
integrations: [combineTelemetry(otlpTelemetry({ endpoint: env.LUNORA_OTLP_ENDPOINT, token: env.LUNORA_OTLP_TOKEN }), consoleTelemetry())],
},
// ...model, tools, loop control
});consoleTelemetry, combineTelemetry, and otlpTelemetry have zero runtime
dependencies; install @sentry/cloudflare / braintrust only if you use those
bridges.
Concurrency + cancellation
Only one run owns a thread at a time: two concurrent runs on the same
threadKey would interleave their messages into the shared per-thread sequence
counter. The onConcurrentRun policy decides what happens when a run starts on
a thread that already has a different instance in flight:
"reject"(default): fail the new run fast with aCONFLICTerror."replace": terminate the in-flight instance and take the thread over."queue": park the new run behind the one in flight (FIFO, up to five deep) and hibernate it until the thread is handed over. Each parked run is a live Workflow instance waiting on an event, so the depth cap is a real resource bound — a sixth start is rejected exactly as"reject"would be, and a parked run waits at most 12 hours for its turn. A"replace"arriving later supersedes the run in flight, not the queue. A dispatch with no instance id (the inbound-email / inbound-channel paths) cannot be parked — nothing later could tell it apart from another such dispatch to wake it — so"queue"rejects those.
A workflow replay re-enters under the same instance id and is never treated as a concurrent run.
Cancel from the server with ctx.agents.<name>.cancel(instanceId), or from the
client through the framework hooks (below). cancel() terminates the in-flight
instance and moves the thread to "cancelled".
Human-in-the-loop approvals
Gate a side-effecting tool behind a human decision with needsApproval (a
boolean, or a per-input predicate; this mirrors the AI SDK's needsApproval):
chargeCard: defineAgentTool({
description: "Charge the customer's saved card.",
execute: async ({ amount }, { run }) => run({ __lunoraRef: "billing:charge" }, { amount }),
inputSchema: jsonSchema({ properties: { amount: { type: "number" } }, required: ["amount"], type: "object" }),
needsApproval: ({ amount }) => amount > 5000,
}),When the gate resolves truthy the run pauses: the thread moves to
"awaiting_input" and the workflow hibernates on approval:<toolCallId>, so
no compute is billed while it waits. A client resolves it via the
auto-registered agents:agentResolveApproval (the framework approve/reject
helpers dispatch it). On approve the tool runs exactly as normal; on reject it
is skipped and a tool result explaining the rejection is persisted so the next
turn recovers. Because needsApproval runs on replay, keep it deterministic
(no Date.now()/Math.random()).
An unanswered approval is not pending forever. The wait carries a timeout —
approvalTimeout on the agent config, default "3 days", clamped to one
week — and when it elapses the call is recorded as rejected ("approval
timed out") and the run continues down the normal rejection path. It has to: a
wait that outlived the thread's abandoned-run horizon would let a new run
reclaim the thread while the approval was still pending, so the approval could
never be answered at all. Size approvalTimeout to your on-call reality, and do
not assume an approval raised before a long weekend is still waiting on
Monday.
Client hooks
Every framework adapter ships two hooks over the same live subscription. The
generated api drives them; you pass the run/send and cancel mutation
references your app exposes.
useAgent(useAgent/createAgent/agent) is the thin driver:{ run, cancel, pending, status, thread }.run(input, args?)dispatches the run mutation with{ input, threadKey }merged overrunArgs;cancel()terminates the in-flight instance (a no-op when nothing is running or no cancel ref was supplied);status/threadflow live fromagents:agentThread.useAgentChatis the batteries-included chat surface:{ send, approve, reject, cancel, messages, status, streamingText }.sendappends an optimistic user turn then starts/continues the run;approve/rejectresolve a pending tool approval;messagesis the live thread.
// React
const chat = useAgentChat({ api, cancel: api.chat.cancelRun, send: api.chat.startRun, threadKey });
chat.send("Refund my last order");
if (chat.status === "awaiting_input") chat.approve(toolCallId);The same hooks exist as Vue composables (useAgent/useAgentChat), Solid
primitives (createAgent/createAgentChat), and Svelte store factories
(agent/agentChat).
Token streaming. The turn seam has a streaming counterpart (streamText
via createStreamGenerate) and the hooks expose a stream option feeding
useAgentChat's live streamingText. The server-side token sink
(onTokenDelta → transport) is a wired-but-dormant seam: until it is threaded
onto a live transport the loop takes the byte-identical non-streaming path, so
streamingText stays empty and messages still arrive per-turn over the thread
subscription. Turn-granular streaming works today; sub-turn token streaming is
the remaining follow-up.
Synced state
Beyond the message log, a run carries a small synced state object: a
setState-style scratchpad (the plan, a step counter, a progress summary) that
a tool writes and a client watches live. Seed it with initialState on the
agent (set once, on thread creation, where first writer wins so a replay never
re-seeds), then read and write it from any tool's ctx:
const planner = defineAgent({
initialState: { applied: [], plan: [], step: 0 },
model: "@cf/meta/llama-3.1-8b-instruct",
tools: {
advance: defineAgentTool({
description: "Record the next step of the plan.",
execute: async ({ note }, { getState, idempotencyKey, setState }) => {
const current = (await getState()) ?? { applied: [], plan: [], step: 0 };
const applied = current.applied as string[];
// Idempotent read-modify-write: skip if this durable step already
// applied. A retried step re-reads the ALREADY-written state, so
// dedupe on `idempotencyKey` or the counter double-advances.
if (applied.includes(idempotencyKey)) {
return "recorded";
}
await setState({
applied: [...applied, idempotencyKey],
plan: [...(current.plan as string[]), note],
step: (current.step as number) + 1,
});
return "recorded";
},
inputSchema: jsonSchema({ properties: { note: { type: "string" } }, required: ["note"], type: "object" }),
}),
},
});setState(state) is an absolute replace (not a patch). The value you pass
must be replay-stable: a constant or derived purely from the tool's
input, never from Date.now() / Math.random(). A step that fails mid-body is
retried at-least-once and re-runs the whole execute against state a prior
attempt may already have written, so re-applying a replay-stable value is a
no-op and a replay converges. A value derived from getState() is not
replay-stable: a naive read-modify-write (setState({ step: (await getState()).step + 1 })) double-advances on a retry because the retry
re-reads the already-written value; dedupe it on idempotencyKey as above.
getState() reads the thread's current state through the owner-gated
agents:agentState query, the same identity gate as the message reads, so only
the thread's owner sees it.
On the client, subscribe with useAgentState, a thin wrapper over the
agents:agentState subscription that re-renders whenever a tool calls
setState (a dedicated query with a per-socket JSON memo suppresses no-op
pushes on unrelated thread writes):
// React — generic over your app's state shape
const { state, error } = useAgentState<PlanState>({ api, threadKey });
// state?.step, state?.plan — live, undefined until first seeded/pusheduseAgentState ships in @lunora/react today; the Vue/Solid/Svelte
equivalents are a follow-up (the server surface and the subscription are already
in place). State sync is opt-in: an agent without initialState and tools
that never touch getState/setState behaves exactly as before, and the
state column stays absent on its threads.
Tools from anywhere
Beyond inline defineAgentTool, three helpers pull tools from other surfaces.
Each returns an AgentToolDefinition you drop into the tools map:
functionTool: expose a Lunora function as a tool.executedispatches the referenced function;inputSchemashould mirror its argument validator.agentAsTool: adapt a declared agent into a tool so a supervisor can delegate to specialists. It starts a child run on the child agent's Workflow binding, polls to completion, and returns the child's final answer. ChildthreadKey+ instance id derive from the parent's (replay-stable), so a retried step reuses the same child run (idempotent). Options:name(child export name →AGENT_<NAME>binding),description,maxPolls(default 600),pollIntervalMs(default 500) — together a five-minute wall-clock budget, past which the child is terminated and the parent is told it did not finish. The child thread inherits the parent run'sowner, so it stays under the same RLS scope. A child that runs out of turns reports its turn cap rather than returning an empty answer.mcpTools: adapt the tools an MCP server lists. Pass a pre-connectedclient(the SSE/HTTP transport cannot run in workerd), optionally filtered withonlyand namespaced withprefix.
import { agentAsTool, functionTool, mcpTools } from "@lunora/agent";
tools: {
lookupOrder: functionTool("orders:byId", { description: "Look up an order by id.", inputSchema }),
research: agentAsTool({ description: "Delegate deep research.", name: "researcher" }),
...mcpTools({ client: docsClient, prefix: "docs_" }),
}Sandbox tools (batteries-included browser + containers + filesystem)
These tools give an agent a headless browser, a sandboxed container, and
an R2-backed filesystem. Import them from @lunora/agent (or the
@lunora/agent/sandbox subpath) and drop them into tools; codegen does the
rest:
import { browserTool, containerTool, defineAgent, fsTool } from "@lunora/agent";
export const operator = defineAgent({
model: "@cf/meta/llama-3.3-70b-instruct-fp8-fast",
tools: {
// One tool, every browser op — the model picks via `op`. `bucket` is
// required for `screenshot`/`pdf`: the render goes to R2 and the model
// gets its key back (see below).
browser: browserTool({ bucket: "SANDBOX_BUCKET", root: "renders" }),
// Talks to `ctx.containers.sandbox` (a `lunora/containers.ts` export).
sandbox: containerTool("sandbox"),
// A persistent R2-backed filesystem scoped to `agents/operator`.
fs: fsTool("SANDBOX_BUCKET", { root: "agents/operator" }),
},
});-
browserTool(opts?)exposes Cloudflare Browser Rendering as one tool. The model setsoptoscreenshot|pdf|content|scrapeplus aurl;content/scrapereturn HTML.screenshot/pdfwrite the render toopts.bucketand return{ bytes, key, mediaType }— not the bytes: a tool result is capped at 4000 characters before it is persisted, and base64 of even a small PNG is several times that, so inline bytes arrived truncated mid-string after the render had already been billed. Each render lands at<root>/<threadKey>/<toolCallId>.<ext>, derived only from replay-stable identifiers, so a retried step overwrites its own object instead of leaking one per attempt. Without abucketthose two ops are refused as a tool result and never dispatched, so nothing is billed. The model picks theurlwith no allowlist, so this is an SSRF surface: a prompt-injected model can aim it at an internal/link-local endpoint. It runs unattended by default; passopts.needsApproval(a boolean or a predicate) to gate it. -
containerTool(name, opts?)talks to the declared container atctx.containers.<name>.op: "fetch"sends an HTTP request (path,method,body);op: "exec"runs a command (routed as aPOST /__lunora/execthe container app serves). Command execution is gated behind a human approval by default: anexec, and anyfetchusing a non-idempotent method (POST/PUT/PATCH/DELETE), since a prompt-injected model could otherwise reach some other mutating route unattended. Afetchcannot reach the exec route at all —@lunora/containerreserves/__lunora/*on the handle — so there is no path spelling that runs a command around the gate. AGET/HEAD/OPTIONSfetch(including a method-omitted one) to any other route runs unattended, so scope the container's read routes accordingly. Passopts.needsApproval(a boolean or a predicate over the input) to widen or disable the gate. Anexecthat fails to run at all is returned to the model as text rather than thrown, so a command whose result was lost is not re-executed by its own throw. That is not exactly-once: the tool call reaches the container as a dispatch, and a dispatch whose reply is lost after the command ran fails the tool's durable step, whose retry execs again. The sandbox action takes noidempotencyKey— there is nothing on an action ctx to dedupe one against, and the approval decision is memoized in its own step, so the gate does not re-ask either. Make a command that cannot afford to run twice idempotent, or have it mark its own completion inside the container.Every call is addressed to one instance per thread, not a random pool member, so
exec("pnpm install")andexec("pnpm test")on the same thread share a filesystem. Size the container'smaxInstancesfor the number of threads you expect to be live at once.Time budgets and the ceiling that remains. A dispatch defaults to a 30s cap, and a dispatch that times out answers 503 — transient, so the tool's durable step rethrows and the host retries it, dispatching the same call again while the first is still running. For an
execthat is the command running twice. The sandbox therefore dispatches a browser op with a 150s budget over a 120s navigation deadline, and a container op with a 150s budget over a 120sexecdeadline, so the inner deadline is always the one that fires and the op ends rather than being abandoned. That widens the window; it does not make anexecexactly-once. A command that outlives its 120s deadline is killed, and one that somehow outlives the dispatch budget on top of that is still re-dispatched by the step retry. Removing that ceiling needs a fire-then-poll shape, which needs somewhere durable to record the op id — which an action ctx does not have.A container
fetchreads at most 1MB of the response before the read is cut and the reader cancelled (the result is marked[truncated at … bytes]), the same boundexechas always had: the body is buffered in a 128MB isolate shared with every other in-flight request. -
containerFsTool(name, opts?)has the samels/read/write/rm/statops asfsTool, but on the container's own disk, the onecontainerTool'sexecruns against. It addresses the same per-thread instance, so a file the model writes is the file its next command sees. The container needssandbox: trueinlunora/containers.ts(see the container docs). Paths are scoped underopts.root(default/workspace), awritecreates missing parent directories, and the writing ops are gated by default. -
fsTool(bucket, opts?)exposes a persistent, R2-backed virtual filesystem.opisls/read/write(withcontent) /rm/stat, over the R2 binding namedbucket. Every path is scoped underopts.root(e.g."agents/coder"); a..that would escape the root is rejected server-side, so the model can only touch its own prefix. The writing ops (write/rm) are gated behind a human approval by default; reads run unattended. workerd has no real shell, so this is object-store file I/O, not a POSIX shell. The app must declare ther2_bucketbinding inwrangler.jsonc(the op throws a directed error until it is wired).
Why it needs little wiring: a tool's execute runs inside the durable tool step,
which has no ctx.browser/ctx.containers. Both tools instead dispatch to a
single internal action, sandbox:invoke, that codegen auto-registers the
moment a lunora/ file imports either helper (no re-export boilerplate). That
action runs on an action ctx, which is where ctx.browser lives (and
ctx.containers rides every ctx once you declare a container). Importing
browserTool also makes codegen provision the BROWSER wrangler binding.
browserTool still needs ctx.browser wired the same way a direct
@lunora/browser user does: supply a config.browser thunk
(createBrowser({ binding: env.BROWSER, launch }), with the optional
@cloudflare/playwright launch peer) to createShardDO(); codegen never
injects that peer, so the browser op throws a directed error until it is wired.
Payloads are replay-stable, so a workflow replay re-encodes identical bytes to
identical text.
Web search (webSearchTool)
webSearchTool(opts?) lets the model search the web through the Cloudflare
Web Search API (ctx.ai.websearch in @lunora/ai). The model passes a
query; the tool returns up to limit results (default 5, max 10), each with
a url, a title and, when the provider has one, a description.
import { defineAgent, webSearchTool } from "@lunora/agent";
export const researcher = defineAgent({
model: "@cf/meta/llama-3.3-70b-instruct-fp8-fast",
tools: { search: webSearchTool({ provider: "exa" }) },
});It calls the AI binding directly from the tool's durable step: a search is a
read, so a step retry only repeats the query, and there is no dispatcher action
to register. A missing binding, a rejected request or an unreadable response
comes back to the model as the tool result, so the run carries on; a throttled
or failed search throws and is retried, and a gateway or credential problem
(UNAUTHORIZED / FORBIDDEN / NOT_FOUND) fails the run so you see it.
A search bills a provider key stored on the gateway (byokAlias, else the
default alias) or, without one, AI Gateway credits; gatewayId overrides
the gateway.
Code mode (codeTool)
Normally the model calls one tool per turn, a round-trip each. codeTool lets
it compose several tool calls in one turn as a script, feeding a later
call's input from an earlier call's output:
import { codeTool, defineAgent, functionTool } from "@lunora/agent";
export const analyst = defineAgent({
model: "@cf/meta/llama-3.3-70b-instruct-fp8-fast",
tools: {
run: codeTool({
findUser: functionTool("users:byEmail", { description: "Look up a user by email.", inputSchema }),
recentOrders: functionTool("orders:recent", { description: "List a user's recent orders.", inputSchema }),
}),
},
});The model writes steps: [{ id, tool, input }]; a later input references an
earlier output with { "$from": "<stepId>", "$path": "optional.dot.path" }. So
in one turn it can findUser, then feed { "$from": "u", "$path": "id" } into
recentOrders. The whole script runs in a single codeTool call and returns
each step's output plus the final one.
Every step's output must be JSON-serializable. Each step runs inside its own
durable step.do, so the workflow host serializes the returned value before the
next step — or the tool result — ever sees it. A Date comes back as an ISO
string, a Map/Set as {}, undefined as null or a dropped key, and a
bigint throws outright. This applies to the value a later step reads through
{ "$from": … } exactly as much as to the returned results: the hand-off is
durable state, not an in-memory one. A composed tool that wants to pass a rich
value should return its JSON form (date.toISOString(), [...map]) and let the
consuming step rebuild it.
This is a safe interpreted data-flow between the whitelisted tools you hand
codeTool, not arbitrary JavaScript, so there is no eval or isolate and it
runs natively in workerd. Each composed tool still dispatches through the same
durable context a normal call gets, keeping its RLS, and each step's resolved
input goes through the same check a top-level call gets — no more and no less.
That check is only as strong as the tool's own inputSchema: a schema carrying
a validator (a Standard Schema such as zod/valibot, or
jsonSchema(schema, { validate })) rejects a wrong-typed argument before the
tool runs, while a bare jsonSchema(schema) validates nothing — the AI SDK
treats a validator-less schema as pass-through, on this path and on the
top-level one alike. Without a validator the argument reaches the tool and its
own rejection (a dispatched function's 400, say) is recorded as a tool error
the model reads and corrects on its next turn, rather than failing the run.
A composed tool's own needsApproval gate is not carried, because a script
runs its steps in one shot and cannot hibernate mid-way for a human decision.
codeTool therefore refuses at construction to compose any tool that declares
one — the agent module throws on load rather than silently bypassing the gate.
Expose that tool as a normal top-level tool, or gate the whole script with
codeTool's own needsApproval.
Scripts are capped at maxSteps (default 16); a
longer script is rejected (code_tool_too_many_steps) rather than run as a
prefix, so the model can split the work instead of silently losing its trailing
steps. To run model-authored code rather than compose tools, use
jsCodeTool.
Sandboxed JavaScript (jsCodeTool)
jsCodeTool runs a script the model writes — the body of an async function —
in a Worker Loader isolate (Dynamic Workers on Cloudflare and on celld), and
returns what it returns plus its console output:
import { defineAgent, jsCodeTool } from "@lunora/agent";
export const analyst = defineAgent({
model: "@cf/meta/llama-3.3-70b-instruct-fp8-fast",
tools: { js: jsCodeTool({ cpuMs: 1000 }) },
});
// The tool result is `{ value, logs }`, or `{ error, logs }` when the script
// threw, did not parse, returned something that is not JSON, or ran out of time.The isolate gets no network (globalOutbound: null), no bindings and
no tools, and cpuMs (default 1000) plus a 30-second wall clock bound each
run. So it is for computation over what the model already has — arithmetic,
parsing, reshaping data — not a way to reach anything. A failure is a result,
not a throw, so the durable step does not retry a script that will fail the same
way again.
Importing jsCodeTool in lunora/ makes codegen provision the
worker_loaders binding (LOADER) and gate the target on
the workerLoaders capability: Cloudflare and celld rate it native, and the
Node host has no isolate sandbox, so target: "node" refuses it with a
platform_unsupported_feature diagnostic.
Limits
- A script cannot call the agent's tools. The only way out of a loaded
isolate is a service binding, and a call through one arrives as a fresh
invocation of the worker: it has neither the run's durable step (so nothing
it does is memoized or replay-safe) nor the caller's identity (so RLS would
not apply). To chain tool calls in one turn, use
codeTool; to compute over their results, let the model pass them intojsCodeToolas literals. - Nothing persists between runs. Each call loads a fresh isolate and discards it, so module-level state does not survive to the next call.
- The result must be JSON. A
bigint, a cycle or a function comes back aserror, and aDateas its ISO string. - Output is capped. A result over 256 KiB comes back as
error(the whole result is memoized as a durable step result, which the host caps at 1 MiB), andconsoleoutput at 100 entries of 2,000 characters each. - The isolate is a V8 boundary, not a process or VM one. It scopes what the
script can address; it is not a claim about V8 escape safety. Run code that
needs a kernel boundary in a container (
containerTool) instead. - Dynamic Workers are in open beta on Cloudflare, so the loader API can
still change. On celld a process holds at most 256 live loaded workers, and a
load past that fails the run with
errorset.
Skills
A skill is a reusable bundle many agents can share: an instruction fragment,
tools, and retrieval knowledge. defineSkill packages them;
defineAgent({ skills: [...] }) composes them in. A skill's tools carry the
same AgentToolDefinition shape agents already use
(functionTool / mcpTools / agentAsTool), and knowledge reuses memory's
retrieval verbatim.
import { defineAgent, defineSkill, functionTool } from "@lunora/agent";
import { jsonSchema } from "@lunora/ai";
import { api } from "./_generated/api";
const billing = defineSkill({
name: "billing",
// Merged into the system prompt, after the agent's own instructions.
instructions: "When asked about an invoice, always cite its invoice id.",
// Retrieved as its own durable step, keyed by the skill name.
knowledge: { source: "rag:searchBillingDocs", topK: 4 },
tools: {
lookupInvoice: functionTool(api.billing.invoiceById, {
description: "Look up an invoice by id.",
inputSchema: jsonSchema({ properties: { id: { type: "string" } }, required: ["id"], type: "object" }),
}),
},
});
export const support = defineAgent({
instructions: "You are a helpful support agent.",
memory: { source: "rag:searchDocs", topK: 5 },
model: "@cf/meta/llama-3.3-70b-instruct-fp8-fast",
skills: [billing],
});defineAgent folds the skills at declaration time:
- Tools merge into one flat namespace. A skill's tools join the agent's own
tools. The model sees a single namespace, so a name collision, between a skill and the agent or between two skills, throws atdefineAgent(name the colliding tool and rename one). This is the strict cousin ofmcpTools'prefix. - Instructions compose in order. The agent's own
instructionscome first, then each skill's fragment inskillsarray order, joined with blank lines. Dynamic fragments (thunks) resolve once at run start, so composition stays replay-stable. - Knowledge retrieves per skill. Each skill's
knowledgebecomes its own memory source, dispatched as a durable step at run start (memory:retrieve:<name>) and injected alongside the agent's ownmemorycontext. The agent'smemorykeeps the historicmemory:retrievestep name, so adding skills never disturbs an in-flight run's replay.
Skills carry no runtime of their own; they are pure config the agent absorbs, so an agent with skills compiles onto the same durable Workflow as one without.
Authoring a skill as markdown
The instructions are the part a non-author reads, reviews and copies between
projects, so they can live in a SKILL.md of their own. skillFromMarkdown
takes the file's source and reads name from its frontmatter:
---
name: triage
description: not read here, but it parses — other tooling uses it
---
Reproduce the report before proposing a cause. Cite the failing test.import { functionTool } from "@lunora/agent";
import { skillFromMarkdown } from "@lunora/agent/skill-markdown";
import { jsonSchema } from "@lunora/ai";
import triage from "../skills/triage/SKILL.md?raw";
import { api } from "./_generated/api";
const triageSkill = skillFromMarkdown(triage, {
tools: {
searchCode: functionTool(api.code.search, {
description: "Search the repository.",
inputSchema: jsonSchema({ properties: { query: { type: "string" } }, required: ["query"], type: "object" }),
}),
},
});skillFromMarkdown lives at its own subpath deliberately: it is the only part
of this package that needs a YAML parser, and the root export is reached by
every app that calls defineAgent. Behind the subpath, only apps that author a
skill in markdown pay for it.
The name travels with the file, so a skill moves between projects without a
matching TypeScript wrapper. Tools and knowledge stay in the second argument,
because they are code. Everything else behaves exactly as the object form,
including the name rules, so a reserved or malformed name fails identically.
Three things worth knowing:
- Pass the source, not a path. A Worker has no filesystem to read one from at
runtime, so the markdown is imported (
?rawunder Vite) and bundled. - The frontmatter is real YAML. Keys this does not read (
description, licences, tool allow-lists meant for other tooling) parse rather than confuse it, including list- and nested-valued ones. - Invalid YAML fails as invalid YAML. It does not degrade to "no
frontmatter", which would report a missing
namefor a file whose real problem is a syntax error and send you to the wrong line.
Scheduling
An agent compiles onto a Cloudflare Workflow, so it can be a cron target
directly. Codegen emits an agents.<name> reference into _generated/api; pass
it to a cronJobs() schedule and each fire starts a fresh durable run with the
trailing object as its AgentRunInput:
import { cronJobs } from "@lunora/scheduler";
import { agents } from "./_generated/api";
const crons = cronJobs();
// Every day at 03:00 UTC, start a fresh `support` run.
crons.daily("nightly sweep", { hourUTC: 3, minuteUTC: 0 }, agents.support, { input: "sweep", threadKey: "cron" });
export default crons;Each fire is an independent run (its own instance id), so give recurring jobs a
stable threadKey only if you want them to share one conversation thread.
For a one-off delay (rather than a recurring cron), pass the same
agents.<name> reference to ctx.scheduler.runAfter/runAt. Each starts a
single durable run when the timer fires, with the trailing object as its
AgentRunInput:
// Kick off a fresh `support` run five minutes from now.
await ctx.scheduler.runAfter(5 * 60_000, agents.support, { input: "follow up", threadKey: "t-1" });Inbound email
Give an agent an onEmail mapper and codegen wires it onto the worker's
top-level email() handler (via @lunora/agent/inbound), so a message routed by
Cloudflare Email Routing
starts a durable run. The mapper turns a parsed InboundEmail into an
AgentRunInput, or returns null/undefined to decline. When several agents
declare onEmail, each received message is offered to them in order and the
first that returns a run claims it:
export const support = defineAgent({
instructions: "You are a helpful support agent.",
model: "@cf/meta/llama-3.3-70b-instruct-fp8-fast",
onEmail: (email) => {
// The handler has already bounced any message whose `From` domain is not
// vouched for by an ALIGNED DMARC, SPF or DKIM pass — that gate is
// `authenticatesFrom` from `@lunora/mail/inbound`, the same helper to reach
// for in a handler of your own. Never hand-roll it as "did any clause
// pass?": an attacker gets a genuine pass for the domain they control.
//
// Layer your own policy on top of that gate — this is NOT the gate. Here:
// narrow to a full DMARC pass, stricter than the handler's
// DMARC-or-SPF-or-DKIM. It needs no alignment check of its own because a
// DMARC pass already checked alignment at the MX.
if (!email.authentication.dmarc.some((entry) => entry.result === "pass")) {
return null; // declined — no run, no bounce
}
return {
input: `${email.subject ?? "(no subject)"}\n\n${email.text ?? ""}`,
// `owner` gates the thread's reads and the run dispatches RLS-bypassed,
// so derive it from a verified signal (a DKIM-checked address mapped to
// an account) — never blindly from the spoofable sender.
owner: accountForVerifiedSender(email.from),
threadKey: email.messageId ?? email.from,
title: email.subject,
};
},
tools: {/* … */},
});Security: inbound mail is untrusted and a run dispatches privileged.
Cloudflare Email Routing authenticates only the recipient domain, not the
sender: the envelope from, subject, and body are trivially spoofable. Before
any mapper runs, the handler bounces a message unless some reported DMARC, SPF or
DKIM clause both passes and names a domain equal to the From domain
(email.authentication.dkim / .spf / .dmarc are lists, because one header
legitimately reports a method more than once) — alignment is strict
(mail.example.com does not vouch for example.com), and a pass with no
reported domain is rejected. Layer any
stricter policy on email.authentication inside the mapper, and treat every
mapped field as attacker-controlled input. A message no agent claims is
dropped silently; a throw while parsing or dispatching bounces the message with a
fixed generic reason (never reflecting internal error detail back to the sender).
The auto-wired handler is registered before any manual .onEmail(...) you add
to the app builder, so a hand-registered handler still wins if you need full
control.
Inbound channels (Slack / GitHub / Discord)
Trigger an agent from a verified inbound webhook. Give it an onInbound
config naming the channel, the verification secret (an env var), and a map
mapper, then mount dispatchAgentChannel(...) (from @lunora/agent/channels) on
an HTTP route:
import { defineAgent } from "@lunora/agent";
export const support = defineAgent({
model: "...",
onInbound: {
channel: "slack",
secret: "SLACK_SIGNING_SECRET", // env var — Slack signing secret / GitHub webhook secret / Discord Ed25519 public key (hex)
map: (event) => {
const payload = event.json() as { event?: { text?: string } };
// SECURITY: derive `owner` from the VERIFIED channel identity, never a payload field.
return { input: payload.event?.text ?? "", owner: "slack-workspace", threadKey: "slack-thread" };
},
},
});// mount on any HTTP route (e.g. an httpRouter POST handler):
import { dispatchAgentChannel } from "@lunora/agent/channels";
// `className` is the agent's generated class — its `ctx.exports` key.
const handler = dispatchAgentChannel([{ agent: support, className: "SupportAgentWorkflow" }]);
// Pass the Worker ctx: on Cloudflare the agent is reached as ctx.exports.SupportAgentWorkflow.
// handler(c.req.raw, c.env, c.executionCtx) → 200 (claimed / Discord PONG), 401 (bad signature), 204 (declined)dispatchAgentChannel detects the channel from the signature headers and
verifies the request over the raw body before any mapper runs: Slack HMAC
over v0:timestamp:body (with a timestamp-freshness replay guard), GitHub HMAC
over the body, Discord Ed25519 over timestamp+body (and it answers the Discord
PING with a PONG). A request that fails verification is rejected 401 and never
reaches map. The verifiers (verifySlack / verifyGithub / verifyDiscord)
are exported too if you wire your own routing.
Security. Trust comes ONLY from the signature check; the payload is attacker-controlled. Derive the run
ownerfrom the verified channel identity (the workspace / installation the secret belongs to), never from an arbitrary payload field; the run dispatches RLS-bypassed under whateverowneryou set.
Access control
Pass owner: ctx.auth.userId when starting a run. An owned thread's public
queries (agents:agentThread, agents:agentMessages) answer only for a caller
with that verified identity; for anyone else the thread is indistinguishable
from one that doesn't exist, so key-guessing leaks nothing. The owner is
immutable after the first run. Omitting owner leaves the thread readable by
anyone who knows its key, which is only appropriate for single-tenant or
anonymous apps.
The agent tables are RLS-exempt (.public()): under a .rls("required")
schema the auto-registered runtime functions cannot engage app RLS policies,
so access control is enforced inside them instead: the owner gate on reads,
and internal-only visibility on every write.
Replay-safety (the correctness core)
- Completed steps never re-run. Each LLM turn is the durable step
llm:turn:N; each tool call istool:NAME:CALL_ID(the provider's stable call id). Cloudflare Workflows memoizes completed steps by name, so a crashed run resumes without re-charging a card or re-paying for a model call. - Failed steps retry at-least-once. A tool that dies mid-body will run
again, so side-effecting tools receive their step name as
ctx.idempotencyKeyand must dedupe on it. - Messages persist idempotently. Every write is keyed
INSTANCE:ROLE:POSITIONand deduped by a unique index, so replays never duplicate the thread. - The loop's own dispatches are at-least-once. They carry no replay-dedup
id: most of them are made from inside a memoized
step.docallback, so any id numbered by call order is re-issued to a different call on the next activation — and the shard's dedup table, keyed(identity, mutationId)with no function path in it, would answer that call with the first one's cached result. What makes a redelivery harmless instead is that everyagents:*function the loop calls is idempotent on its own key: keyed appends, get-or-create threads, an instance-ownership guard on completion, and absolute (never additive) status/usage/state writes.
Memory (RAG)
Point memory.source at an action that wraps
@lunora/ai/rag's retrieve. The loop runs it as a
durable step at run start and injects the returned context into every
turn's prompt:
// lunora/rag.ts
export const searchDocs = action.input({ query: v.string(), topK: v.optional(v.number()) }).action(async ({ args, ctx }) => {
return docs(ctx).retrieve(args.query, { topK: args.topK });
});
// lunora/agents.ts
export const support = defineAgent({ memory: { source: "rag:searchDocs", topK: 5 }, model: "..." });Dispatching to a real action keeps retrieval inside a fully wired ctx (codegen-resolved vector bindings, RLS, observability) instead of re-plumbing vectors into the workflow runtime.
Agentic retrieval (mode: "agentic")
The default (mode: "inject") fetches top-k context once per run. Set
mode: "agentic" and the loop instead mints a searchMemory tool the model
calls mid-reasoning: it decides for itself what to look up, and searches again
with a refined query if the first hits fall short (multi-hop "read what you
need", à la Recursive-LM). No context is auto-injected; nothing reaches the
prompt until the model asks for it.
export const support = defineAgent({
// The model calls `searchMemory({ query, topK? })` when it needs context.
memory: { mode: "agentic", source: "rag:searchDocs", topK: 5 },
model: "...",
});searchMemory returns a compact projection: ranked { id, sourceId, score, snippet } hits plus deduped sources, with the giant joined .context string
dropped (each snippet is truncated to snippetChars, default 240). The model
reads the snippets and decides what, if anything, to pull in full.
Set read to an opt-in fetch-by-id action ({ id } -> string) and a
companion readMemory tool is minted so the model can pull a hit's full
document:
export const support = defineAgent({
memory: { mode: "agentic", read: "rag:getDoc", source: "rag:searchDocs" },
model: "...",
});Because every tool call is already a memoized durable step
(tool:searchMemory:<id>), multi-hop retrieval is crash-safe for free: a
resume replays completed searches from the journal rather than re-querying. A
skill's agentic knowledge mints the same tools namespaced by skill name
(search_<skill> / read_<skill>).
activeToolsand the memory tools. The memory tools are ordinary tools: subject toactiveTools,toolChoice,stopWhen, and bounded bymaxTurns. If you pinactiveTools, you must list every minted memory tool name (e.g.searchMemory);defineAgentthrows at declaration time ifactiveToolsomits one, rather than silently hiding it and leaving the model unable to retrieve. DropactiveToolsentirely to expose all tools.
Graph memory (kind: "graph")
Semantic memory (the default kind: "semantic") retrieves passages. A
graph memory instead builds structured, persistent knowledge: it extracts
the entities and relations stated in each run and traverses them on the
next one. Set kind: "graph"; no source action is needed, and none of your
own infrastructure either:
export const support = defineAgent({
memory: { graph: { depth: 2, maxSeeds: 4 }, kind: "graph" },
model: "...",
});- Owner-scoped, persists across threads. The graph is keyed by the run's
owner(the verified identity), so a fact learned in one conversation is available in every later conversation that same user has. An anonymous run (noowner) no-ops the graph tier for that run, because the tier is owner-scoped by design. - Owner-scoped, not agent-scoped. The
agent_entities/agent_edgesrows carry anownerbut no agent discriminator, so every graph-memory agent in the deployment shares one graph per owner. A fact one agent extracts is visible to another agent's traversal for the same owner. That is often the point (shared long-term knowledge of the user), but if two agents must not see each other's memory, give them separate deployments or don't enable the graph tier on both. - No external graph DB. Entities and relations live in the agent's own
DO-SQLite as the
agent_entities/agent_edgestables (shipped with theagentschema extension). Nothing to provision. - Auto-extraction on write. After a run answers, a memoized
memory:extractstep runs one LLM call over the exchange (user input + final answer) to pull{ entities, relations }, then idempotently upserts them. Because the step is memoized and the upsert is an absolute set (never an increment), a crash + resume never re-extracts or double-counts. A cheapergraph.extractionModelcan be set for this step. Extraction is best-effort: a failure never fails a run whose answer is already persisted. - Bounded traversal on read. A
memory:traversestep tokenizes the input, matches seed entities, then walks the graph with a bounded JS breadth-first search (depth,maxSeeds,fanOut,maxNodes; defaults 2 / 4 / 8 / 32), rendering deterministic- alice —[works_at]→ acmetriples that are injected like any other memory context. Reads are always bounded; storage is not. - Storage grows unbounded, with no eviction or TTL (yet). Every run appends new
entities/edges, deduped by normalized name / triple. Weight is monotonic,
not last-write-wins: an edge's weight becomes
max(prior, confidence)and an entity's weight is never updated after insert, so a later extraction that lowers a relation's confidence has no effect — a stale high-confidence edge keeps out-ranking underfanOutand cannot be corrected by restating it more weakly. Nothing prunes stale ones. Traversal stays fast because it is bounded by the knobs above, but the tables keep growing over an owner's lifetime; a weight/age-based prune is a planned follow-up. Note too that normalization is light (trim / collapse whitespace / lowercase), so distinct senses of the same string (e.g. two people named "alex") merge into one node, so keep names specific if that matters.
Why JS-BFS, not
WITH RECURSIVE. The loop runs inside a Workflow with no DB handle; it reaches SQLite only by dispatching a registered function whose typedctx.dbexposes indexed reads, not raw SQL. Traversal is therefore one dispatch doing bounded, indexed local reads (replay-stable: noDate.now()or randomness), which keeps the whole graph on one owner-stable shard.
Episodic memory (kind: "episodic")
Where semantic and graph memory recall by relevance, episodic memory recalls
by recency: it summarizes each completed run into one line and, on the next
run, injects the owner's most recent episodes as a short timeline. Set
kind: "episodic"; no source action and no infrastructure:
export const support = defineAgent({
memory: { episodic: { recall: 5 }, kind: "episodic" },
model: "...",
});- Owner-scoped, cross-thread. Episodes live in the agent's own DO-SQLite as
the
agent_episodestable (shipped with theagentschema extension), keyed byowner, so a run recalls the timeline of that user's earlier runs across every thread. An anonymous run (noowner) no-ops the tier. - Auto-summarized on write. After a run answers, a memoized
memory:episodestep runs one LLM call (a cheaperepisodic.extractionModelcan be set) to condense the exchange into a one-sentence episode, then idempotently records it, so a crash + resume never double-records the same run. - Recency recall on read. A
memory:recallstep returns the owner's most recentrecallepisodes (default 5, max 20) rendered oldest → newest as- <summary>lines, injected like any other memory context. - Unlike the graph tier, episodes are capped: every
agentEpisodeUpsertdeletes the owner's oldest rows beyond 200, so the table is bounded but is not an audit log — once an owner passes 200 runs, the earliest episodes are gone.
Testing
@lunora/testing ships agentHarness, an in-memory double over the same
AgentGenerate turn seam the production loop runs on. You script the model
turns (never a real model or network) and assert the persisted thread,
dispatched functions, and run result:
import { agentHarness, finalTurn, toolCallTurn } from "@lunora/testing";
const harness = agentHarness(support, {
// Each entry is one scripted LLM turn, consumed in order.
script: [toolCallTurn("call_1", "getWeather", { city: "Berlin" }), finalTurn("It's sunny in Berlin.")],
});
const result = await harness.run({ input: "weather in Berlin?", threadKey: "t1" });
expect(result.stopped).toBe("final");
expect(harness.messages("t1").map((m) => m.role)).toStrictEqual(["user", "assistant", "tool", "assistant"]);
expect(harness.dispatches).toContainEqual({ args: { city: "Berlin" }, path: "weather:lookup" });finalTurn(text, extra?) and toolCallTurn(id, name, input, text?) build turn
scripts; harness.thread(key) reads the persisted thread record; provide functions
runtime stubs to back the functions your tools dispatch to.
See also
- @lunora/ai: the model surface (
ctx.ai) +defineRag - @lunora/workflow: the durable-execution engine underneath
@lunora/agent/telemetry: ready-madetelemetry.integrations(console, Sentry, Braintrust)- @lunora/testing:
agentHarnessand the wider test toolkit - @lunora/mcp: expose an agent to external MCP clients as an
agent_<name>tool (see Expose an agent)