Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 30 additions & 14 deletions packages/chat/src/platform-adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> {
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<void> {
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);
}

/**
Expand Down
134 changes: 134 additions & 0 deletions packages/chat/test/platform-adapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import type { FoldedRunsDeps } from "@corbits/folded-runs";
import {
agentSession,
asset,
sessionAsset,
sessionMail,
workflowDefinition,
workflowDefinitionVersion,
Expand Down Expand Up @@ -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<typeof originalDeploy>[0],
) => {
const result = await originalDeploy(params);
deployCallCount += 1;
if (deployCallCount >= 2) {
await new Promise<void>((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`
Expand Down
Loading