Skip to content

Commit abc5efe

Browse files
committed
fix(approval): reuse observed acceptance across retry
1 parent 0f31491 commit abc5efe

2 files changed

Lines changed: 223 additions & 84 deletions

File tree

‎src/session/approval-resume.test.ts‎

Lines changed: 108 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -595,3 +595,111 @@ describe("approval identity ordering", () => {
595595
expect(ctx.deliver).not.toHaveBeenCalled();
596596
});
597597
});
598+
599+
describe("approval resume retry re-await", () => {
600+
function retryHarness(outcome: { allow: boolean; message?: string }) {
601+
const delivered: InboundMessage[] = [];
602+
const deliver = mock((message: InboundMessage): void => {
603+
delivered.push(message);
604+
});
605+
const resolveSuspended = mock(async () => outcome);
606+
const resume = createApprovalResume({
607+
getAgent: () => ({
608+
deliver,
609+
history: async () => [],
610+
}),
611+
resolveParkedCallId: () => "call-A",
612+
gate: { resolveSuspended } as unknown as PermissionGate,
613+
});
614+
return { resume, delivered, deliver, resolveSuspended };
615+
}
616+
617+
function onlyDecision(delivered: InboundMessage[]): InboundMessage {
618+
const message = delivered[0];
619+
if (message === undefined) throw new Error("expected a delivered decision");
620+
return message;
621+
}
622+
623+
test("retry after a delivered acceptance reuses it exactly once", async () => {
624+
const { resume, delivered, resolveSuspended } = retryHarness({
625+
allow: true,
626+
});
627+
expect(await resume.handle(suspension("corr-A", "echo alpha"))).toBe(true);
628+
expect(resolveSuspended).toHaveBeenCalledTimes(1);
629+
expect(delivered).toHaveLength(1);
630+
expect(decisionBody(onlyDecision(delivered)).outcome).toBe("approved");
631+
632+
expect(await resume.handle(suspension("corr-A", "echo alpha"))).toBe(true);
633+
expect(resolveSuspended).toHaveBeenCalledTimes(1);
634+
expect(delivered).toHaveLength(1);
635+
});
636+
637+
test("late duplicate acceptance after a rejected decision is a no-op", async () => {
638+
const { resume, delivered, resolveSuspended } = retryHarness({
639+
allow: false,
640+
message: "not today",
641+
});
642+
expect(await resume.handle(suspension("corr-A", "echo alpha"))).toBe(true);
643+
expect(delivered).toHaveLength(1);
644+
expect(decisionBody(onlyDecision(delivered)).outcome).toBe("rejected");
645+
646+
expect(await resume.handle(suspension("corr-A", "echo alpha"))).toBe(true);
647+
expect(await resume.handle(suspension("corr-A", "echo alpha"))).toBe(true);
648+
expect(resolveSuspended).toHaveBeenCalledTimes(1);
649+
expect(delivered).toHaveLength(1);
650+
});
651+
652+
test("retry without acceptance still gates", async () => {
653+
const delivered: InboundMessage[] = [];
654+
const deliver = mock((message: InboundMessage): void => {
655+
delivered.push(message);
656+
});
657+
deliver.mockImplementationOnce(() => {
658+
throw new Error("agent is done");
659+
});
660+
const resolveSuspended = mock(async () => ({ allow: true }));
661+
const resume = createApprovalResume({
662+
getAgent: () => ({
663+
deliver,
664+
history: async () => [],
665+
}),
666+
resolveParkedCallId: () => "call-A",
667+
gate: { resolveSuspended } as unknown as PermissionGate,
668+
});
669+
await expect(
670+
resume.handle(suspension("corr-A", "echo alpha")),
671+
).rejects.toThrow("agent is done");
672+
expect(delivered).toHaveLength(0);
673+
expect(resolveSuspended).toHaveBeenCalledTimes(1);
674+
675+
expect(await resume.handle(suspension("corr-A", "echo alpha"))).toBe(true);
676+
expect(resolveSuspended).toHaveBeenCalledTimes(2);
677+
expect(delivered).toHaveLength(1);
678+
});
679+
680+
test("concurrent duplicate handles share one gate and one deliver", async () => {
681+
const gate = Promise.withResolvers<{ allow: boolean }>();
682+
const delivered: InboundMessage[] = [];
683+
const deliver = mock((message: InboundMessage): void => {
684+
delivered.push(message);
685+
});
686+
const resolveSuspended = mock(() => gate.promise);
687+
const resume = createApprovalResume({
688+
getAgent: () => ({
689+
deliver,
690+
history: async () => [],
691+
}),
692+
resolveParkedCallId: () => "call-A",
693+
gate: { resolveSuspended } as unknown as PermissionGate,
694+
});
695+
const first = resume.handle(suspension("corr-A", "echo alpha"));
696+
const second = resume.handle(suspension("corr-A", "echo alpha"));
697+
for (let i = 0; i < 10; i++) await Promise.resolve();
698+
expect(resolveSuspended).toHaveBeenCalledTimes(1);
699+
gate.resolve({ allow: true });
700+
expect(await first).toBe(true);
701+
expect(await second).toBe(true);
702+
expect(resolveSuspended).toHaveBeenCalledTimes(1);
703+
expect(delivered).toHaveLength(1);
704+
});
705+
});

