Durable runs for any agent framework, on any host.
Write your agent with the AI SDK, Mastra, the OpenAI Agents SDK, LangGraph or plain TypeScript. agent-unit makes it durable without a rewrite, so it can pause for a human, sleep for a day, survive a crash and resume where it stopped, and deploys it to Node, Bun, Deno, Cloudflare, Vercel, Netlify or AWS Lambda. Every agent gets the same HTTP API, a resumable event stream, an MCP endpoint and an A2A agent card.
agents/refund.ts ──► agent-unit build --preset vercel ──► .vercel/output
│ ├─ POST /agents/refund/runs
AI SDK · Mastra · OpenAI Agents · LangGraph · plain TS ├─ GET /runs/:id/events (SSE, resumable)
├─ POST /runs/:id/resume (answer a pause)
├─ POST /mcp (every agent is a tool)
└─ GET /.well-known/agent.json
It is to agents what Nitro is to servers, and it is built on Nitro: one build and one runtime contract, with presets for every host.
- Durable without a rewrite. Model calls and tool calls are journaled automatically. When a run continues after a pause, a restart or a crash, completed work replays from the journal instead of running again, so the provider is not billed twice and a refund is not charged twice.
- Pauses that cost nothing.
run.interrupt()parks a run until someone answers, andrun.sleep("2d")until a time passes. A parked run uses no compute, on any host. - One contract for every framework. The same run API, AG-UI events, manifest, MCP and A2A surface, whichever framework an agent uses, so a UI or an approvals inbox works with all of them.
- No deployers to maintain. Framework authors write one small adapter; agent-unit and Nitro handle hosts, storage and scheduling.
npm install agent-unit// agents/refund.ts
import { defineAgent } from "agent-unit";
export default defineAgent({
description: "Refunds an order after a human approves it.",
async run(input, run) {
const order = await run.step("load-order", () => loadOrder(input.orderId));
const decision = await run.interrupt<{ approved: boolean }>("approve-refund", order);
if (!decision.approved) return "declined";
await run.step("charge", () => refund(order)); // journaled: resuming never repeats it
return "refunded";
},
});npx agent-unit dev # http://localhost:3000, reloads on change
npx agent-unit build # .output/server/index.mjs for Node
node .output/server/index.mjscurl -X POST localhost:3000/agents/refund/runs \
-H 'content-type: application/json' -H 'accept: text/event-stream' \
-d '{"input":{"orderId":"o_1"}}'
# … event: RUN_INTERRUPTED {"interrupt":{"name":"approve-refund","payload":{…}}}
curl -X POST localhost:3000/runs/run_…/resume \
-H 'content-type: application/json' -d '{"answer":{"approved":true}}'Every file in agents/ is an agent named after the file (agents/support/index.ts serves as
support). Export a framework agent or a defineAgent agent as the default export.
agent-unit recognises agents from these frameworks when your package.json depends on them:
| Framework | Export | What becomes durable | How it pauses |
|---|---|---|---|
| AI SDK | new ToolLoopAgent(…) or streamText settings |
Every model call and tool call | useRun().interrupt() in a tool |
| Mastra | new Agent(…) |
Every model call and tool call (on a fork; your agent is untouched); memory uses the run's thread | Tool approvals (requireApproval, requireToolApproval), suspend() in a tool |
| OpenAI Agents SDK | new Agent(…) |
Model responses and function tools, per turn; handoffs included, handoff() options too |
Tool approvals (needsApproval: true) |
| LangGraph | graph.compile() |
LangGraph checkpoints, stored in agent-unit storage | interrupt() in a node, resumed with Command({ resume }) |
| Plain TypeScript | defineAgent(…) |
run.step(…) |
run.interrupt(), run.sleep() |
Every framework is tested on every host below (Node, Bun, Deno, Cloudflare Workers with Durable Objects, Vercel, Netlify and AWS Lambda): a run pauses for approval, the host restarts, and the resume repeats no model call or side effect.
Events follow AG-UI: text streams as TEXT_MESSAGE_*, a model's
reasoning as REASONING_*, and RUN_FINISHED carries token usage per provider and model (also on
the run as usage). An agent that asks for structured output (AI SDK output, Mastra
structuredOutput, OpenAI Agents outputType) finishes with that object as its output.
A framework's own pauses become agent-unit interrupts, answered with POST /runs/:id/resume:
- Tool approvals (OpenAI Agents
needsApproval, MastrarequireApprovalorrequireToolApproval) pause astool-approvalwith{ callId, tool, arguments, agent }. Answertrue,falseor{ approved }; a declined call is not run and the model is told so. - Mastra
suspend(payload)in a tool pauses under the tool's name with that payload. The answer becomes the tool'sresumeDatawhen it runs again, as with Mastra's own resume. - LangGraph
interrupt(value)in a node pauses under the node's name; the answer resumes it.
// agents/support.ts: an AI SDK agent, unchanged except for the approval
import { useRun } from "agent-unit";
import { openai } from "@ai-sdk/openai";
import { ToolLoopAgent, tool } from "ai";
import { z } from "zod";
export default new ToolLoopAgent({
model: openai("gpt-5"),
tools: {
refund: tool({
description: "Refund an order",
inputSchema: z.object({ orderId: z.string() }),
execute: async ({ orderId }) => {
const { approved } = await useRun().interrupt<{ approved: boolean }>("approve-refund", { orderId });
return approved ? await refund(orderId) : "declined";
},
}),
},
});When the run resumes, the model's first response comes from the journal, the tool call continues with the answer, and only new work reaches the provider.
Inside any run, useRun() (or the run argument of defineAgent) gives you:
| Primitive | |
|---|---|
run.step(name, fn) |
Run fn once; replays return the journaled result. fn receives { idempotencyKey } |
run.idempotencyKey() |
The key of the step or tool call running now |
run.interrupt(name, payload) |
Park until POST /runs/:id/resume answers, then return the answer |
run.sleep("10m") |
Park until the time passes |
run.state.get/set/delete(key, { scope }) |
State scoped to the thread (default), the agent or the app |
run.secrets.get(name) |
Read a secret from the host environment (Workers bindings included) |
run.emit(name, value) |
Emit a CUSTOM event |
run.signal |
Aborts when the run is cancelled or parked |
The one rule: between steps, code must make the same decisions given the same journaled results. Put anything with side effects or randomness in a step. Model and tool calls already are.
A step is recorded when it finishes. If the process dies while a step is still running, that step runs again when the run recovers. Every step gets an idempotency key that stays the same when it runs again and differs for every other step and run; hand it to the API the step calls, and the repeat has no second effect:
await run.step("charge", ({ idempotencyKey }) =>
stripe.refunds.create({ payment_intent: order.paymentId }, { idempotencyKey }),
);Inside a tool the framework calls (AI SDK, Mastra, OpenAI Agents), useRun().idempotencyKey()
returns the tool call's key. In a LangGraph node, wrap the side effect in useRun().step(…).
npx agent-unit build --preset <preset>| Host | Preset | Storage | Long runs | Sleep and recovery |
|---|---|---|---|---|
| Node | node-server (default) |
Filesystem by default | Unlimited | In-process timers + sweep every minute |
| Bun | bun |
Filesystem by default | Unlimited | In-process timers + sweep |
| Deno | deno-server |
Filesystem by default | Unlimited | In-process timers + sweep |
| Cloudflare Workers | cloudflare-module with runtime: "durable-objects" |
One Durable Object per run (no storage to configure) | Yields every 10 min and continues from an alarm | Each run's own alarm |
| Cloudflare Workers (KV) | cloudflare-module |
Set storage (for example cloudflare-kv-binding) |
Yields every 25s and continues | Cron Trigger runs the sweep |
| Vercel | vercel |
Set storage (Redis, Upstash, Vercel KV, …) |
Yields every 240s and continues | Vercel Cron runs the sweep |
| Netlify | netlify |
Set storage (netlify-blobs, Redis, …) |
Yields every 20s; finishes in the request | Generated scheduled function |
| AWS Lambda | aws-lambda |
Set storage (Redis, Upstash, a database via db0, …) |
Yields every 25s; finishes in the request | Point EventBridge at POST /__agent-unit/sweep |
Cloudflare builds turn on Workers' stub modules for child_process, readline, worker_threads
and tty, and replace the Node ws package with the built-in WebSocket: Mastra and the OpenAI
Agents SDK import these when they load, and one missing module stops the whole Worker from starting.
Your own compatibility_flags and alias settings are kept.
Any other Nitro preset works too. Serverless hosts need shared
storage, because their filesystem does not survive between invocations; the build warns when one is
missing. Set AGENT_UNIT_SECRET in production: it protects the internal continue and sweep
endpoints, which serverless hosts use to continue long runs. A long run continues in a fresh
invocation at the deployment's own URL, read from AGENT_UNIT_URL (or VERCEL_URL, Netlify's URL,
or the origin option) and never from an incoming request; without one it continues in-process.
Already have a server, from your agent framework or written by hand? Point server at it and
agent-unit builds it for every preset above, with none of its own engine, routes or storage added.
Your framework keeps its own runs, memory and storage.
// agent-unit.config.ts
import { defineConfig } from "agent-unit";
export default defineConfig({
server: "./src/server.ts",
cloudflare: { wrangler: { d1_databases: [{ binding: "DB", database_name: "app", database_id: "…" }] } },
});// src/server.ts
import { Hono } from "hono";
import { env } from "agent-unit/env";
const app = new Hono();
app.get("/", (c) => c.text(`hello from ${env.APP_NAME}`));
export default app;The default export handles web requests: a Hono app, export default { fetch }, or a function
(request, env, ctx) => Response. On Workers it receives the Worker's env and ctx; elsewhere ctx
has waitUntil where the host keeps work alive after the response. agent-unit/env reads
configuration the same way on every host. agent-unit dev serves it with reload on change. The
engine settings (agents, storage, runtime, budget, sweep, retention, basePath,
authorize, origin, maxBodyBytes, adapters) do not apply and are rejected with server.
waitUntil from agent-unit/env keeps work running after the response (an email, a webhook, a
cache write). It uses the host's own waitUntil on Node, Bun, Deno, Workers and Vercel; on Netlify
and AWS Lambda, which may freeze after a response, the response waits for the work instead. Hono's
c.executionCtx.waitUntil does the same.
import { waitUntil } from "agent-unit/env";
app.post("/orders", async (c) => {
const order = await createOrder(await c.req.json());
waitUntil(sendReceipt(order));
return c.json(order, 201);
});tasks runs work on a schedule with the host's own scheduler: in-process cron on Node, Bun and
Deno, Cron Triggers on Workers, Vercel Cron Jobs on Vercel. Other presets build with a warning, and
the task runs only when you trigger it there.
export default defineConfig({
server: "./src/server.ts",
tasks: { "send-digest": { schedule: "0 8 * * *", handler: "./src/tasks/digest.ts" } },
});// src/tasks/digest.ts
export default async ({ name, scheduledTime }: { name: string; scheduledTime: number }) => {
await sendDigest(new Date(scheduledTime));
};npx agent-unit task send-digest runs it once from source.
A Mastra app does not need a server file: the Mastra preset serves your Mastra instance with
Mastra's own Hono adapter, so every Mastra route, its auth, memory and workflows work as they do
under mastra dev, on the storage your instance configures.
npm install agent-unit @mastra/hono hono// agent-unit.config.ts
import { defineConfig } from "agent-unit";
import { mastra } from "agent-unit/presets/mastra";
export default defineConfig({ server: mastra() }); // src/mastra/index.ts, export const mastraStorage stays Mastra's. To use a different store per host, let the build pick it with a
package.json import condition: Cloudflare builds resolve workerd, everything else default.
{ "imports": { "#storage": { "workerd": "./src/storage.d1.ts", "default": "./src/storage.libsql.ts" } } }// src/storage.d1.ts: Mastra's D1 store on the Worker's DB binding (declare it in `cloudflare.wrangler`)
import { D1Store } from "@mastra/cloudflare-d1";
import { env } from "agent-unit/env";
export const storage = new D1Store({ id: "app", binding: env.DB });Mastra's D1 store cannot de-duplicate two simultaneous resumes of the same workflow run (Mastra
logs a warning saying so); stores with atomic updates, such as LibSQL and Postgres, can. Other
frameworks can publish a preset the same way: server takes any object with a name and an
entry({ root, genDir, preset }) that returns the path of a server module.
Frameworks without a server of their own run in yours. Give the framework its own durable storage
and point server at the file:
// src/server.ts: LangGraph, paused and resumed from its own checkpointer
import { Hono } from "hono";
import { Command } from "@langchain/langgraph";
import { PostgresSaver } from "@langchain/langgraph-checkpoint-postgres";
import { env } from "agent-unit/env";
import { workflow } from "./graph";
const checkpointer = PostgresSaver.fromConnString(env.DATABASE_URL);
await checkpointer.setup(); // creates LangGraph's tables the first time
const graph = workflow.compile({ checkpointer });
const app = new Hono();
app.post("/threads/:id/runs", async (c) =>
c.json(await graph.invoke(await c.req.json(), { configurable: { thread_id: c.req.param("id") } })));
app.post("/threads/:id/resume", async (c) =>
c.json(await graph.invoke(new Command({ resume: await c.req.json() }), { configurable: { thread_id: c.req.param("id") } })));
export default app;// src/server.ts: the OpenAI Agents SDK, approvals kept as a serialized RunState
import { Hono } from "hono";
import { RunState, run } from "@openai/agents";
import { support } from "./agent";
import { pending } from "./store"; // any key-value store you choose
const app = new Hono();
app.post("/runs", async (c) => {
const result = await run(support, (await c.req.json()).input);
if (!result.interruptions?.length) return c.json({ output: result.finalOutput });
const id = crypto.randomUUID();
await pending.set(id, result.state.toString());
return c.json({ id, approvals: result.interruptions.map((item) => item.rawItem) }, 202);
});
app.post("/runs/:id/approve", async (c) => {
const state = await RunState.fromString(support, await pending.get(c.req.param("id")));
for (const item of state.getInterruptions()) state.approve(item);
return c.json({ output: (await run(support, state)).finalOutput });
});
export default app;Prefer agent-unit to handle durability for you? Leave server out and use its engine.
On Cloudflare, give every run its own Durable Object, the same primitive Cloudflare's own agents run on:
// agent-unit.config.ts
export default defineConfig({ preset: "cloudflare-module", runtime: "durable-objects" });npx agent-unit build && npx wrangler deployYour agent code does not change. Each run's object holds its journal and events in strongly
consistent storage and is the only thing that ever executes it, so there are no leases to race. Its
alarm wakes it from run.sleep, continues it after a long step, and picks it up again if the
isolate running it dies. A shared index object keeps the run list and thread, agent and app state.
The build exports both classes and writes their bindings and migration into wrangler.json.
Your app can bring its own Cloudflare resources beside agent-unit's: D1, KV or R2 for a storage
provider, or its own Durable Objects (Mastra's Cloudflare store ships one, as does a Workflow). Write
them in Cloudflare's own wrangler.json format under cloudflare, and export your classes from a
file:
export default defineConfig({
preset: "cloudflare-module",
runtime: "durable-objects",
cloudflare: {
wrangler: {
d1_databases: [{ binding: "DB", database_name: "memory", database_id: "…" }],
durable_objects: { bindings: [{ name: "STORE", class_name: "StoreObject" }] },
},
exports: "./store.ts", // or exports.cloudflare.ts, picked up on its own
},
});To connect your storage the same way on every host, read your bindings and secrets from
agent-unit/env. Builds for Workers resolve it to the Workers environment (bindings included);
everywhere else it is process.env:
import { env } from "agent-unit/env";
const storage = env.DB ? new D1Store({ binding: env.DB }) : new LibSQLStore({ url: env.DATABASE_URL });They are merged into what agent-unit generates, never in its place: the names AGENT_UNIT_RUNS,
AGENT_UNIT_INDEX, AgentUnitRun, AgentUnitIndex and the migration tag agent-unit-v1 are
reserved, and the build says so if you reuse one. (nitro.cloudflare is the same setting and still
works.)
In a Worker you write yourself:
import { createDurableAgentUnit } from "agent-unit/cloudflare";
import support from "./agents/support";
const unit = createDurableAgentUnit({ agents: { support } });
export const { AgentUnitRun, AgentUnitIndex } = unit;
export default { fetch: unit.fetch };
// wrangler.json: copy unit.wrangler (durable_objects bindings + migrations)Already on Nitro (or something built on it)? Mount your agents with the module:
// nitro.config.ts
import { defineConfig } from "nitro";
import { agentUnit } from "agent-unit/nitro";
export default defineConfig({
modules: [agentUnit({ basePath: "/api/agents" })],
});In a Farm app, mount it with one catch-all API route, which works in
farm dev and every production target (see examples/farm-app):
// src/app/api/ai/[...path]/route.ts
import { after } from "@farm.js/core/after";
import { agents } from "../../../../agents/server"; // createAgentUnit({ …, basePath: "/api/ai" })
const handle = (request: Request) =>
agents.handler(request, { waitUntil: (work) => after(async () => void (await work)) });
export const GET = handle;
export const POST = handle;Anywhere else that speaks Fetch (Hono, Next.js route handlers, Bun.serve, Deno.serve, a Worker):
import { createAgentUnit } from "agent-unit/server";
import { createStorage } from "unstorage";
import redis from "unstorage/drivers/redis";
import support from "./agents/support";
const unit = createAgentUnit({
agents: { support },
storage: createStorage({ driver: redis({ url: process.env.REDIS_URL }) }),
basePath: "/api/agents",
});
export const fetch = (request: Request, ctx?: { waitUntil(p: Promise<unknown>): void }) =>
unit.handler(request, { waitUntil: ctx?.waitUntil.bind(ctx) });Pass the host's waitUntil where it has one. Without it, start and resume requests finish the work
before they answer, which is the safe behaviour on hosts that freeze after a response.
| Example | What it shows |
|---|---|
examples/universal |
The universal build: one agents/ folder built for Node, Bun, Deno, Cloudflare, Vercel, Netlify and AWS Lambda, then every build run through the same pause, restart and resume flow (npm run try). |
examples/cloudflare-ai-sdk |
An AI SDK ToolLoopAgent (OpenAI and Claude) on Cloudflare Workers with the Durable Objects runtime: a refund pauses for approval, and npm run verify checks it against real models across a runtime restart. |
examples/farm-app |
Agents inside a Farm app at /api/ai, with an approvals inbox and live event stream at /agents, tested on the production server across a restart (npm run smoke). |
agent-unit.config.ts is optional:
import { defineConfig } from "agent-unit";
export default defineConfig({
name: "support", // manifest, MCP and A2A name
storage: { driver: "redis", url: process.env.REDIS_URL }, // any unstorage driver
preset: "vercel", // or --preset
runtime: "default", // or "durable-objects" on Cloudflare
budget: "60s", // yield and continue after this long
sweep: "*/5 * * * *", // sweep schedule, or false
basePath: "/api/agents",
authorize: (request) => request.headers.get("authorization") === `Bearer ${process.env.API_TOKEN}`,
retention: "30d", // delete finished runs 30 days after they finish
maxBodyBytes: 1_000_000, // largest accepted request body (default 1 MB)
// server: "./src/server.ts", // your own server instead (see "Your own server")
// agents: { support }, // instead of the agents/ directory
// adapters: [myAdapter], // extra framework adapters
// nitro: { … }, // passed through to Nitro
});The full contract is in spec/. Events follow AG-UI (RUN_STARTED,
TEXT_MESSAGE_CONTENT, TOOL_CALL_START, …) plus RUN_INTERRUPTED, RUN_SLEEPING and
RUN_CANCELLED. Every event has a seq; it is the SSE event id, so EventSource reconnects with
Last-Event-ID and picks up exactly where it left off.
import { createAgentClient } from "agent-unit/client";
const agents = createAgentClient({ baseUrl: "https://agents.example.com" });
for await (const event of agents.run("refund", { orderId: "o_1" })) {
if (event.type === "RUN_INTERRUPTED") console.log("needs approval:", event.interrupt.payload);
}
for await (const event of agents.resume(runId, { approved: true })) console.log(event.type);
await agents.list({ status: "interrupted" }); // an approvals inbox in one call
await agents.delete(runId); // remove a finished run and its historyThe client reconnects dropped streams from the last event it saw.
POST /mcp speaks MCP over streamable HTTP: every agent is a tool, so Claude, Cursor or any MCP
client can call your agents. A run that pauses returns its run id and the pending interrupt.
GET /.well-known/agent.json is an A2A agent card with one skill per agent.
import { createTestUnit } from "agent-unit/testing";
import refund from "../agents/refund";
const unit = createTestUnit({ refund });
const parked = await unit.run("refund", { orderId: "o_1" });
expect(parked.status).toBe("interrupted");
const done = await unit.resume(parked.id, { approved: true });
expect(done.output).toBe("refunded");Pass the same storage to a second createTestUnit to test a restart.
agent-unit is layered. Pick the layer that matches what you are changing:
| You want to… | Layer | Who |
|---|---|---|
| Pause, sleep, keep state, make a call run once | Run primitives: useRun() / run.step, run.interrupt, run.sleep, run.state |
App developers, in agent code |
| Choose host, storage, runtime, auth, time budget | agent-unit.config.ts |
App developers |
| Mount the API in your own server, add routes, compose | createAgentUnit(), createDurableAgentUnit(), createHandler() |
App developers |
| Support a new framework | defineAdapter() with ctx.durable, ctx.kv, ctx.emit |
Framework and library authors |
| Run agents on a new kind of runtime | RunEngine options, or your own RunService |
Runtime and platform authors |
Mount and compose. createAgentUnit returns a Web handler, so it sits beside your own
routes, middleware and auth:
import { createAgentUnit } from "agent-unit/server";
const unit = createAgentUnit({
agents: { support },
storage, // any unstorage instance
basePath: "/api/agents",
// Second argument: what the request does, with the run's agent and thread taken from storage.
authorize: async (request, { action, threadId }) => {
const session = await getSession(request);
return session !== null && (threadId === undefined || threadId.startsWith(`${session.userId}:`));
},
budget: "60s",
});
// Any Fetch-style router, here Hono:
app.all("/api/agents/*", (c) => unit.handler(c.req.raw));Build on the engine. Everything the HTTP API does goes through RunEngine, which you can drive
directly (a queue consumer, a cron job, a test):
import { RunEngine, RunStore, resolveAgent } from "agent-unit/runtime";
const engine = new RunEngine({
store: new RunStore(storage),
agents: [resolveAgent("support", support, [myAdapter])],
budgetMs: 60_000,
// Hand wake-ups to your own scheduler (a delayed queue message, a job runner…) instead of timers:
scheduleWake: (runId, at) => scheduler.at(at, () => engine.handleWake(runId)),
});
const { run, done } = await engine.start("support", { prompt: "refund o_1" });A new runtime. createHandler(service) serves the whole API (runs, SSE, MCP, A2A) from any
object that implements RunService from agent-unit/server. The Durable Objects runtime is one
such implementation: it routes each call to the run's object instead of running it in-process.
AI SDK, Mastra, OpenAI Agents and LangGraph are supported out of the box. Any other framework can publish its own adapter package and maintain it on its own schedule: no change to agent-unit, no pull request. To use one:
npx agent-unit add agent-unit-adapter-crewIt installs the package, checks that it is an adapter, and adds two lines to
agent-unit.config.ts (asking first; --yes skips the question). Adding them by hand works the same:
import crewAdapter from "agent-unit-adapter-crew";
export default defineConfig({ adapters: [crewAdapter()] });Adapters in the config are tried before the built-in ones, so a framework's own adapter replaces agent-unit's.
An adapter teaches agent-unit one framework. Most are a few dozen lines:
// agent-unit-adapter-crew/src/index.ts
import { defineAdapter } from "agent-unit/adapter";
import { CrewAgent } from "crew";
export default function crewAdapter() {
return defineAdapter<CrewAgent>({
name: "crew",
apiVersion: 1, // the adapter interface this was written for
match: (value): value is CrewAgent => value instanceof CrewAgent,
describe: (agent) => ({ description: agent.description, tools: agent.tools.map(({ name }) => ({ name })) }),
async run(agent, ctx) {
const model = ctx.durable.model(agent.model); // journaled, streams text events
const tools = ctx.durable.tools(agent.tools); // journaled, emits tool events
return agent.run(ctx.input, { model, tools, signal: ctx.signal });
},
});
}To publish it:
- Export a function returning the adapter as the default export, so
agent-unit addcan find it. - Name it
agent-unit-adapter-<framework>, or@your-scope/agent-unit, and add theagent-unit-adapterkeyword, so people can find it on npm. - Declare
agent-unitand your framework as peer dependencies, so the app's copies are used. - Check it with the same checks agent-unit runs on its own adapters, in any test runner. Report
every real side effect through
kit.effect; the checks fail when one runs outside a journaled call, or runs again after a restart:
import { assertAdapter } from "agent-unit/testing";
test("agent-unit adapter", () =>
assertAdapter({
adapter: crewAdapter(),
agent: (kit) => new CrewAgent({ model, tools: [refundTool(() => kit.effect("refund"))] }),
// Optional: an agent that pauses through your framework, and the answer that resumes it.
pause: { agent: (kit) => approvalAgent(kit), answer: { approved: true } },
}));See spec/adapters.md for adapters of frameworks that keep their own state.
Run as many instances as you like over one store. With Redis it is exact out of the box:
storage: { driver: "redis", url: process.env.REDIS_URL },With createAgentUnit on your own server, wrap the storage yourself:
import { RunStore } from "agent-unit/runtime";
import { redisCoordination } from "agent-unit/redis";
import { createStorage } from "unstorage";
import redis from "unstorage/drivers/redis";
const driver = redis({ url: process.env.REDIS_URL });
const storage = new RunStore(createStorage({ driver }), redisCoordination({ client: () => driver.getInstance!() }));
// createAgentUnit({ agents, storage })redisCoordination needs one method, eval, so any client with Lua scripting works (ioredis,
node-redis behind a small wrapper). For another store, implement AtomicWrites.compareAndSet (save
a run record only if its version is still the one read) and LeaseBackend with that store's own
conditional write, such as a database UPDATE … WHERE version = ?.
- Leases make one execution run a run at a time, renewed while it works, and every change to a run
record is a versioned save: it lands only if nobody else saved the run since it was read, so two
processes racing to resume, cancel or wake a run cannot both win (the loser gets a 409). On Redis
both are atomic, wired in automatically by
agent-unit buildandagent-unit devwhenstorageuses theredisdriver. On other shared stores the save is check-then-write, which is exact within a process and best effort across processes racing within milliseconds; pass atomicleasesandatomicwrites toRunStorefor those (see below), or use the Durable Objects runtime, where each run's object is its only executor and saves go through the index atomically. - A Mastra agent with memory reads the run's thread (and
input.resourceId, or the thread, as its resource) and stores each finished turn once. With working or observational memory turned on, Mastra saves messages itself as the turn runs, so a run that pauses or recovers from a crash may store that turn twice. - In the Durable Objects runtime, the run list and shared state live in one index object, written when
a run starts, parks or finishes (not on every step).
agent-unit devruns the default runtime locally. - Streamed text and tool-argument deltas that arrive within 50ms of each other are merged into one
event, so a long answer costs a few writes instead of one per token, and live clients see text up
to 50ms later.
deltaBatchMs: 0oncreateAgentUnitkeeps every delta its own event.
MIT, Farming Labs.