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
55 changes: 37 additions & 18 deletions packages/agent-lifecycle/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -169,12 +169,18 @@ export function createAgentLifecycle(
if (typeof interval.unref === "function") interval.unref();

/**
* Races `wake(address)` against `wakeTimeoutMs`. A timeout rejects
* with a distinct error so a caller (and everything coalesced behind
* it) always gets a settled outcome — never a hang — even though the
* underlying `wake` call keeps running unobserved past the deadline.
* Races `pending` against `wakeTimeoutMs`, rejecting with a distinct
* error so a caller always gets a settled outcome — never a hang —
* even though `pending` itself keeps running unobserved past the
* deadline. Called once per `ensureAwake` caller (not once per
* address): every caller coalesced onto the same in-flight `pending`
* gets its own fresh timer, so one caller's timeout never shortens
* another's wait.
*/
function wakeWithTimeout(address: string): Promise<void> {
function boundToTimeout(
pending: Promise<void>,
address: string,
): Promise<void> {
return new Promise((resolve, reject) => {
const timer = setTimeout(() => {
reject(
Expand All @@ -183,7 +189,7 @@ export function createAgentLifecycle(
),
);
}, wakeTimeoutMs);
wake(address).then(
pending.then(
(value) => {
clearTimeout(timer);
resolve(value);
Expand All @@ -196,21 +202,34 @@ export function createAgentLifecycle(
});
}

/**
* `pendingWakes` tracks the raw, untimed `wake(address)` call — never
* the timeout-bound promise a caller sees. Clearing the map entry only
* when the real wake settles (not when a caller's own timeout fires)
* is what keeps a later caller, arriving after an earlier one already
* timed out, coalescing onto the same in-flight wake instead of
* dispatching a second, concurrent one (CL-7217) — the earlier version
* of this function cleared the entry on timeout, which let exactly
* that race through. The trade-off: an address whose injected `wake`
* never settles, ever, stays permanently deduped — no later caller can
* ever trigger a fresh attempt for it — which is an accepted cost
* against a concurrent-wake race that can corrupt host state.
*/
async function ensureAwake(address: string): Promise<void> {
if (isRoutable(address)) return;

const pending = pendingWakes.get(address);
if (pending !== undefined) return pending;

const waking = wakeWithTimeout(address)
.then(() => {
recordActivity(address);
})
.finally(() => {
pendingWakes.delete(address);
});
pendingWakes.set(address, waking);
return waking;
let inFlight = pendingWakes.get(address);
if (inFlight === undefined) {
inFlight = wake(address)
.then(() => {
recordActivity(address);
})
.finally(() => {
pendingWakes.delete(address);
});
pendingWakes.set(address, inFlight);
}
return boundToTimeout(inFlight, address);
}

function stop(): void {
Expand Down
87 changes: 75 additions & 12 deletions packages/agent-lifecycle/test/index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -329,15 +329,26 @@ describe("createAgentLifecycle", () => {
// CL-6643: a run whose deployed record was parked aside (no sidecar
// record, cold wake required) can hang inside the injected `wake` port
// forever — the sidecar never acks a redeploy it has no state for.
// Without a bound, that first hung call never leaves `pendingWakes`
// (its `.finally` never runs because the promise it's attached to
// never settles), so every later `ensureAwake` for the same address
// coalesces onto that same dead promise and hangs too — silently,
// forever, for the rest of the process's life. This is the mechanism
// behind "message accepted, then silence: no fanout, no wake, no
// error" — the caller (`sendMail`) never sees a rejection to report as
// an undelivered notice, because nothing ever rejects.
test("a wake that never settles does not wedge every later call behind it forever", async () => {
// Without a bound, a caller waiting on that hung call would hang too;
// `ensureAwake` bounds each caller's own wait so it always sees a
// settled outcome. This is the mechanism behind "message accepted,
// then silence: no fanout, no wake, no error" — the caller (`sendMail`)
// never sees a rejection to report as an undelivered notice, because
// nothing ever rejects.
//
// CL-7217: a still-hung `wake()` must not let a second caller dispatch
// a SECOND, concurrent `wake()` for the same address — that races both
// redeploys against each other on the host. So every later caller for
// this address keeps coalescing onto the one real `wake()` call for as
// long as it stays in flight, each with its own fresh timeout. The
// trade-off this accepts: an address whose underlying `wake()` never
// settles, ever, is permanently deduped (unwakeable via this dedup
// path) for the rest of the process's life, since nothing ever clears
// `pendingWakes` for it. That is intentional — a silently wedged
// address is a lesser failure than every retry redeploying it
// concurrently forever — and is proven below by a third call still
// seeing no new `wake()` dispatch.
test("a wake that never settles times out every caller but never dispatches a second concurrent wake", async () => {
const routable = new Set<string>();
let wakeCalls = 0;

Expand All @@ -364,11 +375,63 @@ describe("createAgentLifecycle", () => {
expect(wakeCalls).toBe(1);

// A second, later call for the same address — the next message sent
// to the same workbench — must attempt its own wake rather than
// coalescing onto the first call's dead, already-timed-out promise.
// to the same workbench — coalesces onto the still-hung first wake
// rather than dispatching a concurrent one; it gets its own timeout
// and rejects too, but `wake` is not called again.
await expect(lifecycle.ensureAwake("agent-1@t.test")).rejects.toThrow(
/wake/i,
);
expect(wakeCalls).toBe(2);
expect(wakeCalls).toBe(1);

// A third call proves this isn't a one-off race: the dedup holds for
// as long as the underlying wake never settles.
await expect(lifecycle.ensureAwake("agent-1@t.test")).rejects.toThrow(
/wake/i,
);
expect(wakeCalls).toBe(1);
});

// CL-7217: the earlier version of this fix cleared `pendingWakes` the
// moment the per-caller timeout fired, not when `wake()` itself
// settled — so a caller arriving after a timeout, but before the real
// wake finished, dispatched a brand-new concurrent `wake()` call. This
// proves the corrected coalescing: a second caller arriving in that
// window joins the one real wake in flight and observes its outcome,
// instead of triggering its own.
test("a caller arriving after a timeout coalesces onto the still-in-flight wake instead of dispatching a second one", async () => {
const routable = new Set<string>();
let wakeCalls = 0;
let resolveWake: (() => void) | undefined;

const lifecycle = createAgentLifecycle({
idleSleepMs: 1_000,
wakeTimeoutMs: 20,
isRoutable: (address) => routable.has(address),
undeploy: async () => undefined,
wake: async (address) => {
wakeCalls += 1;
await new Promise<void>((resolve) => {
resolveWake = resolve;
});
routable.add(address);
},
log,
});
stop = lifecycle.stop;

const first = lifecycle.ensureAwake("agent-1@t.test");

// The first caller's own 20ms budget expires while the underlying
// wake is still running.
await expect(first).rejects.toThrow(/wake/i);
expect(wakeCalls).toBe(1);

// A second caller arrives after that timeout but before the real
// wake has settled. It must join the one in-flight wake rather than
// starting a new one.
const second = lifecycle.ensureAwake("agent-1@t.test");
resolveWake?.();
await expect(second).resolves.toBeUndefined();
expect(wakeCalls).toBe(1);
});
});
Loading