Skip to content

Commit 6e9d834

Browse files
committed
Reject parked approval on interrupt and hold occupancy until tool start
1 parent 2ed56eb commit 6e9d834

6 files changed

Lines changed: 360 additions & 58 deletions

File tree

src/session/approval-resume.ts

Lines changed: 72 additions & 45 deletions
Original file line numberDiff line numberDiff line change
@@ -123,6 +123,9 @@ export function createApprovalResume(args: {
123123
// TUI: operator-visible notice when an overlay decision is dropped after
124124
// a generation bump.
125125
onDropped?: (text: string) => void;
126+
// TUI: interrupt/clear bump this to reject the parked call on the old
127+
// agent before close/rebuild. Cleared when handle returns.
128+
registerParkedCancel?: (cancel: (() => void) | undefined) => void;
126129
gate: PermissionGate;
127130
}): ApprovalResume {
128131
const { getAgent, gate } = args;
@@ -139,59 +142,83 @@ export function createApprovalResume(args: {
139142
handle: async (result) => {
140143
if (result.type !== "suspended") return false;
141144
const stillCurrent = args.captureGeneration?.() ?? (() => true);
145+
const parkedAgent = requireAgent();
142146
const { correlationId, approvalSnapshot } = result;
143147

144-
const deliverDecision = async (message: InboundMessage): Promise<void> => {
145-
if (!stillCurrent()) return;
146-
if (args.deliver !== undefined) {
147-
await args.deliver(message, stillCurrent);
148-
return;
149-
}
150-
requireAgent().deliver(message);
148+
let cancelled = false;
149+
const cancelParked = (): void => {
150+
if (cancelled) return;
151+
cancelled = true;
152+
parkedAgent.deliver(decisionMessage(correlationId, "rejected", APPROVAL_DROPPED_NOTICE));
151153
};
154+
args.registerParkedCancel?.(cancelParked);
152155

153-
// Turn-count watermark for the settled guard below: a "approval timed
154-
// out" tool result appended after this point means the reactor settled
155-
// this very correlation before our decision lands.
156-
const turnsAtSuspend = (await requireAgent().history()).length;
157-
if (!stillCurrent()) return true;
158-
159-
if (approvalSnapshot === undefined) {
160-
// A suspension without a snapshot cannot be surfaced; fail closed by
161-
// rejecting the parked call so the run does not hang on an invisible
162-
// gate.
163-
await deliverDecision(
164-
decisionMessage(correlationId, "rejected", "approval surface unavailable"),
165-
);
166-
return true;
167-
}
156+
const dropParked = (): void => {
157+
args.onDropped?.(APPROVAL_DROPPED_NOTICE);
158+
cancelParked();
159+
};
168160

169-
const request = requestFromApprovalSnapshot(approvalSnapshot, correlationId);
170-
if (request === null) {
171-
await deliverDecision(
172-
decisionMessage(correlationId, "rejected", "approval surface unavailable"),
173-
);
174-
return true;
175-
}
161+
try {
162+
const deliverDecision = async (message: InboundMessage): Promise<void> => {
163+
if (!stillCurrent()) return;
164+
if (args.deliver !== undefined) {
165+
await args.deliver(message, stillCurrent);
166+
return;
167+
}
168+
requireAgent().deliver(message);
169+
};
170+
171+
// Turn-count watermark for the settled guard below: a "approval timed
172+
// out" tool result appended after this point means the reactor settled
173+
// this very correlation before our decision lands.
174+
const turnsAtSuspend = (await parkedAgent.history()).length;
175+
if (!stillCurrent()) {
176+
dropParked();
177+
return true;
178+
}
176179

177-
const outcome = await gate.resolveSuspended(request, stillCurrent);
178-
if (!stillCurrent()) {
179-
args.onDropped?.(APPROVAL_DROPPED_NOTICE);
180-
return true;
181-
}
182-
if (settledAfterSuspend(await requireAgent().history(), turnsAtSuspend)) {
183-
// The reactor already answered the parked call (its approval timeout
184-
// fired while the surface was still up). Delivering now would append
185-
// the raw decision JSON as an uncorrelated user turn — drop and log.
186-
logger.warn`late approval decision dropped correlation=${correlationId} outcome=${outcome?.allow === true ? "approved" : "rejected"}`;
187-
return true;
188-
}
189-
if (outcome === undefined || !outcome.allow) {
190-
await deliverDecision(decisionMessage(correlationId, "rejected", outcome?.message));
180+
if (approvalSnapshot === undefined) {
181+
// A suspension without a snapshot cannot be surfaced; fail closed by
182+
// rejecting the parked call so the run does not hang on an invisible
183+
// gate.
184+
args.registerParkedCancel?.(undefined);
185+
await deliverDecision(
186+
decisionMessage(correlationId, "rejected", "approval surface unavailable"),
187+
);
188+
return true;
189+
}
190+
191+
const request = requestFromApprovalSnapshot(approvalSnapshot, correlationId);
192+
if (request === null) {
193+
args.registerParkedCancel?.(undefined);
194+
await deliverDecision(
195+
decisionMessage(correlationId, "rejected", "approval surface unavailable"),
196+
);
197+
return true;
198+
}
199+
200+
const outcome = await gate.resolveSuspended(request, stillCurrent);
201+
if (!stillCurrent()) {
202+
dropParked();
203+
return true;
204+
}
205+
args.registerParkedCancel?.(undefined);
206+
if (settledAfterSuspend(await requireAgent().history(), turnsAtSuspend)) {
207+
// The reactor already answered the parked call (its approval timeout
208+
// fired while the surface was still up). Delivering now would append
209+
// the raw decision JSON as an uncorrelated user turn — drop and log.
210+
logger.warn`late approval decision dropped correlation=${correlationId} outcome=${outcome?.allow === true ? "approved" : "rejected"}`;
211+
return true;
212+
}
213+
if (outcome === undefined || !outcome.allow) {
214+
await deliverDecision(decisionMessage(correlationId, "rejected", outcome?.message));
215+
return true;
216+
}
217+
await deliverDecision(decisionMessage(correlationId, "approved"));
191218
return true;
219+
} finally {
220+
args.registerParkedCancel?.(undefined);
192221
}
193-
await deliverDecision(decisionMessage(correlationId, "approved"));
194-
return true;
195222
},
196223
};
197224
}

src/tui/correlation-acceptance.test.ts

Lines changed: 43 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,21 @@
1-
import { describe, test } from "bun:test";
1+
import { describe, expect, test } from "bun:test";
22

33
import { createCorrelationAcceptance } from "./correlation-acceptance.js";
44

5+
function approvedMessage(correlationId: string) {
6+
return {
7+
headers: { interchangeCorrelationId: correlationId },
8+
content: JSON.stringify({ outcome: "approved" }),
9+
};
10+
}
11+
12+
function rejectedMessage(correlationId: string) {
13+
return {
14+
headers: { interchangeCorrelationId: correlationId },
15+
content: JSON.stringify({ outcome: "rejected" }),
16+
};
17+
}
18+
519
describe("createCorrelationAcceptance", () => {
620
test("settle resolves the waiter for that correlation id", async () => {
721
const acceptance = createCorrelationAcceptance();
@@ -22,4 +36,32 @@ describe("createCorrelationAcceptance", () => {
2236
const acceptance = createCorrelationAcceptance();
2337
acceptance.settle("missing");
2438
});
39+
40+
test("an approved correlation does not settle until tool.start", async () => {
41+
const acceptance = createCorrelationAcceptance();
42+
const pending = acceptance.wait("corr-1");
43+
let settled = false;
44+
void pending.then(() => {
45+
settled = true;
46+
});
47+
acceptance.observe({
48+
type: "message.correlated",
49+
data: { correlationId: "corr-1", message: approvedMessage("corr-1") },
50+
});
51+
await Promise.resolve();
52+
expect(settled).toBe(false);
53+
acceptance.observe({ type: "tool.start", data: { call: { id: "call-1" } } });
54+
await pending;
55+
expect(settled).toBe(true);
56+
});
57+
58+
test("a rejected correlation settles at message.correlated", async () => {
59+
const acceptance = createCorrelationAcceptance();
60+
const pending = acceptance.wait("corr-1");
61+
acceptance.observe({
62+
type: "message.correlated",
63+
data: { correlationId: "corr-1", message: rejectedMessage("corr-1") },
64+
});
65+
await pending;
66+
});
2567
});

src/tui/correlation-acceptance.ts

Lines changed: 74 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,19 +2,72 @@
22
* Occupancy wait for a fire-and-forget Agent.deliver of a correlated
33
* approval. The reactor accepts the resume asynchronously after deliver
44
* returns; inFlight must not drop until that acceptance (or an uncorrelated
5-
* pass-through / identity bump) settles the waiter.
5+
* pass-through / identity bump) settles the waiter. An approved re-dispatch
6+
* is not idle at message.correlated — occupancy holds until tool.start.
67
*/
78

9+
import { ApprovalDecision } from "@intx/types";
10+
import { type } from "arktype";
11+
12+
export interface CorrelationStreamEvent {
13+
type: string;
14+
data?: unknown;
15+
}
16+
17+
function isApprovedDecision(content: string | undefined): boolean {
18+
if (content === undefined) return false;
19+
let raw: unknown;
20+
try {
21+
raw = JSON.parse(content);
22+
} catch {
23+
return false;
24+
}
25+
const decision = ApprovalDecision(raw);
26+
if (decision instanceof type.errors) return false;
27+
return decision.outcome === "approved";
28+
}
29+
30+
function correlatedMessage(data: unknown): {
31+
correlationId: string | undefined;
32+
content: string | undefined;
33+
receivedCorrelationId: string | undefined;
34+
} {
35+
if (data === null || typeof data !== "object") {
36+
return { correlationId: undefined, content: undefined, receivedCorrelationId: undefined };
37+
}
38+
const record = data as {
39+
correlationId?: unknown;
40+
message?: {
41+
headers?: { interchangeCorrelationId?: unknown };
42+
content?: unknown;
43+
};
44+
};
45+
return {
46+
correlationId: typeof record.correlationId === "string" ? record.correlationId : undefined,
47+
content: typeof record.message?.content === "string" ? record.message.content : undefined,
48+
receivedCorrelationId:
49+
typeof record.message?.headers?.interchangeCorrelationId === "string"
50+
? record.message.headers.interchangeCorrelationId
51+
: undefined,
52+
};
53+
}
54+
855
export function createCorrelationAcceptance() {
956
const waiters = new Map<string, () => void>();
57+
const holdUntilToolStart = new Set<string>();
1058

1159
const settle = (correlationId: string): void => {
60+
holdUntilToolStart.delete(correlationId);
1261
const resolve = waiters.get(correlationId);
1362
if (resolve === undefined) return;
1463
waiters.delete(correlationId);
1564
resolve();
1665
};
1766

67+
const settleHeldForToolStart = (): void => {
68+
for (const correlationId of [...holdUntilToolStart]) settle(correlationId);
69+
};
70+
1871
return {
1972
wait(correlationId: string): Promise<void> {
2073
const pending = waiters.get(correlationId);
@@ -27,7 +80,27 @@ export function createCorrelationAcceptance() {
2780
},
2881
settle,
2982
settleAll(): void {
83+
holdUntilToolStart.clear();
3084
for (const correlationId of [...waiters.keys()]) settle(correlationId);
3185
},
86+
observe(event: CorrelationStreamEvent): void {
87+
if (event.type === "tool.start") {
88+
settleHeldForToolStart();
89+
return;
90+
}
91+
const fields = correlatedMessage(event.data);
92+
if (event.type === "message.received") {
93+
if (fields.receivedCorrelationId !== undefined) settle(fields.receivedCorrelationId);
94+
return;
95+
}
96+
if (event.type === "message.correlated") {
97+
if (fields.correlationId === undefined) return;
98+
if (isApprovedDecision(fields.content)) {
99+
holdUntilToolStart.add(fields.correlationId);
100+
return;
101+
}
102+
settle(fields.correlationId);
103+
}
104+
},
32105
};
33106
}

src/tui/runner/exit.ts

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -219,11 +219,9 @@ export async function createRunLifecycle(
219219
let eventForSink = event;
220220
if (event.type === "message.received") {
221221
providerFailureAttempts.advanceToNextMessage();
222-
const correlationId = event.data.message.headers.interchangeCorrelationId;
223-
if (correlationId !== undefined) services.correlationAcceptance.settle(correlationId);
224-
} else if (event.type === "message.correlated") {
225-
services.correlationAcceptance.settle(event.data.correlationId);
226-
} else if (event.type === "inference.start" || event.type === "inference.done") {
222+
}
223+
services.correlationAcceptance.observe(event);
224+
if (event.type === "inference.start" || event.type === "inference.done") {
227225
providerFailureAttempts.reset();
228226
} else if (event.type === "inference.error") {
229227
const error = event.data.error;

src/tui/runner/session.ts

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -126,7 +126,11 @@ export async function assembleTUISession(
126126
undefined;
127127

128128
const correlationAcceptance = createCorrelationAcceptance();
129-
const deliveryGeneration = createDeliveryGeneration(() => correlationAcceptance.settleAll());
129+
const parkedApprovalCancel = { fn: undefined as (() => void) | undefined };
130+
const deliveryGeneration = createDeliveryGeneration(() => {
131+
parkedApprovalCancel.fn?.();
132+
correlationAcceptance.settleAll();
133+
});
130134

131135
const { gate: permissionGate } = await assembleSessionGate({
132136
cwd: config.cwd,
@@ -378,9 +382,12 @@ export async function assembleTUISession(
378382
// so a rebuild never races an in-flight deliver.
379383
const sessionOps = createSessionOperationQueue();
380384
const approvalResume = createApprovalResume({
381-
getAgent: () => state.agentProxy ?? state.currentAgent,
385+
getAgent: () => state.currentAgent,
382386
captureGeneration: deliveryGeneration.capture,
383387
onDropped: (text) => state.systemNotice?.(text),
388+
registerParkedCancel: (cancel) => {
389+
parkedApprovalCancel.fn = cancel;
390+
},
384391
deliver: (message, stillCurrent) => {
385392
return sessionOps.enqueue(async () => {
386393
if (!stillCurrent()) return;

0 commit comments

Comments
 (0)