From 284bb834cae08cc6f7aa997dab85f1ec5219a88c Mon Sep 17 00:00:00 2001 From: Ali Mohammad Date: Sun, 4 Oct 2026 12:08:06 +0300 Subject: [PATCH] channels: bind restarted click decisions to the checkpointed session; approvals.get() reads the durable store (#279, #280) A channel button click after a restart resolved the approval but the continuation ran outside the session, so it never reached the transcript; and a function approvers saw only this process's pending list, so it failed closed after a restart. - mountChannels: a decision with inbound that this process did not pause on is checked against the conversation's checkpointed turn (exact approval id, matching kind, no sign-in), then bound via session.resume() (SessionAwaitingApprovalError) so agent.approvals.resolve() continues the turn inside the session and records the transcript. Mismatched, stale or replayed decisions are refused with LOUSHO_APPROVAL_NOT_FOUND before the approval is touched and without an audit entry. - ApprovalStore gains optional load(id) - resolve() without claiming - implemented by the in-memory, StorageService, file, SQLite, KV and agent-forge stores; agent.approvals.get(id) reads through it and ctx.approval() falls back to it, so a function approvers sees the recorded request (toolName, input, sessionId, principal) after a restart and fails closed only when no store knows the id. --- CHANGELOG.md | 2 + apps/agent-forge/server/approvalStore.ts | 11 ++ docs/approvals.md | 13 +- docs/channels.md | 50 ++++--- llms-full.txt | 63 +++++--- src/channels/channelSupport.ts | 2 +- src/channels/channels.test.ts | 141 ++++++++++++++++++ src/channels/defineChannel.ts | 10 +- src/channels/discordChannel.test.ts | 54 +++++++ src/channels/githubChannel.test.ts | 21 +++ src/channels/mountChannels.ts | 76 ++++++++-- src/channels/slackChannel.test.ts | 60 ++++++++ src/channels/teamsChannel.test.ts | 16 ++ src/channels/telegramChannel.test.ts | 16 ++ src/createAgentApprovals.test.ts | 19 +++ src/createAgentApprovals.ts | 13 ++ src/deploy/kvStore.ts | 5 + src/execution/ApprovalGate.ts | 20 +++ src/execution/InMemoryApprovalStore.ts | 5 + src/storage/fileStore.ts | 7 + src/storage/sqlite/SqliteStore.test.ts | 4 + .../sqlite/__fixtures__/storeContracts.ts | 11 ++ src/storage/sqlite/stores.ts | 5 + 23 files changed, 564 insertions(+), 60 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4db31075..954717b6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -67,6 +67,8 @@ This section lists what is on `main` and not yet on npm. - `SlackTriggerAdapter`, `CronTriggerAdapter` and `WebhookTriggerAdapter` (`@lousho/build-ai-agent/triggers`) are deprecated and will be removed in a future major version. Use `slackChannel()` for Slack (sessions per thread, approval buttons), `defineSchedule()` with `startSchedules()` for cron (or `schedules/` in an agent directory, or cron triggers in a spec file), and `webhookChannel()` for webhooks, each mounted with `mountChannels()`. `verifySlackSignature()`, the `WebhookAuth` helpers, `parseCronExpression()`, `TriggerAdapter` and `TriggerRegistry` are not deprecated. Their options types (`SlackTriggerAdapterOptions`, `CronTriggerAdapterOptions` with `CronIntervalOptions` and `CronExpressionOptions`, `WebhookTriggerAdapterOptions`) and `WebhookTriggerHandle` carry `@deprecated` too. Nothing changes at run time. See docs/triggers.md. ### Fixed +- A channel button click (Slack, Discord), inline-keyboard tap (Telegram), Adaptive Card press (Teams) or `/approve` command (GitHub) after a restart is now bound to the conversation's session before the approval is resolved, so its continuation is appended to the session transcript again (#279). The decision names the conversation it came from (`inbound`); with checkpointed sessions (`mountChannels` `store` with `checkpoints`, e.g. a `SqliteStore`) it is accepted only when that session's pending turn waits on exactly that approval id and kind - a decision for another conversation, a replay of a decided one and a decision of the wrong kind are refused (`LOUSHO_APPROVAL_NOT_FOUND`) before the approval is touched, and `onDecision` audits only decisions that pass, so a replayed click cannot run the tool twice and the thread's next message continues the same session. Without a checkpoint store nothing can be checked, so the approval store alone decides, as before. +- A channel's function `approvers` works after a restart (#280): `ctx.approval(id)` asked only this process's `agent.approvals.list()`, so a pause the restarted process did not make looked the same as a denial and the function failed closed. `ApprovalStore` gains an optional `load(id)` - `resolve()` without removing the record; every shipped store implements it, and a custom store without it keeps compiling and behaves as before - and the new `agent.approvals.get(id)` reads through it. `ctx.approval` falls back to that, so the function sees the same request (`toolName`, `input`, `sessionId`, `principal`) it saw before the restart and fails closed only when no store knows the id (a process-local approval store still forgets its pauses at a restart). - `lousho` CLI surface (LOU-R19): `lousho dev`/`lousho chat --traces` now writes traces for a `.ts`/`.js` module that exports an already-built `createAgent()` agent - the module's export could not take the flag's exporter after the fact, so nothing was written. `createAgent()` remembers the config it was built with (readable via the internal `createAgentConfigOf()`), and the loader rebuilds the agent with the flag's exporter (the module's own `exporter` is replaced while the flag is set; a module exporting an agent not made by `createAgent()` prints a warning and keeps its own). A spec path that does not exist (e.g. `lousho chat missing.yaml`) now fails with the coded `LOUSHO_SPEC_NOT_FOUND` instead of a raw ENOENT. `lousho add --dry-run` no longer fails when a target file already exists; it prints the file list with an "(exists - a real install needs --overwrite)" marker and exits 0 (the symlink checks still run). Top-level `lousho --help` lists the implemented flags again: `--no-schedules`/`--traces` on `dev` and `chat`, `eval`'s record/replay/drift/`--config` flags, `init`'s `--sdk-path`, and `build`'s positional form. - A cassette failure keeps its type through `agent.send()` (LOU-R13): a replay mismatch (or a missing cassette raised by the provider-interception seam that `lousho eval` installs) was compacted into a generic `CompactedLLMProviderError`, losing the cassette path, the call number and the re-record hint. `CassetteMismatchError` is now an `SDKError` with code `LOUSHO_CASSETTE_INVALID` (docs/errors.md already claimed cassette errors carry it), and the run loop rethrows cassette errors untouched instead of compacting them, so `send()` rejects with the `CassetteMismatchError` itself (`.cassette`, `.callNumber`). `setProviderInterceptor()` and the `ProviderInterceptor` type are now exported from `@lousho/build-ai-agent/testing`, the public entry for the model-boundary seam docs/evals.md describes. - A tool call approved with `agent.approvals.resolve()` (or `resumeAfterApproval()`) now has its `execute_tool` span (#281). It ran before the continued run started, outside any span, so OpenTelemetry and `fileTraceExporter()` showed the continued run's `invoke_agent` and `chat` spans but not the approved tool. The continued run's `invoke_agent` span now opens first and the approved call runs inside it: its `execute_tool` span is a child of it (a `run_code` call's inner calls and a sub-agent the tool starts are children of the tool span), with the same `captureContent` / `redactContent` rules as any tool span, and no span is made for a rejected call. A pause that happens again (sign-in) leaves the `invoke_agent` span with the tool span and no `chat`. diff --git a/apps/agent-forge/server/approvalStore.ts b/apps/agent-forge/server/approvalStore.ts index f320967f..b4985c19 100644 --- a/apps/agent-forge/server/approvalStore.ts +++ b/apps/agent-forge/server/approvalStore.ts @@ -61,4 +61,15 @@ export class FileApprovalStore implements ApprovalStore { if (!fs.existsSync(file)) return null; return JSON.parse(fs.readFileSync(file, 'utf8')) as ResolvedApproval; } + + /** `ApprovalStore.load` (#280): like `resolve`, it scans every agent's approvals directory for the id, but does not delete. */ + async load(approvalId: string): Promise { + const agentsDir = path.join(this.baseDir, '.lousho', 'agents'); + if (!fs.existsSync(agentsDir)) return null; + for (const agentId of fs.readdirSync(agentsDir)) { + const found = await this.peek(agentId, approvalId); + if (found) return found; + } + return null; + } } diff --git a/docs/approvals.md b/docs/approvals.md index 12378c5c..af280988 100644 --- a/docs/approvals.md +++ b/docs/approvals.md @@ -196,6 +196,10 @@ task. - `agent.approvals.list()` returns the pending calls this agent paused on in this process, oldest first: `{ id, toolCallId, toolName, args, createdAt }`, plus `subagentPath` when the call belongs to a sub-agent. +- `agent.approvals.get(id)` returns one pending call without deciding it: + `list()`'s entries first, then - with an `approvalStore` that implements + `load()` - a pause saved before a restart, `undefined` when the id is + unknown or already resolved. - `agent.approvals.resolve({ id, approved, note? })` runs the call (approved) or gives the model a rejection with your `note` (rejected), continues the run, and resolves with the continued run's result, which may pause again. @@ -219,10 +223,11 @@ const store = new SqliteStore('./.lousho/agent.db'); const agent = createAgent({ provider, tools: [emailTool], store }); // or approvalStore: store.approvals ``` -`list()` only knows the pauses made by this agent object; keep the -`approvalId` (or read the store) to resolve a pause from somewhere else. A -continued run joins a session only when it is resolved through the agent that -owns that session object. +`list()` only knows the pauses made by this agent object; `get(id)` also finds +a pause the durable store still holds (for example one saved before a +restart). Keep the `approvalId` (or read the store) to resolve a pause from +somewhere else. A continued run joins a session only when it is resolved +through the agent that owns that session object. ### Resuming with a changed agent diff --git a/docs/channels.md b/docs/channels.md index 0f2b927a..53ba0acb 100644 --- a/docs/channels.md +++ b/docs/channels.md @@ -205,9 +205,11 @@ and the approval stays pending. message author on Slack, the command's user on Discord) may approve. In earlier versions anyone who could see the message could; `approvers: () => true` restores that, which you should only do in a private channel. A function -`approvers` fails closed when the process does not know the pending call (after -a restart); use the list form or the default for approvals that must survive -one. Answers to an `ask_question` are not restricted: the next message (Slack) +`approvers` sees the pending call after a restart too: `ctx.approval(id)` asks +the durable approval store when this process did not pause on the id, so it +decides the same way it did before. Only with a process-local approval store +(the default `InMemoryApprovalStore` - then nothing survives a restart anyway) +does it still fail closed. Answers to an `ask_question` are not restricted: the next message (Slack) or `/ask` (Discord) in the conversation is the answer. Who decided goes to `mountChannels(agent, channels, { onDecision({ approver, decision, sessionId, channel }) {} })`. The Telegram channel takes the same `approvers` (Telegram user ids as strings; a @@ -281,8 +283,11 @@ with a `SqliteStore`, and the agent's `approvalStore`): the next message in the thread is still the answer, the turn continues, and the question, the answer and the reply are appended to the session transcript. A message in a thread whose session waits on a tool approval does not decide it; only the buttons -do. After a restart, a button click's continuation is posted but not appended -to the session transcript. +do. A button click after a restart is bound to the thread's session first - +the click must name a thread whose checkpointed turn waits on exactly that +approval - so the continuation is appended to the session transcript, a +replayed click does not run the tool twice, and a click in another thread is +refused. ## Discord @@ -348,8 +353,9 @@ Pending approvals are resolved from the button click and the approval store. As on Slack, a pending `ask_question` survives a restart given durable stores for sessions, checkpoints and approvals: the next `/ask` in the channel is the answer and the continued turn is appended to the session transcript. A button -click's continuation after a restart is posted but not appended to the session -transcript. +click after a restart names its conversation, so it is bound to that session +before the approval is resolved - the continuation joins its transcript, and a +stale or copied click decides nothing. An [agent directory](./agent-directories.md)'s `channels/*.ts` files are loaded as channels too, and the node server mounts them. `SlackTriggerAdapter` (see [Triggers](triggers.md)), which replies @@ -434,9 +440,10 @@ messages from the method name and status only: the token never reaches Pending approvals are resolved from the tap and the approval store. As on Slack and Discord, a pending `ask_question` survives a restart given durable stores for sessions, checkpoints and approvals: the next message in the chat is the -answer and the continued turn is appended to the session transcript. A tap's -continuation after a restart is sent but not appended to the session -transcript. +answer and the continued turn is appended to the session transcript. As with a +button click on the other channels, a tap after a restart is checked against +the chat session's checkpointed turn, so the continuation joins the transcript +and a replayed tap does not run the tool twice. ## GitHub @@ -515,8 +522,9 @@ before posting `/approve` is refused. It is a snapshot taken at the command, not a live permission check; if a revocation must take effect inside that window, use `approvers` with a function that asks the GitHub API (for example the collaborator permission) before it returns `true`. A function `approvers`, as on -the other channels, fails closed when the process does not know the pending -call (after a restart). +the other channels, sees the pending call after a restart too (the durable +approval store is asked), and fails closed only with a process-local approval +store. Set up a GitHub App: @@ -568,11 +576,12 @@ call and the HTTP status only (`LOUSHO_CHANNEL_REQUEST_FAILED`). Pending approvals are resolved from the command comment and the approval store. A pending `ask_question` is kept in memory until the next comment, and survives -a restart given durable stores, as on the other channels. After a restart, a -command with an id that is no longer pending (decided already, expired) fails -with the generic "Sorry, that request failed." comment instead of being ignored, -and an id is no longer tied to the thread it was posted in. A continuation after -a restart is posted but not appended to the session transcript. +a restart given durable stores, as on the other channels. After a restart the +command names its thread, so - given durable checkpoints - it is bound to that +thread's session before the approval is resolved: the continuation is appended +to the session transcript, and an id the thread is not waiting on (decided +already, or the approval of another thread) fails with the generic "Sorry, +that request failed." comment. ## Microsoft Teams @@ -658,6 +667,7 @@ http.createServer((req, res) => { The app password and the access token are never put in an error message or log line: a failed call is reported to `onError` with the call and the HTTP status only (`LOUSHO_CHANNEL_REQUEST_FAILED`). A press on a card works on a restarted -process, because it names the conversation and the approval itself; the -continuation after a restart is sent but not appended to the session -transcript. +process, because it names the conversation and the approval itself; given +durable checkpoints the press is bound to that session's pending turn first, +so the continuation is appended to the session transcript and a press a +session is not waiting on decides nothing. diff --git a/llms-full.txt b/llms-full.txt index 3c45cf8c..49eb5f52 100644 --- a/llms-full.txt +++ b/llms-full.txt @@ -2261,6 +2261,10 @@ task. - `agent.approvals.list()` returns the pending calls this agent paused on in this process, oldest first: `{ id, toolCallId, toolName, args, createdAt }`, plus `subagentPath` when the call belongs to a sub-agent. +- `agent.approvals.get(id)` returns one pending call without deciding it: + `list()`'s entries first, then - with an `approvalStore` that implements + `load()` - a pause saved before a restart, `undefined` when the id is + unknown or already resolved. - `agent.approvals.resolve({ id, approved, note? })` runs the call (approved) or gives the model a rejection with your `note` (rejected), continues the run, and resolves with the continued run's result, which may pause again. @@ -2284,10 +2288,11 @@ const store = new SqliteStore('./.lousho/agent.db'); const agent = createAgent({ provider, tools: [emailTool], store }); // or approvalStore: store.approvals ``` -`list()` only knows the pauses made by this agent object; keep the -`approvalId` (or read the store) to resolve a pause from somewhere else. A -continued run joins a session only when it is resolved through the agent that -owns that session object. +`list()` only knows the pauses made by this agent object; `get(id)` also finds +a pause the durable store still holds (for example one saved before a +restart). Keep the `approvalId` (or read the store) to resolve a pause from +somewhere else. A continued run joins a session only when it is resolved +through the agent that owns that session object. ### Resuming with a changed agent @@ -3378,9 +3383,11 @@ and the approval stays pending. message author on Slack, the command's user on Discord) may approve. In earlier versions anyone who could see the message could; `approvers: () => true` restores that, which you should only do in a private channel. A function -`approvers` fails closed when the process does not know the pending call (after -a restart); use the list form or the default for approvals that must survive -one. Answers to an `ask_question` are not restricted: the next message (Slack) +`approvers` sees the pending call after a restart too: `ctx.approval(id)` asks +the durable approval store when this process did not pause on the id, so it +decides the same way it did before. Only with a process-local approval store +(the default `InMemoryApprovalStore` - then nothing survives a restart anyway) +does it still fail closed. Answers to an `ask_question` are not restricted: the next message (Slack) or `/ask` (Discord) in the conversation is the answer. Who decided goes to `mountChannels(agent, channels, { onDecision({ approver, decision, sessionId, channel }) {} })`. The Telegram channel takes the same `approvers` (Telegram user ids as strings; a @@ -3454,8 +3461,11 @@ with a `SqliteStore`, and the agent's `approvalStore`): the next message in the thread is still the answer, the turn continues, and the question, the answer and the reply are appended to the session transcript. A message in a thread whose session waits on a tool approval does not decide it; only the buttons -do. After a restart, a button click's continuation is posted but not appended -to the session transcript. +do. A button click after a restart is bound to the thread's session first - +the click must name a thread whose checkpointed turn waits on exactly that +approval - so the continuation is appended to the session transcript, a +replayed click does not run the tool twice, and a click in another thread is +refused. ## Discord @@ -3521,8 +3531,9 @@ Pending approvals are resolved from the button click and the approval store. As on Slack, a pending `ask_question` survives a restart given durable stores for sessions, checkpoints and approvals: the next `/ask` in the channel is the answer and the continued turn is appended to the session transcript. A button -click's continuation after a restart is posted but not appended to the session -transcript. +click after a restart names its conversation, so it is bound to that session +before the approval is resolved - the continuation joins its transcript, and a +stale or copied click decides nothing. An [agent directory](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/agent-directories.md)'s `channels/*.ts` files are loaded as channels too, and the node server mounts them. `SlackTriggerAdapter` (see [Triggers](https://github.com/LinuxDevil/agent-sdk/blob/main/docs/triggers.md)), which replies @@ -3607,9 +3618,10 @@ messages from the method name and status only: the token never reaches Pending approvals are resolved from the tap and the approval store. As on Slack and Discord, a pending `ask_question` survives a restart given durable stores for sessions, checkpoints and approvals: the next message in the chat is the -answer and the continued turn is appended to the session transcript. A tap's -continuation after a restart is sent but not appended to the session -transcript. +answer and the continued turn is appended to the session transcript. As with a +button click on the other channels, a tap after a restart is checked against +the chat session's checkpointed turn, so the continuation joins the transcript +and a replayed tap does not run the tool twice. ## GitHub @@ -3688,8 +3700,9 @@ before posting `/approve` is refused. It is a snapshot taken at the command, not a live permission check; if a revocation must take effect inside that window, use `approvers` with a function that asks the GitHub API (for example the collaborator permission) before it returns `true`. A function `approvers`, as on -the other channels, fails closed when the process does not know the pending -call (after a restart). +the other channels, sees the pending call after a restart too (the durable +approval store is asked), and fails closed only with a process-local approval +store. Set up a GitHub App: @@ -3741,11 +3754,12 @@ call and the HTTP status only (`LOUSHO_CHANNEL_REQUEST_FAILED`). Pending approvals are resolved from the command comment and the approval store. A pending `ask_question` is kept in memory until the next comment, and survives -a restart given durable stores, as on the other channels. After a restart, a -command with an id that is no longer pending (decided already, expired) fails -with the generic "Sorry, that request failed." comment instead of being ignored, -and an id is no longer tied to the thread it was posted in. A continuation after -a restart is posted but not appended to the session transcript. +a restart given durable stores, as on the other channels. After a restart the +command names its thread, so - given durable checkpoints - it is bound to that +thread's session before the approval is resolved: the continuation is appended +to the session transcript, and an id the thread is not waiting on (decided +already, or the approval of another thread) fails with the generic "Sorry, +that request failed." comment. ## Microsoft Teams @@ -3831,9 +3845,10 @@ http.createServer((req, res) => { The app password and the access token are never put in an error message or log line: a failed call is reported to `onError` with the call and the HTTP status only (`LOUSHO_CHANNEL_REQUEST_FAILED`). A press on a card works on a restarted -process, because it names the conversation and the approval itself; the -continuation after a restart is sent but not appended to the session -transcript. +process, because it names the conversation and the approval itself; given +durable checkpoints the press is bound to that session's pending turn first, +so the continuation is appended to the session transcript and a press a +session is not waiting on decides nothing. # CLI diff --git a/src/channels/channelSupport.ts b/src/channels/channelSupport.ts index a8075252..1a058af6 100644 --- a/src/channels/channelSupport.ts +++ b/src/channels/channelSupport.ts @@ -37,7 +37,7 @@ export function decodeApprovalRef(value: string): ApprovalRef { return { starter: at > 0 ? value.slice(0, at) : undefined, id: value.slice(at + 1) }; } -/** Whether `user` may decide the approval in `ref` (see {@link Approvers}). A function sees the pending call, so it fails closed when this process does not know it. */ +/** Whether `user` may decide the approval in `ref` (see {@link Approvers}). A function sees the pending call - `ctx.approval` answers from the durable approval store after a restart too - and fails closed when no store knows it. */ export async function mayApprove(approvers: Approvers | undefined, user: ChannelUser | undefined, ref: ApprovalRef, ctx: ChannelContext, sessionKey: string): Promise { if (!user) return false; if (approvers === undefined) return ref.starter === user.id; diff --git a/src/channels/channels.test.ts b/src/channels/channels.test.ts index cd9130a7..a56bf559 100644 --- a/src/channels/channels.test.ts +++ b/src/channels/channels.test.ts @@ -255,6 +255,147 @@ describe('defineChannel / mountChannels (LOU-P7)', () => { }); }); + describe('button decisions are checked against the pending turn (#279)', () => { + const ask = { toolCalls: [{ name: 'ask_question', args: { question: 'Which city?' }, id: 'call_q' }] }; + + /** + * A channel with buttons: `{ user, text }` is a message, `{ user, click, approved }` a click on + * approval `click` in `user`'s conversation (the conversation as the click names it). + */ + function clickingChannel() { + let context: ChannelContext | undefined; + const recorded = recordingChannel({ + async parse(req, _respond, ctx) { + context = ctx; + const body = JSON.parse(req.text) as { user: string; text?: string; click?: string; approved?: boolean }; + const inbound = { sessionKey: body.user, input: body.text ?? '', replyTo: `dm:${body.user}` }; + return body.click ? { decision: { id: body.click, approved: body.approved ?? true }, inbound, approver: { id: body.user } } : inbound; + }, + }); + return { ...recorded, ctx: () => context! }; + } + + function mount(responses: Parameters[0], stores: ReturnType, execute = vi.fn(async ({ to }: { to: string }) => `sent to ${to}`)) { + const tool = defineTool({ name: 'send_email', description: 'Sends an email', input: z.object({ to: z.string() }), needsApproval: true, execute }); + const model = mockModel(responses); + const agent = createAgent({ provider: model, tools: [tool], askQuestion: true, approvalStore: stores.approvalStore }); + const onDecision = vi.fn(); + const recorded = clickingChannel(); + return { ...recorded, model, agent, execute, onDecision, handler: mountChannels(agent, [recorded.channel], { store: stores.store, onDecision }) }; + } + + /** Pauses `user`'s conversation on send_email and returns the approval id. */ + async function pauseOn(t: ReturnType, user: string): Promise { + const before = new Set((await t.agent.approvals.list()).map((request) => request.id)); + await post(t.handler, '/channels/test', { user, text: 'Email Sam' }); + return (await t.agent.approvals.list()).find((request) => !before.has(request.id))!.id; + } + + it('after a restart: the continuation of a click is appended to the session transcript', async () => { + const stores = durableStores(); + const first = mount([callEmail], stores); + const ali = await pauseOn(first, 'ali'); + + const second = mount(['Email sent.', 'You are welcome.'], stores, first.execute); + expect(await post(second.handler, '/channels/test', { user: 'ali', click: ali })).toMatchObject({ status: 200 }); + + expect(first.execute).toHaveBeenCalledTimes(1); + expect(second.replies.map((r) => [r.inbound.replyTo, r.text])).toEqual([['dm:ali', 'Email sent.']]); + const transcript = await stores.transcript(); + for (const text of ['Email Sam', '"call_email"', 'sent to sam@example.com', 'Email sent.']) expect(transcript).toContain(text); + + // the conversation has a session now: the next message is its next turn + await post(second.handler, '/channels/test', { user: 'ali', text: 'Thanks' }); + expect(second.replies.at(-1)?.text).toBe('You are welcome.'); + const userTexts = (second.model.calls[1].messages as Message[]).filter((m) => m.role === 'user').map((m) => m.content); + expect(userTexts).toEqual(['Email Sam', 'Thanks']); + expect(JSON.stringify(second.model.calls[1].messages)).toContain('sent to sam@example.com'); + }); + + it('after a restart: a click naming a conversation whose turn waits on another approval is refused (404); the right one still decides', async () => { + const stores = durableStores(); + const first = mount([callEmail, callEmail], stores); + const ali = await pauseOn(first, 'ali'); + const bob = await pauseOn(first, 'bob'); + + const second = mount(['Email sent.'], stores, first.execute); + const forged = await post(second.handler, '/channels/test', { user: 'bob', click: ali }); + + expect(forged).toMatchObject({ status: 404, json: { error: expect.stringContaining(ali) } }); + expect(first.execute).not.toHaveBeenCalled(); + expect(second.onDecision).not.toHaveBeenCalled(); // nothing was decided, so nothing is audited + + expect(await post(second.handler, '/channels/test', { user: 'ali', click: ali })).toMatchObject({ status: 200 }); + expect(first.execute).toHaveBeenCalledTimes(1); + expect(second.replies.map((r) => [r.inbound.replyTo, r.text])).toEqual([['dm:ali', 'Email sent.']]); + expect(second.onDecision).toHaveBeenCalledTimes(1); + expect(bob).not.toBe(ali); + }); + + it('after a restart: a click naming a conversation with nothing pending is refused, and so is a replay', async () => { + const stores = durableStores(); + const first = mount([callEmail], stores); + const ali = await pauseOn(first, 'ali'); + + const second = mount(['Email sent.', 'never'], stores, first.execute); + expect(await post(second.handler, '/channels/test', { user: 'carol', click: ali })).toMatchObject({ status: 404 }); + expect(await post(second.handler, '/channels/test', { user: 'ali', click: ali })).toMatchObject({ status: 200 }); + expect(await post(second.handler, '/channels/test', { user: 'ali', click: ali })).toMatchObject({ status: 404 }); + + expect(first.execute).toHaveBeenCalledTimes(1); + expect(second.model.calls).toHaveLength(1); + }); + + it('after a restart: two clicks at once decide once', async () => { + const stores = durableStores(); + const first = mount([callEmail], stores); + const ali = await pauseOn(first, 'ali'); + + const second = mount(['Email sent.', 'never'], stores, first.execute); + const statuses = await Promise.all([1, 2].map(async () => (await post(second.handler, '/channels/test', { user: 'ali', click: ali })).status)); + + expect(statuses.sort()).toEqual([200, 404]); + expect(first.execute).toHaveBeenCalledTimes(1); + expect(second.replies.map((r) => r.text)).toEqual(['Email sent.']); + }); + + it('after a restart: a button decision is not an answer: a click on a pending question is refused', async () => { + const stores = durableStores(); + const first = mount([ask], stores); + await post(first.handler, '/channels/test', { user: 'ali', text: 'Book a trip' }); + const [question] = await first.agent.approvals.list(); + + const second = mount(['never'], stores); + expect(await post(second.handler, '/channels/test', { user: 'ali', click: question.id })).toMatchObject({ status: 404 }); + expect(second.model.calls).toHaveLength(0); + }); + + it('in process: a click naming another conversation than the one that paused is refused', async () => { + const t = mount([callEmail, 'Email sent.'], durableStores()); + const ali = await pauseOn(t, 'ali'); + + expect(await post(t.handler, '/channels/test', { user: 'bob', click: ali })).toMatchObject({ status: 404 }); + expect(t.execute).not.toHaveBeenCalled(); + expect(await post(t.handler, '/channels/test', { user: 'ali', click: ali })).toMatchObject({ status: 200 }); + expect(t.execute).toHaveBeenCalledTimes(1); + expect(t.replies.at(-1)).toMatchObject({ text: 'Email sent.', inbound: { replyTo: 'dm:ali' } }); + }); + + it('ctx.approval(id) answers after a restart from the approval store, undefined once decided (#280)', async () => { + const stores = durableStores(); + const first = mount([callEmail], stores); + const ali = await pauseOn(first, 'ali'); + + const second = mount(['Hi.', 'Email sent.'], stores, first.execute); + await post(second.handler, '/channels/test', { user: 'bob', text: 'hi' }); // captures the second mount's context + expect(await second.ctx().approval(ali)).toMatchObject({ id: ali, toolName: 'send_email', args: { to: 'sam@example.com' } }); + expect(await second.ctx().approval('nope')).toBeUndefined(); + + await post(second.handler, '/channels/test', { user: 'ali', click: ali }); + expect(await second.ctx().approval(ali)).toBeUndefined(); + }); + }); + it('rejects an invalid channel name and duplicate names', () => { expect(() => recordingChannel({ name: 'has space' })).toThrow(/Invalid channel name/); const { channel } = recordingChannel(); diff --git a/src/channels/defineChannel.ts b/src/channels/defineChannel.ts index 45954f90..822bcdf9 100644 --- a/src/channels/defineChannel.ts +++ b/src/channels/defineChannel.ts @@ -74,7 +74,9 @@ export interface ChannelUser { /** * What `parse` returns for a request that decides a pause (a button click) instead of starting a turn. * With `inbound` (the conversation as the click itself names it), a pause that a restarted process - * no longer remembers is still delivered to the right place; `approver` is who decided. + * no longer remembers is still delivered to the right place; `approver` is who decided. The decision + * must belong to that conversation: with checkpointed sessions, its turn must wait on exactly this + * pause, else the decision is refused (`LOUSHO_APPROVAL_NOT_FOUND`). */ export interface ChannelDecision { decision: ChannelApprovalDecision; @@ -94,7 +96,11 @@ export type ChannelErrorHandler = (error: unknown, context: ChannelErrorContext) /** What `parse` may ask the host about: the agent's own state, so a channel keeps none of its own. */ export interface ChannelContext { - /** The pending approval `id` this process knows, if any. */ + /** + * The pending approval `id`, if still pending: this process's pauses, and - + * through `agent.approvals.get()` - pauses a durable `approvalStore` still + * holds, so a function `approvers` sees the request after a restart too (#280). + */ approval(id: string): Promise; /** The session id `sessionKey` maps to. */ sessionId(sessionKey: string): string; diff --git a/src/channels/discordChannel.test.ts b/src/channels/discordChannel.test.ts index c24d430e..d7df325b 100644 --- a/src/channels/discordChannel.test.ts +++ b/src/channels/discordChannel.test.ts @@ -320,6 +320,60 @@ describe('discordChannel (LOU-P6)', () => { expect(execute).toHaveBeenCalledTimes(1); expect(second.calls.at(-1)).toMatchObject({ method: 'POST', url: 'tok-click', body: { content: 'Email sent.' } }); }); + + it('a click after a restart is appended to the transcript, and the next /ask continues the session (#279)', async () => { + const stores = durableStores(); + const execute = vi.fn(async ({ to }: { to: string }) => `sent to ${to}`); + const agentOptions = { tools: [emailTool(execute)], approvalStore: stores.approvalStore }; + const id = await pause(setup([emailCall], agentOptions, { mount: { store: stores.store } })); + + const second = setup(['Email sent.', 'You are welcome.'], agentOptions, { mount: { store: stores.store } }); + await second.send(click(id)); + + expect(execute).toHaveBeenCalledTimes(1); + expect(second.calls).toEqual([{ method: 'POST', url: 'tok-click', body: expect.objectContaining({ content: 'Email sent.' }) }]); + const transcript = await stores.transcript(); + for (const text of ['Email Sam', '"call_email"', 'sent to sam@example.com', 'Email sent.']) expect(transcript).toContain(text); + + await second.send(command('Thanks')); + expect(second.calls[1]).toMatchObject({ method: 'PATCH', url: 'tok-Thanks/messages/@original', body: { content: 'You are welcome.' } }); + expect(second.userTexts(1)).toEqual(['Email Sam', 'Thanks']); + expect(JSON.stringify(second.model.calls[1].messages)).toContain('sent to sam@example.com'); + }); + + it('a click naming another channel than the one that paused is refused after a restart (#279)', async () => { + const stores = durableStores(); + const execute = vi.fn(async ({ to }: { to: string }) => `sent to ${to}`); + const onError = vi.fn(); + const agentOptions = { tools: [emailTool(execute)], approvalStore: stores.approvalStore }; + const id = await pause(setup([emailCall], agentOptions, { mount: { store: stores.store } })); + + const second = setup(['Email sent.'], agentOptions, { mount: { store: stores.store, onError } }); + await second.send(click(id, 'U1', { channel_id: 'C2', channel: { id: 'C2', type: 0 } })); + expect(execute).not.toHaveBeenCalled(); + expect(onError).toHaveBeenCalledWith(expect.objectContaining({ code: 'LOUSHO_APPROVAL_NOT_FOUND' }), expect.objectContaining({ stage: 'approval' })); + + // still pending: the click in the channel it was posted in decides it + await second.send(click(id)); + expect(execute).toHaveBeenCalledTimes(1); + expect(second.calls.at(-1)).toMatchObject({ method: 'POST', url: 'tok-click', body: { content: 'Email sent.' } }); + }); + + it('a function approver still decides after a restart (#280)', async () => { + const stores = durableStores(); + const execute = vi.fn(async ({ to }: { to: string }) => `sent to ${to}`); + const approvers = vi.fn(async (user: { id: string }) => user.id === 'UADMIN'); + const agentOptions = { tools: [emailTool(execute)], approvalStore: stores.approvalStore }; + const id = await pause(setup([emailCall], agentOptions, { channel: { approvers }, mount: { store: stores.store } })); + + const second = setup(['Email sent.'], agentOptions, { channel: { approvers }, mount: { store: stores.store } }); + await second.send(click(id, 'U1')); + expect(execute).not.toHaveBeenCalled(); + await second.send(click(id, 'UADMIN')); + + expect(execute).toHaveBeenCalledTimes(1); + expect(second.calls.at(-1)).toMatchObject({ method: 'POST', url: 'tok-click', body: { content: 'Email sent.' } }); + }); }); }); diff --git a/src/channels/githubChannel.test.ts b/src/channels/githubChannel.test.ts index ffad8601..bbfcd453 100644 --- a/src/channels/githubChannel.test.ts +++ b/src/channels/githubChannel.test.ts @@ -524,6 +524,27 @@ describe('githubChannel (N11b)', () => { ]); }); + it('after a restart with durable stores, an id is tied to its thread again and the continuation joins the transcript (#279)', async () => { + const stores = durableStores(); + const execute = vi.fn(async ({ to }: { to: string }) => `sent to ${to}`); + const onError = vi.fn(); + const agentOptions = { tools: [emailTool(execute)], approvalStore: stores.approvalStore }; + const id = await pause(setup([emailCall], agentOptions, { mount: { store: stores.store } }), { number: 7 }); + + const second = setup(['Email sent.'], agentOptions, { mount: { store: stores.store }, channel: { onError } }); + await second.send(issueComment(`/approve ${id}`, { login: 'maintainer', association: 'OWNER', number: 8 })); + expect(execute).not.toHaveBeenCalled(); + expect(onError).toHaveBeenCalledWith(expect.objectContaining({ code: 'LOUSHO_APPROVAL_NOT_FOUND' }), expect.objectContaining({ stage: 'approval' })); + + // still pending (this process treats the id as used, so the next one decides it) + const third = setup(['Email sent.'], agentOptions, { mount: { store: stores.store } }); + await third.send(issueComment(`/approve ${id}`, { login: 'maintainer', association: 'OWNER', number: 7 })); + expect(execute).toHaveBeenCalledTimes(1); + expect(third.calls.at(-1)).toMatchObject({ path: '/repos/acme/widgets/issues/7/comments' }); + expect(text(third.calls.at(-1)!)).toBe('Email sent.'); + for (const word of ['email Sam', 'sent to sam@example.com', 'Email sent.']) expect(await stores.transcript()).toContain(word); + }); + it('a restart does not widen who may approve: the default rule still applies', async () => { const approvalStore = new InMemoryApprovalStore(); const store = new MemorySessionStore(); diff --git a/src/channels/mountChannels.ts b/src/channels/mountChannels.ts index 802fe012..1da8adad 100644 --- a/src/channels/mountChannels.ts +++ b/src/channels/mountChannels.ts @@ -125,13 +125,17 @@ export function mountChannels( /** Question ids `pendingQuestion` handed out whose answer has not been processed yet. */ const claimed = new Set(); const answered = new WeakSet(); - const sessions = withDefaultStores({ store }).store as SessionStore | undefined; + const stores = withDefaultStores({ store }); + const sessions = stores.store as SessionStore | undefined; + /** #279: `store` has checkpoints, so a session tells which pause its turn waits on, also after a restart. */ + const checkpointed = stores.checkpointStore !== undefined; /** What `channel.parse` may ask: the agent's pending approvals and saved sessions. */ function contextFor(channel: Channel): ChannelContext { const sessionId = (sessionKey: string) => channelSessionId(channel, { sessionKey, input: '', replyTo: undefined }); return { - approval: async (id) => (await agent.approvals.list()).find((request) => request.id === id), + // #280: this process's pauses first, then the durable approval store, so a function `approvers` still decides after a restart. + approval: async (id) => (await agent.approvals.list()).find((request) => request.id === id) ?? (await agent.approvals.get?.(id)), sessionId, hasSession: async (sessionKey) => (await sessions?.load(sessionId(sessionKey))) !== undefined, pendingQuestion: (sessionKey) => { @@ -218,16 +222,55 @@ export function mountChannels( }); } + /** + * #279: a decision this process did not pause on (a click after a restart). With checkpointed + * sessions it must name a conversation whose turn waits on exactly this pause, of the matching + * kind (a button decides a tool call, an answer a question); `resume()` then binds the pause to + * the session (`SessionAwaitingApprovalError`), so the continuation joins its transcript. Anything + * else - another conversation, a replay, a forged id - is refused before the approval is touched. + * Without a checkpoint store the session cannot tell, and the approval store alone decides (as before). + */ + async function bindToSession(sessionId: string, decision: ChannelApprovalDecision): Promise { + const session = agent.session({ id: sessionId, store }); + const pending = await session.pending(); + if (!pending) return !checkpointed; + const isAnswer = typeof decision.answer === 'string'; + const kind = pending.approvalKind ?? 'tool'; + if (pending.status !== 'awaiting-approval' || pending.approvalId !== decision.id || isAnswer !== (kind === 'question') || kind === 'sign-in') return false; + try { + await session.resume(); + } catch (error) { + if (error instanceof SessionAwaitingApprovalError && error.approvalId === decision.id) return true; + throw error; + } + return false; + } + + /** #279: a decision that names no pending turn of this conversation: 404 while the request is open, else a failure (`onError`). */ + function refuse(channel: Channel, id: string, respond: ChannelRespond | undefined): void { + const error = new SDKError(`No pending approval '${id}' in this conversation on channel '${channel.name}'`, 'LOUSHO_APPROVAL_NOT_FOUND'); + if (!respond || answered.has(respond)) throw error; + respond(404, { error: error.message }); + } + /** * Decides the pause `turn` stopped on and delivers the continuation, as the session's next turn. * N10b: `approver` is recorded as who decided (`ctx.approval.by`); the run keeps its own principal. + * `accept` (#279) runs first, in the session's queue: it checks the decision and audits it, or refuses it. */ - function continueTurn(turn: PausedTurn, decision: ChannelApprovalDecision, respond?: ChannelRespond, approver?: Principal): Promise { + function continueTurn( + turn: PausedTurn, + decision: ChannelApprovalDecision, + respond?: ChannelRespond, + approver?: Principal, + accept?: () => Promise + ): Promise { paused.delete(decision.id); const { id, approved, note, answer } = decision; const decided = { ...(approver && { principal: approver }) }; const run = () => guard(turn, 'approval', respond, async () => { + if (accept && !(await accept())) return refuse(turn.channel, id, respond); const decide = () => typeof answer === 'string' ? agent.approvals.answer({ id, answer }, decided) : agent.approvals.resolve({ id, approved: approved === true, note }, decided); const result = await decide().catch((error: unknown) => { @@ -260,15 +303,30 @@ export function mountChannels( } /** - * Resolves the decision when `channel` paused on it in this process, or when the click - * names the conversation itself (`inbound`: it survives a restart); else answers 404. + * Resolves the decision when `channel` paused on it in this process (a click must name that + * conversation), or when the click names the conversation itself (`inbound`: it survives a + * restart) and, with checkpointed sessions, that conversation's turn waits on it (#279); else 404. */ async function decide(channel: Channel, { decision, inbound, approver }: ChannelDecision, respond: ChannelRespond): Promise { const known = paused.get(decision.id); - const turn = inbound ? { channel, inbound, sessionId: channelSessionId(channel, inbound) } : known?.channel === channel ? known : undefined; - if (!turn) return respond(404, { error: `No pending approval '${decision.id}' on channel '${channel.name}'` }); - if (approver) await options.onDecision?.({ decision, approver, sessionId: turn.sessionId, channel: channel.name }); - await continueTurn(turn, decision, respond, approverPrincipal(channel, approver, inbound)); + const turn = inbound ? { channel, inbound, sessionId: channelSessionId(channel, inbound) } : known; + if (!turn || (known && (known.channel !== channel || known.sessionId !== turn.sessionId))) { + return respond(404, { error: `No pending approval '${decision.id}' on channel '${channel.name}'` }); + } + const audit = async () => { + if (approver) await options.onDecision?.({ decision, approver, sessionId: turn.sessionId, channel: channel.name }); + }; + const principal = approverPrincipal(channel, approver, inbound); + if (known) { + await audit(); + return continueTurn(turn, decision, respond, principal); + } + const accept = async (): Promise => { + if (!(await bindToSession(turn.sessionId, decision))) return false; + await audit(); + return true; + }; + await continueTurn(turn, decision, respond, principal, accept); } const handler = async (req: http.IncomingMessage, res: http.ServerResponse): Promise => { diff --git a/src/channels/slackChannel.test.ts b/src/channels/slackChannel.test.ts index 7d47c2f2..ca42a4e7 100644 --- a/src/channels/slackChannel.test.ts +++ b/src/channels/slackChannel.test.ts @@ -394,6 +394,66 @@ describe('slackChannel (LOU-P5)', () => { expect(execute).toHaveBeenCalledTimes(1); expect(second.posts.at(-1)).toEqual({ channel: 'C1', thread_ts: '100.1', text: 'Email sent.' }); }); + + it('a click after a restart is appended to the transcript, and the next thread message continues the session (#279)', async () => { + const stores = durableStores(); + const execute = vi.fn(async ({ to }: { to: string }) => `sent to ${to}`); + const agentOptions = { tools: [emailTool(execute)], approvalStore: stores.approvalStore }; + const value = await pause(setup([emailCall], agentOptions, { mount: { store: stores.store } })); + + const second = setup(['Email sent.', 'You are welcome.'], agentOptions, { mount: { store: stores.store } }); + await second.send(click(value), { form: true }); + + expect(execute).toHaveBeenCalledTimes(1); + expect(second.posts).toEqual([{ channel: 'C1', thread_ts: '100.1', text: 'Email sent.' }]); + const transcript = await stores.transcript(); + for (const text of ['Email Sam', '"call_email"', 'sent to sam@example.com', 'Email sent.']) expect(transcript).toContain(text); + + // the thread has a session now: a message without a mention is its next turn + await second.send(threadMessage('Thanks', '100.3')); + expect(second.posts[1]).toEqual({ channel: 'C1', thread_ts: '100.1', text: 'You are welcome.' }); + expect(second.userTexts(1)).toEqual(['Email Sam', 'Thanks']); + expect(JSON.stringify(second.model.calls[1].messages)).toContain('sent to sam@example.com'); + }); + + it('a replayed click after a restart decides nothing twice (#279)', async () => { + const stores = durableStores(); + const execute = vi.fn(async ({ to }: { to: string }) => `sent to ${to}`); + const onError = vi.fn(); + const agentOptions = { tools: [emailTool(execute)], approvalStore: stores.approvalStore }; + const value = await pause(setup([emailCall], agentOptions, { mount: { store: stores.store } })); + + const second = setup(['Email sent.', 'never'], agentOptions, { mount: { store: stores.store, onError } }); + await second.send(click(value), { form: true }); + await second.send(click(value), { form: true }); + + expect(execute).toHaveBeenCalledTimes(1); + expect(second.model.calls).toHaveLength(1); + expect(onError).toHaveBeenCalledWith(expect.objectContaining({ code: 'LOUSHO_APPROVAL_NOT_FOUND' }), expect.objectContaining({ stage: 'approval' })); + expect(JSON.parse(await stores.transcript()).filter((m: Message) => m.role === 'tool')).toHaveLength(1); + }); + + it('a function approver still decides after a restart (#280)', async () => { + const stores = durableStores(); + const execute = vi.fn(async ({ to }: { to: string }) => `sent to ${to}`); + const seen: unknown[] = []; + const approvers = vi.fn(async (user: { id: string }, request: { toolName: string }) => (seen.push([user, request]), user.id === 'UADMIN')); + const agentOptions = { tools: [emailTool(execute)], approvalStore: stores.approvalStore }; + const value = await pause(setup([emailCall], agentOptions, { channel: { approvers }, mount: { store: stores.store } })); + + const second = setup(['Email sent.'], agentOptions, { channel: { approvers }, mount: { store: stores.store } }); + await second.send(click(value, 'U1'), { form: true }); // the starter is not the approver + expect(execute).not.toHaveBeenCalled(); + await second.send(click(value, 'UADMIN'), { form: true }); + + expect(execute).toHaveBeenCalledTimes(1); + expect(second.posts.at(-1)).toEqual({ channel: 'C1', thread_ts: '100.1', text: 'Email sent.' }); + // the function saw the request from the store: the paused run's facts, unchanged by the restart + expect(seen[1]).toEqual([ + { id: 'UADMIN', name: 'name-UADMIN' }, + { toolName: 'send_email', input: { to: 'sam@example.com' }, sessionId: expect.stringContaining('slack_T1_C1_100_1'), principal: { id: 'U1', type: 'user', authenticator: 'slack' } }, + ]); + }); }); }); diff --git a/src/channels/teamsChannel.test.ts b/src/channels/teamsChannel.test.ts index b6fac9a0..48d60d26 100644 --- a/src/channels/teamsChannel.test.ts +++ b/src/channels/teamsChannel.test.ts @@ -521,6 +521,22 @@ describe('teamsChannel (N11c)', () => { expect(second.calls.map((c) => c.body.text)).toEqual(['Booked Lisbon.']); expect(JSON.stringify(second.model.calls[0].messages)).toContain('Lisbon'); }); + + it('a card press after a restart is appended to the transcript, and the next message continues the session (#279)', async () => { + const stores = durableStores(); + const execute = vi.fn(async ({ to }: { to: string }) => `sent to ${to}`); + const agentOptions = { tools: [emailTool(execute)], approvalStore: stores.approvalStore }; + const approve = await pause(setup([emailCall], agentOptions, { mount: { store: stores.store } })); + + const second = setup(['Email sent.', 'You are welcome.'], agentOptions, { mount: { store: stores.store } }); + await second.send(click(approve)); + expect(execute).toHaveBeenCalledTimes(1); + for (const text of ['"call_email"', 'sent to sam@example.com', 'Email sent.']) expect(await stores.transcript()).toContain(text); + + await second.send(activity({ text: 'Thanks' })); + expect(second.calls.at(-1)?.body).toMatchObject({ text: 'You are welcome.' }); + expect(JSON.stringify(second.model.calls[1].messages)).toContain('sent to sam@example.com'); + }); }); describe('failures', () => { diff --git a/src/channels/telegramChannel.test.ts b/src/channels/telegramChannel.test.ts index a0f92c0d..6bd35bce 100644 --- a/src/channels/telegramChannel.test.ts +++ b/src/channels/telegramChannel.test.ts @@ -398,6 +398,22 @@ describe('telegramChannel (N11a)', () => { expect(execute).toHaveBeenCalledTimes(1); expect(second.calls.at(-1)?.body.text).toBe('Email sent.'); }); + + it('a tap after a restart is appended to the transcript, and the next message continues the session (#279)', async () => { + const stores = durableStores(); + const execute = vi.fn(async ({ to }: { to: string }) => `sent to ${to}`); + const agentOptions = { tools: [emailTool(execute)], approvalStore: stores.approvalStore }; + const approve = await pause(setup([emailCall], agentOptions, { mount: { store: stores.store } })); + + const second = setup(['Email sent.', 'You are welcome.'], agentOptions, { mount: { store: stores.store } }); + await second.send(tap(approve)); + expect(execute).toHaveBeenCalledTimes(1); + for (const text of ['Email Sam', '"call_email"', 'sent to sam@example.com', 'Email sent.']) expect(await stores.transcript()).toContain(text); + + await second.send(message('Thanks', { message_id: 13 })); + expect(second.calls.at(-1)?.body).toMatchObject({ text: 'You are welcome.' }); + expect(JSON.stringify(second.model.calls[1].messages)).toContain('sent to sam@example.com'); + }); }); it('a failed sendMessage goes to onError and the bot token is not in the logged message', async () => { diff --git a/src/createAgentApprovals.test.ts b/src/createAgentApprovals.test.ts index d1ded5c8..e783636b 100644 --- a/src/createAgentApprovals.test.ts +++ b/src/createAgentApprovals.test.ts @@ -119,6 +119,25 @@ describe('createAgent approvals (LOU-D21)', () => { expect(await approvalStore.resolve(paused.approvalId!)).toMatchObject({ pending: { toolName: 'send_email' } }); }); + it('get() returns a pause this process did not make, through a shared approvalStore, without resolving it (#280)', async () => { + const approvalStore = new InMemoryApprovalStore(); + const { tool, execute } = emailTool(); + const first = createAgent({ provider: mockModel([callEmail, 'Email sent.']), tools: [tool], approvalStore }); + const paused = await first.send('Email Sam'); + + // "after a restart": another agent over the same store does not list the pause, but get() finds it + const second = createAgent({ provider: mockModel(['unused']), tools: [emailTool().tool], approvalStore }); + expect(await second.approvals.list()).toEqual([]); + expect(await second.approvals.get(paused.approvalId!)).toMatchObject({ id: paused.approvalId, toolName: 'send_email', args: { to: 'sam@example.com' } }); + expect(await second.approvals.get('nope')).toBeUndefined(); + + // the read did not resolve it: either agent can still decide it + const result = await first.approvals.resolve({ id: paused.approvalId!, approved: true }); + expect(result.text).toBe('Email sent.'); + expect(execute).toHaveBeenCalledTimes(1); + expect(await second.approvals.get(paused.approvalId!)).toBeUndefined(); + }); + it('resolving a pause from a session continues that session', async () => { const { tool } = emailTool(); const model = mockModel([callEmail, 'Email sent.', 'You asked me to email Sam.']); diff --git a/src/createAgentApprovals.ts b/src/createAgentApprovals.ts index f9c844b8..b6ad44c7 100644 --- a/src/createAgentApprovals.ts +++ b/src/createAgentApprovals.ts @@ -55,6 +55,14 @@ export type ApproveToolCall = (request: PendingApproval) => boolean | string | P export interface AgentApprovals { /** Approvals this agent paused on in this process and that are not decided yet, oldest first. */ list(): Promise; + /** + * The pending approval `id`, without deciding it: this process's pauses + * first, then - through an `approvalStore` that implements `load` - a pause + * saved before a restart (the request's facts as it was recorded, so a + * channel's function `approvers` sees the same input then as now, #280). + * `undefined` when `id` is unknown or already resolved. + */ + get(id: string): Promise; /** * Approves or rejects a paused tool call (`note` is passed to the model * with a rejection) and continues the run, resolving with the continued @@ -259,6 +267,11 @@ export function createAgentApprovals(options: { const approvals: AgentApprovals = { list: async () => [...pending.values()], + // #280: the durable store answers too, so a pause this process did not make is found again. + get: async (id) => { + const found = pending.get(id) ?? (await options.store.load?.(id))?.pending; + return found === undefined ? undefined : describeApproval(found); + }, resolve, answer: ({ id, answer }, resolveOptions) => resolve({ id, approved: true, note: answer }, resolveOptions), streamResolve, diff --git a/src/deploy/kvStore.ts b/src/deploy/kvStore.ts index de6ec3b7..cd8c8eae 100644 --- a/src/deploy/kvStore.ts +++ b/src/deploy/kvStore.ts @@ -103,6 +103,11 @@ class KVApprovalStore implements ApprovalStore { await this.kv.delete(`${this.prefix}${id}`); return JSON.parse(raw) as ResolvedApproval; } + + async load(id: string): Promise { + const raw = await this.kv.get(`${this.prefix}${id}`); + return raw === null ? null : (JSON.parse(raw) as ResolvedApproval); + } } /** diff --git a/src/execution/ApprovalGate.ts b/src/execution/ApprovalGate.ts index 3f692b70..c96ebd26 100644 --- a/src/execution/ApprovalGate.ts +++ b/src/execution/ApprovalGate.ts @@ -223,6 +223,13 @@ export interface ResolvedApproval { export interface ApprovalStore { save(pending: PendingApproval, snapshot: ExecutionSnapshot): Promise; resolve(id: string): Promise; + /** + * The record `resolve(id)` would return, without claiming or deleting it: + * `null` when `id` is unknown or already resolved. Optional (like + * `CheckpointStore.history`): `agent.approvals.get()` reads a pause another + * process saved only through it. All built-in stores implement it. + */ + load?(id: string): Promise; } /** @@ -263,4 +270,17 @@ export class StorageServiceApprovalStore implements ApprovalStore { this.storageService.releaseLock(storageKey); } } + + async load(id: string): Promise { + const storageKey = this.getStorageKey(id); + await this.storageService.acquireLock(storageKey); + try { + if (!this.storageService.fileExists(storageKey)) { + return null; + } + return this.storageService.readPlainJSONAttachment(storageKey); + } finally { + this.storageService.releaseLock(storageKey); + } + } } diff --git a/src/execution/InMemoryApprovalStore.ts b/src/execution/InMemoryApprovalStore.ts index e709c43f..8472c7f0 100644 --- a/src/execution/InMemoryApprovalStore.ts +++ b/src/execution/InMemoryApprovalStore.ts @@ -30,4 +30,9 @@ export class InMemoryApprovalStore implements ApprovalStore { this.records.delete(id); return record; } + + async load(id: string): Promise { + const record = this.records.get(id); + return record ? structuredClone(record) : null; + } } diff --git a/src/storage/fileStore.ts b/src/storage/fileStore.ts index bd8be542..2d33674e 100644 --- a/src/storage/fileStore.ts +++ b/src/storage/fileStore.ts @@ -202,6 +202,13 @@ class FileApprovalStore implements ApprovalStore { const raw = await takeFile(this.fileFor(id)); return raw === undefined ? null : (JSON.parse(raw, decodeBytes) as ResolvedApproval); } + + /** A record a resolver already claimed (its `.claim` file was left by a crash) reads as resolved, like `resolve` sees it. */ + async load(id: string): Promise { + const file = this.fileFor(id); + const [raw, claim] = await Promise.all([readText(file), readText(`${file}.claim`)]); + return raw === undefined || claim !== undefined ? null : (JSON.parse(raw, decodeBytes) as ResolvedApproval); + } } /** The JSON files' names in `dir`, without `.json`; `[]` when `dir` does not exist. */ diff --git a/src/storage/sqlite/SqliteStore.test.ts b/src/storage/sqlite/SqliteStore.test.ts index 7711bc9b..2b41c3dd 100644 --- a/src/storage/sqlite/SqliteStore.test.ts +++ b/src/storage/sqlite/SqliteStore.test.ts @@ -71,6 +71,10 @@ function inMemoryApprovalStore(): ApprovalStore { async save(pending, snapshot) { map.set(pending.id, JSON.stringify({ pending, snapshot })); }, + async load(id) { + const raw = map.get(id); + return raw === undefined ? null : (JSON.parse(raw) as ResolvedApproval); + }, async resolve(id) { const raw = map.get(id); if (raw === undefined) return null; diff --git a/src/storage/sqlite/__fixtures__/storeContracts.ts b/src/storage/sqlite/__fixtures__/storeContracts.ts index 4f8c07a8..f364a73e 100644 --- a/src/storage/sqlite/__fixtures__/storeContracts.ts +++ b/src/storage/sqlite/__fixtures__/storeContracts.ts @@ -157,5 +157,16 @@ export function describeApprovalStoreContract(name: string, factory: Factory { + const store = await factory(); + const pending = makePending('a'); + const snapshot = makeSnapshot(pending); + await store.save(pending, snapshot); + expect(store.load).toBeTypeOf('function'); + expect(await store.load!('a')).toEqual({ pending, snapshot }); + expect(await store.resolve('a')).toEqual({ pending, snapshot }); // the read did not claim it + expect(await store.load!('a')).toBeNull(); + }); }); } diff --git a/src/storage/sqlite/stores.ts b/src/storage/sqlite/stores.ts index f8137e19..05b07041 100644 --- a/src/storage/sqlite/stores.ts +++ b/src/storage/sqlite/stores.ts @@ -155,4 +155,9 @@ export class SqliteApprovalStore implements ApprovalStore { return parse(row) ?? null; }); } + + async load(id: string): Promise { + const row = this.sql.get('SELECT payload FROM approvals WHERE id = ? AND resolved_at IS NULL').get(id); + return parse(row) ?? null; + } }