fix: cancel pending delayed snapshots when the run completes or disconnects
⚒️ Publish Worker (v4) / build (supervisor) (push) Has been cancelled
⚒️ Publish Worker (v4) / build (supervisor) (push) Has been cancelled
The compute suspend flow delays snapshots by snapshotDelayMs to avoid wasted work on short-lived waitpoints, with the intent that a run which continues before the delay expires cancels the pending snapshot. But the only cancel() call site is the /continue workload action, which runners only invoke when restoring from an already-taken snapshot - so a pending snapshot is never actually cancelled (zero snapshot.canceled events in prod). When a run resumes and completes within the delay window, the stale snapshot fires anyway and fcrun pauses the VM for ~6-13s while its controller is mid warm-start long-poll. The frozen guest can't fire its abort timer or send a FIN, so firestarter keeps the connection claimable past the client deadline and dispatches runs into it - each one a ~300s stall (TRI-10293). Cancel the pending snapshot when the attempt completes and when the run socket disconnects. Genuine waitpoint suspensions keep the runner socket connected and the attempt incomplete, so neither hook cancels a snapshot that is still wanted. Cancellation is guarded by runnerId so a stale duplicate runner for a reassigned run can't cancel the new runner's pending snapshot.
This commit is contained in:
@@ -0,0 +1,130 @@
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import { setTimeout as sleep } from "node:timers/promises";
|
||||
import { ComputeSnapshotService } from "./computeSnapshotService.js";
|
||||
import type { ComputeWorkloadManager } from "../workloadManager/compute.js";
|
||||
import type { SupervisorHttpClient } from "@trigger.dev/core/v3/workers";
|
||||
|
||||
// The TimerWheel ticks every 100ms, so a 200ms delay dispatches within ~300ms.
|
||||
const DELAY_MS = 200;
|
||||
// Long enough that a pending snapshot would certainly have dispatched.
|
||||
const SETTLE_MS = 600;
|
||||
|
||||
function createService() {
|
||||
const snapshot = vi.fn(async (_opts: { runnerId: string; metadata: Record<string, string> }) => true);
|
||||
|
||||
const computeManager = {
|
||||
snapshotDelayMs: DELAY_MS,
|
||||
snapshotDispatchLimit: 1,
|
||||
snapshot,
|
||||
} as unknown as ComputeWorkloadManager;
|
||||
|
||||
const service = new ComputeSnapshotService({
|
||||
computeManager,
|
||||
workerClient: {} as SupervisorHttpClient,
|
||||
wideEventOpts: { service: "supervisor-test", env: {}, enabled: false },
|
||||
});
|
||||
|
||||
return { service, snapshot };
|
||||
}
|
||||
|
||||
function delayedSnapshot(runnerId = "runner-1") {
|
||||
return {
|
||||
runnerId,
|
||||
runFriendlyId: "run_1",
|
||||
snapshotFriendlyId: "snapshot_1",
|
||||
};
|
||||
}
|
||||
|
||||
describe("ComputeSnapshotService", () => {
|
||||
it("dispatches a scheduled snapshot after the delay", async () => {
|
||||
const { service, snapshot } = createService();
|
||||
try {
|
||||
service.schedule("run_1", delayedSnapshot());
|
||||
|
||||
await vi.waitFor(() => expect(snapshot).toHaveBeenCalledTimes(1), { timeout: 2_000 });
|
||||
expect(snapshot).toHaveBeenCalledWith({
|
||||
runnerId: "runner-1",
|
||||
metadata: { runId: "run_1", snapshotFriendlyId: "snapshot_1" },
|
||||
});
|
||||
} finally {
|
||||
service.stop();
|
||||
}
|
||||
});
|
||||
|
||||
it("cancel before the delay expires prevents the dispatch", async () => {
|
||||
const { service, snapshot } = createService();
|
||||
try {
|
||||
service.schedule("run_1", delayedSnapshot());
|
||||
|
||||
expect(service.cancel("run_1")).toBe(true);
|
||||
|
||||
await sleep(SETTLE_MS);
|
||||
expect(snapshot).not.toHaveBeenCalled();
|
||||
} finally {
|
||||
service.stop();
|
||||
}
|
||||
});
|
||||
|
||||
it("cancel returns false when nothing is pending", () => {
|
||||
const { service } = createService();
|
||||
try {
|
||||
expect(service.cancel("run_1")).toBe(false);
|
||||
} finally {
|
||||
service.stop();
|
||||
}
|
||||
});
|
||||
|
||||
it("cancel with a matching runnerId cancels the pending snapshot", async () => {
|
||||
const { service, snapshot } = createService();
|
||||
try {
|
||||
service.schedule("run_1", delayedSnapshot("runner-a"));
|
||||
|
||||
expect(service.cancel("run_1", "runner-a")).toBe(true);
|
||||
|
||||
await sleep(SETTLE_MS);
|
||||
expect(snapshot).not.toHaveBeenCalled();
|
||||
} finally {
|
||||
service.stop();
|
||||
}
|
||||
});
|
||||
|
||||
it("cancel with a different runnerId leaves the pending snapshot alone", async () => {
|
||||
const { service, snapshot } = createService();
|
||||
try {
|
||||
service.schedule("run_1", delayedSnapshot("runner-a"));
|
||||
|
||||
// A stale runner for a reassigned run must not cancel the new runner's snapshot.
|
||||
expect(service.cancel("run_1", "runner-b")).toBe(false);
|
||||
|
||||
await vi.waitFor(() => expect(snapshot).toHaveBeenCalledTimes(1), { timeout: 2_000 });
|
||||
expect(snapshot).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ runnerId: "runner-a" })
|
||||
);
|
||||
} finally {
|
||||
service.stop();
|
||||
}
|
||||
});
|
||||
|
||||
it("re-scheduling the same run replaces the pending snapshot", async () => {
|
||||
const { service, snapshot } = createService();
|
||||
try {
|
||||
service.schedule("run_1", delayedSnapshot());
|
||||
service.schedule("run_1", {
|
||||
runnerId: "runner-1",
|
||||
runFriendlyId: "run_1",
|
||||
snapshotFriendlyId: "snapshot_2",
|
||||
});
|
||||
|
||||
await vi.waitFor(() => expect(snapshot).toHaveBeenCalledTimes(1), { timeout: 2_000 });
|
||||
await sleep(SETTLE_MS);
|
||||
|
||||
expect(snapshot).toHaveBeenCalledTimes(1);
|
||||
expect(snapshot).toHaveBeenCalledWith({
|
||||
runnerId: "runner-1",
|
||||
metadata: { runId: "run_1", snapshotFriendlyId: "snapshot_2" },
|
||||
});
|
||||
} finally {
|
||||
service.stop();
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -92,8 +92,19 @@ export class ComputeSnapshotService {
|
||||
});
|
||||
}
|
||||
|
||||
/** Cancel a pending delayed snapshot. Returns true if one was cancelled. */
|
||||
cancel(runFriendlyId: string): boolean {
|
||||
/**
|
||||
* Cancel a pending delayed snapshot. Returns true if one was cancelled.
|
||||
* When `runnerId` is given, only a snapshot scheduled for that same runner
|
||||
* is cancelled - a stale runner for a run that has since been reassigned
|
||||
* must not cancel the new runner's pending snapshot.
|
||||
*/
|
||||
cancel(runFriendlyId: string, runnerId?: string): boolean {
|
||||
if (runnerId) {
|
||||
const pending = this.timerWheel.peek(runFriendlyId);
|
||||
if (pending && pending.data.runnerId !== runnerId) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
const cancelled = this.timerWheel.cancel(runFriendlyId);
|
||||
if (cancelled) {
|
||||
emitOneShot({
|
||||
|
||||
@@ -121,6 +121,12 @@ export class TimerWheel<T> {
|
||||
return true;
|
||||
}
|
||||
|
||||
/** Look up a pending item without removing it. */
|
||||
peek(key: string): TimerWheelItem<T> | undefined {
|
||||
const entry = this.entries.get(key);
|
||||
return entry ? { key, data: entry.data } : undefined;
|
||||
}
|
||||
|
||||
/** Number of pending items in the wheel. */
|
||||
get size(): number {
|
||||
return this.entries.size;
|
||||
|
||||
@@ -303,6 +303,16 @@ export class WorkloadServer extends EventEmitter<WorkloadServerEvents> {
|
||||
return;
|
||||
}
|
||||
|
||||
// A completed attempt invalidates any pending delayed snapshot: the
|
||||
// suspended execution state it was scheduled to capture no longer
|
||||
// exists. Without this, the snapshot fires up to snapshotDelayMs
|
||||
// later and pauses a VM that has long moved on, e.g. mid warm-start
|
||||
// long-poll or already executing the next run.
|
||||
this.snapshotService?.cancel(
|
||||
params.runFriendlyId,
|
||||
this.runnerIdFromRequest(req)
|
||||
);
|
||||
|
||||
reply.json(
|
||||
completeResponse.data satisfies WorkloadRunAttemptCompleteResponseBody
|
||||
);
|
||||
@@ -728,6 +738,14 @@ export class WorkloadServer extends EventEmitter<WorkloadServerEvents> {
|
||||
const runDisconnected = (friendlyId: string, reason: string) => {
|
||||
socketLogger.debug("runDisconnected", { ...getSocketMetadata() });
|
||||
|
||||
// The run is gone from this runner (crash, exit, or replaced by a new
|
||||
// run), so a pending delayed snapshot for it is stale. Genuine
|
||||
// waitpoint suspensions keep the socket connected, so this doesn't
|
||||
// cancel a snapshot that's still wanted; the runnerId match guards
|
||||
// against a stale duplicate runner cancelling a fresh runner's
|
||||
// snapshot after the run was reassigned.
|
||||
this.snapshotService?.cancel(friendlyId, socket.data.runnerId);
|
||||
|
||||
this.runSockets.delete(friendlyId);
|
||||
this.emit("runDisconnected", { run: { friendlyId } });
|
||||
socket.data.runFriendlyId = undefined;
|
||||
|
||||
Reference in New Issue
Block a user