Skip to content

Commit eb04c31

Browse files
1stvampTrigger.dev RepoOps
authored andcommitted
fix(supervisor): replace a requeued cold start's Runner instead of failing every redelivery on it
Supervisor: when a run is requeued after its Runner never started, replace the stale Runner instead of retrying into a name conflict until it expires. Mono-RevId: e402081626e486d4ac4b25b173e65fe2d477802f
1 parent e39f445 commit eb04c31

2 files changed

Lines changed: 184 additions & 16 deletions

File tree

‎apps/supervisor/src/workloadManager/runCrd.test.ts‎

Lines changed: 120 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -516,6 +516,126 @@ describe("run-crd carries every shared create-option or excludes it on purpose",
516516
});
517517
});
518518

519+
describe("RunCrdWorkloadManager.create", () => {
520+
const manager = () =>
521+
new RunCrdWorkloadManager({
522+
workloadApiProtocol: "http",
523+
workloadApiPort: 8020,
524+
namespace: "v4-runs",
525+
runtime: "microvm",
526+
});
527+
528+
beforeEach(() => {
529+
createRunner.mockReset();
530+
createRunner.mockResolvedValue({ metadata: { uid: "uid-created" } });
531+
getRunner.mockReset();
532+
deleteRunner.mockReset();
533+
deleteRunner.mockResolvedValue({});
534+
});
535+
536+
// createOptions() is dequeued at 03:00:00, so the default is an earlier delivery.
537+
const inTheWay = (snapshotFriendlyID: string, dequeuedAt?: string) => ({
538+
metadata: { name: "runner-abc123", uid: "uid-existing" },
539+
spec: {
540+
bootstrap: {
541+
runFriendlyID: "run_abc123",
542+
snapshotFriendlyID,
543+
dequeuedAt: dequeuedAt ?? "2026-08-27T02:59:00.000Z",
544+
},
545+
},
546+
});
547+
548+
// The platform requeued an earlier delivery that never started the attempt.
549+
it("replaces a Runner an earlier delivery left, guarded by its uid", async () => {
550+
createRunner
551+
.mockRejectedValueOnce({ code: 409 })
552+
.mockResolvedValueOnce({ metadata: { uid: "uid-recreated" } });
553+
getRunner.mockResolvedValue(inTheWay("snapshot_earlier"));
554+
555+
await manager().create(createOptions());
556+
557+
expect(deleteRunner).toHaveBeenCalledWith(
558+
expect.objectContaining({
559+
name: "runner-abc123",
560+
body: { preconditions: { uid: "uid-existing" } },
561+
})
562+
);
563+
expect(createRunner).toHaveBeenCalledTimes(2);
564+
});
565+
566+
// The first create found the old Runner's Secret, which the collector takes with it.
567+
it("gives the replacement a token Secret of its own", async () => {
568+
createRunner
569+
.mockRejectedValueOnce({ code: 409 })
570+
.mockResolvedValueOnce({ metadata: { uid: "uid-recreated" } });
571+
getRunner.mockResolvedValue(inTheWay("snapshot_earlier"));
572+
createSecret.mockReset();
573+
createSecret.mockResolvedValue({ metadata: { uid: "uid-secret" } });
574+
patchObject.mockResolvedValue({});
575+
576+
await manager().create(createOptions({ deploymentToken: "tok" }));
577+
578+
const [first, second] = createRunner.mock.calls.map(
579+
([{ body }]) => body.spec.deployment.token.name
580+
);
581+
expect(second).not.toBe(first);
582+
});
583+
584+
it("leaves a Runner this delivery already made", async () => {
585+
createRunner.mockRejectedValueOnce({ code: 409 });
586+
getRunner.mockResolvedValue(inTheWay("snapshot_abc"));
587+
588+
await manager().create(createOptions());
589+
590+
expect(deleteRunner).not.toHaveBeenCalled();
591+
expect(createRunner).toHaveBeenCalledTimes(1);
592+
});
593+
594+
// A delivery held up past the requeue must not remove the one that followed it.
595+
it.each([
596+
["a later delivery", "2026-08-27T03:01:00.000Z"],
597+
["a delivery dequeued at the same time", "2026-08-27T03:00:00.000Z"],
598+
["a delivery of unknown order", ""],
599+
])("leaves a Runner from %s and fails as stale", async (_, dequeuedAt) => {
600+
createRunner.mockRejectedValueOnce({ code: 409 });
601+
getRunner.mockResolvedValue(inTheWay("snapshot_later", dequeuedAt));
602+
603+
await expect(manager().create(createOptions())).rejects.toThrow("stale");
604+
expect(deleteRunner).not.toHaveBeenCalled();
605+
expect(createRunner).toHaveBeenCalledTimes(1);
606+
});
607+
608+
// Another delivery of this snapshot made the replacement in the gap.
609+
it("accepts a replacement for this snapshot that another delivery made first", async () => {
610+
createRunner.mockRejectedValue({ code: 409 });
611+
getRunner
612+
.mockResolvedValueOnce(inTheWay("snapshot_earlier"))
613+
.mockResolvedValueOnce(inTheWay("snapshot_abc", "2026-08-27T03:00:00.000Z"));
614+
615+
await manager().create(createOptions());
616+
617+
expect(createRunner).toHaveBeenCalledTimes(2);
618+
});
619+
620+
it("replaces it only once per delivery", async () => {
621+
createRunner.mockRejectedValue({ code: 409 });
622+
getRunner.mockResolvedValue(inTheWay("snapshot_earlier"));
623+
624+
await expect(manager().create(createOptions())).rejects.toThrow("still in the way");
625+
expect(createRunner).toHaveBeenCalledTimes(2);
626+
});
627+
628+
it("creates again when the Runner went between the create and the read", async () => {
629+
createRunner.mockRejectedValueOnce({ code: 409 });
630+
getRunner.mockRejectedValue({ code: 404 });
631+
632+
await manager().create(createOptions());
633+
634+
expect(deleteRunner).not.toHaveBeenCalled();
635+
expect(createRunner).toHaveBeenCalledTimes(2);
636+
});
637+
});
638+
519639
describe("RunCrdWorkloadManager.restore", () => {
520640
const checkpoint = { id: "checkpoint_abc", location: "node-a/6f1c2a9e-snap" };
521641

@@ -709,12 +829,6 @@ describe("RunCrdWorkloadManager.restore", () => {
709829
expect(createRunner).toHaveBeenCalledTimes(2);
710830
});
711831

712-
it("still fails a cold start that finds a Runner in the way", async () => {
713-
createRunner.mockRejectedValue({ code: 409 });
714-
715-
await expect(manager().create(createOptions())).rejects.toEqual({ code: 409 });
716-
});
717-
718832
it("deletes a failed resume only while it is the Runner the watch saw", async () => {
719833
deleteRunner.mockReset();
720834
deleteRunner.mockResolvedValue({});

‎apps/supervisor/src/workloadManager/runCrd.ts‎

Lines changed: 64 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -339,11 +339,16 @@ export class RunCrdWorkloadManager implements WorkloadManager, RunnerSnapshotter
339339
return outcome ? { ...outcome, createdAt: createdAtOf(runner), uid: uidOf(runner) } : timedOut;
340340
}
341341

342+
/** Deletes a failed resume so its redelivery can create it again. */
343+
async deleteRestoreRunner(runnerId: string, uid: string): Promise<void> {
344+
await this.deleteRunner(runnerId, uid);
345+
}
346+
342347
/**
343-
* Deletes a failed resume so its redelivery can create it again. Preconditioned
344-
* on uid, so a Runner that already replaced it is left alone; one already gone is fine.
348+
* Preconditioned on uid, so a Runner that already replaced this one is left
349+
* alone; one already gone is fine.
345350
*/
346-
async deleteRestoreRunner(runnerId: string, uid: string): Promise<void> {
351+
private async deleteRunner(runnerId: string, uid: string): Promise<void> {
347352
try {
348353
await this.k8s.custom.deleteNamespacedCustomObject({
349354
group: GROUP,
@@ -378,8 +383,54 @@ export class RunCrdWorkloadManager implements WorkloadManager, RunnerSnapshotter
378383
return this.runtime === "microvm" && checkpoint.type === "COMPUTE";
379384
}
380385

386+
/**
387+
* A cold start's Runner is named for its attempt, so one already in the way is
388+
* another delivery of this attempt. One for this snapshot is this delivery's
389+
* own create, made again. One dequeued before this delivery is one the
390+
* platform has since requeued (it never started the attempt), which can't
391+
* start it now, so it is replaced, once and guarded by its uid. Left in place,
392+
* it failed every redelivery's create until the operator's terminal TTL took
393+
* it. One dequeued after this delivery is newer, and this delivery is the
394+
* stale one.
395+
*/
381396
async create(opts: WorkloadManagerCreateOptions) {
382-
await this.createRunner(opts, getRunnerId(opts.runFriendlyId, opts.nextAttemptNumber));
397+
const runnerId = getRunnerId(opts.runFriendlyId, opts.nextAttemptNumber);
398+
if (await this.createRunner(opts, runnerId)) {
399+
return;
400+
}
401+
try {
402+
const existing = (await this.getRunner(runnerId)) as RunnerBootstrapSpec | null;
403+
const bootstrap = existing?.spec?.bootstrap;
404+
if (bootstrap?.snapshotFriendlyID === opts.snapshotFriendlyId) {
405+
return;
406+
}
407+
// Unknown order counts as newer: only a Runner shown to be older is deleted.
408+
if (!(Date.parse(bootstrap?.dequeuedAt ?? "") < opts.dequeuedAt.getTime())) {
409+
throw new Error(
410+
`Runner ${runnerId} belongs to a delivery dequeued after this one, which is stale`
411+
);
412+
}
413+
const uid = uidOf(existing);
414+
if (!uid) {
415+
throw new Error(`Runner ${runnerId} is in the way, with no uid to guard its delete`);
416+
}
417+
await this.deleteRunner(runnerId, uid);
418+
} catch (err: unknown) {
419+
// Deleted between the create and the read.
420+
if (statusCodeOf(err) !== 404) {
421+
throw err;
422+
}
423+
}
424+
// A Secret of its own: the one this create found is the deleted Runner's, and
425+
// the collector takes it.
426+
if (await this.createRunner(opts, runnerId, undefined, true)) {
427+
return;
428+
}
429+
// Another delivery of this snapshot may have made it in the gap.
430+
const raced = (await this.getRunner(runnerId)) as RunnerBootstrapSpec | null;
431+
if (raced?.spec?.bootstrap?.snapshotFriendlyID !== opts.snapshotFriendlyId) {
432+
throw new Error(`Runner ${runnerId} is still in the way after replacing it`);
433+
}
383434
}
384435

385436
/**
@@ -462,13 +513,14 @@ export class RunCrdWorkloadManager implements WorkloadManager, RunnerSnapshotter
462513
}
463514
}
464515

465-
/** Returns false when a resume's Runner was already there. */
516+
/** Returns false when the Runner was already there. */
466517
private async createRunner(
467518
opts: WorkloadManagerCreateOptions,
468519
runnerId: string,
469-
restore?: RunnerRestore
520+
restore?: RunnerRestore,
521+
replacing = false
470522
): Promise<{ uid?: string } | false> {
471-
const token = await this.ensureRunnerToken(opts, runnerId, !!restore);
523+
const token = await this.ensureRunnerToken(opts, runnerId, !!restore || replacing);
472524

473525
const body = runnerBodyFor(opts, {
474526
name: runnerId,
@@ -495,8 +547,8 @@ export class RunCrdWorkloadManager implements WorkloadManager, RunnerSnapshotter
495547
await this.releaseRunnerToken(token, err);
496548
// A resume's name carries the checkpoint, so the Runner in the way is this
497549
// resume, still restoring or held terminal for the operator's TTL. A cold
498-
// start's carries the attempt, so its 409 is an earlier terminal Runner.
499-
if (restore && statusCodeOf(err) === 409) {
550+
// start's carries the attempt: create sorts out which delivery made it.
551+
if (statusCodeOf(err) === 409) {
500552
return false;
501553
}
502554
this.logger.error("[RunCrdWorkloadManager] Create failed", { runnerId, rawError: err });
@@ -971,7 +1023,9 @@ export type UnwatchedRestoreFailure = {
9711023
};
9721024

9731025
type RunnerBootstrapSpec = {
974-
spec?: { bootstrap?: { runFriendlyID?: string; snapshotFriendlyID?: string } };
1026+
spec?: {
1027+
bootstrap?: { runFriendlyID?: string; snapshotFriendlyID?: string; dequeuedAt?: string };
1028+
};
9751029
};
9761030

9771031
/**

0 commit comments

Comments
 (0)