diff --git a/CHANGELOG.md b/CHANGELOG.md index 754bb40d..e7fc5143 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -35,6 +35,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - `http_request` is no longer open to DNS rebinding (N13a). It resolved a host name once to check it against the private-address list and then let `fetch` resolve it again to connect, so a name that answered a public address to the check and `127.0.0.1` (or a cloud metadata or intranet address) to the connection reached the private address; the docs and a comment in the Worker runtime described the check as rebinding-safe. Now the one resolution is the check: the in-process path connects through an undici `Agent` whose `connect.lookup` checks every address and hands the socket the address it checked (each redirect hop included), and the sandboxed path (`sandboxExecute`, which agents use) resolves and checks the host in the agent's process and has the sandboxed process connect to that address (node:http/https with a fixed lookup; the host name still goes in `Host` and TLS SNI). The private-address list is now the shared one (below), which adds `0.0.0.0/8`, `100.64.0.0/10`, `192.0.0.0/24`, `198.18.0.0/15`, multicast and reserved `224.0.0.0/3`, `::`, `ff00::/8`, NAT64 `64:ff9b::/96` and 6to4 `2002::/16` to what `http_request` refused. Behavior changes: a host name that cannot be resolved now fails with the resolver's error instead of being passed to `fetch`; `undici` is loaded on the first `http_request` request (not only with `validateSSL: false`). New option `createHttpTool({ allowPrivate })` for hosts that may resolve to private addresses. Cloudflare Workers have no DNS hook, so neither tool exists in the Worker build; docs/deployment.md now says what a Worker can and cannot check. ### Added +- Two ready-made durable stores (R2). `@lousho/build-ai-agent/kv` is a new subpath exporting `KVStore`, `KVCheckpointStore`, `CHECKPOINT_KV_BINDING` and the types `KVStoreOptions`, `KVBinding`, `KVPutOptions`, so a hand-written Cloudflare Worker can use `createAgent({ provider, store: new KVStore(env.AGENT_KV) })`; nothing in its import graph touches `node:*`, so it bundles without shims (`assertSessionId` moved to the Node-free `src/session/sessionId.ts` and is still exported where it was). `fileStore(dir, { historyLimit? })` (root export) is an `AgentStore` of plain JSON files: `sessions/.json`, `checkpoints/.json`, `checkpoint-history/.json` and `approvals/.json`, each written to a temp file and renamed into place, with no lock files and no `StorageService`; resolving an approval claims it with an exclusive create, so of two processes resolving one approval only one gets it. It replaces combining `FileSessionStore`, `LocalStorageCheckpointStore` and `StorageServiceApprovalStore` by hand. See docs/sessions.md#choosing-a-store and docs/deployment.md. - `web_fetch` built-in tool (N13a): `webFetchTool` / `createWebFetchTool(options)`, and `web-fetch` in spec files. It `GET`s one public web page and returns `{ url, finalUrl, status, contentType, content, truncated }`: HTML converted to text by a small built-in converter (scripts and styles dropped, link URLs kept, entities decoded), JSON and `text/*` as sent, a note for other types; a 4xx/5xx is returned with its body. Loopback, private and link-local destinations are refused on every hop with the connection pinned to the checked address; redirects (10), bytes read (2 MiB), characters returned (50 000) and time (30 s) are capped; `allowedHosts` / `blockedHosts` are checked before DNS and `allowPrivate` exempts named hosts. Node only. See docs/tools.md#built-in-tools. - `src/security/privateAddress.ts` (internal): `isPrivateAddress()` and `pinnedLookup()`, shared by the credential broker, `http_request` and `web_fetch`. The credential broker uses it instead of its own list, so it also refuses `192.0.0.0/24`, `198.18.0.0/15`, NAT64 and 6to4 addresses. - Slack and Discord: a pending `ask_question` survives a restart (M10a). Given durable stores for sessions, checkpoints and approvals, the next message in the Slack thread (or the next `/ask` in the Discord channel) after a restart is still the answer: the turn continues, the reply is posted, and the question, the answer and the reply are appended to the session transcript (before, the answer became a new turn and the question was orphaned). A message in a conversation that waits on a tool approval still does not decide it. New `ChannelContext.pendingQuestion(sessionKey)` for custom channels: the id of the `ask_question` that key's session waits on, also one asked before a restart (it binds the approval to the session, so the continuation is recorded); handed out once until that answer has run. A channel's answer and click continuations now run one at a time with the session's turns. Supporting additions: `Checkpoint.approvalKind`, `PendingTurn.approvalKind` (`session.pending()`) and `SessionAwaitingApprovalError.approvalKind` are `'question'` for a turn paused on an `ask_question`. Slack thread replies that answer a pending question no longer need a saved transcript first (a first turn paused on a question with checkpoints has none yet). See docs/channels.md. diff --git a/docs/api-overview.md b/docs/api-overview.md index 14adf7cd..985573df 100644 --- a/docs/api-overview.md +++ b/docs/api-overview.md @@ -49,6 +49,8 @@ How the pieces fit: | `InMemoryApprovalStore` | Process-local `ApprovalStore`; the default store of `createAgent()` agents. | | `StorageServiceApprovalStore`, `LocalStorageCheckpointStore` | File-backed approval and checkpoint stores over a `StorageService` (see [Approvals](./approvals.md), [Durable execution](./durable-execution.md)). | | `SqliteStore` (from `/sqlite`) | Sessions, checkpoints and approvals in one SQLite file (see [Sessions](./sessions.md#choosing-a-store)). | +| `fileStore(dir)` | Sessions, checkpoints and approvals as plain JSON files under `dir` (see [Sessions](./sessions.md#choosing-a-store)). | +| `KVStore`, `KVCheckpointStore` (from `/kv`) | Stores on a Cloudflare Workers KV binding, for a hand-written Worker (see [Deployment](./deployment.md)). | | `AgentStore`, `memoryStore()` | The `createAgent({ store })` option: `{ sessions?, checkpoints?, approvals? }`, and an in-memory one (see [Sessions](./sessions.md#choosing-a-store)). | | `SessionAwaitingApprovalError` | Thrown by `execute()` when its `sessionId` is paused on an approval (see [Durable execution](./durable-execution.md)). | | `SDKError`, `ERROR_CODES` | Base class of the SDK's errors: a stable `code`, a `hint` and a `docs` link (see [Errors](./errors.md)). | diff --git a/docs/deployment.md b/docs/deployment.md index e1d71bc9..d94fffe3 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -202,10 +202,33 @@ curl -N https://.workers.dev/chat \ -d '{ "sessionId": "alice", "input": "Hello" }' ``` -`KVStore(kvBinding, { prefix?, ttl? })` (`src/deploy/kvStore.ts`) is the -`AgentStore` the generated Worker builds from the binding. It is not exported -from any entry point of the package, so importing it in a hand-written Worker -is not supported yet. Its keys, with an optional `prefix` before each: +`KVStore(kvBinding, { prefix?, ttl?, historyLimit? })` is the `AgentStore` the +generated Worker builds from the binding. A hand-written Worker imports it from +the `/kv` subpath, which has no `node:*` import anywhere in its graph, with the +binding typed as `KVBinding` (the `get`/`put`/`delete` part of Cloudflare's +`KVNamespace`, so `@cloudflare/workers-types` is not needed): + +```ts +import { createAgent } from '@lousho/build-ai-agent'; +import { KVStore, type KVBinding } from '@lousho/build-ai-agent/kv'; + +interface Env { + AGENT_KV: KVBinding; +} + +export default { + async fetch(request: Request, env: Env): Promise { + const agent = createAgent({ provider, store: new KVStore(env.AGENT_KV) }); + const { sessionId, input } = (await request.json()) as { sessionId: string; input: string }; + const { text } = await agent.session({ id: sessionId }).send(input); + return Response.json({ text }); + }, +}; +``` + +`/kv` also exports `KVCheckpointStore` (checkpoints only) and +`CHECKPOINT_KV_BINDING` (`'AGENT_CHECKPOINTS'`, the binding name the generated +Worker reads). `KVStore`'s keys, with an optional `prefix` before each: | Key | Value | | --- | ----- | @@ -306,9 +329,8 @@ one session at the same moment can overwrite each other's turn, since a KV read-modify-write is not atomic. The KV-backed stores (`KVStore`, `KVCheckpointStore` and `CHECKPOINT_KV_BINDING`, -in `src/deploy/kvStore.ts`, `src/deploy/kvCheckpointStore.ts` and -`src/deploy/checkpointBinding.ts`) have no `node:*` references anywhere in their -dependency graph. The Worker runs the spec as a `createAgent()` agent, whose +exported from `@lousho/build-ai-agent/kv`) have no `node:*` references anywhere +in their dependency graph. The Worker runs the spec as a `createAgent()` agent, whose Node-only imports (project instructions, the file session store, guardrail patches, MCP over stdio) the build points at a shim that fails when used (`src/deploy/shims/node.worker.ts`). The built `dist/worker.js` bundle is then diff --git a/docs/durable-execution.md b/docs/durable-execution.md index 4b586036..9a06a812 100644 --- a/docs/durable-execution.md +++ b/docs/durable-execution.md @@ -193,7 +193,8 @@ Which stores keep one: | `memoryStore()` | yes | `memoryStore({ historyLimit })` | | `SqliteStore` | yes (`checkpoint_history` table) | `new SqliteStore(path, { historyLimit })` | | `LocalStorageCheckpointStore` | yes | `new LocalStorageCheckpointStore(storage, { historyLimit })` | -| `KVCheckpointStore` / `KVStore` (Cloudflare Workers KV) | yes (see below) | `new KVStore(kv, { historyLimit })` | +| `fileStore(dir)` (JSON files) | yes (`checkpoint-history/.json`) | `fileStore(dir, { historyLimit })` | +| `KVCheckpointStore` / `KVStore` (Cloudflare Workers KV, from `@lousho/build-ai-agent/kv`) | yes (see below) | `new KVStore(kv, { historyLimit })` | | Agent Forge's `FileCheckpointStore` | yes | `new FileCheckpointStore(dir, { historyLimit })` | | a custom `CheckpointStore` | only if it implements `history()` | - | diff --git a/docs/sessions.md b/docs/sessions.md index cf3864c6..27f026bc 100644 --- a/docs/sessions.md +++ b/docs/sessions.md @@ -227,33 +227,30 @@ implementation by where the process runs: | Store | Sessions | Checkpoints | Approvals | Use it when | | --- | --- | --- | --- | --- | | In memory (`memoryStore()`) | `MemorySessionStore` | in memory | `InMemoryApprovalStore` | Tests, scripts, one process that never restarts | -| Files | `FileSessionStore(dir)` | `LocalStorageCheckpointStore` | `StorageServiceApprovalStore` | One machine, you want plain inspectable files | +| Files (`fileStore(dir)`) | `/sessions/.json` | `/checkpoints/`, `/checkpoint-history/` | `/approvals/.json` | One machine, you want plain inspectable files | | SQLite | `store.sessions` | `store.checkpoints` | `store.approvals` | A Node server: one durable, transactional file, shared safely by several processes | -| Cloudflare KV | - | `KVCheckpointStore` | - | Workers deployments (see [Deployment](deployment.md)) | +| Cloudflare KV (`KVStore` from `@lousho/build-ai-agent/kv`) | `store.sessions` | `store.checkpoints` | `store.approvals` | Workers deployments (see [Deployment](deployment.md)) | Each one is an `AgentStore` part: pass them together as -`createAgent({ store: { sessions, checkpoints, approvals } })`. `memoryStore()` -and `SqliteStore` are ready-made `AgentStore`s; for plain files, combine the -file stores: +`createAgent({ store: { sessions, checkpoints, approvals } })`. `memoryStore()`, +`fileStore(dir)`, `SqliteStore` and `KVStore` are ready-made `AgentStore`s. For +plain files, `fileStore(dir)` writes one JSON file per session, checkpoint and +pending approval under `dir`, each written to a temp file and renamed into +place, with no lock files: ```ts -import { - createAgent, - FileSessionStore, - LocalStorageCheckpointStore, - StorageServiceApprovalStore, - type AgentStore, -} from '@lousho/build-ai-agent'; - -// `storage` is a StorageService rooted where the files should go. -const store: AgentStore = { - sessions: new FileSessionStore('./.lousho/sessions'), - checkpoints: new LocalStorageCheckpointStore(storage), - approvals: new StorageServiceApprovalStore(storage), -}; -const agent = createAgent({ provider, store }); +import { createAgent, fileStore } from '@lousho/build-ai-agent'; + +const agent = createAgent({ provider, store: fileStore('./.lousho') }); +// ./.lousho/sessions/user-42.json, ./.lousho/checkpoints/..., ./.lousho/approvals/... +await agent.session({ id: 'user-42' }).send('Hello'); ``` +Resolving an approval from `fileStore` is safe across processes (only one +caller gets the record); two processes writing one session at the same moment +are not coordinated, so the last write wins. For several processes sharing a +store, use `SqliteStore`. + Any object with the three methods of a part works there too: a Redis `SessionStore`, or a KV-backed `CheckpointStore` on Cloudflare Workers (the generated Worker uses `KVCheckpointStore`, see [Deployment](deployment.md)). diff --git a/llms-full.txt b/llms-full.txt index 286b2fcb..96bbee19 100644 --- a/llms-full.txt +++ b/llms-full.txt @@ -899,6 +899,8 @@ How the pieces fit: | `InMemoryApprovalStore` | Process-local `ApprovalStore`; the default store of `createAgent()` agents. | | `StorageServiceApprovalStore`, `LocalStorageCheckpointStore` | File-backed approval and checkpoint stores over a `StorageService` (see [Approvals](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/approvals.md), [Durable execution](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/durable-execution.md)). | | `SqliteStore` (from `/sqlite`) | Sessions, checkpoints and approvals in one SQLite file (see [Sessions](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/sessions.md#choosing-a-store)). | +| `fileStore(dir)` | Sessions, checkpoints and approvals as plain JSON files under `dir` (see [Sessions](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/sessions.md#choosing-a-store)). | +| `KVStore`, `KVCheckpointStore` (from `/kv`) | Stores on a Cloudflare Workers KV binding, for a hand-written Worker (see [Deployment](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/deployment.md)). | | `AgentStore`, `memoryStore()` | The `createAgent({ store })` option: `{ sessions?, checkpoints?, approvals? }`, and an in-memory one (see [Sessions](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/sessions.md#choosing-a-store)). | | `SessionAwaitingApprovalError` | Thrown by `execute()` when its `sessionId` is paused on an approval (see [Durable execution](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/durable-execution.md)). | | `SDKError`, `ERROR_CODES` | Base class of the SDK's errors: a stable `code`, a `hint` and a `docs` link (see [Errors](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/errors.md)). | @@ -4098,10 +4100,33 @@ curl -N https://.workers.dev/chat \ -d '{ "sessionId": "alice", "input": "Hello" }' ``` -`KVStore(kvBinding, { prefix?, ttl? })` (`src/deploy/kvStore.ts`) is the -`AgentStore` the generated Worker builds from the binding. It is not exported -from any entry point of the package, so importing it in a hand-written Worker -is not supported yet. Its keys, with an optional `prefix` before each: +`KVStore(kvBinding, { prefix?, ttl?, historyLimit? })` is the `AgentStore` the +generated Worker builds from the binding. A hand-written Worker imports it from +the `/kv` subpath, which has no `node:*` import anywhere in its graph, with the +binding typed as `KVBinding` (the `get`/`put`/`delete` part of Cloudflare's +`KVNamespace`, so `@cloudflare/workers-types` is not needed): + +```ts +import { createAgent } from '@lousho/build-ai-agent'; +import { KVStore, type KVBinding } from '@lousho/build-ai-agent/kv'; + +interface Env { + AGENT_KV: KVBinding; +} + +export default { + async fetch(request: Request, env: Env): Promise { + const agent = createAgent({ provider, store: new KVStore(env.AGENT_KV) }); + const { sessionId, input } = (await request.json()) as { sessionId: string; input: string }; + const { text } = await agent.session({ id: sessionId }).send(input); + return Response.json({ text }); + }, +}; +``` + +`/kv` also exports `KVCheckpointStore` (checkpoints only) and +`CHECKPOINT_KV_BINDING` (`'AGENT_CHECKPOINTS'`, the binding name the generated +Worker reads). `KVStore`'s keys, with an optional `prefix` before each: | Key | Value | | --- | ----- | @@ -4202,9 +4227,8 @@ one session at the same moment can overwrite each other's turn, since a KV read-modify-write is not atomic. The KV-backed stores (`KVStore`, `KVCheckpointStore` and `CHECKPOINT_KV_BINDING`, -in `src/deploy/kvStore.ts`, `src/deploy/kvCheckpointStore.ts` and -`src/deploy/checkpointBinding.ts`) have no `node:*` references anywhere in their -dependency graph. The Worker runs the spec as a `createAgent()` agent, whose +exported from `@lousho/build-ai-agent/kv`) have no `node:*` references anywhere +in their dependency graph. The Worker runs the spec as a `createAgent()` agent, whose Node-only imports (project instructions, the file session store, guardrail patches, MCP over stdio) the build points at a shim that fails when used (`src/deploy/shims/node.worker.ts`). The built `dist/worker.js` bundle is then @@ -4422,7 +4446,8 @@ Which stores keep one: | `memoryStore()` | yes | `memoryStore({ historyLimit })` | | `SqliteStore` | yes (`checkpoint_history` table) | `new SqliteStore(path, { historyLimit })` | | `LocalStorageCheckpointStore` | yes | `new LocalStorageCheckpointStore(storage, { historyLimit })` | -| `KVCheckpointStore` / `KVStore` (Cloudflare Workers KV) | yes (see below) | `new KVStore(kv, { historyLimit })` | +| `fileStore(dir)` (JSON files) | yes (`checkpoint-history/.json`) | `fileStore(dir, { historyLimit })` | +| `KVCheckpointStore` / `KVStore` (Cloudflare Workers KV, from `@lousho/build-ai-agent/kv`) | yes (see below) | `new KVStore(kv, { historyLimit })` | | Agent Forge's `FileCheckpointStore` | yes | `new FileCheckpointStore(dir, { historyLimit })` | | a custom `CheckpointStore` | only if it implements `history()` | - | @@ -9009,33 +9034,30 @@ implementation by where the process runs: | Store | Sessions | Checkpoints | Approvals | Use it when | | --- | --- | --- | --- | --- | | In memory (`memoryStore()`) | `MemorySessionStore` | in memory | `InMemoryApprovalStore` | Tests, scripts, one process that never restarts | -| Files | `FileSessionStore(dir)` | `LocalStorageCheckpointStore` | `StorageServiceApprovalStore` | One machine, you want plain inspectable files | +| Files (`fileStore(dir)`) | `/sessions/.json` | `/checkpoints/`, `/checkpoint-history/` | `/approvals/.json` | One machine, you want plain inspectable files | | SQLite | `store.sessions` | `store.checkpoints` | `store.approvals` | A Node server: one durable, transactional file, shared safely by several processes | -| Cloudflare KV | - | `KVCheckpointStore` | - | Workers deployments (see [Deployment](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/deployment.md)) | +| Cloudflare KV (`KVStore` from `@lousho/build-ai-agent/kv`) | `store.sessions` | `store.checkpoints` | `store.approvals` | Workers deployments (see [Deployment](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/deployment.md)) | Each one is an `AgentStore` part: pass them together as -`createAgent({ store: { sessions, checkpoints, approvals } })`. `memoryStore()` -and `SqliteStore` are ready-made `AgentStore`s; for plain files, combine the -file stores: +`createAgent({ store: { sessions, checkpoints, approvals } })`. `memoryStore()`, +`fileStore(dir)`, `SqliteStore` and `KVStore` are ready-made `AgentStore`s. For +plain files, `fileStore(dir)` writes one JSON file per session, checkpoint and +pending approval under `dir`, each written to a temp file and renamed into +place, with no lock files: ```ts -import { - createAgent, - FileSessionStore, - LocalStorageCheckpointStore, - StorageServiceApprovalStore, - type AgentStore, -} from '@lousho/build-ai-agent'; +import { createAgent, fileStore } from '@lousho/build-ai-agent'; -// `storage` is a StorageService rooted where the files should go. -const store: AgentStore = { - sessions: new FileSessionStore('./.lousho/sessions'), - checkpoints: new LocalStorageCheckpointStore(storage), - approvals: new StorageServiceApprovalStore(storage), -}; -const agent = createAgent({ provider, store }); +const agent = createAgent({ provider, store: fileStore('./.lousho') }); +// ./.lousho/sessions/user-42.json, ./.lousho/checkpoints/..., ./.lousho/approvals/... +await agent.session({ id: 'user-42' }).send('Hello'); ``` +Resolving an approval from `fileStore` is safe across processes (only one +caller gets the record); two processes writing one session at the same moment +are not coordinated, so the last write wins. For several processes sharing a +store, use `SqliteStore`. + Any object with the three methods of a part works there too: a Redis `SessionStore`, or a KV-backed `CheckpointStore` on Cloudflare Workers (the generated Worker uses `KVCheckpointStore`, see [Deployment](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/deployment.md)). diff --git a/package.json b/package.json index 5791b6a5..efdb9353 100644 --- a/package.json +++ b/package.json @@ -63,6 +63,11 @@ "import": "./dist/storage/sqlite/index.mjs", "require": "./dist/storage/sqlite/index.js" }, + "./kv": { + "types": "./dist/deploy/kv.d.ts", + "import": "./dist/deploy/kv.mjs", + "require": "./dist/deploy/kv.js" + }, "./traces": { "types": "./dist/traces/index.d.ts", "import": "./dist/traces/index.mjs", diff --git a/src/deploy/kv.test.ts b/src/deploy/kv.test.ts new file mode 100644 index 00000000..693e19a9 --- /dev/null +++ b/src/deploy/kv.test.ts @@ -0,0 +1,47 @@ +/** + * R2: `@lousho/build-ai-agent/kv` must bundle for a hand-written Cloudflare + * Worker, which has no shim for Node builtins. This guards the barrel's + * import graph: no module it reaches may import a `node:` path. + */ +import { afterAll, describe, expect, it } from 'vitest'; +import { build, stop } from 'esbuild'; +import { join } from 'node:path'; +import * as kv from './kv'; + +// esbuild keeps a service process alive after build(); end it so it does not outlive this file. +afterAll(() => stop()); + +describe('@lousho/build-ai-agent/kv (R2)', () => { + it('bundles for the browser platform with no node: import in its graph', async () => { + const result = await build({ + entryPoints: [join(__dirname, 'kv.ts')], + bundle: true, + write: false, + platform: 'browser', + format: 'esm', + external: ['ai', 'zod', '@opentelemetry/api'], + metafile: true, + logLevel: 'silent', + }); + + expect(result.errors).toEqual([]); + const nodeImports = Object.entries(result.metafile.inputs).flatMap(([file, input]) => + input.imports.filter((entry) => entry.path.startsWith('node:')).map((entry) => `${file} -> ${entry.path}`) + ); + expect(nodeImports).toEqual([]); + expect(Object.keys(result.metafile.inputs).some((file) => file.endsWith('kvStore.ts'))).toBe(true); + }, 30_000); + + it('exports KVStore, KVCheckpointStore and CHECKPOINT_KV_BINDING', () => { + expect(Object.keys(kv).sort()).toEqual(['CHECKPOINT_KV_BINDING', 'KVCheckpointStore', 'KVStore']); + expect(kv.CHECKPOINT_KV_BINDING).toBe('AGENT_CHECKPOINTS'); + const memory = new Map(); + const binding: kv.KVBinding = { + get: async (key) => memory.get(key) ?? null, + put: async (key, value, _options?: kv.KVPutOptions) => void memory.set(key, value), + delete: async (key) => void memory.delete(key), + }; + const options: kv.KVStoreOptions = { prefix: 'app/' }; + expect(new kv.KVStore(binding, options).checkpoints).toBeInstanceOf(kv.KVCheckpointStore); + }); +}); diff --git a/src/deploy/kv.ts b/src/deploy/kv.ts new file mode 100644 index 00000000..e22e1506 --- /dev/null +++ b/src/deploy/kv.ts @@ -0,0 +1,9 @@ +/** + * `@lousho/build-ai-agent/kv` (R2): the Cloudflare Workers KV stores, for a + * hand-written Worker. Node-free on purpose (no `node:*` import anywhere in + * its graph, guarded by kv.test.ts) so it bundles for the Workers runtime + * without shims. + */ +export { KVStore, type KVStoreOptions } from './kvStore'; +export { KVCheckpointStore, type KVBinding, type KVPutOptions } from './kvCheckpointStore'; +export { CHECKPOINT_KV_BINDING } from './checkpointBinding'; diff --git a/src/deploy/kvStore.ts b/src/deploy/kvStore.ts index fadd58a9..ce5cb60f 100644 --- a/src/deploy/kvStore.ts +++ b/src/deploy/kvStore.ts @@ -16,7 +16,8 @@ */ import type { ApprovalStore, ExecutionSnapshot, PendingApproval, ResolvedApproval } from '../execution/ApprovalGate'; import type { Message } from '../providers/llm'; -import { assertSessionId, type SessionStore } from '../session/sessionStore'; +import { assertSessionId } from '../session/sessionId'; +import type { SessionStore } from '../session/sessionStore'; import type { AgentStore } from '../storage/agentStore'; import { DEFAULT_KV_KEY_PREFIX, KVCheckpointStore, type KVBinding } from './kvCheckpointStore'; diff --git a/src/execution/checkpoint.ts b/src/execution/checkpoint.ts index ba38716c..031fc4db 100644 --- a/src/execution/checkpoint.ts +++ b/src/execution/checkpoint.ts @@ -3,8 +3,8 @@ * Backend-agnostic durable-execution checkpointing for AgentExecutor runs. */ -import { Message } from '../providers'; -import { StorageService } from '../storage'; +import type { Message } from '../providers'; +import type { StorageService } from '../storage'; import type { StepUsage } from '../models/usage'; import type { CheckpointUsage } from './runUsage'; import type { AgentFingerprint } from './agentFingerprint'; diff --git a/src/execution/checkpointHistory.contract.test.ts b/src/execution/checkpointHistory.contract.test.ts index c8d186db..c8420d71 100644 --- a/src/execution/checkpointHistory.contract.test.ts +++ b/src/execution/checkpointHistory.contract.test.ts @@ -1,10 +1,13 @@ /** * LOU-D43: one contract suite for `CheckpointStore.history()`, run against * every store that implements it (the in-memory store, `SqliteStore`, - * `LocalStorageCheckpointStore` and `KVCheckpointStore`), plus the `getCheckpointHistory()` helper on + * `LocalStorageCheckpointStore`, `KVCheckpointStore` and `fileStore()`), plus the `getCheckpointHistory()` helper on * stores that do not. */ -import { describe, it, expect, afterEach } from 'vitest'; +import { describe, it, expect, afterEach, afterAll, vi } from 'vitest'; +import { mkdtempSync, rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; import { getCheckpointHistory, LocalStorageCheckpointStore, @@ -12,6 +15,7 @@ import { } from './checkpoint'; import { StorageService } from '../storage/StorageService'; import { memoryStore } from '../storage/agentStore'; +import { fileStore } from '../storage/fileStore'; import { SqliteStore } from '../storage/sqlite'; import { KVCheckpointStore } from '../deploy/kvCheckpointStore'; import { createFakeFs } from './__fixtures__/fakeFs'; @@ -23,6 +27,14 @@ afterEach(() => { while (sqliteStores.length) sqliteStores.pop()?.close(); }); +// fileStore() writes 55 checkpoints in one test; real files are slow on a loaded Windows machine. +vi.setConfig({ testTimeout: 20_000 }); + +const tempDirs: string[] = []; +afterAll(() => { + for (const dir of tempDirs) rmSync(dir, { recursive: true, force: true }); +}); + const implementations: Array<[string, CheckpointHistoryStoreFactory]> = [ ['memoryStore()', (options) => memoryStore(options).checkpoints], [ @@ -41,6 +53,14 @@ const implementations: Array<[string, CheckpointHistoryStoreFactory]> = [ return new KVCheckpointStore(kv, undefined, undefined, options); }, ], + [ + 'fileStore()', + (options) => { + const dir = mkdtempSync(join(tmpdir(), 'lousho-history-')); + tempDirs.push(dir); + return fileStore(dir, options).checkpoints; + }, + ], [ 'LocalStorageCheckpointStore', (options) => { diff --git a/src/execution/guardrails.test.ts b/src/execution/guardrails.test.ts index e4eed189..445948c9 100644 --- a/src/execution/guardrails.test.ts +++ b/src/execution/guardrails.test.ts @@ -118,9 +118,19 @@ describe('command guardrails (LOU-E11)', () => { fs.cpSync(fixtureRepo, scratchDir, { recursive: true }); }); - afterEach(() => { - fs.rmSync(scratchDir, { recursive: true, force: true }); - }); + afterEach(async () => { + // The timeout test's child is killed by a taskkill nobody waits for; on + // Windows its working directory stays locked (EPERM) until it exits. + for (let attempt = 0; ; attempt++) { + try { + fs.rmSync(scratchDir, { recursive: true, force: true }); + return; + } catch (error) { + if ((error as NodeJS.ErrnoException).code !== 'EPERM' || attempt >= 50) throw error; + await new Promise((done) => setTimeout(done, 200)); + } + } + }, 15_000); it('createTestRunGuardrail passes against the unmutated fixture repo', async () => { const guardrail = createTestRunGuardrail(scratchDir); diff --git a/src/index.ts b/src/index.ts index bd027723..82de0098 100644 --- a/src/index.ts +++ b/src/index.ts @@ -54,6 +54,7 @@ export * from './createAgent'; export type { AgentApprovals, ApproveToolCall } from './createAgentApprovals'; // One store for sessions, checkpoints and approvals: createAgent({ store }) (LOU-D30) export { memoryStore, type AgentStore, type MemoryStoreOptions } from './storage/agentStore'; +export { fileStore, type FileStoreOptions } from './storage/fileStore'; // Sessions: multi-turn conversations for createAgent() (LOU-W4) export * from './session'; diff --git a/src/session/sessionId.ts b/src/session/sessionId.ts new file mode 100644 index 00000000..9c4c4f8d --- /dev/null +++ b/src/session/sessionId.ts @@ -0,0 +1,24 @@ +/** + * Session id validation, in a module with no Node builtins so Worker code + * (`KVStore`, `@lousho/build-ai-agent/kv`) can import it. `sessionStore.ts` + * re-exports `assertSessionId`. + */ + +import { ConfigurationError } from '../execution/errors'; + +const SESSION_ID_PATTERN = /^[A-Za-z0-9_-]{1,128}$/; + +/** + * Throws unless `id` is 1-128 characters of letters, digits, `_` or `-`. + * Session ids become file names, so anything else (`../`, `/`, `.`) is refused. + */ +export function assertSessionId(id: string): void { + if (typeof id !== 'string' || !SESSION_ID_PATTERN.test(id)) { + throw new ConfigurationError( + `Invalid session id ${JSON.stringify(id)}: use 1-128 characters from A-Z, a-z, 0-9, '_' and '-' ` + + "(e.g. 'user-42'). Omit the id to get a generated one.", + 'id', + 'LOUSHO_SESSION_ID_INVALID' + ); + } +} diff --git a/src/session/sessionStore.ts b/src/session/sessionStore.ts index 538b2f20..b4b855ac 100644 --- a/src/session/sessionStore.ts +++ b/src/session/sessionStore.ts @@ -7,7 +7,8 @@ import { mkdir, readFile, rename, rm, writeFile } from 'node:fs/promises'; import { join, resolve } from 'node:path'; import { randomUUID } from 'node:crypto'; import type { Message } from '../providers/llm'; -import { ConfigurationError, SDKError } from '../execution/errors'; +import { SDKError } from '../execution/errors'; +import { assertSessionId } from './sessionId'; /** How bytes (image and file parts, LOU-V11) are saved in a JSON transcript: `{ "$bytes": "" }`. */ const BYTES_KEY = '$bytes'; @@ -45,22 +46,7 @@ export interface SessionStore { delete(id: string): Promise; } -const SESSION_ID_PATTERN = /^[A-Za-z0-9_-]{1,128}$/; - -/** - * Throws unless `id` is 1-128 characters of letters, digits, `_` or `-`. - * Session ids become file names, so anything else (`../`, `/`, `.`) is refused. - */ -export function assertSessionId(id: string): void { - if (typeof id !== 'string' || !SESSION_ID_PATTERN.test(id)) { - throw new ConfigurationError( - `Invalid session id ${JSON.stringify(id)}: use 1-128 characters from A-Z, a-z, 0-9, '_' and '-' ` + - "(e.g. 'user-42'). Omit the id to get a generated one.", - 'id', - 'LOUSHO_SESSION_ID_INVALID' - ); - } -} +export { assertSessionId }; /** * In-memory store: the default. Transcripts live as long as the process diff --git a/src/storage/fileStore.test.ts b/src/storage/fileStore.test.ts new file mode 100644 index 00000000..14c0fb5d --- /dev/null +++ b/src/storage/fileStore.test.ts @@ -0,0 +1,187 @@ +/** + * R2: `fileStore(dir)`, an AgentStore of JSON files. Runs the shared store + * contracts, then the file-specific behavior: layout, atomic claims of + * approvals, id validation, a corrupt history file, and two agents on one dir. + */ +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { mkdtempSync, readdirSync, rmSync, writeFileSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { z } from 'zod'; +import { fileStore } from './fileStore'; +import { createAgent } from '../createAgent'; +import { defineTool } from '../tools/defineTool'; +import { mockModel } from '../testing'; +import type { Message } from '../providers/llm'; +import { + describeApprovalStoreContract, + describeCheckpointStoreContract, + describeSessionStoreContract, + makeCheckpoint, + makePending, + makeSnapshot, +} from './sqlite/__fixtures__/storeContracts'; + +// Real files: slow on a loaded Windows machine (antivirus scans every temp file). +vi.setConfig({ testTimeout: 30_000 }); + +const dirs: string[] = []; +function tempDir(): string { + const dir = mkdtempSync(join(tmpdir(), 'lousho-filestore-')); + dirs.push(dir); + return dir; +} +afterEach(() => { + while (dirs.length) rmSync(dirs.pop()!, { recursive: true, force: true }); +}); + +describeSessionStoreContract('fileStore().sessions', () => fileStore(tempDir()).sessions); +describeCheckpointStoreContract('fileStore().checkpoints', () => fileStore(tempDir()).checkpoints); +describeApprovalStoreContract('fileStore().approvals', () => fileStore(tempDir()).approvals); + +const save = (store: ReturnType, sessionId: string, stepIndex: number) => + store.checkpoints.save(sessionId, makeCheckpoint({ sessionId, stepIndex })); + +describe('fileStore(dir) (R2)', () => { + it('writes sessions, checkpoints, history and approvals at the documented layout, with no temp or lock files left', async () => { + const dir = tempDir(); + const store = fileStore(dir); + await store.sessions.save('chat', [{ role: 'user', content: 'hi' }]); + await save(store, 'chat.turn-0', 1); + const pending = makePending('appr_1'); + await store.approvals.save(pending, makeSnapshot(pending)); + + expect(readdirSync(dir).sort()).toEqual(['approvals', 'checkpoint-history', 'checkpoints', 'sessions']); + expect(readdirSync(join(dir, 'sessions'))).toEqual(['chat.json']); + expect(readdirSync(join(dir, 'checkpoints'))).toEqual(['chat.turn-0.json']); + expect(readdirSync(join(dir, 'checkpoint-history'))).toEqual(['chat.turn-0.json']); + expect(readdirSync(join(dir, 'approvals'))).toEqual(['appr_1.json']); + }); + + it('round-trips a transcript with a Uint8Array file part, across store instances', async () => { + const dir = tempDir(); + const bytes = new Uint8Array([0, 1, 2, 250, 255]); + const messages: Message[] = [ + { role: 'user', content: [{ type: 'text', text: 'read this' }, { type: 'file', data: bytes, mimeType: 'application/octet-stream' }] }, + ]; + await fileStore(dir).sessions.save('doc', messages); + + const loaded = await fileStore(dir).sessions.load('doc'); + const part = (loaded![0].content as Array<{ type: string; data?: unknown }>)[1]; + expect(part.data).toBeInstanceOf(Uint8Array); + expect(Array.from(part.data as Uint8Array)).toEqual(Array.from(bytes)); + }); + + it('checkpoints: save, load, delete and history() newest first', async () => { + const store = fileStore(tempDir()); + await save(store, 's', 1); + await save(store, 's', 2); + expect((await store.checkpoints.load('s'))?.stepIndex).toBe(2); + expect((await store.checkpoints.history('s')).map((entry) => entry.step)).toEqual([2, 1]); + await store.checkpoints.delete('s'); + expect(await store.checkpoints.load('s')).toBeNull(); + expect(await store.checkpoints.history('s')).toEqual([]); + }); + + it('honors historyLimit: 3 keeps the newest three, 0 keeps none, keepHistory keeps the ring', async () => { + const small = fileStore(tempDir(), { historyLimit: 3 }); + for (let step = 0; step < 5; step++) await save(small, 's', step); + expect((await small.checkpoints.history('s')).map((entry) => entry.step)).toEqual([4, 3, 2]); + await small.checkpoints.delete('s', { keepHistory: true }); + expect(await small.checkpoints.load('s')).toBeNull(); + expect((await small.checkpoints.history('s')).map((entry) => entry.step)).toEqual([4, 3, 2]); + + const dir = tempDir(); + const off = fileStore(dir, { historyLimit: 0 }); + await save(off, 's', 1); + expect(await off.checkpoints.history('s')).toEqual([]); + expect(readdirSync(dir)).not.toContain('checkpoint-history'); + }); + + it('reads a truncated history file as empty, and the next save starts a new ring', async () => { + const dir = tempDir(); + const store = fileStore(dir); + await save(store, 's', 1); + writeFileSync(join(dir, 'checkpoint-history', 's.json'), '[{"step":1,"savedAt":"20'); + expect(await store.checkpoints.history('s')).toEqual([]); + await save(store, 's', 2); + expect((await store.checkpoints.history('s')).map((entry) => entry.step)).toEqual([2]); + }); + + it('approvals: resolve returns the record once, then null', async () => { + const store = fileStore(tempDir()); + const pending = makePending('appr_once'); + await store.approvals.save(pending, makeSnapshot(pending)); + expect((await store.approvals.resolve('appr_once'))?.pending.id).toBe('appr_once'); + expect(await store.approvals.resolve('appr_once')).toBeNull(); + }); + + it('two concurrent resolves of one approval (two store instances on one dir) give exactly one record', async () => { + const dir = tempDir(); + for (let round = 0; round < 10; round++) { + const id = `appr_race_${round}`; + const pending = makePending(id); + await fileStore(dir).approvals.save(pending, makeSnapshot(pending)); + const results = await Promise.all([fileStore(dir).approvals.resolve(id), fileStore(dir).approvals.resolve(id)]); + expect(results.filter((result) => result !== null)).toHaveLength(1); + } + expect(readdirSync(join(dir, 'approvals'))).toEqual([]); + }); + + it('a claim left by a crashed resolver makes resolve return null until the approval is saved again', async () => { + const dir = tempDir(); + const store = fileStore(dir); + const pending = makePending('appr_stale'); + await store.approvals.save(pending, makeSnapshot(pending)); + writeFileSync(join(dir, 'approvals', 'appr_stale.json.claim'), '1234'); + expect(await store.approvals.resolve('appr_stale')).toBeNull(); + await store.approvals.save(pending, makeSnapshot(pending)); + expect((await store.approvals.resolve('appr_stale'))?.pending.id).toBe('appr_stale'); + }); + + it('rejects session, checkpoint and approval ids that would escape the directory', async () => { + const store = fileStore(tempDir()); + await expect(store.sessions.save('../evil', [])).rejects.toThrow(/Invalid session id/); + await expect(store.sessions.load('a/../../b')).rejects.toThrow(/Invalid session id/); + await expect(save(store, '../evil', 1)).rejects.toThrow(/Invalid session id/); + await expect(store.checkpoints.load('..')).rejects.toThrow(/Invalid session id/); + await expect(store.approvals.resolve('../evil')).rejects.toThrow(/Invalid approval id/); + const pending = makePending('../evil'); + await expect(store.approvals.save(pending, makeSnapshot(pending))).rejects.toThrow(/Invalid approval id/); + }); +}); + +describe('createAgent({ store: fileStore(dir) }) (R2)', () => { + it('a session survives a second createAgent() on the same dir', async () => { + const dir = tempDir(); + await createAgent({ provider: mockModel(['Hi Ali.']), store: fileStore(dir) }).session({ id: 'chat' }).send('My name is Ali.'); + + const model = mockModel(['Ali.']); + const { text } = await createAgent({ provider: model, store: fileStore(dir) }).session({ id: 'chat' }).send('What is my name?'); + + expect(text).toBe('Ali.'); + expect(model.calls[0].messages.some((m) => m.content === 'My name is Ali.')).toBe(true); + }); + + it('an approval-gated run paused on one agent resumes on a second agent over the same dir', async () => { + const dir = tempDir(); + let sent = 0; + const tools = [ + defineTool({ name: 'send_email', description: 'send', input: z.object({}), needsApproval: true, execute: async () => `sent ${++sent}` }), + ]; + const paused = await createAgent({ provider: mockModel([{ toolCalls: [{ name: 'send_email', id: 'call_1' }] }]), tools, store: fileStore(dir) }).send( + 'Email Sam', + { sessionId: 'job-1' } + ); + expect(paused.finishReason).toBe('awaiting-approval'); + expect(readdirSync(join(dir, 'approvals'))).toEqual([`${paused.approvalId}.json`]); + + const agent = createAgent({ provider: mockModel(['Email sent.']), tools, store: fileStore(dir) }); + const result = await agent.approvals.resolve({ id: paused.approvalId!, approved: true }); + + expect(result.text).toBe('Email sent.'); + expect(sent).toBe(1); + expect(await fileStore(dir).checkpoints.load('job-1')).toMatchObject({ status: 'finished' }); + expect(readdirSync(join(dir, 'approvals'))).toEqual([]); + }); +}); diff --git a/src/storage/fileStore.ts b/src/storage/fileStore.ts new file mode 100644 index 00000000..95ffaeac --- /dev/null +++ b/src/storage/fileStore.ts @@ -0,0 +1,207 @@ +/** + * `fileStore(dir)` (R2): an `AgentStore` of plain JSON files, for a Node + * process that should keep sessions, checkpoints and paused approvals across + * restarts without a database: + * + * `/sessions/.json` a transcript (`FileSessionStore`) + * `/checkpoints/.json` the latest checkpoint of a run + * `/checkpoint-history/.json` its bounded history, oldest first + * `/approvals/.json` a pending approval and its snapshot + * + * Every write goes to a temp file and is renamed into place, so a crash never + * leaves half a file. No lock files: a crashed writer cannot block anyone. + * Two processes saving one session's checkpoint at the same time can lose one + * history entry (the ring is a read-modify-write); resolving an approval is + * safe across processes (see `FileApprovalStore`). + */ + +import { mkdir, readFile, rename, rm, writeFile } from 'node:fs/promises'; +import { dirname, join, resolve } from 'node:path'; +import { randomUUID } from 'node:crypto'; +import type { ApprovalStore, ExecutionSnapshot, PendingApproval, ResolvedApproval } from '../execution/ApprovalGate'; +import { + appendToRing, + newestFirst, + resolveHistoryLimit, + toHistoryEntry, + type Checkpoint, + type CheckpointDeleteOptions, + type CheckpointHistoryEntry, + type CheckpointHistoryOptions, + type CheckpointStore, +} from '../execution/checkpoint'; +import { ConfigurationError } from '../execution/errors'; +import { decodeBytes, encodeBytes, FileSessionStore } from '../session/sessionStore'; +import type { AgentStore } from './agentStore'; + +/** Options of {@link fileStore}. */ +export interface FileStoreOptions { + /** Checkpoints kept per session in `checkpoints.history()` (default 50, `0` keeps none). */ + historyLimit?: number; +} + +/** + * Checkpoint ids are session ids plus the `.turn-` / `.fork-` suffixes + * the SDK appends, so a `.` is allowed, but not first (no `.`, `..` or hidden + * files) and never a path separator. + */ +const CHECKPOINT_ID_PATTERN = /^[A-Za-z0-9_-][A-Za-z0-9_.-]{0,199}$/; + +/** Approval ids arrive from HTTP input and become file names. */ +const APPROVAL_ID_PATTERN = /^[A-Za-z0-9_-]{1,128}$/; + +function assertId(id: string, pattern: RegExp, what: string): void { + if (typeof id !== 'string' || !pattern.test(id)) { + throw new ConfigurationError(`Invalid ${what} ${JSON.stringify(id)}: it becomes a file name, so use letters, digits, '_' and '-'.`, 'id'); + } +} + +function isMissing(error: unknown): boolean { + return (error as NodeJS.ErrnoException).code === 'ENOENT'; +} + +/** `undefined` when `file` does not exist. */ +async function readText(file: string): Promise { + try { + return await readFile(file, 'utf8'); + } catch (error) { + if (isMissing(error)) return undefined; + throw error; + } +} + +/** Write to a unique temp file next to `file` and rename it over `file`; creates the directory when it is missing. */ +async function writeAtomic(file: string, value: unknown): Promise { + const temp = `${file}.${process.pid}.${randomUUID()}.tmp`; + const content = JSON.stringify(value, encodeBytes); + try { + try { + await writeFile(temp, content, 'utf8'); + } catch (error) { + if (!isMissing(error)) throw error; + await mkdir(dirname(file), { recursive: true }); + await writeFile(temp, content, 'utf8'); + } + await rename(temp, file); + } catch (error) { + await rm(temp, { force: true }); + throw error; + } +} + +/** Checkpoints and their history as JSON files; created on first write. */ +class FileCheckpointStore implements CheckpointStore { + private readonly historyLimit: number; + + constructor( + private readonly checkpointDir: string, + private readonly historyDir: string, + options: FileStoreOptions + ) { + this.historyLimit = resolveHistoryLimit(options.historyLimit); + } + + private fileFor(dir: string, sessionId: string): string { + assertId(sessionId, CHECKPOINT_ID_PATTERN, 'session id'); + return join(dir, `${sessionId}.json`); + } + + /** The oldest-first ring; a missing or unreadable (half-written) file counts as empty. */ + private async readRing(sessionId: string): Promise { + const raw = await readText(this.fileFor(this.historyDir, sessionId)); + if (raw === undefined) return []; + try { + const parsed: unknown = JSON.parse(raw, decodeBytes); + return Array.isArray(parsed) ? (parsed as CheckpointHistoryEntry[]) : []; + } catch { + return []; + } + } + + async save(sessionId: string, checkpoint: Checkpoint): Promise { + const file = this.fileFor(this.checkpointDir, sessionId); + await writeAtomic(file, checkpoint); + if (this.historyLimit === 0) return; + const ring = appendToRing(await this.readRing(sessionId), toHistoryEntry(checkpoint), this.historyLimit); + await writeAtomic(this.fileFor(this.historyDir, sessionId), ring); + } + + async load(sessionId: string): Promise { + const raw = await readText(this.fileFor(this.checkpointDir, sessionId)); + return raw === undefined ? null : (JSON.parse(raw, decodeBytes) as Checkpoint); + } + + async delete(sessionId: string, options: CheckpointDeleteOptions = {}): Promise { + await rm(this.fileFor(this.checkpointDir, sessionId), { force: true }); + if (!options.keepHistory) await rm(this.fileFor(this.historyDir, sessionId), { force: true }); + } + + async history(sessionId: string, options?: CheckpointHistoryOptions): Promise { + return newestFirst(await this.readRing(sessionId), options); + } +} + +/** + * One JSON file per pending approval. `resolve` first creates `.json.claim` + * with an exclusive create (`wx`, atomic on POSIX and NTFS), so of two callers + * resolving one approval, in one process or two, exactly one gets the record. + * (A rename is not a safe claim on Windows: two renames of one file can both + * succeed.) A resolver that crashes after claiming leaves the claim file, and + * the approval then resolves to `null`; saving it again clears the claim. + */ +class FileApprovalStore implements ApprovalStore { + constructor(private readonly dir: string) {} + + private fileFor(id: string): string { + assertId(id, APPROVAL_ID_PATTERN, 'approval id'); + return join(this.dir, `${id}.json`); + } + + async save(pending: PendingApproval, snapshot: ExecutionSnapshot): Promise { + const file = this.fileFor(pending.id); + const record: ResolvedApproval = { pending, snapshot }; + await writeAtomic(file, record); + await rm(`${file}.claim`, { force: true }); + } + + async resolve(id: string): Promise { + const file = this.fileFor(id); + const claim = `${file}.claim`; + try { + await writeFile(claim, String(process.pid), { flag: 'wx' }); + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if (code === 'EEXIST' || code === 'ENOENT') return null; // another caller has it, or nothing was ever saved + throw error; + } + try { + const raw = await readText(file); + if (raw === undefined) return null; + await rm(file, { force: true }); + return JSON.parse(raw, decodeBytes) as ResolvedApproval; + } finally { + await rm(claim, { force: true }); + } + } +} + +/** + * An {@link AgentStore} of plain, inspectable JSON files under `dir` + * (`sessions/`, `checkpoints/`, `checkpoint-history/`, `approvals/`), for + * `createAgent({ store })` in a Node process. Directories are created on + * first write. + * + * @example + * ```ts + * const agent = createAgent({ provider, store: fileStore('./.lousho') }); + * await agent.session({ id: 'user-42' }).send('Hello'); + * ``` + */ +export function fileStore(dir: string, options: FileStoreOptions = {}): Required { + const root = resolve(dir); + return { + sessions: new FileSessionStore(join(root, 'sessions')), + checkpoints: new FileCheckpointStore(join(root, 'checkpoints'), join(root, 'checkpoint-history'), options), + approvals: new FileApprovalStore(join(root, 'approvals')), + }; +} diff --git a/tsup.config.ts b/tsup.config.ts index 8d176190..bc03b2b7 100644 --- a/tsup.config.ts +++ b/tsup.config.ts @@ -40,6 +40,7 @@ export default defineConfig({ 'execution/otel': 'src/execution/otel.ts', 'execution/hooks': 'src/execution/hooks.ts', 'storage/sqlite/index': 'src/storage/sqlite/index.ts', + 'deploy/kv': 'src/deploy/kv.ts', 'traces/index': 'src/traces/index.ts', 'triggers/index': 'src/triggers/index.ts', 'react/index': 'src/react/index.ts',