Skip to content

Commit 158d646

Browse files
cxbAsDevsteipete
andauthored
fix(signal): retry inbound flushes on reply session init conflict (#103218)
* fix(signal): retry inbound flushes on reply session init conflict * fix(signal): bind conflict retries to monitor lifecycle --------- Co-authored-by: Peter Steinberger <[email protected]>
1 parent 20bc552 commit 158d646

5 files changed

Lines changed: 554 additions & 33 deletions

extensions/signal/src/monitor.tool-result.pairs-uuid-only-senders-uuid-allowlist-entry.test.ts

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -126,6 +126,49 @@ describe("monitorSignalProvider tool results", () => {
126126
}
127127
});
128128

129+
it("cancels a pending reply-session conflict retry when the monitor stops", async () => {
130+
vi.useFakeTimers();
131+
const abortController = new AbortController();
132+
replyMock.mockRejectedValue(
133+
new Error(
134+
"reply session initialization conflicted for agent:main:signal:direct:+15550001111",
135+
),
136+
);
137+
streamMock.mockImplementation(async ({ onEvent, abortSignal }) => {
138+
onEvent({
139+
event: "receive",
140+
data: JSON.stringify({
141+
envelope: {
142+
sourceNumber: "+15550001111",
143+
sourceName: "Ada",
144+
timestamp: 1,
145+
dataMessage: { message: "hello after the prior turn" },
146+
},
147+
}),
148+
});
149+
await new Promise<void>((resolve) => {
150+
abortSignal?.addEventListener("abort", () => resolve(), { once: true });
151+
});
152+
});
153+
154+
try {
155+
const monitorPromise = monitorSignalProvider({
156+
autoStart: false,
157+
baseUrl: "http://127.0.0.1:8080",
158+
abortSignal: abortController.signal,
159+
});
160+
161+
await vi.waitFor(() => expect(replyMock).toHaveBeenCalledTimes(1));
162+
abortController.abort(new Error("monitor stopped"));
163+
await monitorPromise;
164+
await vi.advanceTimersByTimeAsync(10_000);
165+
166+
expect(replyMock).toHaveBeenCalledTimes(1);
167+
} finally {
168+
vi.useRealTimers();
169+
}
170+
});
171+
129172
it("sizes attachment RPC response caps from mediaMaxMb", async () => {
130173
const abortController = new AbortController();
131174
const maxBytes = 2 * 1024 * 1024;

extensions/signal/src/monitor.ts

Lines changed: 13 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -92,10 +92,10 @@ function resolveRuntime(opts: MonitorSignalOpts): RuntimeEnv {
9292
function createSignalMonitorTaskRunner(runtime: RuntimeEnv) {
9393
const inFlight = new Set<Promise<void>>();
9494
return {
95-
runEventTask(task: () => Promise<void>): void {
95+
runTask(task: () => Promise<void>): void {
9696
const trackedTask = Promise.resolve()
9797
.then(task)
98-
.catch((err: unknown) => runtime.error?.(`event handler failed: ${String(err)}`))
98+
.catch((err: unknown) => runtime.error?.(`signal monitor task failed: ${String(err)}`))
9999
.finally(() => inFlight.delete(trackedTask));
100100
inFlight.add(trackedTask);
101101
},
@@ -155,6 +155,11 @@ function createSignalDaemonLifecycle(params: { abortSignal?: AbortSignal }) {
155155
const mergedAbort = mergeAbortSignals(params.abortSignal, daemonAbortController.signal);
156156
const stop = () => {
157157
daemonStopRequested = true;
158+
if (!daemonAbortController.signal.aborted) {
159+
daemonAbortController.abort(
160+
params.abortSignal?.reason ?? new Error("Signal monitor stopped"),
161+
);
162+
}
158163
daemonHandle?.stop();
159164
};
160165
const attach = (handle: SignalDaemonHandle) => {
@@ -678,6 +683,8 @@ export async function monitorSignalProvider(opts: MonitorSignalOpts = {}): Promi
678683

679684
const handleEvent = createSignalEventHandler({
680685
runtime,
686+
abortSignal: daemonLifecycle.abortSignal,
687+
runTrackedTask: (task) => monitorTaskRunner.runTask(task),
681688
cfg,
682689
baseUrl,
683690
account,
@@ -715,7 +722,7 @@ export async function monitorSignalProvider(opts: MonitorSignalOpts = {}): Promi
715722
apiMode: configuredApiMode,
716723
policy: opts.reconnectPolicy,
717724
onEvent: (event) => {
718-
monitorTaskRunner.runEventTask(() => handleEvent(event));
725+
monitorTaskRunner.runTask(() => handleEvent(event));
719726
},
720727
});
721728
const daemonExitError = daemonLifecycle.getExitError();
@@ -729,9 +736,11 @@ export async function monitorSignalProvider(opts: MonitorSignalOpts = {}): Promi
729736
}
730737
throw err;
731738
} finally {
739+
// Stop first: pending retry delays observe abort, while retries that already
740+
// started remain in the task runner and drain before monitor teardown.
741+
daemonLifecycle.stop();
732742
await monitorTaskRunner.waitForIdle();
733743
daemonLifecycle.dispose();
734744
opts.abortSignal?.removeEventListener("abort", onAbort);
735-
daemonLifecycle.stop();
736745
}
737746
}

0 commit comments

Comments
 (0)