From 960d9a868381c6904c6bd035c848fef2c55ace76 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 21:15:27 -0700 Subject: [PATCH 1/2] Add tests for bounding the orchestrator's process-lifetime collections MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Proves eviction actually happens for postedApprovalIds and pendingDelegationThreads, and that eviction is safe: a redelivered gate-blocked event for a resolved approval still posts nothing once its guard entry expires, and an abandoned delegation's reply still posts — just unthreaded — once its entry expires. --- packages/chat/test/chat-orchestrator.test.ts | 210 ++++++++++++++++++- 1 file changed, 209 insertions(+), 1 deletion(-) diff --git a/packages/chat/test/chat-orchestrator.test.ts b/packages/chat/test/chat-orchestrator.test.ts index b816f068d..9d531f35a 100644 --- a/packages/chat/test/chat-orchestrator.test.ts +++ b/packages/chat/test/chat-orchestrator.test.ts @@ -22,8 +22,12 @@ import { createSidecarEmitter } from "@intx/hub-sessions"; import { createArtifactDeliveryHandler, createChatOrchestrator, + POSTED_APPROVAL_GUARD_TTL_MS, } from "../src/chat-orchestrator"; -import { createInMemoryAgentTurnStore } from "../src/agent-turns"; +import { + AGENT_TURN_STALE_MS, + createInMemoryAgentTurnStore, +} from "../src/agent-turns"; import { parseBlock } from "../src/blocks"; import type { ChatPlatform, ChatWorkbenchEvent } from "../src/platform-port"; import { @@ -973,6 +977,135 @@ describe("createChatOrchestrator", () => { orchestrator.dispose(); }); + // CL-7229: `pendingDelegationThreads` is bounded by `AGENT_TURN_STALE_MS` + // rather than growing one entry per delegation for the life of the hub + // process. A specialist that never replies (never woke, crashed, or was + // mentioned by mistake) has its entry aged out once the rest of the + // system would already consider that occurrence stale — an eviction + // this test proves happens, and proves is safe: the specialist's reply + // still posts (nothing is lost), it simply lands unthreaded instead of + // nested under the delegating message. + test("a delegation nobody ever replies to is evicted after AGENT_TURN_STALE_MS, not held forever", async () => { + const room = fakeRoom(); + const openedThreadFor: string[] = []; + const events = createSidecarEmitter(); + let now = 0; + + // Two distinct addresses resolve in event order (host first, + // specialist second) — see the "delegated specialist" test above + // for why a single-run `createFakeDb` can't tell them apart. + const runsByCallOrder = [ + { id: "ins_myra1", tenantId: "ten_1" }, + { id: "ins_echo1", tenantId: "ten_1" }, + ]; + let dbCallIndex = 0; + let launchCallIndex = 0; + const db = { + query: { + workflowRun: { + findFirst: async () => { + const run = runsByCallOrder[dbCallIndex]; + dbCallIndex += 1; + return run === undefined + ? undefined + : { ...run, principalId: null }; + }, + }, + }, + select: () => ({ + from: () => ({ + where: () => ({ + limit: async () => { + const run = runsByCallOrder[launchCallIndex]; + launchCallIndex += 1; + return run === undefined + ? [] + : [launchRowFor(run.id, run.tenantId)]; + }, + }), + }), + }), + }; + + const orchestrator = createChatOrchestrator( + { + db: db as never, + store: { + listWorkbenchSettings: async () => [ + workbenchRow("ins_workbench1", [ + "ins_myra1@ten1.workbench.test", + "ins_echo1@ten1.workbench.test", + ]), + ], + }, + roomMessages: room.roomMessages, + publish: room.publish, + platform: fakeMail().platform, + threads: { + openReplyThread: async (input) => { + openedThreadFor.push(input.parentMessageId); + return { + id: "thr_1", + tenantId: input.tenantId, + workbenchId: input.workbenchId, + kind: "reply", + parentMessageId: input.parentMessageId, + parentThreadId: "thr_root", + runRef: null, + title: null, + createdAt: new Date(), + }; + }, + assignMessage: async () => {}, + }, + events, + claims: fakeClaims(), + approvals: { findByCorrelationId: async () => null }, + }, + { now: () => now }, + ); + + // The host delegates to the specialist by @mention. + events.emit("agent.event", { + agentAddress: "ins_myra1@ten1.workbench.test", + sessionId: "ses_1", + event: { + type: "connector.reply", + data: { content: "@ins_echo1 can you take this one" }, + }, + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(room.posted).toHaveLength(1); + + // The specialist never wakes up. Time passes well beyond + // `AGENT_TURN_STALE_MS` — the rest of the system has already given + // up on that occurrence ever finishing. `createExpiringMap` evicts + // lazily, so the entry is gone the moment anything next reads it. + now += AGENT_TURN_STALE_MS + 1; + + // The specialist finally does reply, long after being abandoned. + events.emit("agent.event", { + agentAddress: "ins_echo1@ten1.workbench.test", + sessionId: "ses_2", + event: { + type: "connector.reply", + data: { content: "Sorry for the wait — filed the ticket." }, + }, + }); + await new Promise((resolve) => setTimeout(resolve, 0)); + + // The reply is never lost — it still posts... + const specialistReply = room.posted.find( + (message) => message.sender.address === "ins_echo1@ten1.workbench.test", + ); + expect(specialistReply).toMatchObject({ workbenchId: "ins_workbench1" }); + // ...but the evicted entry means it's no longer threaded under the + // (long-abandoned) delegating message. + expect(openedThreadFor).toHaveLength(0); + + orchestrator.dispose(); + }); + // CL-6137 / turn-drop notice (stress round 3): a turn that never // emits `connector.reply` content is no longer invisible — the // bracket close posts an honest in-workbench notice — but a reply this @@ -1408,6 +1541,81 @@ describe("createChatOrchestrator", () => { orchestrator.dispose(); }); + + // CL-7229: `postedApprovalIds` is bounded by `POSTED_APPROVAL_GUARD_TTL_MS` + // rather than growing one entry per approval ever carded for the life of + // the hub process. Proves both halves: the entry is actually evicted + // after the TTL, and — the safety property that matters — a redelivery + // that arrives after eviction still posts nothing for an approval that + // has since resolved, because the live status re-read in + // `postApproveBlock` is the real guard; the in-memory guard is only ever + // an optimization on top of it. + test("an evicted approval guard entry still can't cause a duplicate card once the approval has resolved", async () => { + const room = fakeRoom(); + const events = createSidecarEmitter(); + let status: "pending" | "approved" = "pending"; + let now = 0; + + const orchestrator = createChatOrchestrator( + { + db: createFakeDb({ id: "ins_echo1", tenantId: "ten_1" }) as never, + store: { + listWorkbenchSettings: async () => [ + workbenchRow("ins_workbench1", ["ins_echo1@ten1.workbench.test"]), + ], + }, + roomMessages: room.roomMessages, + publish: room.publish, + platform: fakeMail().platform, + events, + claims: fakeClaims(), + approvals: { + findByCorrelationId: async () => approvalRow({ status }), + }, + }, + { now: () => now }, + ); + + const emitGateBlocked = () => + events.emit("agent.event", { + agentAddress: "ins_echo1@ten1.workbench.test", + sessionId: "ses_1", + event: { + type: "reactor.gate.blocked", + data: { + reason: "approval", + gateId: "gate_1", + correlationId: "cor_1", + }, + }, + }); + + emitGateBlocked(); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(room.posted).toHaveLength(1); + + // A redelivery well within the TTL is still deduped by the guard. + now += 1_000; + emitGateBlocked(); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(room.posted).toHaveLength(1); + + // The approval resolves — a human acted on the card already posted — + // and enough time passes for the guard entry to be evicted. + status = "approved"; + now += POSTED_APPROVAL_GUARD_TTL_MS + 1; + + // A very late, redelivered gate-blocked event arrives (a stale + // wire-layer replay) for the now-evicted, now-resolved approval. + emitGateBlocked(); + await new Promise((resolve) => setTimeout(resolve, 0)); + + // No second card: the live status re-read still catches it even + // though the in-memory guard no longer remembers this approval. + expect(room.posted).toHaveLength(1); + + orchestrator.dispose(); + }); }); describe("createArtifactDeliveryHandler", () => { From 7b98c36fc85266c85337d32cab63c4d7dc837d9b Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 21:15:36 -0700 Subject: [PATCH 2/2] Bound the orchestrator's postedApprovalIds and pendingDelegationThreads MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Both grew one entry per approval/delegation for the life of the hub process, never pruned. Bound each with @corbits/collections' createExpiringMap rather than a plain Set/Map: - pendingDelegationThreads: already deleted on the event that matters (the specialist's first reply consumes it); now also TTL-bounded at AGENT_TURN_STALE_MS, the same threshold agent-turns.ts already uses to call a running turn no longer believable. A delegation nobody answers within that window is already considered abandoned elsewhere in the system, so evicting it can't drop a live reply — a very late reply still posts, just unthreaded instead of nested under the delegating message. - postedApprovalIds: every gate the reactor blocks on resolves out of "pending" on its own within one hour (DEFAULT_GATE_TIMEOUT_MS), so a guard entry older than that is guaranteed to belong to a non-pending approval. TTL is set well past that bound; eviction can never reopen the duplicate-card hole this guard exists to close, because postApproveBlock re-reads the approval's live status independently of the guard on every call. --- bun.lock | 1 + packages/chat/package.json | 1 + packages/chat/src/chat-orchestrator.ts | 77 ++++++++++++++++++++++---- 3 files changed, 69 insertions(+), 10 deletions(-) diff --git a/bun.lock b/bun.lock index ebb446c1a..265ac646c 100644 --- a/bun.lock +++ b/bun.lock @@ -440,6 +440,7 @@ "@corbits/agent-lifecycle": "workspace:*", "@corbits/agent-runtime": "workspace:*", "@corbits/approvals": "workspace:*", + "@corbits/collections": "workspace:*", "@corbits/commands": "workspace:*", "@corbits/error-sink": "workspace:*", "@corbits/folded-runs": "workspace:*", diff --git a/packages/chat/package.json b/packages/chat/package.json index 20ec8cb22..170d72e07 100644 --- a/packages/chat/package.json +++ b/packages/chat/package.json @@ -29,6 +29,7 @@ "@corbits/agent-lifecycle": "workspace:*", "@corbits/agent-runtime": "workspace:*", "@corbits/approvals": "workspace:*", + "@corbits/collections": "workspace:*", "@corbits/commands": "workspace:*", "@corbits/error-sink": "workspace:*", "@corbits/folded-runs": "workspace:*", diff --git a/packages/chat/src/chat-orchestrator.ts b/packages/chat/src/chat-orchestrator.ts index 6deb22877..3c866956a 100644 --- a/packages/chat/src/chat-orchestrator.ts +++ b/packages/chat/src/chat-orchestrator.ts @@ -29,6 +29,7 @@ import { toolDoneResult, type ReplyContentBlock, } from "@corbits/agent-events"; +import { createExpiringMap, type ExpiringMap } from "@corbits/collections"; import { reportError } from "@corbits/error-sink"; import type { Memory } from "@corbits/memory"; import { @@ -58,7 +59,7 @@ import { parseParticipants, type ParticipantRecord } from "./participants"; import type { Part, TextPart } from "./parts"; import type { ChatPlatform } from "./platform-port"; import { postRoomMessage, type RoomMessageStore } from "./room-messages"; -import type { AgentTurnStore } from "./agent-turns"; +import { AGENT_TURN_STALE_MS, type AgentTurnStore } from "./agent-turns"; import type { ChatStore } from "./store"; import type { ThreadStore } from "./threads"; import type { WorkbenchSubscriberRegistry } from "./workbench-events"; @@ -420,6 +421,24 @@ async function resolveMemberWorkbenches( * "thread machinery that already exists" scope for CL-5879 rather than * tracking an open-ended delegation session. A specialist mentioned * again gets a fresh entry from that later delegating message. + * + * Bounded with a TTL (CL-7229) rather than a plain `Map`: a specialist + * that never replies (crashed, never woke, or was mentioned by mistake) + * left its entry here forever, one per abandoned delegation for the + * life of the hub process. The TTL is `AGENT_TURN_STALE_MS` — the same + * threshold `./agent-turns.ts` already uses to decide a `running` turn + * row is no longer believable and fails it outright. That threshold is + * exactly the point past which this entry stops mattering too: once a + * specialist's occurrence has run longer than any turn is allowed to, + * the rest of the system has already given up on it ever finishing, so + * evicting the pending-thread entry at the same age can't drop a reply + * a running turn still needs. If a reply somehow still lands after + * that (a very late, already-abandoned occurrence), it just posts + * unthreaded into the main feed instead of nested under the delegating + * message — a cosmetic degrade, never a lost message. The common case + * (event-driven delete) is unchanged: `threadDelegatedReply` below + * still deletes the entry the moment the specialist's first reply + * consumes it, long before the TTL would ever fire. */ type PendingDelegationThread = { readonly tenantId: string; @@ -429,7 +448,7 @@ type PendingDelegationThread = { async function threadDelegatedReply( deps: ChatOrchestratorDeps, - pendingDelegationThreads: Map, + pendingDelegationThreads: ExpiringMap, agentAddress: string, workbenchId: string, messageId: string, @@ -529,7 +548,7 @@ function correlateOriginatingWorkbench( async function postReply( deps: ChatOrchestratorDeps, - pendingDelegationThreads: Map, + pendingDelegationThreads: ExpiringMap, agentAddress: string, parts: readonly Part[], outcome: TurnOutcome, @@ -642,6 +661,25 @@ async function postReply( } } +/** + * How long `postApproveBlock`'s `postedApprovalIds` guard remembers a + * carded approval (CL-7229). Every gate the reactor blocks on — approval + * gates included — resolves out of `"pending"` on its own within + * `DEFAULT_GATE_TIMEOUT_MS` (one hour: `vendor/intx/inference/src/reactor.ts`, + * not exported — `./chat-orchestrator.ts`'s own copy of the same number + * `./authz-extension.ts` already keeps for the same reason), either by a + * human's decision or by timing out to a terminal `"timeout"`/`"expired"` + * status; the gate never stays blocked, and stays open, past that bound. + * A guard entry older than that bound is therefore guaranteed to belong + * to an approval that is no longer pending, so evicting it can never + * cause a duplicate card: the very next line of `postApproveBlock` below + * re-reads the approval's live status and bails out on anything but + * `"pending"` regardless of whether this guard remembers it. Set well + * past the gate timeout (double it, plus a margin) so an evicted entry's + * safety never depends on exact timing. + */ +export const POSTED_APPROVAL_GUARD_TTL_MS = 2 * 60 * 60 * 1000 + 5 * 60 * 1000; + /** * Posts the platform-minted approve block for a gate-blocked run into every * workbench the parked agent is a member of. `postedApprovalIds` is the @@ -652,20 +690,25 @@ async function postReply( * (a race with `POST .../resolve`, or a *very* stale replay) is a no-op * too, since a card for an already-resolved approval would render terminal * state a human never got to act on — nothing to add to the workbench. + * + * Bounded with `POSTED_APPROVAL_GUARD_TTL_MS` (CL-7229) instead of a plain + * `Set` that grew one entry per approval ever carded for the life of the + * hub process — see that constant's own doc comment for why an eviction + * can never reopen the duplicate-card hole this guard exists to close. */ async function postApproveBlock( deps: ChatOrchestratorDeps, agentAddress: string, correlationId: string, - postedApprovalIds: Set, + postedApprovalIds: ExpiringMap, ): Promise { const approval = await deps.approvals.findByCorrelationId(correlationId); if (approval === null || approval.status !== "pending") return; - if (postedApprovalIds.has(approval.id)) return; + if (postedApprovalIds.get(approval.id) !== undefined) return; // Marked before the awaits below: two redelivered events racing this // function must not both pass the guard while the first resolves // workbenches. - postedApprovalIds.add(approval.id); + postedApprovalIds.set(approval.id, true); const resolved = await resolveMemberWorkbenches(deps, agentAddress); if (resolved === undefined) return; @@ -1004,10 +1047,21 @@ export function createArtifactDeliveryHandler(deps: ChatOrchestratorDeps): ( export function createChatOrchestrator( deps: ChatOrchestratorDeps, + options?: { + /** Injectable clock, for tests that age the two TTL-bounded + * collections below without waiting on a real timer. */ + readonly now?: () => number; + }, ): ChatOrchestrator { - // Process-lifetime idempotency guard for `postApproveBlock` — see its own - // doc comment for what this does and doesn't cover. - const postedApprovalIds = new Set(); + const now = options?.now ?? Date.now; + + // Bounded idempotency guard for `postApproveBlock` (CL-7229) — see + // `POSTED_APPROVAL_GUARD_TTL_MS`'s own doc comment for what this does + // and doesn't cover, and why an eviction is safe. + const postedApprovalIds = createExpiringMap({ + ttlMs: POSTED_APPROVAL_GUARD_TTL_MS, + now, + }); // Every turn with a `connector.reply` pending delivery — added the // moment reply content is seen, cleared the moment that same turn's own @@ -1036,7 +1090,10 @@ export function createChatOrchestrator( const notifiedDropTurns = new Set(); // See `PendingDelegationThread`'s own doc comment above `postReply`. - const pendingDelegationThreads = new Map(); + const pendingDelegationThreads = createExpiringMap< + string, + PendingDelegationThread + >({ ttlMs: AGENT_TURN_STALE_MS, now }); // See `createReplyPartsAccumulator`'s own doc comment above. const replyParts = createReplyPartsAccumulator();