From 727154d22eb77496dbe95ad6fea577ef78322579 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Fri, 28 Aug 2026 08:54:20 -0700 Subject: [PATCH 1/2] Isolate sidecar conversation per originating workbench Same agent principal. Durable turns nested under agent-state///. A new room starts empty. https://linear.app/abklabs/issue/CL-7168 --- apps/sidecar/package.json | 1 + apps/sidecar/src/conversation-state.test.ts | 191 +++++++++++++++ apps/sidecar/src/conversation-state.ts | 217 ++++++++++++++---- .../sidecar/src/originating-workbench.test.ts | 106 +++++++++ apps/sidecar/src/originating-workbench.ts | 102 ++++++++ .../sidecar/src/workflow-host-wiring/index.ts | 14 +- .../src/workflow-substrate-factory/index.ts | 41 +++- .../parked-approvals.ts | 30 ++- .../workflow-substrate-factory/step-env.ts | 13 +- bun.lock | 1 + 10 files changed, 641 insertions(+), 75 deletions(-) create mode 100644 apps/sidecar/src/conversation-state.test.ts create mode 100644 apps/sidecar/src/originating-workbench.test.ts create mode 100644 apps/sidecar/src/originating-workbench.ts diff --git a/apps/sidecar/package.json b/apps/sidecar/package.json index 73ad0f766..3ec85acec 100644 --- a/apps/sidecar/package.json +++ b/apps/sidecar/package.json @@ -28,6 +28,7 @@ "@intx/inference": "workspace:*", "@intx/log": "0.3.0", "@intx/mail-memory": "workspace:*", + "@intx/mime": "workspace:*", "@intx/storage-isogit": "0.3.0", "@intx/tool-packaging": "0.3.0", "@intx/types": "workspace:*", diff --git a/apps/sidecar/src/conversation-state.test.ts b/apps/sidecar/src/conversation-state.test.ts new file mode 100644 index 000000000..e09c6487b --- /dev/null +++ b/apps/sidecar/src/conversation-state.test.ts @@ -0,0 +1,191 @@ +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +import { afterEach, expect, test } from "bun:test"; + +import type { Principal, RepoId, RepoStore } from "@intx/hub-sessions/substrate"; +import type { ConversationTurn } from "@intx/types/runtime"; + +import { + createDurableConversationStore, + durableConversationAgentStateDir, + durableConversationAgentStatePrefix, + prepareConversationForOriginatingWorkbench, + reconstructDurableConversation, +} from "./conversation-state"; + +const tmpDirs: string[] = []; + +afterEach(async () => { + await Promise.all( + tmpDirs + .splice(0) + .map((dir) => fs.promises.rm(dir, { recursive: true, force: true })), + ); +}); + +function tmpDir(prefix: string): string { + const dir = fs.mkdtempSync(path.join(os.tmpdir(), prefix)); + tmpDirs.push(dir); + return dir; +} + +const WORKFLOW_RUN_REPO_ID: RepoId = { kind: "workflow-run", id: "wfr_iso" }; + +function userTurn(text: string): ConversationTurn { + return { + role: "user", + content: [{ type: "text", text }], + } as ConversationTurn; +} + +function peekTurns(store: { storage: unknown }): ConversationTurn[] { + const storage = store.storage as { peekTurns: () => ConversationTurn[] }; + return storage.peekTurns(); +} + +function fsBackedSubstrate(repoDir: string): RepoStore { + return { + getRepoDir: () => repoDir, + writeTreePreservingPrefix: async ( + _principal: Principal, + _repoId: RepoId, + _ref: string, + args: Parameters[3], + ) => { + const prefixDir = path.join( + repoDir, + ...args.preservePrefix.split("/").filter((seg: string) => seg.length > 0), + ); + const existing = new Map(); + let names: string[] = []; + try { + names = await fs.promises.readdir(prefixDir); + } catch (cause) { + if ( + !( + typeof cause === "object" && + cause !== null && + "code" in cause && + cause.code === "ENOENT" + ) + ) { + throw cause; + } + } + for (const name of names) { + const full = path.join(prefixDir, name); + const st = await fs.promises.stat(full); + if (st.isFile()) { + existing.set(`${args.preservePrefix}${name}`, await fs.promises.readFile(full)); + } + } + const files = await args.merge(existing); + await fs.promises.rm(prefixDir, { recursive: true, force: true }); + for (const [rel, content] of Object.entries(files)) { + const dest = path.join( + repoDir, + ...rel.split("/").filter((seg: string) => seg.length > 0), + ); + await fs.promises.mkdir(path.dirname(dest), { recursive: true }); + await fs.promises.writeFile(dest, content); + } + return { commitSha: "test" }; + }, + } as unknown as RepoStore; +} + +async function makeStore(root: string) { + const repoDir = path.join(root, "repo"); + await fs.promises.mkdir(repoDir, { recursive: true }); + return createDurableConversationStore({ + localStoreDir: path.join(root, "local"), + signer: () => Promise.resolve("sig"), + substrate: fsBackedSubstrate(repoDir), + workflowRunRepoId: WORKFLOW_RUN_REPO_ID, + workflowRunRef: "refs/heads/main", + principal: { kind: "workflow-process", anchorRunId: "run_1" } as Principal, + agentKey: "default", + }).then((store) => ({ store, repoDir })); +} + +test("durable conversation prefix nests workbench under agentKey", () => { + expect(durableConversationAgentStatePrefix("default", "chan_a")).toBe( + "agent-state/default/chan_a/", + ); + expect(durableConversationAgentStatePrefix("step_1", "chan_a")).not.toBe( + durableConversationAgentStatePrefix("step_2", "chan_a"), + ); +}); + +test("two rooms of one agent keep isolated turns; a new room starts empty", async () => { + const root = tmpDir("conv-iso-"); + const { store, repoDir } = await makeStore(root); + + expect(await store.bindOriginatingWorkbench("chan_a")).toBe(true); + await store.storage.writeTurns([userTurn("room A tool result")]); + await store.mirrorToSubstrate(); + + const roomADir = durableConversationAgentStateDir(repoDir, "default", "chan_a"); + const reconstructedA = await reconstructDurableConversation( + roomADir, + "default/chan_a", + ); + expect(reconstructedA?.turns).toHaveLength(1); + + expect(await store.bindOriginatingWorkbench("chan_b")).toBe(true); + expect(peekTurns(store)).toEqual([]); + + const mixedLegacy = path.join(repoDir, "agent-state", "default", "checkpoint.json"); + expect(fs.existsSync(mixedLegacy)).toBe(false); + + await store.storage.writeTurns([userTurn("room B first infer")]); + await store.mirrorToSubstrate(); + + expect(await store.bindOriginatingWorkbench("chan_a")).toBe(true); + expect(peekTurns(store)).toEqual([userTurn("room A tool result")]); + + expect(await store.bindOriginatingWorkbench("chan_b")).toBe(true); + expect(peekTurns(store)).toEqual([userTurn("room B first infer")]); + + expect(await store.bindOriginatingWorkbench("chan_b")).toBe(false); +}); + +test("prepareConversationForOriginatingWorkbench evicts the warm agent on room change only", async () => { + const evictions: string[] = []; + const store = { + bindOriginatingWorkbench: (id: string) => + Promise.resolve(id === "chan_b"), + }; + const registry = { + acquire: () => Promise.resolve(store), + get: () => store, + peek: () => store, + }; + const warmCache = { + evictAll: (reason: string) => { + evictions.push(reason); + return Promise.resolve(); + }, + }; + + const swappedA = await prepareConversationForOriginatingWorkbench({ + registry: registry as never, + agentKey: "default", + originatingWorkbenchId: "chan_a", + warmCache, + }); + expect(swappedA).toBe(false); + expect(evictions).toEqual([]); + + const swappedB = await prepareConversationForOriginatingWorkbench({ + registry: registry as never, + agentKey: "default", + originatingWorkbenchId: "chan_b", + warmCache, + }); + expect(swappedB).toBe(true); + expect(evictions).toHaveLength(1); + expect(evictions[0]).toContain("chan_b"); +}); diff --git a/apps/sidecar/src/conversation-state.ts b/apps/sidecar/src/conversation-state.ts index 4674d8baa..60092c53c 100644 --- a/apps/sidecar/src/conversation-state.ts +++ b/apps/sidecar/src/conversation-state.ts @@ -10,10 +10,13 @@ // This module makes the warm agent's conversation DURABLE in the // workflow-run substrate (the single-writer proxy `RepoStore`, written // through the supervisor). The durable copy lives under the workflow-run -// repo at a stable per-agent path (`agent-state//...`), sibling -// to the per-run event log under `runs//...` and NOT confused with -// it. On a new run (and after a child respawn, once the warm agent is -// rebuilt lazily) the conversation is restored from the substrate into +// repo at a per-agent, per-originating-workbench path +// (`agent-state///...`), sibling to the per-run +// event log under `runs//...` and NOT confused with it. `agentKey` +// stays the warm stepId; `workbenchId` is the inbound mail From local-part +// so two rooms sharing one principal do not share turns. On a new run (and +// after a child respawn, once the warm agent is rebuilt lazily) the +// conversation is restored from the substrate into // the agent's local store BEFORE the agent's reactor loads, so multi-turn // continuity holds across runs and across respawn. // @@ -24,7 +27,7 @@ // over a conversation. It is replaced here with an append-only, // bucket-sharded write-ahead log plus a periodic compacted checkpoint: // -// agent-state// +// agent-state/// // checkpoint.json compacted full snapshot (turns + metadata) // checkpoint.meta.json { checkpointSeq: , turnCount, // tokenUsage, pendingOperations, @@ -68,18 +71,18 @@ // prefix pass through untouched). A WAL blob is two levels below // `agent-state//`, so: // -// - WAL append uses `preservePrefix = agent-state//wal//`. +// - WAL append uses `preservePrefix = agent-state///wal//`. // The bucket's existing blobs ARE direct children, so the merge // pre-image is exactly that bucket and the append adds one entry -- -// no isogit side-read, and the checkpoint / other buckets are -// untouched (outside the prefix). +// no isogit side-read, and the checkpoint / other buckets / sibling +// rooms are untouched (outside the prefix). // - Checkpoint write + WAL truncate uses `preservePrefix = -// agent-state//`. The top-level checkpoint files are direct -// children; the merge returns ONLY those files and NO `wal/...` -// paths, so the recursive `clearPrefix` at `agent-state//` drops -// the entire WAL subtree in the same atomic commit. The truncate -// needs no nested read: omitting the WAL paths from the returned set -// IS the truncate. +// agent-state///`. The top-level checkpoint files are +// direct children; the merge returns ONLY those files and NO `wal/...` +// paths, so the recursive `clearPrefix` at the workbench dir drops +// that room's WAL subtree in the same atomic commit without touching +// sibling rooms. The truncate needs no nested read: omitting the WAL +// paths from the returned set IS the truncate. // // Persistence sink (the riskiest part). The connector router's // `snapshot()` / `restore()` surface and the harness's @@ -159,6 +162,39 @@ const CHECKPOINT_FILE = "checkpoint.json"; const CHECKPOINT_META_FILE = "checkpoint.meta.json"; const WAL_DIR = "wal"; +const EMPTY_TOKEN_USAGE: TokenUsage = { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + thinking: 0, +}; + +/** + * Substrate prefix for one agent's conversation in one originating + * workbench. Nested under `agent-state//` so two agents in the + * same room cannot collide, and two rooms of one agent cannot share turns. + */ +export function durableConversationAgentStatePrefix( + agentKey: string, + originatingWorkbenchId: string, +): string { + return `${WORKFLOW_RUN_AGENT_STATE_PREFIX}/${encodeURIComponent(agentKey)}/${encodeURIComponent(originatingWorkbenchId)}/`; +} + +export function durableConversationAgentStateDir( + repoDir: string, + agentKey: string, + originatingWorkbenchId: string, +): string { + return path.join( + repoDir, + WORKFLOW_RUN_AGENT_STATE_PREFIX, + encodeURIComponent(agentKey), + encodeURIComponent(originatingWorkbenchId), + ); +} + /** * Compaction interval: fold the WAL into a fresh checkpoint once it holds * this many turns since the last checkpoint. Bounds the WAL tail (and so @@ -272,9 +308,10 @@ export interface DurableConversationStoreOpts { principal: Principal; /** * Stable per-agent key the snapshot is filed under - * (`agent-state//`). The warm single-step agent's stepId is - * the natural key: it is stable across that agent's whole lifetime and - * disjoint from any runId. + * (`agent-state///`). The warm single-step agent's + * stepId is the natural agentKey: it is stable across that agent's whole + * lifetime and disjoint from any runId. Originating workbench is bound + * separately so one warm agent can swap rooms without cloning. */ agentKey: string; } @@ -309,7 +346,8 @@ export interface DurableConversationStore { * Commit the local store's new turn(s) to the substrate as O(1) WAL * appends, folding into a fresh checkpoint when the WAL reaches the * compaction interval. Called synchronously at the run boundary (after - * the agent's send settles). A write failure surfaces. + * the agent's send settles). A write failure surfaces. Flushes the + * currently bound originating-workbench prefix (the composite key). */ mirrorToSubstrate(): Promise; /** @@ -344,6 +382,16 @@ export interface DurableConversationStore { * seeded thread. A metadata write failure surfaces. */ onReplySent(receipt: SendReceipt): Promise; + /** + * Point this store at `originatingWorkbenchId`'s nested substrate + * snapshot. Same room is a no-op. On a change: mirror the current room, + * retarget `agent-state///`, restore that + * snapshot. Missing snapshot starts empty -- the prior mixed + * `agent-state//` blob is never migrated. Returns true when + * the bound room changed (caller must rebuild the warm agent so + * `reactor.start()` loads the restored turns). + */ + bindOriginatingWorkbench(originatingWorkbenchId: string): Promise; } export async function createDurableConversationStore( @@ -370,8 +418,27 @@ export async function createDurableConversationStore( }, }); - const agentStatePrefix = `${WORKFLOW_RUN_AGENT_STATE_PREFIX}/${encodeURIComponent(opts.agentKey)}/`; + // Nested substrate prefix is retargeted by bindOriginatingWorkbench. + // Unbound until the first inbound origin is known -- mirror/restore + // require a bound room so they cannot write the mixed agent-state// + // blob. + let originatingWorkbenchId: string | null = null; + + function requireBoundWorkbenchId(): string { + if (originatingWorkbenchId === null) { + throw new Error( + `sidecar conversation-state: originating workbench is not bound for ${opts.agentKey}; bindOriginatingWorkbench must run before restore or mirror`, + ); + } + return originatingWorkbenchId; + } + function agentStatePrefix(): string { + return durableConversationAgentStatePrefix( + opts.agentKey, + requireBoundWorkbenchId(), + ); + } // The number of mirror boundaries already durably committed (the // checkpoint's folded boundaries plus every appended WAL entry). It is // the seq of the NEXT WAL entry. `null` until learned -- lazily from the @@ -436,11 +503,10 @@ export async function createDurableConversationStore( } function substrateAgentStateFsDir(): string { - const repoDir = opts.substrate.getRepoDir(opts.workflowRunRepoId); - return path.join( - repoDir, - WORKFLOW_RUN_AGENT_STATE_PREFIX, - encodeURIComponent(opts.agentKey), + return durableConversationAgentStateDir( + opts.substrate.getRepoDir(opts.workflowRunRepoId), + opts.agentKey, + requireBoundWorkbenchId(), ); } @@ -449,7 +515,7 @@ export async function createDurableConversationStore( } function walBucketPrefix(bucket: number): string { - return `${agentStatePrefix}${WAL_DIR}/${String(bucket)}/`; + return `${agentStatePrefix()}${WAL_DIR}/${String(bucket)}/`; } function walEntryPath(seq: number): string { @@ -457,25 +523,38 @@ export async function createDurableConversationStore( } function checkpointPath(): string { - return `${agentStatePrefix}${CHECKPOINT_FILE}`; + return `${agentStatePrefix()}${CHECKPOINT_FILE}`; } function checkpointMetaPath(): string { - return `${agentStatePrefix}${CHECKPOINT_META_FILE}`; + return `${agentStatePrefix()}${CHECKPOINT_META_FILE}`; + } + + async function applyEmptyLocalConversation(reason: string): Promise { + await baseStorage.writeTurns([]); + baseStorage.setConnectorState(null); + connectorRouter.restore(null); + await baseStorage.writeMetadata({ + pendingOperations: [], + tokenUsage: EMPTY_TOKEN_USAGE, + }); + await baseStorage.commit({ message: reason }); + mirroredBoundaryCount = 0; + mirroredTurnCount = 0; + checkpointBoundarySeq = 0; } async function runRestore(): Promise { const reconstructed = await reconstructDurableConversation( substrateAgentStateFsDir(), - opts.agentKey, + `${opts.agentKey}/${requireBoundWorkbenchId()}`, ); if (reconstructed === null) { - // No durable state yet: the next mirror starts the WAL from an empty - // checkpoint. Record the (empty) committed counts so the first - // mirror appends from boundary seq 0. - mirroredBoundaryCount = 0; - mirroredTurnCount = 0; - checkpointBoundarySeq = 0; + // No snapshot for this room: start empty. Do not migrate a mixed + // agent-state// blob from before per-room nesting. + await applyEmptyLocalConversation( + `reset conversation for ${opts.agentKey} workbench ${requireBoundWorkbenchId()} (no prior snapshot)`, + ); return false; } // Write the reconstructed turns + metadata into the local store's @@ -587,7 +666,7 @@ export async function createDurableConversationStore( opts.workflowRunRepoId, opts.workflowRunRef, { - preservePrefix: agentStatePrefix, + preservePrefix: agentStatePrefix(), merge: async () => ({ [checkpointPath()]: JSON.stringify(snapshot), [checkpointMetaPath()]: JSON.stringify(meta), @@ -620,7 +699,7 @@ export async function createDurableConversationStore( if (mirroredBoundaryCount === null) { const reconstructed = await reconstructDurableConversation( substrateAgentStateFsDir(), - opts.agentKey, + `${opts.agentKey}/${requireBoundWorkbenchId()}`, ); checkpointBoundarySeq = reconstructed?.checkpointBoundarySeq ?? 0; mirroredBoundaryCount = reconstructed?.boundaryCount ?? 0; @@ -727,6 +806,23 @@ export async function createDurableConversationStore( }); } + function bindOriginatingWorkbench(nextWorkbenchId: string): Promise { + return serializeStateOp(async () => { + if (nextWorkbenchId.length === 0) { + throw new Error( + `sidecar conversation-state: originating workbench id must be non-empty for ${opts.agentKey}`, + ); + } + if (originatingWorkbenchId === nextWorkbenchId) return false; + if (originatingWorkbenchId !== null) { + await runMirror(); + } + originatingWorkbenchId = nextWorkbenchId; + await runRestore(); + return true; + }); + } + return { storage: baseStorage, restoreFromSubstrate, @@ -734,6 +830,7 @@ export async function createDurableConversationStore( seedInbound, composeReply, onReplySent, + bindOriginatingWorkbench, }; } @@ -754,11 +851,9 @@ export interface DurableConversationRegistryOpts { /** * Per-agent durable-conversation store registry. One store - * per warm agent key, built lazily and reused across runs in the same - * child. The first `acquire` for a key builds the store and restores its - * prior conversation snapshot from the substrate -- the path that runs on - * the lazy first build AND on the respawn rebuild, so the warm agent - * resumes its conversation across child respawn. The registry is empty + * per warm agent key (stepId), built lazily and reused across runs in the + * same child. Originating-workbench snapshots nest under that key; bind + * retargets the live store before each send. The registry is empty * after a respawn (it lives in the child's address space); the substrate * is the durable mirror that survives. */ @@ -807,14 +902,9 @@ export function createDurableConversationRegistry( principal: opts.principal, agentKey: key, }); - // Restore the prior conversation BEFORE the store is observable (and - // before the warm agent's reactor `load()` reads it). On a genuine - // first-ever run this is a no-op (no checkpoint/WAL yet); on a - // respawn rebuild it pulls the pre-respawn conversation back from the - // substrate (checkpoint + WAL replay). A restore failure surfaces -- - // a lost conversation on respawn is a correctness failure, not a - // silently-fresh start. - await store.restoreFromSubstrate(); + // Restore happens in bindOriginatingWorkbench once the inbound + // From local-part is known. Restoring here would load a mixed + // agent-state// blob (or throw unbound). stores.set(key, store); building.delete(key); return store; @@ -843,6 +933,30 @@ export function createDurableConversationRegistry( return { acquire, get, peek }; } +/** + * Bind the warm agent's live store to the originating workbench and, when + * the room changed, evict the cached agent so the next send rebuilds + * against the restored snapshot. The cache stays keyed by stepId -- one + * warm agent, not one per room. + */ +export async function prepareConversationForOriginatingWorkbench(args: { + registry: DurableConversationRegistry; + agentKey: string; + originatingWorkbenchId: string; + warmCache?: { evictAll: (reason: string) => Promise }; +}): Promise { + const store = await args.registry.acquire(args.agentKey); + const swapped = await store.bindOriginatingWorkbench( + args.originatingWorkbenchId, + ); + if (swapped && args.warmCache !== undefined) { + await args.warmCache.evictAll( + `originating workbench changed to ${args.originatingWorkbenchId}`, + ); + } + return swapped; +} + interface SnapshotMetadataValue { pendingOperations: unknown[]; tokenUsage: TokenUsage; @@ -877,8 +991,9 @@ function castValidatedEnvelope(items: unknown[]): T[] { /** * Reconstruct the warm agent's conversation from the two-tier on-disk - * layout under `agentStateDir` (`/agent-state//`): the - * compacted `checkpoint.json` turns followed by the replayed WAL tail. + * layout under `agentStateDir` + * (`/agent-state///`): the compacted + * `checkpoint.json` turns followed by the replayed WAL tail. * Pure read against the substrate working tree -- no inference, no commit. * Returns `null` when neither a checkpoint nor any WAL exists (the genuine * first-ever run). The latest metadata source wins (the last replayed WAL diff --git a/apps/sidecar/src/originating-workbench.test.ts b/apps/sidecar/src/originating-workbench.test.ts new file mode 100644 index 000000000..c35a13dc5 --- /dev/null +++ b/apps/sidecar/src/originating-workbench.test.ts @@ -0,0 +1,106 @@ +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +import { afterEach, expect, test } from "bun:test"; + +import { + extractOriginatingWorkbenchId, + readOriginatingWorkbenchId, + recordOriginatingWorkbench, + resolveOriginatingWorkbenchId, + UNSCOPED_ORIGINATING_WORKBENCH_ID, +} from "./originating-workbench"; + +const tmpDirs: string[] = []; + +afterEach(async () => { + await Promise.all( + tmpDirs + .splice(0) + .map((dir) => fs.promises.rm(dir, { recursive: true, force: true })), + ); +}); + +function mailFrom(from: string): Uint8Array { + return new TextEncoder().encode( + [`From: ${from}`, "To: myra@alice.localhost", "", "hello"].join("\r\n"), + ); +} + +test("extractOriginatingWorkbenchId reads the From local-part", () => { + expect(extractOriginatingWorkbenchId(mailFrom("chan_room_a@alice.localhost"))).toBe( + "chan_room_a", + ); + expect( + extractOriginatingWorkbenchId( + mailFrom("Workbench "), + ), + ).toBe("ins_workbench1"); +}); + +test("extractOriginatingWorkbenchId is undefined without a From", () => { + const raw = new TextEncoder().encode("To: myra@alice.localhost\r\n\r\nhello"); + expect(extractOriginatingWorkbenchId(raw)).toBeUndefined(); +}); + +test("resolveOriginatingWorkbenchId uses the unscoped sentinel when From is missing", () => { + expect(resolveOriginatingWorkbenchId(undefined)).toBe( + UNSCOPED_ORIGINATING_WORKBENCH_ID, + ); + expect(resolveOriginatingWorkbenchId("chan_a")).toBe("chan_a"); +}); + +test("record and read round-trip the originating workbench id", async () => { + const dataDir = fs.mkdtempSync(path.join(os.tmpdir(), "origin-wb-")); + tmpDirs.push(dataDir); + await recordOriginatingWorkbench({ + dataDir, + mailboxAddress: "ins_myra@alice.localhost", + raw: mailFrom("chan_room_b@alice.localhost"), + }); + expect( + await readOriginatingWorkbenchId({ + dataDir, + mailboxAddress: "ins_myra@alice.localhost", + }), + ).toBe("chan_room_b"); +}); + +test("a later mail overwrites a stale room id", async () => { + const dataDir = fs.mkdtempSync(path.join(os.tmpdir(), "origin-wb-")); + tmpDirs.push(dataDir); + const mailboxAddress = "ins_myra@alice.localhost"; + await recordOriginatingWorkbench({ + dataDir, + mailboxAddress, + raw: mailFrom("chan_a@alice.localhost"), + }); + await recordOriginatingWorkbench({ + dataDir, + mailboxAddress, + raw: mailFrom("chan_b@alice.localhost"), + }); + expect(await readOriginatingWorkbenchId({ dataDir, mailboxAddress })).toBe( + "chan_b", + ); +}); + +test("unparseable mail records the unscoped sentinel rather than leaving a stale room", async () => { + const dataDir = fs.mkdtempSync(path.join(os.tmpdir(), "origin-wb-")); + tmpDirs.push(dataDir); + const mailboxAddress = "ins_myra@alice.localhost"; + await recordOriginatingWorkbench({ + dataDir, + mailboxAddress, + raw: mailFrom("chan_a@alice.localhost"), + }); + await recordOriginatingWorkbench({ + dataDir, + mailboxAddress, + raw: new TextEncoder().encode("not-a-mail"), + }); + expect(await readOriginatingWorkbenchId({ dataDir, mailboxAddress })).toBe( + UNSCOPED_ORIGINATING_WORKBENCH_ID, + ); +}); diff --git a/apps/sidecar/src/originating-workbench.ts b/apps/sidecar/src/originating-workbench.ts new file mode 100644 index 000000000..2a4a1d8c0 --- /dev/null +++ b/apps/sidecar/src/originating-workbench.ts @@ -0,0 +1,102 @@ +// Originating workbench for the warm agent's durable conversation. +// +// Chat mail names the room it speaks for in the RFC 2822 From local-part +// (`fromWorkbenchId`). The supervisor records that id when inbound mail +// arrives; the workflow-process child reads it before each send and binds +// the agent's durable conversation to `agent-state///`. +// Missing or unparseable From is `_unscoped` so a later room never inherits +// a mixed blob. + +import fs from "node:fs"; +import path from "node:path"; + +import { type } from "arktype"; +import { extractAddrSpec, parseHeaderSection } from "@intx/mime"; + +import { writeFileAtomicDurable } from "./atomic-write"; +import { isErrnoNotFound } from "./conversation-state"; + +export const UNSCOPED_ORIGINATING_WORKBENCH_ID = "_unscoped"; + +const OriginRecord = type({ + fromWorkbenchId: "string", +}); + +function originRecordPath(dataDir: string, mailboxAddress: string): string { + return path.join( + dataDir, + "inbound-origin", + `${encodeURIComponent(mailboxAddress)}.json`, + ); +} + +/** + * Local-part of the inbound MIME From header. Chat encodes the originating + * workbench id there; anything else (no From, unparseable) is undefined. + */ +export function extractOriginatingWorkbenchId( + raw: Uint8Array, +): string | undefined { + try { + const { headers } = parseHeaderSection(raw); + const from = headers.get("from"); + if (from === undefined || from.trim() === "") return undefined; + const spec = extractAddrSpec(from); + const at = spec.lastIndexOf("@"); + if (at <= 0) return undefined; + const local = spec.slice(0, at); + return local.length > 0 ? local : undefined; + } catch { + return undefined; + } +} + +export function resolveOriginatingWorkbenchId( + fromWorkbenchId: string | undefined, +): string { + return fromWorkbenchId !== undefined && fromWorkbenchId.length > 0 + ? fromWorkbenchId + : UNSCOPED_ORIGINATING_WORKBENCH_ID; +} + +/** + * Persist the originating workbench id from inbound mail so the child can + * bind conversation state before send. Never throws: a parse failure records + * the unscoped sentinel rather than leaving a stale room id on disk. + */ +export async function recordOriginatingWorkbench(args: { + dataDir: string; + mailboxAddress: string; + raw: Uint8Array; +}): Promise { + const fromWorkbenchId = resolveOriginatingWorkbenchId( + extractOriginatingWorkbenchId(args.raw), + ); + const target = originRecordPath(args.dataDir, args.mailboxAddress); + await fs.promises.mkdir(path.dirname(target), { recursive: true }); + await writeFileAtomicDurable( + target, + JSON.stringify({ fromWorkbenchId }), + { mode: 0o600 }, + ); +} + +export async function readOriginatingWorkbenchId(args: { + dataDir: string; + mailboxAddress: string; +}): Promise { + const target = originRecordPath(args.dataDir, args.mailboxAddress); + let raw: string; + try { + raw = await fs.promises.readFile(target, "utf8"); + } catch (cause) { + if (isErrnoNotFound(cause)) return UNSCOPED_ORIGINATING_WORKBENCH_ID; + throw cause; + } + const parsed: unknown = JSON.parse(raw); + const validated = OriginRecord(parsed); + if (validated instanceof type.errors || validated.fromWorkbenchId.length === 0) { + return UNSCOPED_ORIGINATING_WORKBENCH_ID; + } + return validated.fromWorkbenchId; +} diff --git a/apps/sidecar/src/workflow-host-wiring/index.ts b/apps/sidecar/src/workflow-host-wiring/index.ts index 9d4319fa2..7e37cd0aa 100644 --- a/apps/sidecar/src/workflow-host-wiring/index.ts +++ b/apps/sidecar/src/workflow-host-wiring/index.ts @@ -70,6 +70,7 @@ import { snapshotAgentIdentity, } from "../hibernated-agent-identity-vault"; import { isErrnoNotFound } from "../conversation-state"; +import { recordOriginatingWorkbench } from "../originating-workbench"; import { computeWireDefinitionHash, validateWorkflowProjection, @@ -1042,7 +1043,18 @@ export function createSidecarDeployRouter(deps: { // supervisor's mail-bus subscription. Registration happens after // `spawn` succeeds so a spawn-time rejection leaves the registry // untouched. - deps.multistepMailRouter?.register(spec.agentAddress, (message) => { + deps.multistepMailRouter?.register(spec.agentAddress, async (message) => { + if (stepStateDataDir !== undefined) { + try { + await recordOriginatingWorkbench({ + dataDir: stepStateDataDir, + mailboxAddress: spec.agentAddress, + raw: message, + }); + } catch (cause) { + logger.warn`record originating workbench failed for ${spec.agentAddress}: ${cause instanceof Error ? cause.message : String(cause)}`; + } + } return wired.routeInbound(message); }); // Register the signal-delivery handler so a hub `signal.deliver` frame diff --git a/apps/sidecar/src/workflow-substrate-factory/index.ts b/apps/sidecar/src/workflow-substrate-factory/index.ts index bc1c2d4c8..cb65bc667 100644 --- a/apps/sidecar/src/workflow-substrate-factory/index.ts +++ b/apps/sidecar/src/workflow-substrate-factory/index.ts @@ -70,8 +70,10 @@ import { } from "../step-agent-tools"; import { createDurableConversationRegistry, + prepareConversationForOriginatingWorkbench, type DurableConversationRegistry, } from "../conversation-state"; +import { readOriginatingWorkbenchId } from "../originating-workbench"; import { parseToolRegistries } from "../tool-materialization"; import { deriveHubHttpUrl, @@ -483,6 +485,17 @@ export function createSidecarSubstrateFactory( credentialWiring, mailPartReader, ) => { + const bodyStepId = req.authzContext.stepId; + if (durableConversation !== undefined && bodyStepId !== undefined) { + await prepareConversationForOriginatingWorkbench({ + registry: durableConversation, + agentKey: bodyStepId, + originatingWorkbenchId: await readOriginatingWorkbenchId({ + dataDir: validated.SIDECAR_DATA_DIR, + mailboxAddress: env.spawn.mailboxAddress, + }), + }); + } try { return await createWorkflowStepInvoker({ workflowAuthorize: authorize, @@ -494,11 +507,10 @@ export function createSidecarSubstrateFactory( ...(mailPartReader !== undefined ? { mailPartReader } : {}), })(req); } finally { - const bodyStepId = req.authzContext.stepId; if (durableConversation !== undefined && bodyStepId !== undefined) { // `peek`, not `get`: a build failure before the env's acquire // must surface as itself, not as the registry's missing-store - // throw. + // throw. Flushes the currently bound workbench prefix. await durableConversation.peek(bodyStepId)?.mirrorToSubstrate(); } } @@ -541,10 +553,10 @@ export function createSidecarSubstrateFactory( // lifecycle. // Run-boundary durability flush. When the deployment is // warm-kept, mirror the warm agent's conversation snapshot to the - // workflow-run substrate after each message's send settles. The key - // is the step identity, the same key the env builder filed the - // durable store under, so the hook resolves the right per-agent - // store. Absent for a multi-step deploy (no durable registry). + // currently bound originating-workbench prefix after each message's + // send settles. The registry key is still the step identity (one warm + // agent); the store flushes `agent-state///`. + // Absent for a multi-step deploy (no durable registry). const onRunBoundary: ((key: string) => Promise) | undefined = durableConversation !== undefined ? async (key: string) => { @@ -603,8 +615,20 @@ export function createSidecarSubstrateFactory( sourcesRef, credentialWiring, mailPartReader, - ) => - createWorkflowStepInvoker({ + ) => { + const stepId = req.authzContext.stepId; + if (durableConversation !== undefined && stepId !== undefined) { + await prepareConversationForOriginatingWorkbench({ + registry: durableConversation, + agentKey: stepId, + originatingWorkbenchId: await readOriginatingWorkbenchId({ + dataDir: validated.SIDECAR_DATA_DIR, + mailboxAddress: env.spawn.mailboxAddress, + }), + ...(warmCache !== undefined ? { warmCache } : {}), + }); + } + return createWorkflowStepInvoker({ workflowAuthorize: authorize, buildEnv: (buildReq: Parameters[0]) => buildStepEnv(buildReq, sourcesRef, credentialWiring), @@ -617,6 +641,7 @@ export function createSidecarSubstrateFactory( ...(seedInbound !== undefined ? { seedInbound } : {}), ...(driveReplies !== undefined ? { driveReplies } : {}), })(req); + }; const evaluateGrantsAdapter: GrantEvaluator = async ({ resource, diff --git a/apps/sidecar/src/workflow-substrate-factory/parked-approvals.ts b/apps/sidecar/src/workflow-substrate-factory/parked-approvals.ts index 3c18efd75..f558ae2de 100644 --- a/apps/sidecar/src/workflow-substrate-factory/parked-approvals.ts +++ b/apps/sidecar/src/workflow-substrate-factory/parked-approvals.ts @@ -92,7 +92,7 @@ export async function readColdParkedApprovalSnapshot(args: { * Read a warm (single-step) parked agent's durable pending operations * from substrate state. A warm agent's pending operations live in its * durable conversation store, mirrored to the workflow-run substrate - * under `agent-state//`. + * under `agent-state///`. * * Reconstructs that state read-only — deliberately NOT through * `DurableConversationRegistry.acquire`, whose first acquire writes and @@ -101,7 +101,8 @@ export async function readColdParkedApprovalSnapshot(args: { * rebuilt the live store when re-registration runs (resume re-parks * without re-invoking the step), so the substrate is the only place the * pending operations live at that moment. Returns an empty list when no - * durable state exists for the agent. + * durable state exists for the agent. Walks every nested workbench dir + * so a park in any room is visible. */ export async function readWarmParkedPendingOperations(args: { substrate: RepoStore; @@ -113,12 +114,25 @@ export async function readWarmParkedPendingOperations(args: { WORKFLOW_RUN_AGENT_STATE_PREFIX, encodeURIComponent(args.stepId), ); - const reconstructed = await reconstructDurableConversation( - agentStateDir, - args.stepId, - ); - if (reconstructed === null) return []; - return reconstructed.pendingOperations; + let children: fs.Dirent[]; + try { + children = await fs.promises.readdir(agentStateDir, { withFileTypes: true }); + } catch (cause) { + if (isErrnoNotFound(cause)) return []; + throw cause; + } + const pending: PendingOperation[] = []; + for (const child of children) { + if (!child.isDirectory()) continue; + const reconstructed = await reconstructDurableConversation( + path.join(agentStateDir, child.name), + `${args.stepId}/${child.name}`, + ); + if (reconstructed !== null) { + pending.push(...reconstructed.pendingOperations); + } + } + return pending; } export async function readWarmParkedApprovalSnapshot(args: { diff --git a/apps/sidecar/src/workflow-substrate-factory/step-env.ts b/apps/sidecar/src/workflow-substrate-factory/step-env.ts index abeb25139..138a0a217 100644 --- a/apps/sidecar/src/workflow-substrate-factory/step-env.ts +++ b/apps/sidecar/src/workflow-substrate-factory/step-env.ts @@ -291,13 +291,12 @@ export function createSidecarStepBuildEnv( }); // Conversation storage. For the warm single-step agent the // conversation must survive child respawn, so it is backed by a - // per-agent durable store whose content is mirrored to the - // workflow-run substrate; building it here restores the - // prior conversation before the agent's reactor loads. A multi-step - // deploy (no durable registry) keeps the per-run isogit store: its - // per-step agents are not warm/long-lived and have no cross-run - // conversation to carry. The workdir + tools stay per-run in both - // cases -- only the conversation context is durable across runs. + // per-agent durable store (keyed by stepId) whose content is mirrored + // to the workflow-run substrate under + // `agent-state///`. Bind to the originating + // workbench happens before this builder runs so restore is already + // applied. A multi-step deploy (no durable registry) keeps the + // per-run isogit store. const storage: ContextStore & AuditStore = deps.durableConversation !== undefined ? (await deps.durableConversation.acquire(stepId)).storage diff --git a/bun.lock b/bun.lock index 4db1d162e..872b272d8 100644 --- a/bun.lock +++ b/bun.lock @@ -120,6 +120,7 @@ "@intx/inference": "workspace:*", "@intx/log": "0.3.0", "@intx/mail-memory": "workspace:*", + "@intx/mime": "workspace:*", "@intx/storage-isogit": "0.3.0", "@intx/tool-packaging": "0.3.0", "@intx/types": "workspace:*", From e290dbd0ab26beb05e61c2b822d2c908f72358e3 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Fri, 28 Aug 2026 09:13:37 -0700 Subject: [PATCH 2/2] Read the originating workbench from this request's mail, not a latest-wins file Supervisor enqueue overwrote one origin file per mailbox, so two queued mails bound the first under the second's workbench. The step request already carries this message's Mail as its input (the object the workflow-host step invoker sends to the agent); read From from it. A mail-less request (an approval resume) keeps the current binding. No vendored workflow-host delta. --- apps/sidecar/src/conversation-state.test.ts | 98 ++++++++++++--- apps/sidecar/src/conversation-state.ts | 44 ++++--- .../sidecar/src/originating-workbench.test.ts | 115 +++++++----------- apps/sidecar/src/originating-workbench.ts | 98 ++++----------- .../sidecar/src/workflow-host-wiring/index.ts | 12 -- .../src/workflow-substrate-factory/index.ts | 12 +- .../parked-approvals.test.ts | 66 +++++++++- .../parked-approvals.ts | 4 +- 8 files changed, 247 insertions(+), 202 deletions(-) diff --git a/apps/sidecar/src/conversation-state.test.ts b/apps/sidecar/src/conversation-state.test.ts index e09c6487b..b905f68ed 100644 --- a/apps/sidecar/src/conversation-state.test.ts +++ b/apps/sidecar/src/conversation-state.test.ts @@ -4,7 +4,11 @@ import path from "node:path"; import { afterEach, expect, test } from "bun:test"; -import type { Principal, RepoId, RepoStore } from "@intx/hub-sessions/substrate"; +import type { + Principal, + RepoId, + RepoStore, +} from "@intx/hub-sessions/substrate"; import type { ConversationTurn } from "@intx/types/runtime"; import { @@ -14,6 +18,7 @@ import { prepareConversationForOriginatingWorkbench, reconstructDurableConversation, } from "./conversation-state"; +import { originatingWorkbenchIdFromRequest } from "./originating-workbench"; const tmpDirs: string[] = []; @@ -56,21 +61,21 @@ function fsBackedSubstrate(repoDir: string): RepoStore { ) => { const prefixDir = path.join( repoDir, - ...args.preservePrefix.split("/").filter((seg: string) => seg.length > 0), + ...args.preservePrefix + .split("/") + .filter((seg: string) => seg.length > 0), ); const existing = new Map(); let names: string[] = []; try { names = await fs.promises.readdir(prefixDir); } catch (cause) { - if ( - !( - typeof cause === "object" && - cause !== null && - "code" in cause && - cause.code === "ENOENT" - ) - ) { + if (!( + typeof cause === "object" && + cause !== null && + "code" in cause && + cause.code === "ENOENT" + )) { throw cause; } } @@ -78,7 +83,10 @@ function fsBackedSubstrate(repoDir: string): RepoStore { const full = path.join(prefixDir, name); const st = await fs.promises.stat(full); if (st.isFile()) { - existing.set(`${args.preservePrefix}${name}`, await fs.promises.readFile(full)); + existing.set( + `${args.preservePrefix}${name}`, + await fs.promises.readFile(full), + ); } } const files = await args.merge(existing); @@ -127,7 +135,11 @@ test("two rooms of one agent keep isolated turns; a new room starts empty", asyn await store.storage.writeTurns([userTurn("room A tool result")]); await store.mirrorToSubstrate(); - const roomADir = durableConversationAgentStateDir(repoDir, "default", "chan_a"); + const roomADir = durableConversationAgentStateDir( + repoDir, + "default", + "chan_a", + ); const reconstructedA = await reconstructDurableConversation( roomADir, "default/chan_a", @@ -137,7 +149,12 @@ test("two rooms of one agent keep isolated turns; a new room starts empty", asyn expect(await store.bindOriginatingWorkbench("chan_b")).toBe(true); expect(peekTurns(store)).toEqual([]); - const mixedLegacy = path.join(repoDir, "agent-state", "default", "checkpoint.json"); + const mixedLegacy = path.join( + repoDir, + "agent-state", + "default", + "checkpoint.json", + ); expect(fs.existsSync(mixedLegacy)).toBe(false); await store.storage.writeTurns([userTurn("room B first infer")]); @@ -152,11 +169,62 @@ test("two rooms of one agent keep isolated turns; a new room starts empty", asyn expect(await store.bindOriginatingWorkbench("chan_b")).toBe(false); }); +test("two queued mails bind and mirror under each message's From not a later origin", async () => { + const root = tmpDir("conv-queued-"); + const { store } = await makeStore(root); + const registry = { + acquire: () => Promise.resolve(store), + get: () => store, + peek: () => store, + }; + const mailRequest = (from: string) => ({ + input: { + headers: { from, to: ["myra@alice.localhost"] }, + rawHeaders: {}, + parts: [], + }, + }); + // Both inbound mails exist (enqueued) before the first invokeStep. A + // latest-wins origin file would now name chan_b for both binds. + const fromA = originatingWorkbenchIdFromRequest( + mailRequest("chan_a@alice.localhost"), + ); + const fromB = originatingWorkbenchIdFromRequest( + mailRequest("chan_b@alice.localhost"), + ); + expect(fromA).toBe("chan_a"); + expect(fromB).toBe("chan_b"); + + await prepareConversationForOriginatingWorkbench({ + registry: registry as never, + agentKey: "default", + originatingWorkbenchId: fromA, + }); + await store.storage.writeTurns([userTurn("from A")]); + await store.mirrorToSubstrate(); + + await prepareConversationForOriginatingWorkbench({ + registry: registry as never, + agentKey: "default", + originatingWorkbenchId: fromB, + }); + expect(peekTurns(store)).toEqual([]); + await store.storage.writeTurns([userTurn("from B")]); + await store.mirrorToSubstrate(); + + await prepareConversationForOriginatingWorkbench({ + registry: registry as never, + agentKey: "default", + originatingWorkbenchId: fromA, + }); + expect(peekTurns(store)).toEqual([userTurn("from A")]); +}); + test("prepareConversationForOriginatingWorkbench evicts the warm agent on room change only", async () => { const evictions: string[] = []; const store = { - bindOriginatingWorkbench: (id: string) => - Promise.resolve(id === "chan_b"), + bindOriginatingWorkbench: (id: string) => Promise.resolve(id === "chan_b"), + boundOriginatingWorkbenchId: () => null, }; const registry = { acquire: () => Promise.resolve(store), diff --git a/apps/sidecar/src/conversation-state.ts b/apps/sidecar/src/conversation-state.ts index 60092c53c..b411ef8b4 100644 --- a/apps/sidecar/src/conversation-state.ts +++ b/apps/sidecar/src/conversation-state.ts @@ -143,6 +143,8 @@ import type { RepoStore, } from "@intx/hub-sessions/substrate"; import { WORKFLOW_RUN_AGENT_STATE_PREFIX } from "@intx/hub-sessions/substrate"; + +import { UNSCOPED_ORIGINATING_WORKBENCH_ID } from "./originating-workbench"; import { ConnectorThreadState, TokenUsage, @@ -229,8 +231,8 @@ const SnapshotMetadata = type({ /** * On-disk shape of the compacted checkpoint blob committed at - * `agent-state//checkpoint.json`. Carries the folded turn - * history (turns 0..checkpointSeq-1) plus the non-turn reactor metadata. + * `agent-state///checkpoint.json`. Carries the folded + * turn history (turns 0..checkpointSeq-1) plus the non-turn reactor metadata. * Validated on read because it crosses back into the program from the * substrate working tree -- a corrupt or partially-written checkpoint must * surface at the boundary, never be half-applied into the agent. @@ -243,8 +245,9 @@ const CheckpointSnapshot = type({ }); /** - * On-disk shape of `agent-state//checkpoint.meta.json`. The - * checkpoint pointer the restore path reads first to learn the boundary seq + * On-disk shape of + * `agent-state///checkpoint.meta.json`. The checkpoint + * pointer the restore path reads first to learn the boundary seq * the checkpoint folded to (`checkpointSeq`) -- and therefore which WAL * boundary seqs remain to replay -- plus the folded turn count * (`turnCount`, = `checkpoint.json`'s turn array length) and the freshest @@ -263,8 +266,8 @@ const CheckpointMeta = type({ /** * On-disk shape of one WAL entry blob at - * `agent-state//wal//.json`. One entry per MIRROR - * BOUNDARY (keyed by boundary `seq`, not turn index). Records the 0-or-more + * `agent-state///wal//.json`. One entry + * per MIRROR BOUNDARY (keyed by boundary `seq`, not turn index). Records the 0-or-more * new turns that boundary added (the O(1) append payload -- it never * carries prior turns) plus the latest non-turn metadata snapshot. The * append is UNCONDITIONAL: a turnless boundary still writes one entry with @@ -392,6 +395,8 @@ export interface DurableConversationStore { * `reactor.start()` loads the restored turns). */ bindOriginatingWorkbench(originatingWorkbenchId: string): Promise; + /** Room this store is bound to, or null before the first bind. */ + boundOriginatingWorkbenchId(): string | null; } export async function createDurableConversationStore( @@ -627,12 +632,12 @@ export async function createDurableConversationStore( } /** - * Fold the full conversation into a fresh checkpoint and truncate the - * WAL in one atomic commit at `preservePrefix = agent-state//`. The - * merge returns ONLY the two checkpoint files and NO `wal/...` paths; - * because the substrate's `clearPrefix` recursively removes the whole - * `agent-state//` subtree before writing the returned set, omitting - * the WAL paths IS the truncate. + * Fold the full conversation into a fresh checkpoint and truncate the WAL in + * one atomic commit at `preservePrefix = agent-state///`. + * The merge returns ONLY the two checkpoint files and NO `wal/...` paths; + * because the substrate's `clearPrefix` recursively removes the whole room + * subtree before writing the returned set, omitting the WAL paths IS the + * truncate. */ async function writeCheckpoint( boundarySeq: number, @@ -831,6 +836,7 @@ export async function createDurableConversationStore( composeReply, onReplySent, bindOriginatingWorkbench, + boundOriginatingWorkbenchId: () => originatingWorkbenchId, }; } @@ -942,17 +948,17 @@ export function createDurableConversationRegistry( export async function prepareConversationForOriginatingWorkbench(args: { registry: DurableConversationRegistry; agentKey: string; - originatingWorkbenchId: string; + /** Room named by this request's mail; undefined keeps the current binding. */ + originatingWorkbenchId: string | undefined; warmCache?: { evictAll: (reason: string) => Promise }; }): Promise { const store = await args.registry.acquire(args.agentKey); - const swapped = await store.bindOriginatingWorkbench( - args.originatingWorkbenchId, - ); + const requested = args.originatingWorkbenchId; + const current = store.boundOriginatingWorkbenchId(); + const target = requested ?? current ?? UNSCOPED_ORIGINATING_WORKBENCH_ID; + const swapped = await store.bindOriginatingWorkbench(target); if (swapped && args.warmCache !== undefined) { - await args.warmCache.evictAll( - `originating workbench changed to ${args.originatingWorkbenchId}`, - ); + await args.warmCache.evictAll(`originating workbench changed to ${target}`); } return swapped; } diff --git a/apps/sidecar/src/originating-workbench.test.ts b/apps/sidecar/src/originating-workbench.test.ts index c35a13dc5..e56009008 100644 --- a/apps/sidecar/src/originating-workbench.test.ts +++ b/apps/sidecar/src/originating-workbench.test.ts @@ -1,47 +1,34 @@ -import fs from "node:fs"; -import os from "node:os"; -import path from "node:path"; - -import { afterEach, expect, test } from "bun:test"; +import { expect, test } from "bun:test"; import { extractOriginatingWorkbenchId, - readOriginatingWorkbenchId, - recordOriginatingWorkbench, + originatingWorkbenchIdFromRequest, resolveOriginatingWorkbenchId, UNSCOPED_ORIGINATING_WORKBENCH_ID, } from "./originating-workbench"; -const tmpDirs: string[] = []; - -afterEach(async () => { - await Promise.all( - tmpDirs - .splice(0) - .map((dir) => fs.promises.rm(dir, { recursive: true, force: true })), - ); -}); - -function mailFrom(from: string): Uint8Array { - return new TextEncoder().encode( - [`From: ${from}`, "To: myra@alice.localhost", "", "hello"].join("\r\n"), - ); +function mailRequest(from: string) { + return { + input: { + headers: { from, to: ["myra@alice.localhost"] }, + rawHeaders: {}, + parts: [], + }, + }; } test("extractOriginatingWorkbenchId reads the From local-part", () => { - expect(extractOriginatingWorkbenchId(mailFrom("chan_room_a@alice.localhost"))).toBe( + expect(extractOriginatingWorkbenchId("chan_room_a@alice.localhost")).toBe( "chan_room_a", ); expect( - extractOriginatingWorkbenchId( - mailFrom("Workbench "), - ), + extractOriginatingWorkbenchId("Workbench "), ).toBe("ins_workbench1"); }); -test("extractOriginatingWorkbenchId is undefined without a From", () => { - const raw = new TextEncoder().encode("To: myra@alice.localhost\r\n\r\nhello"); - expect(extractOriginatingWorkbenchId(raw)).toBeUndefined(); +test("extractOriginatingWorkbenchId is undefined for an empty or bare From", () => { + expect(extractOriginatingWorkbenchId("")).toBeUndefined(); + expect(extractOriginatingWorkbenchId("@alice.localhost")).toBeUndefined(); }); test("resolveOriginatingWorkbenchId uses the unscoped sentinel when From is missing", () => { @@ -51,56 +38,36 @@ test("resolveOriginatingWorkbenchId uses the unscoped sentinel when From is miss expect(resolveOriginatingWorkbenchId("chan_a")).toBe("chan_a"); }); -test("record and read round-trip the originating workbench id", async () => { - const dataDir = fs.mkdtempSync(path.join(os.tmpdir(), "origin-wb-")); - tmpDirs.push(dataDir); - await recordOriginatingWorkbench({ - dataDir, - mailboxAddress: "ins_myra@alice.localhost", - raw: mailFrom("chan_room_b@alice.localhost"), - }); +test("originatingWorkbenchIdFromRequest is a function of this request's mail", () => { expect( - await readOriginatingWorkbenchId({ - dataDir, - mailboxAddress: "ins_myra@alice.localhost", - }), - ).toBe("chan_room_b"); + originatingWorkbenchIdFromRequest(mailRequest("chan_a@alice.localhost")), + ).toBe("chan_a"); + expect( + originatingWorkbenchIdFromRequest(mailRequest("chan_b@alice.localhost")), + ).toBe("chan_b"); }); -test("a later mail overwrites a stale room id", async () => { - const dataDir = fs.mkdtempSync(path.join(os.tmpdir(), "origin-wb-")); - tmpDirs.push(dataDir); - const mailboxAddress = "ins_myra@alice.localhost"; - await recordOriginatingWorkbench({ - dataDir, - mailboxAddress, - raw: mailFrom("chan_a@alice.localhost"), - }); - await recordOriginatingWorkbench({ - dataDir, - mailboxAddress, - raw: mailFrom("chan_b@alice.localhost"), - }); - expect(await readOriginatingWorkbenchId({ dataDir, mailboxAddress })).toBe( - "chan_b", - ); +test("a resumed input delivers the decision's mail, not the original input", () => { + expect( + originatingWorkbenchIdFromRequest({ + ...mailRequest("chan_a@alice.localhost"), + resume: { + correlationId: "c1", + kind: "input", + decision: mailRequest("chan_b@alice.localhost").input, + }, + }), + ).toBe("chan_b"); }); -test("unparseable mail records the unscoped sentinel rather than leaving a stale room", async () => { - const dataDir = fs.mkdtempSync(path.join(os.tmpdir(), "origin-wb-")); - tmpDirs.push(dataDir); - const mailboxAddress = "ins_myra@alice.localhost"; - await recordOriginatingWorkbench({ - dataDir, - mailboxAddress, - raw: mailFrom("chan_a@alice.localhost"), - }); - await recordOriginatingWorkbench({ - dataDir, - mailboxAddress, - raw: new TextEncoder().encode("not-a-mail"), - }); - expect(await readOriginatingWorkbenchId({ dataDir, mailboxAddress })).toBe( - UNSCOPED_ORIGINATING_WORKBENCH_ID, +test("a request without mail names no room", () => { + expect(originatingWorkbenchIdFromRequest({ input: "plain text" })).toBe( + undefined, ); + expect( + originatingWorkbenchIdFromRequest({ + input: mailRequest("chan_a@alice.localhost").input, + resume: { correlationId: "c1", kind: "approval", decision: "approved" }, + }), + ).toBe(undefined); }); diff --git a/apps/sidecar/src/originating-workbench.ts b/apps/sidecar/src/originating-workbench.ts index 2a4a1d8c0..db2094258 100644 --- a/apps/sidecar/src/originating-workbench.ts +++ b/apps/sidecar/src/originating-workbench.ts @@ -1,51 +1,34 @@ // Originating workbench for the warm agent's durable conversation. // // Chat mail names the room it speaks for in the RFC 2822 From local-part -// (`fromWorkbenchId`). The supervisor records that id when inbound mail -// arrives; the workflow-process child reads it before each send and binds -// the agent's durable conversation to `agent-state///`. -// Missing or unparseable From is `_unscoped` so a later room never inherits -// a mixed blob. - -import fs from "node:fs"; -import path from "node:path"; - -import { type } from "arktype"; -import { extractAddrSpec, parseHeaderSection } from "@intx/mime"; - -import { writeFileAtomicDurable } from "./atomic-write"; -import { isErrnoNotFound } from "./conversation-state"; +// (`fromWorkbenchId`). The step request carries THIS message's Mail as its +// input (the same object the workflow-host step invoker sends to the agent), +// so the room is read from the request itself -- never from a per-mailbox +// side-channel, which binds an earlier mail under a later room when two +// inbound messages enqueue before the first invokeStep. A request that +// carries no mail (an approval resume, a synthetic input) names no room and +// leaves the current binding alone; a mail without a usable From is +// `_unscoped` so a later room never inherits a mixed blob. + +import { extractAddrSpec } from "@intx/mime"; +import { isMail } from "@intx/types/runtime"; +import type { StepInvokeRequest } from "@intx/workflow"; export const UNSCOPED_ORIGINATING_WORKBENCH_ID = "_unscoped"; -const OriginRecord = type({ - fromWorkbenchId: "string", -}); - -function originRecordPath(dataDir: string, mailboxAddress: string): string { - return path.join( - dataDir, - "inbound-origin", - `${encodeURIComponent(mailboxAddress)}.json`, - ); -} - /** - * Local-part of the inbound MIME From header. Chat encodes the originating - * workbench id there; anything else (no From, unparseable) is undefined. + * Local-part of a From header value. Chat encodes the originating workbench + * id there; anything else (empty, unparseable) is undefined. */ export function extractOriginatingWorkbenchId( - raw: Uint8Array, + from: string, ): string | undefined { + if (from.trim() === "") return undefined; try { - const { headers } = parseHeaderSection(raw); - const from = headers.get("from"); - if (from === undefined || from.trim() === "") return undefined; const spec = extractAddrSpec(from); const at = spec.lastIndexOf("@"); if (at <= 0) return undefined; - const local = spec.slice(0, at); - return local.length > 0 ? local : undefined; + return spec.slice(0, at); } catch { return undefined; } @@ -60,43 +43,16 @@ export function resolveOriginatingWorkbenchId( } /** - * Persist the originating workbench id from inbound mail so the child can - * bind conversation state before send. Never throws: a parse failure records - * the unscoped sentinel rather than leaving a stale room id on disk. + * Room named by the mail this step request delivers, or undefined when the + * request carries no mail. Mirrors the step invoker's own input selection: + * a resume delivers its decision, anything else delivers `input`. */ -export async function recordOriginatingWorkbench(args: { - dataDir: string; - mailboxAddress: string; - raw: Uint8Array; -}): Promise { - const fromWorkbenchId = resolveOriginatingWorkbenchId( - extractOriginatingWorkbenchId(args.raw), - ); - const target = originRecordPath(args.dataDir, args.mailboxAddress); - await fs.promises.mkdir(path.dirname(target), { recursive: true }); - await writeFileAtomicDurable( - target, - JSON.stringify({ fromWorkbenchId }), - { mode: 0o600 }, +export function originatingWorkbenchIdFromRequest( + req: Pick, +): string | undefined { + const delivered = req.resume === undefined ? req.input : req.resume.decision; + if (!isMail(delivered)) return undefined; + return resolveOriginatingWorkbenchId( + extractOriginatingWorkbenchId(delivered.headers.from), ); } - -export async function readOriginatingWorkbenchId(args: { - dataDir: string; - mailboxAddress: string; -}): Promise { - const target = originRecordPath(args.dataDir, args.mailboxAddress); - let raw: string; - try { - raw = await fs.promises.readFile(target, "utf8"); - } catch (cause) { - if (isErrnoNotFound(cause)) return UNSCOPED_ORIGINATING_WORKBENCH_ID; - throw cause; - } - const parsed: unknown = JSON.parse(raw); - const validated = OriginRecord(parsed); - if (validated instanceof type.errors || validated.fromWorkbenchId.length === 0) { - return UNSCOPED_ORIGINATING_WORKBENCH_ID; - } - return validated.fromWorkbenchId; -} diff --git a/apps/sidecar/src/workflow-host-wiring/index.ts b/apps/sidecar/src/workflow-host-wiring/index.ts index 7e37cd0aa..86ea0b7b6 100644 --- a/apps/sidecar/src/workflow-host-wiring/index.ts +++ b/apps/sidecar/src/workflow-host-wiring/index.ts @@ -70,7 +70,6 @@ import { snapshotAgentIdentity, } from "../hibernated-agent-identity-vault"; import { isErrnoNotFound } from "../conversation-state"; -import { recordOriginatingWorkbench } from "../originating-workbench"; import { computeWireDefinitionHash, validateWorkflowProjection, @@ -1044,17 +1043,6 @@ export function createSidecarDeployRouter(deps: { // `spawn` succeeds so a spawn-time rejection leaves the registry // untouched. deps.multistepMailRouter?.register(spec.agentAddress, async (message) => { - if (stepStateDataDir !== undefined) { - try { - await recordOriginatingWorkbench({ - dataDir: stepStateDataDir, - mailboxAddress: spec.agentAddress, - raw: message, - }); - } catch (cause) { - logger.warn`record originating workbench failed for ${spec.agentAddress}: ${cause instanceof Error ? cause.message : String(cause)}`; - } - } return wired.routeInbound(message); }); // Register the signal-delivery handler so a hub `signal.deliver` frame diff --git a/apps/sidecar/src/workflow-substrate-factory/index.ts b/apps/sidecar/src/workflow-substrate-factory/index.ts index cb65bc667..089e76606 100644 --- a/apps/sidecar/src/workflow-substrate-factory/index.ts +++ b/apps/sidecar/src/workflow-substrate-factory/index.ts @@ -73,7 +73,7 @@ import { prepareConversationForOriginatingWorkbench, type DurableConversationRegistry, } from "../conversation-state"; -import { readOriginatingWorkbenchId } from "../originating-workbench"; +import { originatingWorkbenchIdFromRequest } from "../originating-workbench"; import { parseToolRegistries } from "../tool-materialization"; import { deriveHubHttpUrl, @@ -490,10 +490,7 @@ export function createSidecarSubstrateFactory( await prepareConversationForOriginatingWorkbench({ registry: durableConversation, agentKey: bodyStepId, - originatingWorkbenchId: await readOriginatingWorkbenchId({ - dataDir: validated.SIDECAR_DATA_DIR, - mailboxAddress: env.spawn.mailboxAddress, - }), + originatingWorkbenchId: originatingWorkbenchIdFromRequest(req), }); } try { @@ -621,10 +618,7 @@ export function createSidecarSubstrateFactory( await prepareConversationForOriginatingWorkbench({ registry: durableConversation, agentKey: stepId, - originatingWorkbenchId: await readOriginatingWorkbenchId({ - dataDir: validated.SIDECAR_DATA_DIR, - mailboxAddress: env.spawn.mailboxAddress, - }), + originatingWorkbenchId: originatingWorkbenchIdFromRequest(req), ...(warmCache !== undefined ? { warmCache } : {}), }); } diff --git a/apps/sidecar/src/workflow-substrate-factory/parked-approvals.test.ts b/apps/sidecar/src/workflow-substrate-factory/parked-approvals.test.ts index d3b35cefc..914cf1ef0 100644 --- a/apps/sidecar/src/workflow-substrate-factory/parked-approvals.test.ts +++ b/apps/sidecar/src/workflow-substrate-factory/parked-approvals.test.ts @@ -4,15 +4,17 @@ // was parked on an ask-approval: the child's re-registration enumeration // found the park but had no `loadParkedApproval` binding wired and threw, // killing the deployment at boot. -import { mkdtemp, rm } from "node:fs/promises"; +import { mkdtemp, mkdir, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import path from "node:path"; import { afterEach, describe, expect, test } from "bun:test"; +import type { RepoStore } from "@intx/hub-sessions"; import type { PendingOperation } from "@intx/types/runtime"; import { findApprovalSnapshot, readColdParkedApprovalSnapshot, + readWarmParkedPendingOperations, toParkedApprovalOps, } from "./parked-approvals"; @@ -73,3 +75,65 @@ describe("readColdParkedApprovalSnapshot", () => { expect(snapshot).toBeUndefined(); }); }); + +const EMPTY_TOKEN_USAGE = { + input: 0, + output: 0, + cacheRead: 0, + cacheWrite: 0, + thinking: 0, +}; + +async function writeRoomCheckpoint( + repoDir: string, + stepId: string, + workbenchId: string, + pendingOperations: unknown[], +): Promise { + const dir = path.join(repoDir, "agent-state", stepId, workbenchId); + await mkdir(dir, { recursive: true }); + const snapshot = { + turns: [], + pendingOperations, + tokenUsage: EMPTY_TOKEN_USAGE, + connectorState: null, + }; + const meta = { + checkpointSeq: 1, + turnCount: 0, + pendingOperations, + tokenUsage: EMPTY_TOKEN_USAGE, + connectorState: null, + }; + await writeFile(path.join(dir, "checkpoint.json"), JSON.stringify(snapshot)); + await writeFile(path.join(dir, "checkpoint.meta.json"), JSON.stringify(meta)); +} + +describe("readWarmParkedPendingOperations", () => { + test("collects pending ops from nested rooms and ignores a mixed agent-state blob", async () => { + const repoDir = await mkdtemp(path.join(tmpdir(), "parked-warm-")); + tempDirs.push(repoDir); + await writeRoomCheckpoint(repoDir, "default", "chan_a", [ + approvalOp("corr-a"), + ]); + await writeRoomCheckpoint(repoDir, "default", "chan_b", [ + approvalOp("corr-b"), + ]); + await writeFile( + path.join(repoDir, "agent-state", "default", "checkpoint.json"), + JSON.stringify({ pendingOperations: [approvalOp("corr-legacy")] }), + ); + + const substrate = { + getRepoDir: () => repoDir, + } as unknown as RepoStore; + + const pending = await readWarmParkedPendingOperations({ + substrate, + workflowRunRepoId: REPO_ID, + stepId: "default", + }); + const ids = pending.map((op) => op.correlationId).sort(); + expect(ids).toEqual(["corr-a", "corr-b"]); + }); +}); diff --git a/apps/sidecar/src/workflow-substrate-factory/parked-approvals.ts b/apps/sidecar/src/workflow-substrate-factory/parked-approvals.ts index f558ae2de..0c1894d0e 100644 --- a/apps/sidecar/src/workflow-substrate-factory/parked-approvals.ts +++ b/apps/sidecar/src/workflow-substrate-factory/parked-approvals.ts @@ -116,7 +116,9 @@ export async function readWarmParkedPendingOperations(args: { ); let children: fs.Dirent[]; try { - children = await fs.promises.readdir(agentStateDir, { withFileTypes: true }); + children = await fs.promises.readdir(agentStateDir, { + withFileTypes: true, + }); } catch (cause) { if (isErrnoNotFound(cause)) return []; throw cause;