FN-8812: fix workspace merge retry timer cleanup
Ensure workspace contention retries are canceled when the project engine stops. - Track workspace busy re-enqueue timers separately from other merge timers - Clear pending workspace retry callbacks during engine shutdown - Strengthen retry ladder and timer cleanup coverage Files changed: .changeset/workspace-busy-timer-cleanup.md | 7 +++ .../engine/src/__tests__/project-engine.test.ts | 66 +++++++++++----------- packages/engine/src/project-engine.ts | 23 +++++++- 3 files changed, 59 insertions(+), 37 deletions(-) Fusion-Task-Id: FN-8812 Fusion-Task-Lineage: dd4c9581-8738-4082-8d6b-d91d4c463a5a Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/workspace-busy-timer-cleanup.md
Normal file
7
.changeset/workspace-busy-timer-cleanup.md
Normal file
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
summary: Stop cancels pending workspace merge contention retries.
|
||||
category: fix
|
||||
dev: Tracks workspace busy re-enqueue timers separately from merge maintenance timers.
|
||||
@@ -1715,8 +1715,15 @@ describe("ProjectEngine workspace merge dispatch hardening (Phase C review)", ()
|
||||
"FN-WSH",
|
||||
expect.objectContaining({ mergeRetries: expect.anything(), status: null }),
|
||||
);
|
||||
// The isolated child-process seam proves this fail-closed ordering never falls through
|
||||
// to the root-cwd reachability probe (`git remote` was the original guard timeout).
|
||||
const gitRemoteCalls = (mocks.execFile.mock.calls as Array<[string, string[]]>).filter(
|
||||
(call) => Array.isArray(call[1]) && call[1][0] === "remote",
|
||||
);
|
||||
expect(gitRemoteCalls).toHaveLength(0);
|
||||
|
||||
await engine.stop();
|
||||
await expect(engine.stop()).resolves.toBeUndefined();
|
||||
expect(vi.getTimerCount()).toBe(0);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
@@ -1801,6 +1808,10 @@ describe("ProjectEngine workspace merge dispatch hardening (Phase C review)", ()
|
||||
engine as unknown as { internalEnqueueMerge: (id: string) => void },
|
||||
"internalEnqueueMerge",
|
||||
);
|
||||
const scheduleBusyRetrySpy = vi.spyOn(
|
||||
engine as unknown as { scheduleWorkspaceBusyReenqueue: (id: string, delayMs: number) => void },
|
||||
"scheduleWorkspaceBusyReenqueue",
|
||||
);
|
||||
engine.enqueueMerge("FN-WSH");
|
||||
|
||||
// The busy catch logs a WorkspaceRepoLandBusy entry then schedules a backoff timer.
|
||||
@@ -1821,44 +1832,31 @@ describe("ProjectEngine workspace merge dispatch hardening (Phase C review)", ()
|
||||
expect(burnedRetries).toBe(false);
|
||||
|
||||
/*
|
||||
FNXC:Workspace 2026-06-22-09:30 (Phase C review B5b — assert the 60s CAP, not just the first retry):
|
||||
Advancing 60s once only proves the first 5s timer fired; an UNcapped exponential
|
||||
(5s,10s,20s,40s,80s,160s,…) would still pass that. Capture EVERY scheduled busy backoff delay
|
||||
across enough cycles to pass the cap point (busyCount=4 → 5000*2^4 = 80_000ms, clamped to 60_000)
|
||||
and assert no delay exceeds 60_000 AND the cap is actually reached. Each advance fires the pending
|
||||
timer → re-enqueue → landWorkspaceTask rejects busy again → next backoff is scheduled.
|
||||
FNXC:WorkspaceMergeDispatch 2026-08-05-23:56:
|
||||
Attribute delays at the workspace-owned scheduler, never global setTimeout: merge dispatch also
|
||||
schedules body-settle and maintenance timers, which made an unrelated 120s timer look like a
|
||||
workspace backoff. Drive the first six busy attempts through the cap and require its exact ladder.
|
||||
*/
|
||||
const scheduledBusyDelays: number[] = [];
|
||||
// `globalThis.setTimeout` is already the fake-timer impl here (vi.useFakeTimers above).
|
||||
// Wrap it to record the requested delay, then delegate to the SAME fake timer so the
|
||||
// fake clock still drives the callback — no real-timer leakage.
|
||||
const fakeSetTimeout = globalThis.setTimeout;
|
||||
const setTimeoutSpy = vi
|
||||
.spyOn(globalThis, "setTimeout")
|
||||
.mockImplementation(((cb: (...a: unknown[]) => void, ms?: number, ...rest: unknown[]) => {
|
||||
if (typeof ms === "number") scheduledBusyDelays.push(ms);
|
||||
return (fakeSetTimeout as (...a: unknown[]) => unknown)(cb, ms, ...rest);
|
||||
}) as typeof setTimeout);
|
||||
|
||||
try {
|
||||
// Drive enough busy cycles to climb past the cap point (busyCount 0..5 = 6 cycles).
|
||||
for (let i = 0; i < 6; i++) {
|
||||
await vi.advanceTimersByTimeAsync(60_000);
|
||||
}
|
||||
} finally {
|
||||
setTimeoutSpy.mockRestore();
|
||||
const expectedBusyDelays = [5_000, 10_000, 20_000, 40_000, 60_000, 60_000] as const;
|
||||
await vi.waitFor(() => {
|
||||
expect(scheduleBusyRetrySpy).toHaveBeenCalledWith("FN-WSH", expectedBusyDelays[0]);
|
||||
});
|
||||
for (const [index, delayMs] of expectedBusyDelays.slice(0, -1).entries()) {
|
||||
await vi.advanceTimersByTimeAsync(delayMs);
|
||||
await vi.waitFor(() => {
|
||||
expect(scheduleBusyRetrySpy).toHaveBeenCalledTimes(index + 2);
|
||||
});
|
||||
}
|
||||
|
||||
// The exponential climbed (more than one distinct delay) AND every delay is capped at 60s.
|
||||
expect(scheduledBusyDelays.length).toBeGreaterThanOrEqual(5);
|
||||
expect(Math.max(...scheduledBusyDelays)).toBe(60_000);
|
||||
expect(scheduledBusyDelays.every((d) => d <= 60_000)).toBe(true);
|
||||
// The cap was actually exercised: at least one delay sits at the 60s ceiling.
|
||||
expect(scheduledBusyDelays).toContain(60_000);
|
||||
// Each fired backoff re-enqueued the merge (the contention retry loop is live).
|
||||
const scheduledBusyDelays = scheduleBusyRetrySpy.mock.calls.map(([, delayMs]) => delayMs);
|
||||
expect(scheduledBusyDelays).toEqual(expectedBusyDelays);
|
||||
expect(scheduledBusyDelays).not.toContain(120_000);
|
||||
// Each fired backoff reaches one fresh merge body; the queue remains single-flight.
|
||||
expect(mocks.landWorkspaceTask).toHaveBeenCalledTimes(expectedBusyDelays.length);
|
||||
expect(enqueueSpy).toHaveBeenCalledWith("FN-WSH");
|
||||
|
||||
await engine.stop();
|
||||
await expect(engine.stop()).resolves.toBeUndefined();
|
||||
expect(vi.getTimerCount()).toBe(0);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
|
||||
@@ -497,8 +497,23 @@ export class ProjectEngine {
|
||||
Cleared on the first non-busy outcome (success path resets it).
|
||||
*/
|
||||
private workspaceBusyReenqueues = new Map<string, number>();
|
||||
private readonly workspaceBusyReenqueueTimers = new Set<ReturnType<typeof setTimeout>>();
|
||||
private static readonly WORKSPACE_BUSY_MAX_REENQUEUES = 10;
|
||||
|
||||
/*
|
||||
FNXC:WorkspaceMergeDispatch 2026-08-05-23:56:
|
||||
Workspace lease-contention retries are engine-owned lifecycle work, not detached callbacks.
|
||||
Track only these busy re-enqueue timers so stop() can cancel them and tests can measure the
|
||||
capped workspace ladder without confusing it with merge-body or maintenance timers.
|
||||
*/
|
||||
private scheduleWorkspaceBusyReenqueue(taskId: string, delayMs: number): void {
|
||||
const timer = setTimeout(() => {
|
||||
this.workspaceBusyReenqueueTimers.delete(timer);
|
||||
if (!this.shuttingDown) this.internalEnqueueMerge(taskId);
|
||||
}, delayMs);
|
||||
this.workspaceBusyReenqueueTimers.add(timer);
|
||||
}
|
||||
|
||||
/**
|
||||
* Pending manual merge resolvers — keyed by taskId.
|
||||
* When `onMerge` is called, the task is enqueued like auto-merge but a
|
||||
@@ -1328,6 +1343,10 @@ export class ProjectEngine {
|
||||
clearInterval(this.mergeActiveReconcileTimer);
|
||||
this.mergeActiveReconcileTimer = null;
|
||||
}
|
||||
for (const timer of this.workspaceBusyReenqueueTimers) {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
this.workspaceBusyReenqueueTimers.clear();
|
||||
for (const [taskId, deferred] of this.capacityDeferredMerges) {
|
||||
clearTimeout(deferred.timer);
|
||||
for (const resolver of deferred.resolvers) {
|
||||
@@ -4323,9 +4342,7 @@ export class ProjectEngine {
|
||||
runtimeLog.log(
|
||||
`Workspace land busy re-enqueue ${busyCount + 1}/${ProjectEngine.WORKSPACE_BUSY_MAX_REENQUEUES} for ${taskId} in ${delayMs / 1000}s (no mergeRetry consumed — pure lease contention)`,
|
||||
);
|
||||
setTimeout(() => {
|
||||
if (!this.shuttingDown) this.internalEnqueueMerge(taskId);
|
||||
}, delayMs);
|
||||
this.scheduleWorkspaceBusyReenqueue(taskId, delayMs);
|
||||
} else {
|
||||
// Pathological sustained contention — surface but do NOT burn mergeRetries; park as
|
||||
// failed so the cooldown sweep stops re-attempting and an operator can intervene.
|
||||
|
||||
Reference in New Issue
Block a user