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..b905f68ed --- /dev/null +++ b/apps/sidecar/src/conversation-state.test.ts @@ -0,0 +1,259 @@ +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"; +import { originatingWorkbenchIdFromRequest } 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 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("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"), + boundOriginatingWorkbenchId: () => null, + }; + 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..b411ef8b4 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 @@ -140,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, @@ -159,6 +164,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 @@ -193,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. @@ -207,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 @@ -227,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 @@ -272,9 +311,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 +349,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 +385,18 @@ 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; + /** Room this store is bound to, or null before the first bind. */ + boundOriginatingWorkbenchId(): string | null; } export async function createDurableConversationStore( @@ -370,8 +423,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 +508,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 +520,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 +528,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 @@ -548,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, @@ -587,7 +671,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 +704,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 +811,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 +835,8 @@ export async function createDurableConversationStore( seedInbound, composeReply, onReplySent, + bindOriginatingWorkbench, + boundOriginatingWorkbenchId: () => originatingWorkbenchId, }; } @@ -754,11 +857,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 +908,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 +939,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; + /** 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 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 ${target}`); + } + return swapped; +} + interface SnapshotMetadataValue { pendingOperations: unknown[]; tokenUsage: TokenUsage; @@ -877,8 +997,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..e56009008 --- /dev/null +++ b/apps/sidecar/src/originating-workbench.test.ts @@ -0,0 +1,73 @@ +import { expect, test } from "bun:test"; + +import { + extractOriginatingWorkbenchId, + originatingWorkbenchIdFromRequest, + resolveOriginatingWorkbenchId, + UNSCOPED_ORIGINATING_WORKBENCH_ID, +} from "./originating-workbench"; + +function mailRequest(from: string) { + return { + input: { + headers: { from, to: ["myra@alice.localhost"] }, + rawHeaders: {}, + parts: [], + }, + }; +} + +test("extractOriginatingWorkbenchId reads the From local-part", () => { + expect(extractOriginatingWorkbenchId("chan_room_a@alice.localhost")).toBe( + "chan_room_a", + ); + expect( + extractOriginatingWorkbenchId("Workbench "), + ).toBe("ins_workbench1"); +}); + +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", () => { + expect(resolveOriginatingWorkbenchId(undefined)).toBe( + UNSCOPED_ORIGINATING_WORKBENCH_ID, + ); + expect(resolveOriginatingWorkbenchId("chan_a")).toBe("chan_a"); +}); + +test("originatingWorkbenchIdFromRequest is a function of this request's mail", () => { + expect( + originatingWorkbenchIdFromRequest(mailRequest("chan_a@alice.localhost")), + ).toBe("chan_a"); + expect( + originatingWorkbenchIdFromRequest(mailRequest("chan_b@alice.localhost")), + ).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("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 new file mode 100644 index 000000000..db2094258 --- /dev/null +++ b/apps/sidecar/src/originating-workbench.ts @@ -0,0 +1,58 @@ +// 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 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"; + +/** + * Local-part of a From header value. Chat encodes the originating workbench + * id there; anything else (empty, unparseable) is undefined. + */ +export function extractOriginatingWorkbenchId( + from: string, +): string | undefined { + if (from.trim() === "") return undefined; + try { + const spec = extractAddrSpec(from); + const at = spec.lastIndexOf("@"); + if (at <= 0) return undefined; + return spec.slice(0, at); + } catch { + return undefined; + } +} + +export function resolveOriginatingWorkbenchId( + fromWorkbenchId: string | undefined, +): string { + return fromWorkbenchId !== undefined && fromWorkbenchId.length > 0 + ? fromWorkbenchId + : UNSCOPED_ORIGINATING_WORKBENCH_ID; +} + +/** + * 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 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), + ); +} diff --git a/apps/sidecar/src/workflow-host-wiring/index.ts b/apps/sidecar/src/workflow-host-wiring/index.ts index 9d4319fa2..86ea0b7b6 100644 --- a/apps/sidecar/src/workflow-host-wiring/index.ts +++ b/apps/sidecar/src/workflow-host-wiring/index.ts @@ -1042,7 +1042,7 @@ 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) => { 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..089e76606 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 { originatingWorkbenchIdFromRequest } from "../originating-workbench"; import { parseToolRegistries } from "../tool-materialization"; import { deriveHubHttpUrl, @@ -483,6 +485,14 @@ export function createSidecarSubstrateFactory( credentialWiring, mailPartReader, ) => { + const bodyStepId = req.authzContext.stepId; + if (durableConversation !== undefined && bodyStepId !== undefined) { + await prepareConversationForOriginatingWorkbench({ + registry: durableConversation, + agentKey: bodyStepId, + originatingWorkbenchId: originatingWorkbenchIdFromRequest(req), + }); + } try { return await createWorkflowStepInvoker({ workflowAuthorize: authorize, @@ -494,11 +504,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 +550,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 +612,17 @@ export function createSidecarSubstrateFactory( sourcesRef, credentialWiring, mailPartReader, - ) => - createWorkflowStepInvoker({ + ) => { + const stepId = req.authzContext.stepId; + if (durableConversation !== undefined && stepId !== undefined) { + await prepareConversationForOriginatingWorkbench({ + registry: durableConversation, + agentKey: stepId, + originatingWorkbenchId: originatingWorkbenchIdFromRequest(req), + ...(warmCache !== undefined ? { warmCache } : {}), + }); + } + return createWorkflowStepInvoker({ workflowAuthorize: authorize, buildEnv: (buildReq: Parameters[0]) => buildStepEnv(buildReq, sourcesRef, credentialWiring), @@ -617,6 +635,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.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 3c18efd75..0c1894d0e 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,27 @@ 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:*",