From d42cebadca5ec554c090c607fb550b62dc8b6dd7 Mon Sep 17 00:00:00 2001 From: NarwhalChen Date: Mon, 3 Aug 2026 00:02:26 +0800 Subject: [PATCH] Add durable Agent Host crash recovery --- docs/agent-collaboration.md | 49 +- .../src/electron/agent-host/executor.test.ts | 64 +++ .../app/src/electron/agent-host/executor.ts | 102 +++- .../app/src/electron/agent-host/host.test.ts | 104 ++++ packages/app/src/electron/agent-host/host.ts | 287 +++++++++- .../electron/agent-host/repository.test.ts | 27 +- .../app/src/electron/agent-host/repository.ts | 38 +- .../src/electron/agent-host/workflow.test.ts | 160 ++++++ .../app/src/electron/agent-host/workflow.ts | 506 ++++++++++++++++++ .../electron/ai/structured-task-contracts.ts | 1 + packages/app/src/electron/trace/agent-host.ts | 5 +- .../components/chat/TypingIndicator.tsx | 4 +- .../libs/__tests__/agent-speech.test.ts | 64 +++ .../src/renderer/libs/agent-host-service.ts | 55 +- .../app/src/renderer/libs/agent-speech.ts | 84 ++- packages/app/src/renderer/libs/agent-tasks.ts | 15 +- packages/app/src/renderer/libs/db/database.ts | 22 + packages/app/src/renderer/libs/db/hooks.ts | 2 + packages/app/src/shared/types/agent-host.ts | 85 +++ .../src/shared/types/workspace-perception.ts | 4 + 20 files changed, 1590 insertions(+), 88 deletions(-) create mode 100644 packages/app/src/electron/agent-host/workflow.test.ts create mode 100644 packages/app/src/electron/agent-host/workflow.ts diff --git a/docs/agent-collaboration.md b/docs/agent-collaboration.md index b7d9563d..8ee68cd7 100644 --- a/docs/agent-collaboration.md +++ b/docs/agent-collaboration.md @@ -12,12 +12,12 @@ first — this is the mechanism behind "agents are colleagues, not features". `src/electron/ai/workspace-tools.ts` and answered in the renderer by `resolveWorkspaceQuery` (`src/renderer/libs/workspace-perception.ts`). Visibility is enforced in exactly one function, `canViewChannel` in that same - file, and an invisible channel returns the *same* error as a nonexistent one, + file, and an invisible channel returns the _same_ error as a nonexistent one, so agents cannot probe for hidden rooms. Channel isolation lands by editing that one function. - **What a read returns is what a person sees.** A `read_channel` message carries `replyTo`, `replyCount`, and `reactions` as `{emoji, reactors}` with - reactor *names* resolved, not member ids — a tally an agent cannot read is + reactor _names_ resolved, not member ids — a tally an agent cannot read is not perception. Names are resolved past the current roster: `readChannel` bulk-gets any sender or reactor id missing from it, so people who left still print as names. `replyCount` is counted over the whole transcript, not the @@ -40,7 +40,7 @@ first — this is the mechanism behind "agents are colleagues, not features". - **Typing is read off the stream, never assumed.** An agent counts as typing from `tool-input-start` for `send_message` until that call closes, keyed by the stream's own `toolCallId`. Being offered a message, reading a channel and - thinking are all *working*, not composing — and showing them as typing claims + thinking are all _working_, not composing — and showing them as typing claims a reply is coming when most colleagues will stay quiet. Both the channel path (`agent-host-service.ts`) and the 1:1 path (`use-local-ai-chat.ts`) decide this through one shared function, `typingTransition` in `typing-store.ts`. @@ -51,6 +51,21 @@ first — this is the mechanism behind "agents are colleagues, not features". prepare and execute each turn over `agent-host:request`/`respond`. The renderer half is `RendererAgentHostService` (`src/renderer/libs/agent-host-service.ts`). +- **A run advances through durable semantic checkpoints.** The job snapshot + owns a small `prepare-turn -> provider-turn -> provider-retry? -> finalize` + workflow with immutable checkpoints, node attempts, pending writes and a + provider-effect ledger (`src/electron/agent-host/workflow.ts`). A restart may + resume a node that stopped before an external effect, or commit a node whose + provider receipt was already durable. A provider effect recorded as started + without a receipt becomes `uncertain` and is never replayed automatically. + Older running jobs without this workflow stay `interrupted` rather than + being guessed forward. +- **Speech delivery is idempotent at the Dexie boundary.** Agent Host injects + an effect key and payload hash outside model-controlled input. The renderer + commits the channel message and `agentEffectReceipts` inbox row in one Dexie + transaction; replay returns the original message id and a key reused with + different content is rejected. Messages remain the channel-text authority; + the receipt is only delivery evidence. - **Turns are concurrent.** A turn that reserves no transcript row sets `concurrent: true` and serializes per (conversation, actor) rather than per conversation (`sessionKey` in `src/electron/ai/runtime.ts`), so colleagues @@ -67,7 +82,7 @@ These were decided deliberately. Change them on purpose, not by accident. `src/renderer/libs/agent-templates.ts` are Elena, Mika, Omar, Noah, Vera, Hana, Ivan and Zoe — people, not mascots. Each `systemPrompt` ends in a temperament, not just a job, because a room of eight competent assistants - reads as one assistant eight times. (Template *ids* are still the old mascot + reads as one assistant eight times. (Template _ids_ are still the old mascot names — `sage`, `patch`, `atlas` — and renaming them would orphan every hired row; leave them.) - **Everyone may answer, but the room is not a chorus.** Several agents @@ -82,18 +97,18 @@ These were decided deliberately. Change them on purpose, not by accident. (`elena 最近在忙什么` names Elena as surely as `@Elena`) because the mechanical instruction is what models follow; a description of the principle was not enough. A message that names someone else is theirs - alone: don't answer it, don't answer *for* them, don't rephrase the + alone: don't answer it, don't answer _for_ them, don't rephrase the question back at them. - **Read your own sentence back.** If any colleague here could have written it, it's an echo — say it in your own voice or say nothing. This lives in the same block, plus an unconditional clause forbidding checklists and process descriptions in every channel context. - **Tasks: visible everywhere, controllable only in private.** `manage_task` - (`src/electron/ai/agent-host-tools.ts`) is attached to *every* Agent + (`src/electron/ai/agent-host-tools.ts`) is attached to _every_ Agent Host turn, so an agent can answer "what are you working on?" from any room. Outside a DM, `canControl` is false and only `list`/`inspect` run; `pause`, `resume`, `cancel` and `redirect` return `TASK_CONTROL_UNAVAILABLE`. Two - structural reasons, both load-bearing: a channel turn *is itself* one of the + structural reasons, both load-bearing: a channel turn _is itself_ one of the listed tasks, so a cancel from a channel could kill the run that is speaking; and redirect guidance is stored under a never-reveal-in-public contract (`formatTaskGuidance` in `agent-host-service.ts`) that authoring it from a @@ -104,7 +119,7 @@ These were decided deliberately. Change them on purpose, not by accident. addressed in. `MAX_CHAIN_HOPS` (counting every agent message) contains the blast radius. - **Reply markers are rare.** The boxed quote is a one-line strip (`↩ Name - text`, `message-row.tsx`). The tool description tells agents a marker is for +text`, `message-row.tsx`). The tool description tells agents a marker is for pulling an older or ambiguous message back into view — answering the latest message needs none. - **Reasoning stays on.** With reasoning off, models answer into the void @@ -114,12 +129,12 @@ These were decided deliberately. Change them on purpose, not by accident. ## Providers -| id | transport | notes | -| --- | --- | --- | -| `openai-api` | OpenAI Responses API | `gpt-5.6-luna`, the only model in its catalog, so nothing can pick an expensive one by accident. Default for hired agents. | -| `fireworks-api` | Fireworks Responses API | DeepSeek V4 Flash. `store: false`, `reasoningEffort: "max"`. | -| `claude-code` | Claude Agent SDK | Brings its own tools. | -| `codex-cli` | Codex app-server (local process) | Brings its own tools. | +| id | transport | notes | +| --------------- | -------------------------------- | -------------------------------------------------------------------------------------------------------------------------- | +| `openai-api` | OpenAI Responses API | `gpt-5.6-luna`, the only model in its catalog, so nothing can pick an expensive one by accident. Default for hired agents. | +| `fireworks-api` | Fireworks Responses API | DeepSeek V4 Flash. `store: false`, `reasoningEffort: "max"`. | +| `claude-code` | Claude Agent SDK | Brings its own tools. | +| `codex-cli` | Codex app-server (local process) | Brings its own tools. | Every OpenAI-compatible HTTP provider goes through the **Responses API**, not `/chat/completions` — Fireworks implements `/v1/responses` at the same baseURL, @@ -154,7 +169,7 @@ plus `run_command`. macOS and bubblewrap+seccomp on Linux. On an unsupported platform it falls back to a bare shell, unsandboxed — a known ceiling, marked in place. - **`resolveInSandbox` (`src/electron/ai/sandbox.ts`) is the boundary.** pi's - own cwd only resolves *relative* paths; it does not confine absolute ones. So + own cwd only resolves _relative_ paths; it does not confine absolute ones. So every pi tool call goes through `sandboxedPiTool`, which rewrites `path` through `resolveInSandbox`: both sides are realpath'd (walking up to the nearest existing ancestor, so a not-yet-created file is still symlink-checked), @@ -243,7 +258,7 @@ remember to increment. localStorage, falling back to the workspace `epoch` so pre-install history stays quiet). Own messages never count. The scan floor is the earliest boundary among conversations that actually exist — computing it as - `min(epoch, ...lastSeen)` made the floor *always* the epoch, i.e. all of + `min(epoch, ...lastSeen)` made the floor _always_ the epoch, i.e. all of history re-scanned on every message write. - **The member card is the roster.** `MemberCard` (`src/renderer/components/common/member-card.tsx`) opens from a channel @@ -271,7 +286,7 @@ remember to increment. `ready` gates whether a turn pinned to that provider can run at all. - **Search is visibility-gated and speaks Chinese.** `use-conversation-search.ts` filters through `useVisibleChannelsStrict` — - the same `canViewChannel` — and drops hidden conversations *before* matching; + the same `canViewChannel` — and drops hidden conversations _before_ matching; "strict" means an unresolved check searches nothing rather than leaking. A CJK needle (`hasCJK`) short-circuits to case-folded substring matching, because uFuzzy's tokenizer treats every non-Latin codepoint as a separator, diff --git a/packages/app/src/electron/agent-host/executor.test.ts b/packages/app/src/electron/agent-host/executor.test.ts index d2345545..198070a5 100644 --- a/packages/app/src/electron/agent-host/executor.test.ts +++ b/packages/app/src/electron/agent-host/executor.test.ts @@ -7,6 +7,8 @@ import type { } from "@/shared/types/local-ai"; import type { AgentHostRendererBridge } from "./renderer-bridge"; import { LocalAiAgentHostExecutor } from "./executor"; +import { AgentHost } from "./host"; +import { InMemoryAgentHostJobRepository } from "./repository"; /** Prepares a turn whose transcript is one direct question, unanswered. */ function silentBridge(): AgentHostRendererBridge { @@ -201,4 +203,66 @@ describe("LocalAiAgentHostExecutor", () => { { role: "user", content: "@trusted are you there?" }, ]); }); + + it("persists the normal prepare, provider-effect, and finalize path", async () => { + const runtime = { + startChat: async ( + request: LocalAIChatRequest, + emit: (event: LocalAIStreamEvent) => void, + ) => { + emit({ + type: "interaction", + requestId: request.requestId, + interactionId: "interaction", + kind: "input", + name: "workspace:query", + prompt: "Send a message", + input: { + kind: "send_message", + viewerMemberId: "agent:trusted", + channelId: "channel", + content: "done", + }, + }); + }, + } as unknown as LocalAIRuntimeService; + const host = new AgentHost({ + repository: new InMemoryAgentHostJobRepository(), + executor: new LocalAiAgentHostExecutor(runtime, silentBridge()), + createId: () => "durable-job", + }); + + await host.enqueue({ + channelId: "channel", + channelKind: "channel", + conversationId: "trusted-conversation", + triggerMessageId: "message", + contextMessageIds: ["message"], + mode: "direct", + offeredAgentMemberIds: ["agent:trusted"], + targets: [{ agentId: "trusted", memberId: "agent:trusted" }], + chain: { hops: 0, invoked: ["agent:trusted"] }, + }); + await vi.waitFor(async () => + expect((await host.listJobs())[0].status).toBe("completed"), + ); + + const completed = (await host.listJobs())[0]; + expect(completed.workflow).toMatchObject({ + checkpoint: { step: 3, next: [] }, + effects: [ + expect.objectContaining({ + kind: "provider-turn", + status: "committed", + receipt: { spoke: true, next: ["finalize"] }, + }), + ], + }); + expect(completed.workflow?.checkpoints.map(({ next }) => next)).toEqual([ + ["prepare-turn"], + ["provider-turn"], + ["finalize"], + [], + ]); + }); }); diff --git a/packages/app/src/electron/agent-host/executor.ts b/packages/app/src/electron/agent-host/executor.ts index 69bc3c6c..918c1e7a 100644 --- a/packages/app/src/electron/agent-host/executor.ts +++ b/packages/app/src/electron/agent-host/executor.ts @@ -11,7 +11,11 @@ import type { } from "@/shared/types/local-ai"; import { WORKSPACE_QUERY_INTERACTION } from "@/shared/types/workspace-perception"; import type { AgentHostTraceRecorder } from "@/electron/trace/agent-host"; -import type { AgentHostExecutor } from "./host"; +import type { + AgentHostExecutionControl, + AgentHostExecutor, + AgentHostProviderOutcome, +} from "./host"; import type { AgentHostRendererBridge } from "./renderer-bridge"; /** @@ -75,6 +79,7 @@ export class LocalAiAgentHostExecutor implements AgentHostExecutor { async execute( job: AgentHostJob, emit: (event: AgentHostEvent) => void, + control?: AgentHostExecutionControl, ): Promise { try { const prepared = await this.bridge.request({ @@ -96,8 +101,44 @@ export class LocalAiAgentHostExecutor implements AgentHostExecutor { memberId: job.agentMemberId, }, }; + await control?.prepared(request); this.activeRequests.set(job.id, request.requestId); - if (await this.runTurn(request, job, emit)) return; + if (control?.currentNode() === "provider-retry") { + emit({ type: "retrying", jobId: job.id }); + const retry = remind(request); + const retrySpoke = await this.runDurableTurn( + retry, + job, + emit, + control, + "provider-retry", + request.turnId, + ); + const retryOutcome: AgentHostProviderOutcome = retrySpoke + ? { spoke: true, disposition: "complete" } + : { + spoke: false, + disposition: "fail", + error: + "The agent completed a direct offer without sending a message.", + }; + await control.providerCompleted(retry, "provider-retry", retryOutcome); + if (retrySpoke) return; + throw new Error(retryOutcome.error); + } + + const spoke = await this.runDurableTurn( + request, + job, + emit, + control, + "provider-turn", + ); + await control?.providerCompleted(request, "provider-turn", { + spoke, + disposition: spoke || job.mode !== "direct" ? "complete" : "retry", + }); + if (spoke) return; // Silence is a complete answer on an open floor; only a direct question // left unanswered is worth a second ask. if (job.mode !== "direct") { @@ -107,17 +148,66 @@ export class LocalAiAgentHostExecutor implements AgentHostExecutor { // typing indicator across the gap rather than retiring it and lighting // it again for the same reply. emit({ type: "retrying", jobId: job.id }); - if (await this.runTurn(remind(request), job, emit, request.turnId)) { + const retry = remind(request); + const retrySpoke = await this.runDurableTurn( + retry, + job, + emit, + control, + "provider-retry", + request.turnId, + ); + const retryOutcome: AgentHostProviderOutcome = retrySpoke + ? { spoke: true, disposition: "complete" } + : { + spoke: false, + disposition: "fail", + error: + "The agent completed a direct offer without sending a message.", + }; + await control?.providerCompleted(retry, "provider-retry", retryOutcome); + if (retrySpoke) { return; } - throw new Error( - "The agent completed a direct offer without sending a message.", - ); + throw new Error(retryOutcome.error); } finally { this.activeRequests.delete(job.id); } } + private async runDurableTurn( + request: LocalAIChatRequest, + job: AgentHostJob, + emit: (event: AgentHostEvent) => void, + control: AgentHostExecutionControl | undefined, + node: "provider-turn" | "provider-retry", + previousTurnId?: string, + ): Promise { + await control?.providerStarted(request, node); + try { + return await this.runTurn(request, job, emit, previousTurnId); + } catch (error) { + if (control) { + const state = await Promise.resolve( + this.runtime.getTurnRuntimeState({ + conversationId: request.conversationId, + turnId: request.turnId, + }), + ).catch(() => null); + const message = + error instanceof Error ? error.message : "The provider turn failed."; + await control.providerFailed( + request, + state?.status === "uncertain" || state?.status === "pending" + ? "uncertain" + : "failed", + message, + ); + } + throw error; + } + } + /** Runs one turn; resolves to whether the agent actually spoke. */ private async runTurn( request: LocalAIChatRequest, diff --git a/packages/app/src/electron/agent-host/host.test.ts b/packages/app/src/electron/agent-host/host.test.ts index 70ce1925..b1a11f00 100644 --- a/packages/app/src/electron/agent-host/host.test.ts +++ b/packages/app/src/electron/agent-host/host.test.ts @@ -5,6 +5,16 @@ import type { } from "@/shared/types/agent-host"; import { AgentHost, type AgentHostExecutor } from "./host"; import { InMemoryAgentHostJobRepository } from "./repository"; +import { + beginWorkflowNode, + commitProviderEffect, + commitWorkflowNode, + createAgentHostWorkflow, + prepareProviderEffect, + providerEffectInputHash, + recordWorkflowNodeResult, + startProviderEffect, +} from "./workflow"; function deferred() { let resolve!: (value: T) => void; @@ -59,6 +69,47 @@ function storedJob(status: AgentHostJob["status"]): AgentHostJob { }; } +function runningProviderJob(committed: boolean): AgentHostJob { + const job = storedJob("running"); + let workflow = createAgentHostWorkflow(job, job.updatedAt); + workflow = beginWorkflowNode(workflow, "prepare-turn", job.updatedAt); + workflow = recordWorkflowNodeResult( + workflow, + "prepare-turn", + { + next: ["provider-turn"], + values: { requestId: "request-1", turnId: "turn-1" }, + }, + job.updatedAt, + ); + workflow = commitWorkflowNode(workflow, "prepare-turn", job.updatedAt); + workflow = beginWorkflowNode(workflow, "provider-turn", job.updatedAt); + workflow = prepareProviderEffect( + workflow, + { + node: "provider-turn", + requestId: "request-1", + turnId: "turn-1", + inputHash: providerEffectInputHash({ + providerId: "codex-cli", + requestId: "request-1", + turnId: "turn-1", + }), + }, + job.updatedAt, + ); + workflow = startProviderEffect(workflow, "turn-1", job.updatedAt); + if (committed) { + workflow = commitProviderEffect( + workflow, + "turn-1", + { spoke: true, next: ["finalize"] }, + job.updatedAt, + ); + } + return { ...job, requestId: "request-1", turnId: "turn-1", workflow }; +} + describe("AgentHost", () => { it("deduplicates the same trigger and actor", async () => { const execute = vi.fn(async () => undefined); @@ -202,6 +253,59 @@ describe("AgentHost", () => { expect(execute).not.toHaveBeenCalled(); }); + it("requeues a checkpointed node that stopped before an external effect", async () => { + const running = storedJob("running"); + running.workflow = beginWorkflowNode( + createAgentHostWorkflow(running, running.updatedAt), + "prepare-turn", + running.updatedAt, + ); + const repository = new InMemoryAgentHostJobRepository([running]); + const execute = vi.fn(async () => undefined); + const host = new AgentHost({ + repository, + executor: { execute }, + startPaused: true, + }); + + await host.initialize(); + + expect((await host.listJobs())[0]).toMatchObject({ status: "queued" }); + expect(execute).not.toHaveBeenCalled(); + }); + + it("quarantines a started provider effect instead of replaying it", async () => { + const repository = new InMemoryAgentHostJobRepository([ + runningProviderJob(false), + ]); + const execute = vi.fn(async () => undefined); + const host = new AgentHost({ repository, executor: { execute } }); + + await host.initialize(); + + expect((await host.listJobs())[0]).toMatchObject({ + status: "uncertain", + workflow: { effects: [expect.objectContaining({ status: "uncertain" })] }, + }); + expect(execute).not.toHaveBeenCalled(); + }); + + it("finalizes a committed provider receipt after restart", async () => { + const repository = new InMemoryAgentHostJobRepository([ + runningProviderJob(true), + ]); + const execute = vi.fn(async () => undefined); + const host = new AgentHost({ repository, executor: { execute } }); + + await host.initialize(); + + expect((await host.listJobs())[0]).toMatchObject({ + status: "completed", + workflow: { checkpoint: { next: [] } }, + }); + expect(execute).not.toHaveBeenCalled(); + }); + it("cancels the exact running provider request", async () => { const gate = deferred(); const cancel = vi.fn(async () => { diff --git a/packages/app/src/electron/agent-host/host.ts b/packages/app/src/electron/agent-host/host.ts index a1869e9f..c19e32ad 100644 --- a/packages/app/src/electron/agent-host/host.ts +++ b/packages/app/src/electron/agent-host/host.ts @@ -7,13 +7,55 @@ import type { AgentHostStructuredTaskBrief, AgentHostTarget, AgentHostTaskSummary, + AgentHostWorkflowNode, } from "@/shared/types/agent-host"; +import type { LocalAIChatRequest } from "@/shared/types/local-ai"; import type { AgentHostJobRepository } from "./repository"; +import { + beginWorkflowNode, + commitProviderEffect, + commitWorkflowNode, + createAgentHostWorkflow, + currentWorkflowNode, + failProviderEffect, + finalizeAgentHostWorkflow, + prepareProviderEffect, + providerEffectInputHash, + recordWorkflowNodeResult, + recoverAgentHostWorkflow, + startProviderEffect, +} from "./workflow"; + +export interface AgentHostProviderOutcome { + spoke: boolean; + disposition: "retry" | "complete" | "fail"; + error?: string; +} + +export interface AgentHostExecutionControl { + currentNode(): AgentHostWorkflowNode | undefined; + prepared(request: LocalAIChatRequest): Promise; + providerStarted( + request: LocalAIChatRequest, + node: "provider-turn" | "provider-retry", + ): Promise; + providerCompleted( + request: LocalAIChatRequest, + node: "provider-turn" | "provider-retry", + outcome: AgentHostProviderOutcome, + ): Promise; + providerFailed( + request: LocalAIChatRequest, + status: "failed" | "uncertain", + error: string, + ): Promise; +} export interface AgentHostExecutor { execute( job: AgentHostJob, emit: (event: AgentHostEvent) => void, + control?: AgentHostExecutionControl, ): Promise; cancel?(job: AgentHostJob): Promise | boolean; } @@ -242,11 +284,41 @@ export class AgentHost { for (const stored of jobs) { const job = structuredClone(stored); if (job.status === "running") { - job.status = "interrupted"; - job.error = - "Agent work was interrupted by an application restart and was not replayed to avoid duplicating tool actions."; - job.completedAt = this.now().toISOString(); - job.updatedAt = job.completedAt; + const now = this.now().toISOString(); + if (!job.workflow) { + job.status = "interrupted"; + job.error = + "Agent work was interrupted by an application restart and was not replayed because this run predates durable checkpoints."; + job.completedAt = now; + } else { + const recovery = recoverAgentHostWorkflow(job.workflow, now); + job.workflow = recovery.workflow; + if (recovery.disposition === "uncertain") { + job.status = "uncertain"; + job.error = recovery.error; + job.completedAt = undefined; + } else if (recovery.disposition === "failed") { + job.status = "failed"; + job.error = recovery.error; + job.completedAt = now; + } else if (recovery.disposition === "complete") { + job.workflow = finalizeAgentHostWorkflow(job.workflow, now); + job.status = "completed"; + job.error = undefined; + job.completedAt = now; + } else { + job.status = "queued"; + job.error = undefined; + job.completedAt = undefined; + } + } + job.updatedAt = now; + await this.repository.put(job); + } else if ( + !job.workflow && + (job.status === "queued" || job.status === "paused") + ) { + job.workflow = createAgentHostWorkflow(job, job.updatedAt); await this.repository.put(job); } this.jobs.set(job.id, job); @@ -328,6 +400,7 @@ export class AgentHost { createdAt: now, updatedAt: now, }; + job.workflow = createAgentHostWorkflow(job, now); this.jobs.set(job.id, job); await this.repository.put(job); this.emit({ type: "job", job }); @@ -433,6 +506,7 @@ export class AgentHost { createdAt: now, updatedAt: now, }; + job.workflow = createAgentHostWorkflow(job, now); this.jobs.set(id, job); await this.repository.put(job); this.emit({ type: "job", job }); @@ -532,6 +606,7 @@ export class AgentHost { createdAt: now, updatedAt: now, }; + successor.workflow = createAgentHostWorkflow(successor, now); this.jobs.set(id, successor); await this.repository.put(successor); this.emit({ type: "job", job: successor }); @@ -654,7 +729,12 @@ export class AgentHost { ): Promise { await this.ready; const job = this.taskCurrentJob(taskId, agentMemberId); - if (!job || TERMINAL.has(job.status) || job.status === "paused") { + if ( + !job || + TERMINAL.has(job.status) || + job.status === "paused" || + job.status === "uncertain" + ) { return false; } if (job.status === "running") { @@ -782,6 +862,7 @@ export class AgentHost { createdAt: now, updatedAt: now, }; + successor.workflow = createAgentHostWorkflow(successor, now); this.jobs.set(id, successor); await this.repository.put(successor); this.emit({ type: "job", job: successor }); @@ -911,19 +992,180 @@ export class AgentHost { } } - private async run(job: AgentHostJob): Promise { - const startedAt = this.now().toISOString(); - job.status = "running"; - job.attempts += 1; - job.startedAt = startedAt; - job.updatedAt = startedAt; - await this.repository.put(job); - this.emit({ type: "job", job }); + private executionControl(job: AgentHostJob): AgentHostExecutionControl { + const persist = async () => { + job.updatedAt = this.now().toISOString(); + await this.repository.put(job); + }; + const workflow = () => { + job.workflow ??= createAgentHostWorkflow(job, this.now().toISOString()); + return job.workflow; + }; + return { + currentNode: () => currentWorkflowNode(workflow()), + prepared: async (request) => { + if (currentWorkflowNode(workflow()) !== "prepare-turn") return; + job.requestId = request.requestId; + job.turnId = request.turnId; + job.workflow = recordWorkflowNodeResult( + workflow(), + "prepare-turn", + { + next: ["provider-turn"], + values: { + requestId: request.requestId, + turnId: request.turnId, + }, + }, + this.now().toISOString(), + ); + await persist(); + job.workflow = commitWorkflowNode( + workflow(), + "prepare-turn", + this.now().toISOString(), + ); + await persist(); + }, + providerStarted: async (request, node) => { + if (currentWorkflowNode(workflow()) !== node) { + throw new Error( + `Agent Host workflow expected ${currentWorkflowNode(workflow()) ?? "completion"}, not ${node}.`, + ); + } + job.requestId = request.requestId; + job.turnId = request.turnId; + job.workflow = beginWorkflowNode( + workflow(), + node, + this.now().toISOString(), + ); + await persist(); + job.workflow = prepareProviderEffect( + workflow(), + { + node, + requestId: request.requestId, + turnId: request.turnId, + inputHash: providerEffectInputHash(request), + }, + this.now().toISOString(), + ); + await persist(); + // This durable fence intentionally precedes runtime.startChat. A crash + // after it is conservative: the provider may have advanced, so the + // effect is reconciled rather than replayed. + job.workflow = startProviderEffect( + workflow(), + request.turnId, + this.now().toISOString(), + ); + await persist(); + }, + providerCompleted: async (request, node, outcome) => { + const next: AgentHostWorkflowNode[] = + outcome.disposition === "retry" + ? ["provider-retry"] + : outcome.disposition === "complete" + ? ["finalize"] + : []; + job.workflow = commitProviderEffect( + workflow(), + request.turnId, + { + spoke: outcome.spoke, + next, + terminalStatus: + outcome.disposition === "fail" ? "failed" : undefined, + terminalError: outcome.error, + }, + this.now().toISOString(), + ); + await persist(); + job.workflow = recordWorkflowNodeResult( + workflow(), + node, + { + next, + values: { + requestId: request.requestId, + turnId: request.turnId, + providerTurnCount: + workflow().checkpoint.values.providerTurnCount + 1, + spoke: outcome.spoke, + terminalStatus: + outcome.disposition === "fail" ? "failed" : undefined, + terminalError: outcome.error, + }, + }, + this.now().toISOString(), + ); + await persist(); + job.workflow = commitWorkflowNode( + workflow(), + node, + this.now().toISOString(), + ); + await persist(); + }, + providerFailed: async (request, status, error) => { + job.workflow = failProviderEffect( + workflow(), + request.turnId, + status, + error, + this.now().toISOString(), + ); + await persist(); + }, + }; + } + private async run(job: AgentHostJob): Promise { + const release = () => { + this.running.delete(job.id); + this.activeActors.delete(actorKey(job)); + this.scheduleDrain(); + }; try { - await this.executor.execute(structuredClone(job), (event) => - this.emit(event), - ); + const startedAt = this.now().toISOString(); + job.workflow ??= createAgentHostWorkflow(job, startedAt); + const recovery = recoverAgentHostWorkflow(job.workflow, startedAt); + job.workflow = recovery.workflow; + if (recovery.disposition === "uncertain") { + job.status = "uncertain"; + job.error = recovery.error; + job.updatedAt = startedAt; + await this.repository.put(job); + this.emit({ type: "job", job }); + return; + } + if (recovery.disposition === "failed") { + await this.finish(job, "failed", recovery.error); + return; + } + if (recovery.disposition === "complete") { + await this.finish(job, "completed"); + return; + } + const node = currentWorkflowNode(job.workflow); + if (node === "prepare-turn") { + job.workflow = beginWorkflowNode(job.workflow, node, startedAt); + } + job.status = "running"; + job.attempts += 1; + job.startedAt = startedAt; + job.updatedAt = startedAt; + await this.repository.put(job); + this.emit({ type: "job", job }); + + if (node !== "finalize") { + await this.executor.execute( + structuredClone(job), + (event) => this.emit(event), + this.executionControl(job), + ); + } if (job.status !== "running" || this.pausing.has(job.id)) return; await this.finish(job, "completed"); } catch (error) { @@ -940,9 +1182,7 @@ export class AgentHost { ); } } finally { - this.running.delete(job.id); - this.activeActors.delete(actorKey(job)); - this.scheduleDrain(); + release(); } } @@ -955,6 +1195,13 @@ export class AgentHost { error?: string, ): Promise { const completedAt = this.now().toISOString(); + if ( + status === "completed" && + job.workflow && + currentWorkflowNode(job.workflow) === "finalize" + ) { + job.workflow = finalizeAgentHostWorkflow(job.workflow, completedAt); + } job.status = status; job.error = error; job.completedAt = completedAt; diff --git a/packages/app/src/electron/agent-host/repository.test.ts b/packages/app/src/electron/agent-host/repository.test.ts index b582b15d..79908467 100644 --- a/packages/app/src/electron/agent-host/repository.test.ts +++ b/packages/app/src/electron/agent-host/repository.test.ts @@ -1,4 +1,4 @@ -import { mkdtemp, rm, writeFile } from "node:fs/promises"; +import { mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; import { join } from "node:path"; import { tmpdir } from "node:os"; import { afterEach, describe, expect, it } from "vitest"; @@ -44,13 +44,17 @@ describe("JsonAgentHostJobRepository", () => { it("round-trips jobs and updates by stable id", async () => { const directory = await mkdtemp(join(tmpdir(), "convera-agent-host-")); temporaryDirectories.push(directory); + const path = join(directory, "jobs.json"); const repository = new JsonAgentHostJobRepository({ - path: join(directory, "jobs.json"), + path, }); await repository.put(job("1", "queued")); await repository.put(job("1", "completed")); expect(await repository.list()).toEqual([job("1", "completed")]); + expect(JSON.parse(await readFile(path, "utf8"))).toMatchObject({ + schemaVersion: 4, + }); }); it("persists structured task provenance and Dexie result receipts", async () => { @@ -162,4 +166,23 @@ describe("JsonAgentHostJobRepository", () => { }), ]); }); + + it("migrates version three without making old running work replayable", async () => { + const directory = await mkdtemp(join(tmpdir(), "convera-agent-host-")); + temporaryDirectories.push(directory); + const path = join(directory, "jobs.json"); + await writeFile( + path, + JSON.stringify({ schemaVersion: 3, jobs: [job("1", "running")] }), + "utf8", + ); + + const [migrated] = await new JsonAgentHostJobRepository({ path }).list(); + + expect(migrated).toMatchObject({ status: "running" }); + expect(migrated.workflow).toBeUndefined(); + expect(JSON.parse(await readFile(path, "utf8"))).toMatchObject({ + schemaVersion: 4, + }); + }); }); diff --git a/packages/app/src/electron/agent-host/repository.ts b/packages/app/src/electron/agent-host/repository.ts index cd827b5d..5b76dcb4 100644 --- a/packages/app/src/electron/agent-host/repository.ts +++ b/packages/app/src/electron/agent-host/repository.ts @@ -2,6 +2,7 @@ import { AtomicJsonFile } from "../memory/json-file"; import { SerialTaskQueue } from "../memory/serial-queue"; import type { AgentHostJob } from "@/shared/types/agent-host"; import { z } from "zod"; +import { agentHostWorkflowSchema } from "./workflow"; export interface AgentHostJobRepository { list(): Promise; @@ -19,6 +20,7 @@ const statusSchema = z.enum([ "queued", "running", "paused", + "uncertain", "completed", "failed", "cancelled", @@ -75,7 +77,7 @@ const versionTwoJobSchema = legacyJobSchema.extend({ agentId: z.string().min(1), }); -const jobSchema = versionTwoJobSchema.extend({ +const versionThreeJobSchema = versionTwoJobSchema.extend({ taskId: z.string().min(1), parentTaskId: z.string().min(1).optional(), parentJobId: z.string().min(1).optional(), @@ -86,6 +88,10 @@ const jobSchema = versionTwoJobSchema.extend({ maxOutputTokens: z.number().int().min(1).max(1_000_000).optional(), }); +const jobSchema = versionThreeJobSchema.extend({ + workflow: agentHostWorkflowSchema.optional(), +}); + const legacyStateSchema = z.object({ schemaVersion: z.literal(1), jobs: z.array(legacyJobSchema), @@ -98,6 +104,11 @@ const stateSchema = z.object({ const currentStateSchema = z.object({ schemaVersion: z.literal(3), + jobs: z.array(versionThreeJobSchema), +}); + +const durableStateSchema = z.object({ + schemaVersion: z.literal(4), jobs: z.array(jobSchema), }); @@ -185,21 +196,21 @@ export class JsonAgentHostJobRepository implements AgentHostJobRepository { private async readState(): Promise<{ state: { - schemaVersion: 3; + schemaVersion: 4; jobs: AgentHostJob[]; }; migrated: boolean; }> { const value = await this.file.read(); if (value === undefined) { - return { state: { schemaVersion: 3, jobs: [] }, migrated: false }; + return { state: { schemaVersion: 4, jobs: [] }, migrated: false }; } const version = z.object({ schemaVersion: z.number() }).parse(value); if (version.schemaVersion === 1) { const legacy = legacyStateSchema.parse(value); return { state: { - schemaVersion: 3, + schemaVersion: 4, jobs: legacy.jobs.map(migrateLegacyJob), }, migrated: true, @@ -209,15 +220,22 @@ export class JsonAgentHostJobRepository implements AgentHostJobRepository { const prior = stateSchema.parse(value); return { state: { - schemaVersion: 3, + schemaVersion: 4, jobs: prior.jobs.map(migrateVersionTwoJob), }, migrated: true, }; } + if (version.schemaVersion === 3) { + const prior = currentStateSchema.parse(value); + return { + state: { schemaVersion: 4, jobs: prior.jobs }, + migrated: true, + }; + } return { - state: currentStateSchema.parse(value) as { - schemaVersion: 3; + state: durableStateSchema.parse(value) as { + schemaVersion: 4; jobs: AgentHostJob[]; }, migrated: false, @@ -225,10 +243,10 @@ export class JsonAgentHostJobRepository implements AgentHostJobRepository { } private async writeState(state: { - schemaVersion: 3; + schemaVersion: 4; jobs: AgentHostJob[]; }): Promise { - await this.file.write(currentStateSchema.parse(state)); + await this.file.write(durableStateSchema.parse(state)); } async list(): Promise { @@ -236,7 +254,7 @@ export class JsonAgentHostJobRepository implements AgentHostJobRepository { const { state, migrated } = await this.readState(); const jobs = prune(state.jobs, this.maxTerminalJobs); if (migrated || jobs.length !== state.jobs.length) { - await this.writeState({ schemaVersion: 3, jobs }); + await this.writeState({ schemaVersion: 4, jobs }); } return structuredClone(jobs); }); diff --git a/packages/app/src/electron/agent-host/workflow.test.ts b/packages/app/src/electron/agent-host/workflow.test.ts new file mode 100644 index 00000000..11bf6408 --- /dev/null +++ b/packages/app/src/electron/agent-host/workflow.test.ts @@ -0,0 +1,160 @@ +import { describe, expect, it } from "vitest"; +import type { + AgentHostJob, + AgentHostWorkflowState, +} from "@/shared/types/agent-host"; +import { + beginWorkflowNode, + commitProviderEffect, + commitWorkflowNode, + createAgentHostWorkflow, + currentWorkflowNode, + prepareProviderEffect, + providerEffectInputHash, + recordWorkflowNodeResult, + recoverAgentHostWorkflow, + startProviderEffect, +} from "./workflow"; + +const T0 = "2026-08-02T12:00:00.000Z"; +const T1 = "2026-08-02T12:00:01.000Z"; +const T2 = "2026-08-02T12:00:02.000Z"; +const T3 = "2026-08-02T12:00:03.000Z"; + +function job(): AgentHostJob { + return { + id: "job-1", + taskId: "task-1", + channelId: "channel-1", + channelKind: "channel", + conversationId: "conversation-1", + triggerMessageId: "message-1", + contextMessageIds: ["message-1"], + mode: "direct", + offeredAgentMemberIds: ["agent:a"], + agentId: "a", + agentMemberId: "agent:a", + chain: { hops: 0, invoked: ["agent:a"] }, + controlInstructions: [], + status: "queued", + attempts: 0, + createdAt: T0, + updatedAt: T0, + }; +} + +function atProviderNode(): AgentHostWorkflowState { + let workflow = createAgentHostWorkflow(job(), T0); + workflow = beginWorkflowNode(workflow, "prepare-turn", T0); + workflow = recordWorkflowNodeResult( + workflow, + "prepare-turn", + { + next: ["provider-turn"], + values: { requestId: "request-1", turnId: "turn-1" }, + }, + T1, + ); + return commitWorkflowNode(workflow, "prepare-turn", T1); +} + +function startedProviderEffect(): AgentHostWorkflowState { + const request = { + providerId: "codex-cli", + requestId: "request-1", + turnId: "turn-1", + }; + let workflow = beginWorkflowNode(atProviderNode(), "provider-turn", T1); + workflow = prepareProviderEffect( + workflow, + { + node: "provider-turn", + requestId: request.requestId, + turnId: request.turnId, + inputHash: providerEffectInputHash(request), + }, + T1, + ); + return startProviderEffect(workflow, request.turnId, T2); +} + +describe("Agent Host durable workflow", () => { + it("commits semantic node writes into immutable checkpoints", () => { + const workflow = atProviderNode(); + + expect(workflow.checkpoints).toHaveLength(2); + expect(workflow.checkpoint).toMatchObject({ + parentId: "job-1:checkpoint:0", + step: 1, + next: ["provider-turn"], + values: { requestId: "request-1", turnId: "turn-1" }, + }); + expect(workflow.pendingWrites).toEqual([ + expect.objectContaining({ status: "committed" }), + ]); + }); + + it("commits a durable pending write after a process restart", () => { + let workflow = createAgentHostWorkflow(job(), T0); + workflow = beginWorkflowNode(workflow, "prepare-turn", T0); + workflow = recordWorkflowNodeResult( + workflow, + "prepare-turn", + { next: ["provider-turn"], values: { requestId: "request-1" } }, + T1, + ); + + const recovered = recoverAgentHostWorkflow(workflow, T2); + + expect(recovered.disposition).toBe("resume"); + expect(currentWorkflowNode(recovered.workflow)).toBe("provider-turn"); + expect(recovered.workflow.pendingWrites[0].status).toBe("committed"); + }); + + it("never replays a provider effect that was started without a receipt", () => { + const recovered = recoverAgentHostWorkflow(startedProviderEffect(), T3); + + expect(recovered.disposition).toBe("uncertain"); + expect(recovered.workflow.effects[0]).toMatchObject({ + status: "uncertain", + turnId: "turn-1", + }); + expect(currentWorkflowNode(recovered.workflow)).toBe("provider-turn"); + }); + + it("finishes a node from a committed provider receipt without rerunning it", () => { + let workflow = startedProviderEffect(); + workflow = commitProviderEffect( + workflow, + "turn-1", + { spoke: true, next: ["finalize"] }, + T2, + ); + + const recovered = recoverAgentHostWorkflow(workflow, T3); + + expect(recovered.disposition).toBe("complete"); + expect(currentWorkflowNode(recovered.workflow)).toBe("finalize"); + expect(recovered.workflow.checkpoint.values).toMatchObject({ + providerTurnCount: 1, + spoke: true, + }); + }); + + it("rejects one effect key being reused for different provider input", () => { + const workflow = startedProviderEffect(); + + expect(() => + prepareProviderEffect( + workflow, + { + node: "provider-turn", + requestId: "request-1", + turnId: "turn-1", + inputHash: "f".repeat(64), + }, + T3, + ), + ).toThrow("different input"); + }); +}); diff --git a/packages/app/src/electron/agent-host/workflow.ts b/packages/app/src/electron/agent-host/workflow.ts new file mode 100644 index 00000000..3293bb47 --- /dev/null +++ b/packages/app/src/electron/agent-host/workflow.ts @@ -0,0 +1,506 @@ +import { createHash } from "node:crypto"; +import type { + AgentHostJob, + AgentHostWorkflowEffect, + AgentHostWorkflowNode, + AgentHostWorkflowNodeAttempt, + AgentHostWorkflowPendingWrite, + AgentHostWorkflowState, +} from "@/shared/types/agent-host"; +import { z } from "zod"; + +export const AGENT_HOST_WORKFLOW_GRAPH_VERSION = "agent-host-turn-v1" as const; + +const nodeSchema = z.enum([ + "prepare-turn", + "provider-turn", + "provider-retry", + "finalize", +]); + +const checkpointValuesSchema = z.object({ + inputHash: z.string().length(64), + requestId: z.string().min(1).max(256).optional(), + turnId: z.string().min(1).max(256).optional(), + providerTurnCount: z.number().int().min(0), + spoke: z.boolean().optional(), + terminalStatus: z.enum(["completed", "failed"]).optional(), + terminalError: z.string().max(16_000).optional(), +}); + +const checkpointSchema = z.object({ + id: z.string().min(1).max(512), + parentId: z.string().min(1).max(512).optional(), + step: z.number().int().min(0), + next: z.array(nodeSchema).max(16), + values: checkpointValuesSchema, + committedWriteIds: z.array(z.string().min(1).max(512)).max(256), + createdAt: z.string().datetime(), +}); + +const nodeAttemptSchema = z.object({ + id: z.string().min(1).max(512), + node: nodeSchema, + attempt: z.number().int().min(1), + status: z.enum(["running", "completed", "failed", "interrupted"]), + inputHash: z.string().length(64), + startedAt: z.string().datetime(), + completedAt: z.string().datetime().optional(), + error: z.string().max(16_000).optional(), +}); + +const pendingWriteSchema = z.object({ + id: z.string().min(1).max(512), + checkpointId: z.string().min(1).max(512), + attemptId: z.string().min(1).max(512), + channel: z.literal("node-result"), + value: z.object({ + next: z.array(nodeSchema).max(16), + values: checkpointValuesSchema.partial(), + }), + status: z.enum(["pending", "committed"]), + createdAt: z.string().datetime(), + committedAt: z.string().datetime().optional(), +}); + +const effectSchema = z.object({ + id: z.string().min(1).max(768), + attemptId: z.string().min(1).max(512), + kind: z.literal("provider-turn"), + idempotencyKey: z.string().min(1).max(768), + inputHash: z.string().length(64), + requestId: z.string().min(1).max(256), + turnId: z.string().min(1).max(256), + status: z.enum(["prepared", "started", "committed", "failed", "uncertain"]), + preparedAt: z.string().datetime(), + startedAt: z.string().datetime().optional(), + completedAt: z.string().datetime().optional(), + receipt: z + .object({ + spoke: z.boolean(), + next: z.array(nodeSchema).max(16), + terminalStatus: z.enum(["completed", "failed"]).optional(), + terminalError: z.string().max(16_000).optional(), + }) + .optional(), + error: z.string().max(16_000).optional(), +}); + +export const agentHostWorkflowSchema = z.object({ + schemaVersion: z.literal(1), + graphVersion: z.literal(AGENT_HOST_WORKFLOW_GRAPH_VERSION), + stateSchemaVersion: z.literal(1), + threadId: z.string().min(1).max(256), + checkpoint: checkpointSchema, + checkpoints: z.array(checkpointSchema).min(1).max(256), + attempts: z.array(nodeAttemptSchema).max(512), + pendingWrites: z.array(pendingWriteSchema).max(512), + effects: z.array(effectSchema).max(256), +}); + +function hash(value: unknown): string { + return createHash("sha256").update(JSON.stringify(value)).digest("hex"); +} + +export function agentHostJobInputHash( + job: Pick< + AgentHostJob, + | "taskId" + | "conversationId" + | "triggerMessageId" + | "contextMessageIds" + | "agentId" + | "agentMemberId" + | "controlInstructions" + >, +): string { + return hash({ + taskId: job.taskId, + conversationId: job.conversationId, + triggerMessageId: job.triggerMessageId, + contextMessageIds: job.contextMessageIds, + agentId: job.agentId, + agentMemberId: job.agentMemberId, + controlInstructions: job.controlInstructions, + }); +} + +export function providerEffectInputHash(input: { + providerId: string; + modelId?: string; + requestId: string; + turnId: string; +}): string { + return hash(input); +} + +export function createAgentHostWorkflow( + job: Pick< + AgentHostJob, + | "id" + | "taskId" + | "conversationId" + | "triggerMessageId" + | "contextMessageIds" + | "agentId" + | "agentMemberId" + | "controlInstructions" + >, + now: string, +): AgentHostWorkflowState { + const inputHash = agentHostJobInputHash(job); + const checkpoint = { + id: `${job.id}:checkpoint:0`, + step: 0, + next: ["prepare-turn" as const], + values: { inputHash, providerTurnCount: 0 }, + committedWriteIds: [], + createdAt: now, + }; + return agentHostWorkflowSchema.parse({ + schemaVersion: 1, + graphVersion: AGENT_HOST_WORKFLOW_GRAPH_VERSION, + stateSchemaVersion: 1, + threadId: job.id, + checkpoint, + checkpoints: [checkpoint], + attempts: [], + pendingWrites: [], + effects: [], + }) as AgentHostWorkflowState; +} + +function copy(workflow: AgentHostWorkflowState): AgentHostWorkflowState { + return structuredClone(workflow); +} + +export function currentWorkflowNode( + workflow: AgentHostWorkflowState, +): AgentHostWorkflowNode | undefined { + return workflow.checkpoint.next[0]; +} + +function currentAttempt( + workflow: AgentHostWorkflowState, + node = currentWorkflowNode(workflow), +): AgentHostWorkflowNodeAttempt | undefined { + if (!node) return undefined; + return [...workflow.attempts] + .reverse() + .find((attempt) => attempt.node === node); +} + +export function beginWorkflowNode( + workflow: AgentHostWorkflowState, + node: AgentHostWorkflowNode, + now: string, +): AgentHostWorkflowState { + const next = currentWorkflowNode(workflow); + if (next !== node) { + throw new Error(`Workflow expected ${next ?? "completion"}, not ${node}.`); + } + const latest = currentAttempt(workflow, node); + if (latest?.status === "running") return copy(workflow); + const result = copy(workflow); + const attempt = + result.attempts.filter((candidate) => candidate.node === node).length + 1; + result.attempts.push({ + id: `${result.threadId}:attempt:${node}:${attempt}`, + node, + attempt, + status: "running", + inputHash: result.checkpoint.values.inputHash, + startedAt: now, + }); + return agentHostWorkflowSchema.parse(result) as AgentHostWorkflowState; +} + +export function recordWorkflowNodeResult( + workflow: AgentHostWorkflowState, + node: AgentHostWorkflowNode, + result: AgentHostWorkflowPendingWrite["value"], + now: string, +): AgentHostWorkflowState { + const next = copy(workflow); + const attempt = currentAttempt(next, node); + if (!attempt || attempt.status !== "running") { + throw new Error(`Workflow node ${node} has no running attempt.`); + } + attempt.status = "completed"; + attempt.completedAt = now; + const existing = next.pendingWrites.find( + (write) => + write.attemptId === attempt.id && write.channel === "node-result", + ); + if (!existing) { + next.pendingWrites.push({ + id: `${attempt.id}:write:node-result`, + checkpointId: next.checkpoint.id, + attemptId: attempt.id, + channel: "node-result", + value: structuredClone(result), + status: "pending", + createdAt: now, + }); + } + return agentHostWorkflowSchema.parse(next) as AgentHostWorkflowState; +} + +export function commitWorkflowNode( + workflow: AgentHostWorkflowState, + node: AgentHostWorkflowNode, + now: string, +): AgentHostWorkflowState { + const next = copy(workflow); + const attempt = currentAttempt(next, node); + if (!attempt || attempt.status !== "completed") { + throw new Error(`Workflow node ${node} has no completed attempt.`); + } + const write = next.pendingWrites.find( + (candidate) => + candidate.attemptId === attempt.id && + candidate.channel === "node-result" && + candidate.status === "pending", + ); + if (!write) throw new Error(`Workflow node ${node} has no pending result.`); + write.status = "committed"; + write.committedAt = now; + const checkpoint = { + id: `${next.threadId}:checkpoint:${next.checkpoint.step + 1}`, + parentId: next.checkpoint.id, + step: next.checkpoint.step + 1, + next: [...write.value.next], + values: { ...next.checkpoint.values, ...write.value.values }, + committedWriteIds: [write.id], + createdAt: now, + }; + next.checkpoint = checkpoint; + next.checkpoints.push(checkpoint); + return agentHostWorkflowSchema.parse(next) as AgentHostWorkflowState; +} + +export function prepareProviderEffect( + workflow: AgentHostWorkflowState, + input: { + node: "provider-turn" | "provider-retry"; + requestId: string; + turnId: string; + inputHash: string; + }, + now: string, +): AgentHostWorkflowState { + const next = copy(workflow); + const attempt = currentAttempt(next, input.node); + if (!attempt || attempt.status !== "running") { + throw new Error(`Provider node ${input.node} has no running attempt.`); + } + const idempotencyKey = `${next.threadId}:${input.node}:${input.turnId}`; + const existing = next.effects.find( + (effect) => effect.idempotencyKey === idempotencyKey, + ); + if (existing) { + if (existing.inputHash !== input.inputHash) { + throw new Error("Provider effect idempotency key has different input."); + } + return next; + } + next.effects.push({ + id: `${next.threadId}:effect:${input.turnId}`, + attemptId: attempt.id, + kind: "provider-turn", + idempotencyKey, + inputHash: input.inputHash, + requestId: input.requestId, + turnId: input.turnId, + status: "prepared", + preparedAt: now, + }); + return agentHostWorkflowSchema.parse(next) as AgentHostWorkflowState; +} + +function findEffect( + workflow: AgentHostWorkflowState, + turnId: string, +): AgentHostWorkflowEffect { + const effect = workflow.effects.find( + (candidate) => candidate.turnId === turnId, + ); + if (!effect) + throw new Error(`Provider effect for turn ${turnId} was not prepared.`); + return effect; +} + +export function startProviderEffect( + workflow: AgentHostWorkflowState, + turnId: string, + now: string, +): AgentHostWorkflowState { + const next = copy(workflow); + const effect = findEffect(next, turnId); + if (effect.status === "started") return next; + if (effect.status !== "prepared") { + throw new Error( + `Provider effect ${turnId} cannot start from ${effect.status}.`, + ); + } + effect.status = "started"; + effect.startedAt = now; + return agentHostWorkflowSchema.parse(next) as AgentHostWorkflowState; +} + +export function commitProviderEffect( + workflow: AgentHostWorkflowState, + turnId: string, + receipt: NonNullable, + now: string, +): AgentHostWorkflowState { + const next = copy(workflow); + const effect = findEffect(next, turnId); + if (effect.status === "committed") return next; + if (effect.status !== "started") { + throw new Error( + `Provider effect ${turnId} cannot commit from ${effect.status}.`, + ); + } + effect.status = "committed"; + effect.receipt = structuredClone(receipt); + effect.completedAt = now; + return agentHostWorkflowSchema.parse(next) as AgentHostWorkflowState; +} + +export function failProviderEffect( + workflow: AgentHostWorkflowState, + turnId: string, + status: "failed" | "uncertain", + error: string, + now: string, +): AgentHostWorkflowState { + const next = copy(workflow); + const effect = findEffect(next, turnId); + if (["committed", "failed", "uncertain"].includes(effect.status)) return next; + effect.status = status; + effect.error = error; + effect.completedAt = now; + const attempt = next.attempts.find( + (candidate) => candidate.id === effect.attemptId, + ); + if (attempt?.status === "running") { + attempt.status = status === "failed" ? "failed" : "interrupted"; + attempt.error = error; + attempt.completedAt = now; + } + return agentHostWorkflowSchema.parse(next) as AgentHostWorkflowState; +} + +export type AgentHostWorkflowRecovery = { + workflow: AgentHostWorkflowState; + disposition: "resume" | "complete" | "failed" | "uncertain"; + error?: string; +}; + +export function recoverAgentHostWorkflow( + workflow: AgentHostWorkflowState, + now: string, +): AgentHostWorkflowRecovery { + let next = copy(workflow); + const node = currentWorkflowNode(next); + if (!node) { + return next.checkpoint.values.terminalStatus === "failed" + ? { + workflow: next, + disposition: "failed", + error: + next.checkpoint.values.terminalError ?? + "The agent did not produce a deliverable output.", + } + : { workflow: next, disposition: "complete" }; + } + if (node === "finalize") { + return { workflow: next, disposition: "complete" }; + } + const attempt = currentAttempt(next, node); + if (!attempt) return { workflow: next, disposition: "resume" }; + + const effect = [...next.effects] + .reverse() + .find((candidate) => candidate.attemptId === attempt.id); + if (effect?.status === "started" || effect?.status === "uncertain") { + const error = + effect.error ?? + "The provider may have advanced before the application stopped. This effect was not replayed."; + next = failProviderEffect(next, effect.turnId, "uncertain", error, now); + return { workflow: next, disposition: "uncertain", error }; + } + if (effect?.status === "failed") { + return { + workflow: next, + disposition: "failed", + error: effect.error ?? "The provider effect failed.", + }; + } + if (effect?.status === "committed" && attempt.status === "running") { + next = recordWorkflowNodeResult( + next, + node, + { + next: effect.receipt?.next ?? ["finalize"], + values: { + requestId: effect.requestId, + turnId: effect.turnId, + providerTurnCount: next.checkpoint.values.providerTurnCount + 1, + spoke: effect.receipt?.spoke, + terminalStatus: effect.receipt?.terminalStatus, + terminalError: effect.receipt?.terminalError, + }, + }, + now, + ); + } + const refreshedAttempt = currentAttempt(next, node); + if (refreshedAttempt?.status === "completed") { + next = commitWorkflowNode(next, node, now); + if (next.checkpoint.values.terminalStatus === "failed") { + return { + workflow: next, + disposition: "failed", + error: + next.checkpoint.values.terminalError ?? + "The agent did not produce a deliverable output.", + }; + } + return { + workflow: next, + disposition: + currentWorkflowNode(next) === "finalize" ? "complete" : "resume", + }; + } + if (refreshedAttempt?.status === "running") { + refreshedAttempt.status = "interrupted"; + refreshedAttempt.completedAt = now; + refreshedAttempt.error = + "The node stopped before starting an external effect."; + } + return { + workflow: agentHostWorkflowSchema.parse(next) as AgentHostWorkflowState, + disposition: "resume", + }; +} + +export function finalizeAgentHostWorkflow( + workflow: AgentHostWorkflowState, + now: string, +): AgentHostWorkflowState { + if (currentWorkflowNode(workflow) === undefined) return copy(workflow); + let next = workflow; + if (currentWorkflowNode(next) !== "finalize") { + throw new Error("Agent Host workflow is not ready to finalize."); + } + next = beginWorkflowNode(next, "finalize", now); + next = recordWorkflowNodeResult( + next, + "finalize", + { next: [], values: {} }, + now, + ); + return commitWorkflowNode(next, "finalize", now); +} diff --git a/packages/app/src/electron/ai/structured-task-contracts.ts b/packages/app/src/electron/ai/structured-task-contracts.ts index 1e9e40bb..d2a0639d 100644 --- a/packages/app/src/electron/ai/structured-task-contracts.ts +++ b/packages/app/src/electron/ai/structured-task-contracts.ts @@ -264,6 +264,7 @@ export type DelegateTaskToolResult = | "queued" | "running" | "paused" + | "uncertain" | "completed" | "failed" | "cancelled" diff --git a/packages/app/src/electron/trace/agent-host.ts b/packages/app/src/electron/trace/agent-host.ts index dee383c3..76fecfbd 100644 --- a/packages/app/src/electron/trace/agent-host.ts +++ b/packages/app/src/electron/trace/agent-host.ts @@ -73,6 +73,7 @@ function terminalStatus(job: AgentHostJob): TraceStatus | undefined { if (job.status === "failed") return "error"; if (job.status === "cancelled") return "cancelled"; if (job.status === "interrupted") return "interrupted"; + if (job.status === "uncertain") return "interrupted"; return undefined; } @@ -263,7 +264,9 @@ export class AgentHostTraceRecorder { }, }); await this.safeAppend(inputs); - if (job.status === "interrupted") await this.recoverInterruptedRun(job); + if (job.status === "interrupted" || job.status === "uncertain") { + await this.recoverInterruptedRun(job); + } } async beginTurn( diff --git a/packages/app/src/renderer/components/chat/TypingIndicator.tsx b/packages/app/src/renderer/components/chat/TypingIndicator.tsx index 6c378899..92ad4630 100644 --- a/packages/app/src/renderer/components/chat/TypingIndicator.tsx +++ b/packages/app/src/renderer/components/chat/TypingIndicator.tsx @@ -43,7 +43,9 @@ export function AgentTurnFailureNotice({ (job) => job.conversationId === conversationId && job.error && - (job.status === "failed" || job.status === "interrupted"), + (job.status === "failed" || + job.status === "interrupted" || + job.status === "uncertain"), ); if (!failure || failure.id === dismissed) return null; diff --git a/packages/app/src/renderer/libs/__tests__/agent-speech.test.ts b/packages/app/src/renderer/libs/__tests__/agent-speech.test.ts index 8a695487..40a1316d 100644 --- a/packages/app/src/renderer/libs/__tests__/agent-speech.test.ts +++ b/packages/app/src/renderer/libs/__tests__/agent-speech.test.ts @@ -76,6 +76,70 @@ describe("agent speech", () => { }); }); + it("returns the original message when an AgentHost effect is delivered twice", async () => { + const { channelId, conversationId } = await seedRoom(); + const query = { + kind: "send_message" as const, + viewerMemberId: sageMemberId, + channelId, + content: "durable result", + agentHost: { + jobId: "job-1", + effectId: "effect-1", + payloadHash: "a".repeat(64), + triggerMessageId: "trigger-1", + contextMessageIds: ["trigger-1"], + chain: { hops: 0, invoked: [sageMemberId] }, + }, + }; + + const first = await resolveWorkspaceQuery(query); + const replay = await resolveWorkspaceQuery(query); + + expect(first).toMatchObject({ ok: true, kind: "send_message" }); + expect(replay).toEqual(first); + expect( + await db.messages.where("conversationId").equals(conversationId).count(), + ).toBe(1); + expect(await db.agentEffectReceipts.get("effect-1")).toMatchObject({ + messageId: + first.ok && first.kind === "send_message" ? first.messageId : "missing", + payloadHash: "a".repeat(64), + }); + }); + + it("rejects an AgentHost effect key reused with different content", async () => { + const { channelId } = await seedRoom(); + const agentHost = { + jobId: "job-1", + effectId: "effect-1", + payloadHash: "a".repeat(64), + triggerMessageId: "trigger-1", + contextMessageIds: ["trigger-1"], + chain: { hops: 0, invoked: [sageMemberId] }, + }; + await resolveWorkspaceQuery({ + kind: "send_message", + viewerMemberId: sageMemberId, + channelId, + content: "first", + agentHost, + }); + + expect( + await resolveWorkspaceQuery({ + kind: "send_message", + viewerMemberId: sageMemberId, + channelId, + content: "different", + agentHost: { ...agentHost, payloadHash: "b".repeat(64) }, + }), + ).toMatchObject({ + ok: false, + error: { message: expect.stringContaining("different content") }, + }); + }); + it("re-parses mentions from the posted text", async () => { // A mention is what routes the next turn, so it has to come from the text // that actually landed rather than anything the model claimed alongside it. diff --git a/packages/app/src/renderer/libs/agent-host-service.ts b/packages/app/src/renderer/libs/agent-host-service.ts index 2c208712..d8a9209c 100644 --- a/packages/app/src/renderer/libs/agent-host-service.ts +++ b/packages/app/src/renderer/libs/agent-host-service.ts @@ -9,6 +9,7 @@ import type { LocalAIStreamEvent, } from "@/shared/types/local-ai"; import type { Member } from "@/shared/types/workspace"; +import type { WorkspaceSendMessageQuery } from "@/shared/types/workspace-perception"; import { db, type Agent, type Channel, type Message } from "./db/database"; import { useSelectionStore } from "./db/ui-state"; import { @@ -27,6 +28,26 @@ import { typingTransition, useTypingStore } from "./stores/typing-store"; import { useUserInputStore } from "./stores/user-input-store"; import { beginTrace, endTrace, noteRetry, recordTrace } from "./agent-trace"; +async function workspaceMessagePayloadHash(input: { + viewerMemberId: string; + channelId: string; + content: string; + replyToMessageId?: string; +}): Promise { + const bytes = new TextEncoder().encode( + JSON.stringify({ + viewerMemberId: input.viewerMemberId, + channelId: input.channelId, + content: input.content, + replyToMessageId: input.replyToMessageId ?? null, + }), + ); + const digest = await crypto.subtle.digest("SHA-256", bytes); + return [...new Uint8Array(digest)] + .map((byte) => byte.toString(16).padStart(2, "0")) + .join(""); +} + interface ActiveOffer { job: AgentHostJob; requestId: string; @@ -457,25 +478,33 @@ export class RendererAgentHostService { if (event.type === "interaction") { const respond = (response: LocalAIInteractionResponse) => this.respondToInteraction(event, response); - const workspaceEvent = + const messageInput = event.name === "workspace:query" && typeof event.input === "object" && event.input !== null && "kind" in event.input && event.input.kind === "send_message" - ? { - ...event, - input: { - ...event.input, - agentHost: { - jobId: active.job.id, - triggerMessageId: active.job.triggerMessageId, - contextMessageIds: active.job.contextMessageIds, - chain: active.job.chain, - }, + ? (event.input as WorkspaceSendMessageQuery) + : undefined; + const payloadHash = messageInput + ? await workspaceMessagePayloadHash(messageInput) + : undefined; + const workspaceEvent = messageInput + ? { + ...event, + input: { + ...messageInput, + agentHost: { + jobId: active.job.id, + effectId: `agent-host:${active.job.id}:${event.requestId}:${event.interactionId}`, + payloadHash, + triggerMessageId: active.job.triggerMessageId, + contextMessageIds: active.job.contextMessageIds, + chain: active.job.chain, }, - } - : event; + }, + } + : event; const workspaceRespond = async (response: LocalAIInteractionResponse) => { if ( workspaceEvent.name === "workspace:query" && diff --git a/packages/app/src/renderer/libs/agent-speech.ts b/packages/app/src/renderer/libs/agent-speech.ts index 0acddaed..f436461c 100644 --- a/packages/app/src/renderer/libs/agent-speech.ts +++ b/packages/app/src/renderer/libs/agent-speech.ts @@ -1,4 +1,4 @@ -import { addMessage, db, type Channel } from "./db"; +import { db, type Channel } from "./db"; import { registerWorkspaceSendMessage } from "./workspace-perception"; import type { WorkspaceQueryResult } from "@/shared/types/workspace-perception"; import { parseMentions } from "./mention-parser"; @@ -24,18 +24,61 @@ async function speak( senderId: string, content: string, replyToMessageId?: string, -): Promise { + effect?: { idempotencyKey: string; payloadHash: string; jobId: string }, +): Promise<{ messageId: string; created: boolean }> { const members = await db.members.toArray(); - return addMessage(channel.conversationId, { - role: "assistant", - content, - senderId, - // Parsed here rather than trusted from the model: a mention is what routes - // the next turn, so it has to reflect the text that was actually posted. - mentions: parseMentions(content, members), - ...(replyToMessageId ? { replyToMessageId } : {}), - status: "completed", - }); + return db.transaction( + "rw", + [db.messages, db.conversations, db.agentEffectReceipts], + async () => { + if (effect) { + const receipt = await db.agentEffectReceipts.get(effect.idempotencyKey); + if (receipt) { + if (receipt.payloadHash !== effect.payloadHash) { + throw new Error( + "Agent speech idempotency key was reused with different content.", + ); + } + if (!(await db.messages.get(receipt.messageId))) { + throw new Error( + "Agent speech receipt points to a missing Dexie message.", + ); + } + return { messageId: receipt.messageId, created: false }; + } + } + + const messageId = crypto.randomUUID(); + const createdAt = new Date(); + await db.messages.add({ + id: messageId, + conversationId: channel.conversationId, + role: "assistant", + content, + senderId, + // Parsed here rather than trusted from the model: a mention is what + // routes the next turn, so it reflects the text that actually landed. + mentions: parseMentions(content, members), + ...(replyToMessageId ? { replyToMessageId } : {}), + status: "completed", + createdAt, + }); + await db.conversations.update(channel.conversationId, { + updatedAt: createdAt, + }); + if (effect) { + await db.agentEffectReceipts.add({ + idempotencyKey: effect.idempotencyKey, + payloadHash: effect.payloadHash, + jobId: effect.jobId, + conversationId: channel.conversationId, + messageId, + createdAt, + }); + } + return { messageId, created: true }; + }, + ); } /** @@ -62,13 +105,26 @@ export function installAgentSpeech(): void { }, }; } - const messageId = await speak( + if ( + agentHost && + Boolean(agentHost.effectId) !== Boolean(agentHost.payloadHash) + ) { + throw new Error("Agent speech effect metadata is incomplete."); + } + const { messageId, created } = await speak( channel, viewerMemberId, content, replyToMessageId, + agentHost?.effectId && agentHost.payloadHash + ? { + idempotencyKey: agentHost.effectId, + payloadHash: agentHost.payloadHash, + jobId: agentHost.jobId, + } + : undefined, ); - if (agentHost) { + if (agentHost && created) { const memberRows = await db.members.bulkGet(channel.memberIds); const members = memberRows.filter( (member): member is Member => member !== undefined, diff --git a/packages/app/src/renderer/libs/agent-tasks.ts b/packages/app/src/renderer/libs/agent-tasks.ts index 3be1abab..de85e316 100644 --- a/packages/app/src/renderer/libs/agent-tasks.ts +++ b/packages/app/src/renderer/libs/agent-tasks.ts @@ -11,7 +11,12 @@ const TERMINAL = new Set([ ]); /** What is happening now, then what is waiting, then what is being held. */ -const OPEN_ORDER: AgentHostJobStatus[] = ["running", "queued", "paused"]; +const OPEN_ORDER: AgentHostJobStatus[] = [ + "running", + "uncertain", + "queued", + "paused", +]; /** * What a colleague is currently carrying. @@ -37,7 +42,9 @@ export function openAgentTasks( export function taskStatusLabel(status: string): string { return status === "running" ? "Working" - : status === "queued" - ? "Queued" - : status.charAt(0).toUpperCase() + status.slice(1); + : status === "uncertain" + ? "Needs review" + : status === "queued" + ? "Queued" + : status.charAt(0).toUpperCase() + status.slice(1); } diff --git a/packages/app/src/renderer/libs/db/database.ts b/packages/app/src/renderer/libs/db/database.ts index df816319..1fb0b719 100644 --- a/packages/app/src/renderer/libs/db/database.ts +++ b/packages/app/src/renderer/libs/db/database.ts @@ -118,6 +118,20 @@ export interface PendingTurnJournal { updatedAt: Date; } +/** + * A renderer-side inbox receipt for an AgentHost effect. The message remains + * authoritative in `messages`; this row only makes response-loss retries + * return the same message id instead of posting a duplicate. + */ +export interface AgentEffectReceipt { + idempotencyKey: string; + payloadHash: string; + jobId: string; + conversationId: string; + messageId: string; + createdAt: Date; +} + export type PendingConversationDeletionState = | "pending" | "deleting" @@ -250,6 +264,7 @@ export class ConveraDB extends Dexie { conversations!: EntityTable; messages!: EntityTable; pendingTurns!: EntityTable; + agentEffectReceipts!: EntityTable; pendingConversationDeletions!: EntityTable< PendingConversationDeletion, "conversationId" @@ -482,6 +497,13 @@ export class ConveraDB extends Dexie { agentTraces: "id, conversationId, memberId, startedAt", }); + // v11: endpoint idempotency for AgentHost speech. A message and its + // receipt are committed in the same Dexie transaction by agent-speech. + this.version(11).stores({ + agentEffectReceipts: + "idempotencyKey, jobId, conversationId, messageId, createdAt", + }); + // A database created fresh at the latest version never runs upgrade hooks, // so the local human member, default workspace and built-in tag must also // be seeded on populate. diff --git a/packages/app/src/renderer/libs/db/hooks.ts b/packages/app/src/renderer/libs/db/hooks.ts index 862e5a65..025e3cb3 100644 --- a/packages/app/src/renderer/libs/db/hooks.ts +++ b/packages/app/src/renderer/libs/db/hooks.ts @@ -276,11 +276,13 @@ export async function deleteConversation(id: string): Promise { db.conversations, db.messages, db.pendingTurns, + db.agentEffectReceipts, db.pendingConversationDeletions, ], async () => { await db.messages.where("conversationId").equals(id).delete(); await db.pendingTurns.where("conversationId").equals(id).delete(); + await db.agentEffectReceipts.where("conversationId").equals(id).delete(); await db.pendingConversationDeletions.delete(id); await db.conversations.delete(id); }, diff --git a/packages/app/src/shared/types/agent-host.ts b/packages/app/src/shared/types/agent-host.ts index 82627426..9c30a754 100644 --- a/packages/app/src/shared/types/agent-host.ts +++ b/packages/app/src/shared/types/agent-host.ts @@ -4,11 +4,94 @@ export type AgentHostJobStatus = | "queued" | "running" | "paused" + | "uncertain" | "completed" | "failed" | "cancelled" | "interrupted"; +export type AgentHostWorkflowNode = + | "prepare-turn" + | "provider-turn" + | "provider-retry" + | "finalize"; + +export interface AgentHostWorkflowCheckpoint { + id: string; + parentId?: string; + step: number; + next: AgentHostWorkflowNode[]; + values: { + inputHash: string; + requestId?: string; + turnId?: string; + providerTurnCount: number; + spoke?: boolean; + terminalStatus?: "completed" | "failed"; + terminalError?: string; + }; + committedWriteIds: string[]; + createdAt: string; +} + +export interface AgentHostWorkflowNodeAttempt { + id: string; + node: AgentHostWorkflowNode; + attempt: number; + status: "running" | "completed" | "failed" | "interrupted"; + inputHash: string; + startedAt: string; + completedAt?: string; + error?: string; +} + +export interface AgentHostWorkflowPendingWrite { + id: string; + checkpointId: string; + attemptId: string; + channel: "node-result"; + value: { + next: AgentHostWorkflowNode[]; + values: Partial; + }; + status: "pending" | "committed"; + createdAt: string; + committedAt?: string; +} + +export interface AgentHostWorkflowEffect { + id: string; + attemptId: string; + kind: "provider-turn"; + idempotencyKey: string; + inputHash: string; + requestId: string; + turnId: string; + status: "prepared" | "started" | "committed" | "failed" | "uncertain"; + preparedAt: string; + startedAt?: string; + completedAt?: string; + receipt?: { + spoke: boolean; + next: AgentHostWorkflowNode[]; + terminalStatus?: "completed" | "failed"; + terminalError?: string; + }; + error?: string; +} + +export interface AgentHostWorkflowState { + schemaVersion: 1; + graphVersion: "agent-host-turn-v1"; + stateSchemaVersion: 1; + threadId: string; + checkpoint: AgentHostWorkflowCheckpoint; + checkpoints: AgentHostWorkflowCheckpoint[]; + attempts: AgentHostWorkflowNodeAttempt[]; + pendingWrites: AgentHostWorkflowPendingWrite[]; + effects: AgentHostWorkflowEffect[]; +} + export type AgentHostChannelKind = "channel" | "dm"; export interface AgentHostChain { @@ -87,6 +170,8 @@ export interface AgentHostJob { collaboration?: AgentHostCollaboration; /** Renderer/Dexie-owned result messages posted by this run. */ outputMessageIds?: string[]; + /** Main-owned execution metadata. It stores references and receipts, never transcript text. */ + workflow?: AgentHostWorkflowState; maxOutputTokens?: number; status: AgentHostJobStatus; attempts: number; diff --git a/packages/app/src/shared/types/workspace-perception.ts b/packages/app/src/shared/types/workspace-perception.ts index 71c7b1ee..6b5e5a5b 100644 --- a/packages/app/src/shared/types/workspace-perception.ts +++ b/packages/app/src/shared/types/workspace-perception.ts @@ -61,6 +61,10 @@ export interface WorkspaceSendMessageQuery { /** Renderer-owned lifecycle context; never accepted from a model tool input. */ agentHost?: { jobId: string; + /** Renderer-generated endpoint key; never accepted from model input. */ + effectId?: string; + /** SHA-256 of the destination, author, body, and reply target. */ + payloadHash?: string; triggerMessageId: string; contextMessageIds: string[]; chain: { hops: number; invoked: string[] };