From 3511e4ba61d39ac44aed5d9d4afdbb06a0aa2898 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 06:23:04 -0700 Subject: [PATCH 1/2] Add test for wakeByAddressBounded reclaim-retry coalescing (CL-7214) sendFoldedMailWithReclaimRetry's reclaim path calls wakeByAddressBounded directly, bypassing lifecycle.ensureAwake's per-address dedup even when a lifecycle is configured. Drives a sendMail call whose first delivery attempt fails "agent is unreachable" (forcing the reclaim retry), gates its redeploy open, and fires a concurrent ensureAwake call for the same address -- asserting only one wakeFoldedRun delete+redeploy cycle runs, not two racing ones. Confirmed the test fails (deployCallCount reaches 3) against the unfixed platform-adapter.ts. --- packages/chat/test/platform-adapter.test.ts | 134 ++++++++++++++++++++ 1 file changed, 134 insertions(+) diff --git a/packages/chat/test/platform-adapter.test.ts b/packages/chat/test/platform-adapter.test.ts index 480768214..f99c00053 100644 --- a/packages/chat/test/platform-adapter.test.ts +++ b/packages/chat/test/platform-adapter.test.ts @@ -27,6 +27,7 @@ import type { FoldedRunsDeps } from "@corbits/folded-runs"; import { agentSession, asset, + sessionAsset, sessionMail, workflowDefinition, workflowDefinitionVersion, @@ -2471,6 +2472,139 @@ describe("createHubChatPlatform", () => { }); }); + // CL-7214: `sendFoldedMailWithReclaimRetry`'s reclaim-retry loop used + // to call `wakeByAddress` directly, bypassing `lifecycle.ensureAwake`'s + // per-address coalescing entirely. Proves that a reclaim-retry wake and + // an independent, concurrent `ensureAwake` call for the same address + // now coalesce onto the same in-flight wake — never dispatching a + // second, concurrent `wakeFoldedRun` that would race the first on the + // same `session_asset` primary key and git ref. + describe("wakeByAddressBounded reclaim-retry coalescing", () => { + test("a reclaim-retry wake and a concurrent ensureAwake call for the same address never redeploy it twice", async () => { + resolveDefinitionSourcesResult = { + ok: true, + sources: [ + { + id: "off_1", + provider: "anthropic", + baseURL: "https://inference.invalid", + apiKey: "placeholder", + model: "claude-sonnet-5", + }, + ], + defaultSource: "off_1", + }; + const address = "ins_workbench1@ten1.workbench.test"; + const db = createFakeDb({ + assetRow: { + tenantId: "ten_1", + creatorPrincipalId: "prin_creator", + name: "workbench-1", + displayName: null, + }, + definitionId: "wfd_workbench1", + workflowRunRow: { + id: "ins_workbench1", + address, + principalId: "prin_run1", + }, + workbenchLaunchRow: { + tenantId: "ten_1", + instanceId: "ins_workbench1", + foldedBody: { + systemPrompt: "host prompt", + model: "claude-sonnet-5", + toolPackagePins: [], + grantRequirements: [], + credentialBindings: [], + }, + }, + }); + db.inserted.push({ + table: agentSession, + values: { id: "ses_run1", principalId: "prin_run1" }, + }); + + const sessionService = createFakeSessionService(); + let sendAttempts = 0; + sessionService.sendUserMessage = async (params: unknown) => { + sessionService.sendUserMessageCalls.push(params); + sendAttempts += 1; + // The first delivery attempt fails as "agent is unreachable", + // forcing `sendFoldedMailWithReclaimRetry` down its reclaim + // path — the call site that used to bypass coalescing. + if (sendAttempts === 1) { + throw new Error("agent is unreachable"); + } + return new TextEncoder().encode("raw-mime-bytes"); + }; + + // The first deploy (the cold wake before the first send attempt) + // resolves immediately; every deploy after it is held open until + // the test releases it, so the reclaim-retry's redeploy is still + // in flight when the concurrent `ensureAwake` call joins it. + let deployCallCount = 0; + const gatedDeployReleases: (() => void)[] = []; + const originalDeploy = + sessionService.deployAdoptedWorkflowFromSource.bind(sessionService); + sessionService.deployAdoptedWorkflowFromSource = async ( + params: Parameters[0], + ) => { + const result = await originalDeploy(params); + deployCallCount += 1; + if (deployCallCount >= 2) { + await new Promise((resolve) => { + gatedDeployReleases.push(resolve); + }); + } + return result; + }; + + const platform = createHubChatPlatform({ + toolGrantsForPins: () => [], + db: db as never, + sessionService, + assetService: createFakeAssetService(), + sidecarRouter: createFakeSidecarRouter({ routableAddresses: [] }), + eventCollectors: createFakeEventCollectors(), + lifecycle: { idleSleepMs: 60_000 }, + reclaimRetryDelaysMs: [1], + mailDeliveryTimeoutMs: 5_000, + }); + + const sendMailPromise = platform.sendMail({ + tenantId: "ten_1", + workbenchId: "ins_workbench1", + principalId: "prin_sender", + content: { content: "hello workbench" }, + }); + + // Let the cold wake (deploy #1) resolve, the first send attempt + // fail, and the reclaim retry's own wake (deploy #2) start and + // gate. + await new Promise((resolve) => setTimeout(resolve, 50)); + expect(deployCallCount).toBe(2); + + // A second, independent caller asks for the same address while + // deploy #2 is still in flight — the concurrent-wake race CL-7214 + // closes. + const ensureAwakePromise = platform.ensureAwake(address); + await new Promise((resolve) => setTimeout(resolve, 50)); + expect(deployCallCount).toBe(2); + expect( + db.deleted.filter((row) => row.table === sessionAsset).length, + ).toBe(2); + + gatedDeployReleases.forEach((release) => release()); + await Promise.all([sendMailPromise, ensureAwakePromise]); + + expect(deployCallCount).toBe(2); + expect( + db.deleted.filter((row) => row.table === sessionAsset).length, + ).toBe(2); + }); + }); + // Proves the actual lever an edited system prompt reaches a running // instance through: `wakeFoldedRun` (exercised via `sendMail`'s // wake-on-send path above) replays `workbench_launch.foldedBody` From 840950d45cdacd3f1182f84d191ff502a4ea0311 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sun, 30 Aug 2026 06:23:17 -0700 Subject: [PATCH 2/2] Route wakeByAddressBounded through lifecycle.ensureAwake (CL-7214) wakeFoldedRun clears a run's session_asset rows and redeploys with no lock spanning the two steps, so two concurrent calls for the same instance race on the same primary key and git ref. sendFoldedMailWithReclaimRetry's reclaim path called wakeByAddress directly, bypassing lifecycle.ensureAwake's per-address pendingWakes coalescing even when a lifecycle was configured -- the one caller that could still race a lifecycle-driven wake for the same address. wakeByAddressBounded now routes through lifecycle.ensureAwake whenever a lifecycle is configured, converging every wake path in this adapter (sendMail, the exported ensureAwake hook, and the reclaim retry) onto the one coalescing map @corbits/agent-lifecycle already owns, rather than adding a second one here. reconcileDriftedRun runs unconditionally afterward, matching the CL-6588 pattern already used at the other two call sites, since lifecycle.ensureAwake no-ops on an address that is already routable. Only the no-lifecycle fallback still calls wakeByAddress directly, unchanged. --- packages/chat/src/platform-adapter.ts | 44 ++++++++++++++++++--------- 1 file changed, 30 insertions(+), 14 deletions(-) diff --git a/packages/chat/src/platform-adapter.ts b/packages/chat/src/platform-adapter.ts index 8f15235bc..32db90fb4 100644 --- a/packages/chat/src/platform-adapter.ts +++ b/packages/chat/src/platform-adapter.ts @@ -675,21 +675,37 @@ export function createHubChatPlatform( /** * `wakeByAddress`, bounded to `DEFAULT_WAKE_TIMEOUT_MS` — the same * bound `@corbits/agent-lifecycle`'s `ensureAwake` puts on this exact - * call when `lifecycle` is configured (CL-6643). Every call site that - * invokes `wakeByAddress` directly rather than through - * `lifecycle.ensureAwake` — `sendFoldedMailWithReclaimRetry`'s reclaim - * retry, and the exported `ensureAwake` hook's no-`lifecycle` fallback - * — used to await the underlying sidecar deploy round-trip with no - * bound at all, so a deploy the sidecar never acked wedged the caller - * (and, for the reclaim retry, every later send) forever instead of - * failing loud (CL-6644). + * call when `lifecycle` is configured (CL-6643), so a deploy the + * sidecar never acked fails loud instead of wedging the caller forever + * (CL-6644). + * + * When `lifecycle` is configured, this routes through + * `lifecycle.ensureAwake` rather than calling `wakeByAddress` itself — + * `sendFoldedMailWithReclaimRetry`'s reclaim retry used to call + * `wakeByAddress` directly, bypassing `lifecycle.ensureAwake`'s + * per-address coalescing entirely. Two wakes for the same instance + * racing in through this bypass could both pass `wakeFoldedRun`'s + * `session_asset` delete and both redeploy, colliding on the same + * primary key and git ref (CL-7214). Every wake path now funnels + * through the one coalescing map `@corbits/agent-lifecycle` owns, + * rather than this package growing a second one beside it. + * `reconcileDriftedRun` mirrors the pattern `sendMail` and the + * exported `ensureAwake` hook already use: `lifecycle.ensureAwake` + * no-ops on an address that is already routable, so a staleness check + * needs to run unconditionally alongside it (CL-6588). Only the + * no-`lifecycle` fallback still calls `wakeByAddress` directly. */ - function wakeByAddressBounded(address: string): Promise { - return withTimeout( - wakeByAddress(address), - DEFAULT_WAKE_TIMEOUT_MS, - `wake for "${address}" did not settle within ${String(DEFAULT_WAKE_TIMEOUT_MS)}ms`, - ); + async function wakeByAddressBounded(address: string): Promise { + if (lifecycle === undefined) { + await withTimeout( + wakeByAddress(address), + DEFAULT_WAKE_TIMEOUT_MS, + `wake for "${address}" did not settle within ${String(DEFAULT_WAKE_TIMEOUT_MS)}ms`, + ); + return; + } + await lifecycle.ensureAwake(address); + await reconcileDriftedRun(address); } /**