From d92b91ae70b470b36c13726513451fc4e7b238c8 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sun, 16 Aug 2026 21:52:30 -0700 Subject: [PATCH] fix(matrix): preserve thread binding activity on shutdown (#125039) * fix(matrix): await binding flush on shutdown Amp-Thread-ID: https://ampcode.com/threads/T-01a00b6b-e4e9-74af-bb31-30363fae6c89 * test(gateway): split provisioning intent coverage Amp-Thread-ID: https://ampcode.com/threads/T-01a00b6b-e4e9-74af-bb31-30363fae6c89 * chore: drop superseded CI unblock Amp-Thread-ID: https://ampcode.com/threads/T-01a00b6b-e4e9-74af-bb31-30363fae6c89 --------- Co-authored-by: Amp --- .../matrix/src/matrix/monitor/index.test.ts | 18 +++++++- extensions/matrix/src/matrix/monitor/index.ts | 4 +- .../src/matrix/thread-bindings-shared.ts | 2 +- .../matrix/src/matrix/thread-bindings.test.ts | 42 +++++++------------ .../matrix/src/matrix/thread-bindings.ts | 10 +++-- 5 files changed, 41 insertions(+), 35 deletions(-) diff --git a/extensions/matrix/src/matrix/monitor/index.test.ts b/extensions/matrix/src/matrix/monitor/index.test.ts index 2d46e707aa59..5601f1027205 100644 --- a/extensions/matrix/src/matrix/monitor/index.test.ts +++ b/extensions/matrix/src/matrix/monitor/index.test.ts @@ -129,7 +129,7 @@ describe("monitorMatrixProvider", () => { hoisted.getMemberDisplayName.mockReset().mockResolvedValue("Bot"); hoisted.registeredOnRoomMessage = null; hoisted.registeredHealthySyncGetter = undefined; - hoisted.stopThreadBindingManager.mockReset(); + hoisted.stopThreadBindingManager.mockReset().mockResolvedValue(undefined); hoisted.client.removeAllListeners(); hoisted.client.hasPersistedSyncState.mockReset().mockReturnValue(false); hoisted.client.drainPendingDecryptions.mockReset().mockResolvedValue(undefined); @@ -592,6 +592,7 @@ describe("monitorMatrixProvider", () => { it("detaches listeners, closes admission, waits for handlers, then releases", async () => { const abortController = new AbortController(); const pendingHandlers = new Map void>(); + let finishManagerStop: (() => void) | undefined; hoisted.createMatrixRoomMessageHandler.mockReturnValue( vi.fn((_roomId: string, event: unknown) => { @@ -608,8 +609,14 @@ describe("monitorMatrixProvider", () => { hoisted.client.drainPendingDecryptions.mockImplementation(async () => { hoisted.callOrder.push("drain-decrypts"); }); - hoisted.stopThreadBindingManager.mockImplementation(() => { + hoisted.stopThreadBindingManager.mockImplementation(async () => { hoisted.callOrder.push("stop-manager"); + await new Promise((resolve) => { + finishManagerStop = () => { + hoisted.callOrder.push("manager-stopped"); + resolve(); + }; + }); }); hoisted.releaseSharedClientInstance.mockImplementation(async () => { await hoisted.client.drainPendingDecryptions(); @@ -632,6 +639,10 @@ describe("monitorMatrixProvider", () => { pendingHandlers.get("$event")?.(); await roomMessagePromise; + await waitForCallOrderEntry("stop-manager"); + expect(hoisted.callOrder).not.toContain("release-client"); + + finishManagerStop?.(); await monitorPromise; expect(hoisted.callOrder.indexOf("drain-decrypts")).toBeLessThan( @@ -647,6 +658,9 @@ describe("monitorMatrixProvider", () => { hoisted.callOrder.indexOf("stop-manager"), ); expect(hoisted.callOrder.indexOf("stop-manager")).toBeLessThan( + hoisted.callOrder.indexOf("manager-stopped"), + ); + expect(hoisted.callOrder.indexOf("manager-stopped")).toBeLessThan( hoisted.callOrder.indexOf("release-client"), ); }); diff --git a/extensions/matrix/src/matrix/monitor/index.ts b/extensions/matrix/src/matrix/monitor/index.ts index 8fd8c73aca99..63dfabc12f0f 100644 --- a/extensions/matrix/src/matrix/monitor/index.ts +++ b/extensions/matrix/src/matrix/monitor/index.ts @@ -204,7 +204,7 @@ export async function monitorMatrixProvider(opts: MonitorMatrixOpts = {}): Promi let client: MatrixClient | null = null; let clientLease: SharedMatrixClientLease | null = null; let monitorLifecycleSignal = opts.abortSignal; - let threadBindingManager: { accountId: string; stop: () => void } | null = null; + let threadBindingManager: { accountId: string; stop: () => Promise } | null = null; const monitorTaskRunner = createMatrixMonitorTaskRunner({ logger, logVerboseMessage, @@ -444,7 +444,7 @@ export async function monitorMatrixProvider(opts: MonitorMatrixOpts = {}): Promi logVerboseMessage, }); if (monitorSetupClosed) { - createdThreadBindingManager.stop(); + await createdThreadBindingManager.stop(); await cleanup("stop"); return; } diff --git a/extensions/matrix/src/matrix/thread-bindings-shared.ts b/extensions/matrix/src/matrix/thread-bindings-shared.ts index 8e0d0decc26b..efac64f4a0d8 100644 --- a/extensions/matrix/src/matrix/thread-bindings-shared.ts +++ b/extensions/matrix/src/matrix/thread-bindings-shared.ts @@ -42,7 +42,7 @@ export type MatrixThreadBindingManager = { maxAgeMs: number; }) => MatrixThreadBindingRecord[]; persist: () => Promise; - stop: () => void; + stop: () => Promise; }; type MatrixThreadBindingManagerCacheEntry = { diff --git a/extensions/matrix/src/matrix/thread-bindings.test.ts b/extensions/matrix/src/matrix/thread-bindings.test.ts index 2406f86a7be7..e3c900cca78f 100644 --- a/extensions/matrix/src/matrix/thread-bindings.test.ts +++ b/extensions/matrix/src/matrix/thread-bindings.test.ts @@ -52,10 +52,8 @@ describe("matrix thread bindings", () => { const matrixClient = {} as never; const trackedManagers = new Set(); - function resetThreadBindingAdapters() { - for (const manager of trackedManagers) { - manager.stop(); - } + async function resetThreadBindingAdapters() { + await Promise.all([...trackedManagers].map((manager) => manager.stop())); trackedManagers.clear(); testing.resetSessionBindingAdaptersForTests(); } @@ -197,9 +195,9 @@ describe("matrix thread bindings", () => { return call; } - beforeEach(() => { + beforeEach(async () => { stateDir = fsSync.mkdtempSync(path.join(os.tmpdir(), "matrix-thread-bindings-")); - resetThreadBindingAdapters(); + await resetThreadBindingAdapters(); resetPluginStateStoreForTests(); sendMessageMatrixMock.mockClear(); setMatrixRuntime({ @@ -213,8 +211,8 @@ describe("matrix thread bindings", () => { } as PluginRuntime); }); - afterEach(() => { - resetThreadBindingAdapters(); + afterEach(async () => { + await resetThreadBindingAdapters(); resetPluginStateStoreForTests(); vi.restoreAllMocks(); vi.useRealTimers(); @@ -433,8 +431,8 @@ describe("matrix thread bindings", () => { }); writeAuthStorageMeta(initialAuth, initialStoragePaths); - initialManager.stop(); - resetThreadBindingAdapters(); + await initialManager.stop(); + await resetThreadBindingAdapters(); await createBindingManager({ auth: rotatedAuth }); @@ -484,8 +482,8 @@ describe("matrix thread bindings", () => { targetSessionKey: "agent:ops:subagent:child", }); - initialManager.stop(); - resetThreadBindingAdapters(); + await initialManager.stop(); + await resetThreadBindingAdapters(); await createBindingManager({ auth: rotatedAuth }); @@ -549,7 +547,7 @@ describe("matrix thread bindings", () => { targetSessionKey: "agent:ops:subagent:child", }); - initialManager.stop(); + await initialManager.stop(); expect( replacementManager.getByConversation({ @@ -635,13 +633,8 @@ describe("matrix thread bindings", () => { await vi.advanceTimersByTimeAsync(1_000); vi.useRealTimers(); - manager.stop(); - await vi.waitFor( - async () => { - expect(await readPersistedLastActivityAt(bindingsPath)).toBe(secondTouchedAt); - }, - { interval: 1, timeout: 5_000 }, - ); + await manager.stop(); + expect(await readPersistedLastActivityAt(bindingsPath)).toBe(secondTouchedAt); } finally { vi.useRealTimers(); } @@ -656,16 +649,11 @@ describe("matrix thread bindings", () => { const touchedAt = Date.parse("2026-03-06T12:00:00.000Z"); getSessionBindingService().touch(binding.bindingId, touchedAt); - manager.stop(); + await manager.stop(); vi.useRealTimers(); const bindingsPath = resolveBindingsFilePath(); - await vi.waitFor( - async () => { - expect(await readPersistedLastActivityAt(bindingsPath)).toBe(touchedAt); - }, - { interval: 1, timeout: 1_000 }, - ); + expect(await readPersistedLastActivityAt(bindingsPath)).toBe(touchedAt); } finally { vi.useRealTimers(); } diff --git a/extensions/matrix/src/matrix/thread-bindings.ts b/extensions/matrix/src/matrix/thread-bindings.ts index f1d1d31de6c1..f653bd75c4b8 100644 --- a/extensions/matrix/src/matrix/thread-bindings.ts +++ b/extensions/matrix/src/matrix/thread-bindings.ts @@ -316,7 +316,7 @@ export async function createMatrixThreadBindingManager(params: { if (existingEntry.storageKey === storageKey) { return existingEntry.manager; } - existingEntry.manager.stop(); + await existingEntry.manager.stop(); } const pluginLoaded = await loadBindingsFromPluginState({ accountId: params.accountId, @@ -482,15 +482,18 @@ export async function createMatrixThreadBindingManager(params: { }), }); }, - stop: () => { + stop: async () => { if (sweepTimer) { clearInterval(sweepTimer); } + let finalPersist = persistQueue; if (persistTimer) { clearTimeout(persistTimer); persistTimer = null; - persistSafely("shutdown-flush"); + finalPersist = enqueuePersist(); } + // Retire the live generation now, but settle its captured persistence before + // shutdown can close the shared Matrix state store. unregisterSessionBindingAdapter({ channel: "matrix", accountId: params.accountId, @@ -503,6 +506,7 @@ export async function createMatrixThreadBindingManager(params: { removeBindingRecord(record); } } + await finalPersist; }, };