Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions bun.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions packages/chat/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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:*",
Expand Down
77 changes: 67 additions & 10 deletions packages/chat/src/chat-orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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";
Expand Down Expand Up @@ -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;
Expand All @@ -429,7 +448,7 @@ type PendingDelegationThread = {

async function threadDelegatedReply(
deps: ChatOrchestratorDeps,
pendingDelegationThreads: Map<string, PendingDelegationThread>,
pendingDelegationThreads: ExpiringMap<string, PendingDelegationThread>,
agentAddress: string,
workbenchId: string,
messageId: string,
Expand Down Expand Up @@ -529,7 +548,7 @@ function correlateOriginatingWorkbench(

async function postReply(
deps: ChatOrchestratorDeps,
pendingDelegationThreads: Map<string, PendingDelegationThread>,
pendingDelegationThreads: ExpiringMap<string, PendingDelegationThread>,
agentAddress: string,
parts: readonly Part[],
outcome: TurnOutcome,
Expand Down Expand Up @@ -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
Expand All @@ -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<string>,
postedApprovalIds: ExpiringMap<string, true>,
): Promise<void> {
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;
Expand Down Expand Up @@ -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<string>();
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<string, true>({
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
Expand Down Expand Up @@ -1036,7 +1090,10 @@ export function createChatOrchestrator(
const notifiedDropTurns = new Set<string>();

// See `PendingDelegationThread`'s own doc comment above `postReply`.
const pendingDelegationThreads = new Map<string, PendingDelegationThread>();
const pendingDelegationThreads = createExpiringMap<
string,
PendingDelegationThread
>({ ttlMs: AGENT_TURN_STALE_MS, now });

// See `createReplyPartsAccumulator`'s own doc comment above.
const replyParts = createReplyPartsAccumulator();
Expand Down
210 changes: 209 additions & 1 deletion packages/chat/test/chat-orchestrator.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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", () => {
Expand Down
Loading