‎src/session/approval-resume.ts‎

Lines changed: 115 additions & 84 deletions
Original file line numberDiff line numberDiff line change
@@ -142,100 +142,131 @@ export function createApprovalResume(args: {
142142
) => string | undefined | Promise<string | undefined>;
143143
gate: PermissionGate;
144144
}): ApprovalResume {
145-
return {
146-
handle: async (result) => {
147-
if (result.type !== "suspended") return false;
148-
const generationCurrent = args.captureGeneration?.() ?? (() => true);
149-
const parkedAgent = args.getAgent();
150-
if (parkedAgent === undefined)
151-
throw new Error("approval resume: no live agent");
152-
const { correlationId, approvalSnapshot } = result;
153-
let cancelled = false;
154-
let canReject = false;
155-
const stillCurrent = (): boolean => !cancelled && generationCurrent();
156-
const cancelParked = (): void => {
157-
if (cancelled) return;
158-
cancelled = true;
159-
if (canReject) {
160-
parkedAgent.deliver(
161-
decisionMessage(correlationId, "rejected", APPROVAL_DROPPED_NOTICE),
162-
);
163-
}
164-
};
165-
const dropParked = (): void => {
166-
args.onDropped?.(APPROVAL_DROPPED_NOTICE);
167-
cancelParked();
168-
};
169-
args.registerParkedCancel?.(cancelParked);
170-
try {
171-
// The resolver captures the paired store synchronously before its first await.
172-
const parkedCallId = await args.resolveParkedCallId(correlationId);
173-
if (!stillCurrent()) {
174-
dropParked();
175-
return true;
176-
}
177-
if (parkedCallId === undefined) return true;
178-
const initialHistory = await parkedAgent.history();
179-
if (!stillCurrent()) {
180-
dropParked();
181-
return true;
182-
}
183-
if (timeoutResult(initialHistory, parkedCallId)) return true;
184-
canReject = true;
145+
// Retry re-await wiring: correlation ids whose decision was handed to the
146+
// reactor reuse that acceptance. A retry after an observed acceptance
147+
// returns without opening the gate or delivering again, so the parked call
148+
// resumes exactly once and a late duplicate acceptance is a no-op. Ids are
149+
// recorded only when the decision is actually handed over — a failed send
150+
// (deliver threw, so nothing reached the reactor) retries as before.
151+
const handedOver = new Set<string>();
152+
// Settlements currently gating-and-delivering, keyed by correlation id. A
153+
// concurrent duplicate handle shares the one in-flight outcome instead of
154+
// opening a second gate, so no waiter is lost and none double-resumes.
155+
const inflight = new Map<string, Promise<boolean>>();
185156

186-
const deliverDecision = async (
187-
message: InboundMessage,
188-
): Promise<void> => {
189-
if (!stillCurrent()) return;
190-
if (args.deliver !== undefined) {
191-
await args.deliver(message, stillCurrent);
192-
} else {
193-
parkedAgent.deliver(message);
194-
}
195-
};
196-
const request =
197-
approvalSnapshot === undefined
198-
? null
199-
: requestFromApprovalSnapshot(approvalSnapshot, correlationId);
200-
if (request === null) {
201-
args.registerParkedCancel?.(undefined);
202-
await deliverDecision(
203-
decisionMessage(
204-
correlationId,
205-
"rejected",
206-
"approval surface unavailable",
207-
),
208-
);
209-
return true;
210-
}
211-
const outcome = await args.gate.resolveSuspended(request, stillCurrent);
212-
if (!stillCurrent()) {
213-
dropParked();
214-
return true;
157+
const settleSuspended = async (
158+
result: Extract<SendResult, { type: "suspended" }>,
159+
): Promise<boolean> => {
160+
const generationCurrent = args.captureGeneration?.() ?? (() => true);
161+
const parkedAgent = args.getAgent();
162+
if (parkedAgent === undefined)
163+
throw new Error("approval resume: no live agent");
164+
const { correlationId, approvalSnapshot } = result;
165+
let cancelled = false;
166+
let canReject = false;
167+
const stillCurrent = (): boolean => !cancelled && generationCurrent();
168+
const cancelParked = (): void => {
169+
if (cancelled) return;
170+
cancelled = true;
171+
if (canReject) {
172+
parkedAgent.deliver(
173+
decisionMessage(correlationId, "rejected", APPROVAL_DROPPED_NOTICE),
174+
);
175+
}
176+
};
177+
const dropParked = (): void => {
178+
args.onDropped?.(APPROVAL_DROPPED_NOTICE);
179+
cancelParked();
180+
};
181+
args.registerParkedCancel?.(cancelParked);
182+
try {
183+
// The resolver captures the paired store synchronously before its first await.
184+
const parkedCallId = await args.resolveParkedCallId(correlationId);
185+
if (!stillCurrent()) {
186+
dropParked();
187+
return true;
188+
}
189+
if (parkedCallId === undefined) return true;
190+
const initialHistory = await parkedAgent.history();
191+
if (!stillCurrent()) {
192+
dropParked();
193+
return true;
194+
}
195+
if (timeoutResult(initialHistory, parkedCallId)) return true;
196+
canReject = true;
197+
198+
const deliverDecision = async (
199+
message: InboundMessage,
200+
): Promise<boolean> => {
201+
if (!stillCurrent()) return false;
202+
if (args.deliver !== undefined) {
203+
await args.deliver(message, stillCurrent);
204+
} else {
205+
parkedAgent.deliver(message);
215206
}
207+
handedOver.add(correlationId);
208+
return true;
209+
};
210+
const request =
211+
approvalSnapshot === undefined
212+
? null
213+
: requestFromApprovalSnapshot(approvalSnapshot, correlationId);
214+
if (request === null) {
216215
args.registerParkedCancel?.(undefined);
217-
const history = await parkedAgent.history();
218-
const timedOut = timeoutResult(history, parkedCallId);
219-
if (timedOut) canReject = false;
220-
if (!stillCurrent()) {
221-
dropParked();
222-
return true;
223-
}
224-
if (timedOut) {
225-
logger.warn`late approval decision dropped correlation=${correlationId} timeoutCall=${parkedCallId} outcome=${outcome?.allow === true ? "approved" : "rejected"}`;
226-
return true;
227-
}
228216
await deliverDecision(
229217
decisionMessage(
230218
correlationId,
231-
outcome?.allow === true ? "approved" : "rejected",
232-
outcome?.allow === true ? undefined : outcome?.message,
219+
"rejected",
220+
"approval surface unavailable",
233221
),
234222
);
235223
return true;
236-
} finally {
237-
args.registerParkedCancel?.(undefined);
238224
}
225+
const outcome = await args.gate.resolveSuspended(request, stillCurrent);
226+
if (!stillCurrent()) {
227+
dropParked();
228+
return true;
229+
}
230+
args.registerParkedCancel?.(undefined);
231+
const history = await parkedAgent.history();
232+
const timedOut = timeoutResult(history, parkedCallId);
233+
if (timedOut) canReject = false;
234+
if (!stillCurrent()) {
235+
dropParked();
236+
return true;
237+
}
238+
if (timedOut) {
239+
logger.warn`late approval decision dropped correlation=${correlationId} timeoutCall=${parkedCallId} outcome=${outcome?.allow === true ? "approved" : "rejected"}`;
240+
return true;
241+
}
242+
await deliverDecision(
243+
decisionMessage(
244+
correlationId,
245+
outcome?.allow === true ? "approved" : "rejected",
246+
outcome?.allow === true ? undefined : outcome?.message,
247+
),
248+
);
249+
return true;
250+
} finally {
251+
args.registerParkedCancel?.(undefined);
252+
}
253+
};
254+
255+
return {
256+
handle: (result) => {
257+
if (result.type !== "suspended") return Promise.resolve(false);
258+
const { correlationId } = result;
259+
if (handedOver.has(correlationId)) return Promise.resolve(true);
260+
const ongoing = inflight.get(correlationId);
261+
if (ongoing !== undefined) return ongoing;
262+
const task = settleSuspended(result);
263+
inflight.set(correlationId, task);
264+
const forget = (): void => {
265+
if (inflight.get(correlationId) === task)
266+
inflight.delete(correlationId);
267+
};
268+
task.then(forget, forget);
269+
return task;
239270
},
240271
};
241272
}

0 commit comments

Comments
 (0)