Skip to content

Commit f049c34

Browse files
1stvampTrigger.dev RepoOps
authored andcommitted
fix(supervisor): watch a restore Runner and report a failed restore
Supervisor: with the Runner CRD backend, a restore that fails on its node is detected and reported with the reason, instead of only surfacing when the platform's heartbeat timeout fires. Mono-RevId: c95a9ebc98d5b336717fba5cae03510a46c5b19c
1 parent 0dd2a32 commit f049c34

3 files changed

Lines changed: 250 additions & 5 deletions

File tree

‎apps/supervisor/src/index.ts‎

Lines changed: 34 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -98,6 +98,7 @@ class ManagedSupervisor {
9898
private readonly failedPodHandler?: FailedPodHandler;
9999
private readonly tracing?: OtlpTraceService;
100100
private readonly backpressureMonitors: BackpressureMonitor[] = [];
101+
private readonly watchedRestores = new Set<string>();
101102
private readonly backpressureRedis?: Redis;
102103

103104
private readonly isKubernetes = isKubernetesEnvironment(env.KUBERNETES_FORCE_ENABLED);
@@ -709,10 +710,12 @@ class ManagedSupervisor {
709710
) {
710711
const restoreStart = performance.now();
711712
try {
712-
await manager.restore(await this.createOptionsFor(message), checkpoint);
713+
const runnerId = await manager.restore(await this.createOptionsFor(message), checkpoint);
713714
recordPhaseSince("restore", restoreStart, undefined);
714715
setExtra(fromContext(), "did_restore", true);
715716
this.logger.debug("Runner restore created", { runId: message.run.id });
717+
// Not awaited: the resume can take minutes, and the dequeue is done.
718+
void this.watchRestore(manager, message.run.friendlyId, runnerId);
716719
} catch (error) {
717720
recordPhaseSince(
718721
"restore",
@@ -724,6 +727,36 @@ class ManagedSupervisor {
724727
}
725728
}
726729

730+
/** A resume that fails on the node is otherwise silent until the run's heartbeat stalls. */
731+
private async watchRestore(
732+
manager: RunCrdWorkloadManager,
733+
runFriendlyId: string,
734+
runnerId: string
735+
) {
736+
// A redelivered restore finds the same Runner, which needs only one watch.
737+
if (this.watchedRestores.has(runnerId)) {
738+
return;
739+
}
740+
this.watchedRestores.add(runnerId);
741+
const outcome = await manager
742+
.awaitRestore(runnerId)
743+
.finally(() => this.watchedRestores.delete(runnerId));
744+
if (outcome.ok) {
745+
this.logger.debug("Runner restore started", { runFriendlyId, runnerId });
746+
return;
747+
}
748+
this.logger.error("Runner restore failed (run-crd)", {
749+
runFriendlyId,
750+
runnerId,
751+
error: outcome.error,
752+
});
753+
await this.workerSession.httpClient.sendDebugLog(runFriendlyId, {
754+
time: new Date(),
755+
message: "restore failed on the node",
756+
properties: { runnerId, error: outcome.error },
757+
});
758+
}
759+
727760
private async createWorkload(message: DequeuedMessage, timings: WarmStartTimings) {
728761
const createStart = performance.now();
729762
try {

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

Lines changed: 129 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,9 +5,11 @@ import {
55
SUSPEND_ANNOTATION,
66
SUSPEND_RUN_ANNOTATION,
77
SUSPEND_SUBMITTED_ANNOTATION,
8+
awaitRestoreOf,
89
checkpointLocation,
910
parseCheckpointLocation,
1011
publishedSuspend,
12+
restoreOutcome,
1113
runnerBodyFor,
1214
runnerTokenSecretName,
1315
suspendOutcome,
@@ -517,7 +519,9 @@ describe("RunCrdWorkloadManager.restore", () => {
517519
}
518520

519521
it("creates a Runner named from the checkpoint that restores its location", async () => {
520-
await manager().restore(createOptions(), checkpoint);
522+
await expect(manager().restore(createOptions(), checkpoint)).resolves.toBe(
523+
getRestoreRunnerId("run_abc123", "checkpoint_abc")
524+
);
521525

522526
const { body } = createRunner.mock.calls[0]![0];
523527
expect(body.metadata.name).toBe(getRestoreRunnerId("run_abc123", "checkpoint_abc"));
@@ -538,7 +542,9 @@ describe("RunCrdWorkloadManager.restore", () => {
538542
createRunner.mockRejectedValue({ code: 409 });
539543
getRunner.mockResolvedValue(existing(phase));
540544

541-
await expect(manager().restore(createOptions(), checkpoint)).resolves.toBeUndefined();
545+
await expect(manager().restore(createOptions(), checkpoint)).resolves.toBe(
546+
getRestoreRunnerId("run_abc123", "checkpoint_abc")
547+
);
542548
expect(getRunner).toHaveBeenCalledWith(
543549
expect.objectContaining({ name: getRestoreRunnerId("run_abc123", "checkpoint_abc") })
544550
);
@@ -589,6 +595,127 @@ describe("RunCrdWorkloadManager.restore", () => {
589595
});
590596
});
591597

598+
describe("restoreOutcome", () => {
599+
it.each([undefined, "Pending", "Admitted", "Scheduling", "Restoring"])(
600+
"is still waiting while the Runner is %s",
601+
(phase) => {
602+
expect(restoreOutcome({ status: phase ? { phase } : undefined })).toBeUndefined();
603+
}
604+
);
605+
606+
it.each(["Running", "Suspending", "Succeeded"])("is a success once the Runner is %s", (phase) => {
607+
expect(restoreOutcome({ status: { phase } })).toEqual({ ok: true });
608+
});
609+
610+
it("is a failure carrying the operator's reason", () => {
611+
const runner = {
612+
status: {
613+
phase: "Failed",
614+
conditions: [
615+
{
616+
type: "Failed",
617+
status: "True",
618+
reason: "StartError",
619+
message: "pulling the image: 401",
620+
},
621+
],
622+
},
623+
};
624+
expect(restoreOutcome(runner)).toEqual({
625+
ok: false,
626+
error: "StartError: pulling the image: 401",
627+
});
628+
});
629+
630+
it("is a failure even when the operator recorded no reason", () => {
631+
expect(restoreOutcome({ status: { phase: "Failed" } })).toEqual({
632+
ok: false,
633+
error: "the Runner failed with no reason recorded",
634+
});
635+
});
636+
});
637+
638+
describe("awaitRestoreOf", () => {
639+
/** Answers each read with the next response, repeating the last. `{ throw: err }` fails that read. */
640+
function reads(...responses: unknown[]) {
641+
const log: unknown[] = [];
642+
let i = 0;
643+
const readRunner = async () => {
644+
const response = responses[Math.min(i++, responses.length - 1)];
645+
log.push(response);
646+
if (response && typeof response === "object" && "throw" in response) {
647+
throw response.throw;
648+
}
649+
return response;
650+
};
651+
return { readRunner, log };
652+
}
653+
654+
function awaitWith(readRunner: () => Promise<unknown>, timeoutMs = 1_000) {
655+
return awaitRestoreOf(readRunner, { pollMs: 1, timeoutMs, onReadError: () => {} });
656+
}
657+
658+
it("waits through Restoring until the Runner is Running", async () => {
659+
const { readRunner, log } = reads(
660+
{ status: { phase: "Pending" } },
661+
{ status: { phase: "Restoring" } },
662+
{ status: { phase: "Running" } }
663+
);
664+
665+
await expect(awaitWith(readRunner)).resolves.toEqual({ ok: true });
666+
expect(log).toHaveLength(3);
667+
});
668+
669+
it("reports a failed restore with the operator's reason", async () => {
670+
const { readRunner } = reads(
671+
{ status: { phase: "Restoring" } },
672+
{
673+
status: {
674+
phase: "Failed",
675+
conditions: [{ type: "Failed", reason: "SnapshotNodeGone", message: "node a is gone" }],
676+
},
677+
}
678+
);
679+
680+
await expect(awaitWith(readRunner)).resolves.toEqual({
681+
ok: false,
682+
error: "SnapshotNodeGone: node a is gone",
683+
});
684+
});
685+
686+
it("keeps polling through a read that may succeed next time", async () => {
687+
const readErrors: unknown[] = [];
688+
const { readRunner } = reads({ throw: { code: 500 } }, { status: { phase: "Running" } });
689+
690+
await expect(
691+
awaitRestoreOf(readRunner, {
692+
pollMs: 1,
693+
timeoutMs: 1_000,
694+
onReadError: (err) => readErrors.push(err),
695+
})
696+
).resolves.toEqual({ ok: true });
697+
expect(readErrors).toEqual([{ code: 500 }]);
698+
});
699+
700+
it("fails when the Runner is gone", async () => {
701+
const { readRunner } = reads({ throw: { code: 404 } });
702+
703+
await expect(awaitWith(readRunner)).resolves.toEqual({
704+
ok: false,
705+
error: "the Runner no longer exists",
706+
});
707+
});
708+
709+
it("gives up after its timeout", async () => {
710+
const { readRunner } = reads({ status: { phase: "Restoring" } });
711+
712+
await expect(awaitWith(readRunner, 20)).resolves.toEqual({
713+
ok: false,
714+
error: "the Runner did not start within 20ms",
715+
});
716+
});
717+
});
718+
592719
describe("RunCrdWorkloadManager.suspend", () => {
593720
function manager(snapshots = { enabled: true, delayMs: 5_000, dispatchLimit: 10 }) {
594721
return new RunCrdWorkloadManager({

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

Lines changed: 87 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -60,8 +60,15 @@ export type RunCrdWorkloadManagerOptions = WorkloadManagerOptions & {
6060
* than its own snapshot budget, so its failure arrives first and carries the reason.
6161
*/
6262
suspendTimeoutMs?: number;
63+
/**
64+
* How long a resume has to reach Running. Longer than the operator's pod start
65+
* deadline, so its failure arrives first and carries the reason.
66+
*/
67+
restoreTimeoutMs?: number;
6368
};
6469

70+
export type RunnerRestoreResult = { ok: true } | { ok: false; error: string };
71+
6572
/**
6673
* Creates a Runner and stops; the operator builds the pod, so uid, node selection
6774
* and labels are decided in one place. A Runner serves every later run warm start
@@ -75,6 +82,7 @@ export class RunCrdWorkloadManager implements WorkloadManager, RunnerSnapshotter
7582
private readonly snapshots?: RunCrdWorkloadManagerOptions["snapshots"];
7683
private readonly suspendPollMs: number;
7784
private readonly suspendTimeoutMs: number;
85+
private readonly restoreTimeoutMs: number;
7886

7987
constructor(opts: RunCrdWorkloadManagerOptions) {
8088
this.k8s = createK8sApi();
@@ -83,6 +91,7 @@ export class RunCrdWorkloadManager implements WorkloadManager, RunnerSnapshotter
8391
this.snapshots = opts.snapshots;
8492
this.suspendPollMs = opts.suspendPollMs ?? 1_000;
8593
this.suspendTimeoutMs = opts.suspendTimeoutMs ?? 6 * 60_000;
94+
this.restoreTimeoutMs = opts.restoreTimeoutMs ?? 16 * 60_000;
8695
}
8796

8897
get snapshotsEnabled(): boolean {
@@ -229,6 +238,22 @@ export class RunCrdWorkloadManager implements WorkloadManager, RunnerSnapshotter
229238
};
230239
}
231240

241+
/**
242+
* Waits for a resume to start or to fail. Nothing else watches one: a resume
243+
* that fails never has a runner to connect and say so.
244+
*/
245+
async awaitRestore(runnerId: string): Promise<RunnerRestoreResult> {
246+
return awaitRestoreOf(() => this.getRunner(runnerId), {
247+
pollMs: this.suspendPollMs,
248+
timeoutMs: this.restoreTimeoutMs,
249+
onReadError: (err) =>
250+
this.logger.warn("[RunCrdWorkloadManager] Runner read failed during restore", {
251+
runnerId,
252+
rawError: err,
253+
}),
254+
});
255+
}
256+
232257
private getRunner(name: string): Promise<unknown> {
233258
return this.k8s.custom.getNamespacedCustomObject({
234259
group: GROUP,
@@ -255,15 +280,20 @@ export class RunCrdWorkloadManager implements WorkloadManager, RunnerSnapshotter
255280
* Creates a resume: a Runner the operator restores from the checkpoint's
256281
* snapshot, on the node holding it, instead of cold-starting. Named from the
257282
* checkpoint, so a redelivered restore finds the first Runner and leaves it
258-
* rather than restoring twice.
283+
* rather than restoring twice. Returns the Runner's name, also for a resume
284+
* found in the way, which a restarted supervisor has no other watch on.
259285
*/
260-
async restore(opts: WorkloadManagerCreateOptions, checkpoint: { id: string; location: string }) {
286+
async restore(
287+
opts: WorkloadManagerCreateOptions,
288+
checkpoint: { id: string; location: string }
289+
): Promise<string> {
261290
const runnerId = getRestoreRunnerId(opts.runFriendlyId, checkpoint.id);
262291
const restore = parseCheckpointLocation(checkpoint.location);
263292
const created = await this.createRunner(opts, runnerId, restore);
264293
if (!created) {
265294
await this.checkExistingRestore(runnerId, restore);
266295
}
296+
return runnerId;
267297
}
268298

269299
/**
@@ -528,6 +558,61 @@ export function suspendOutcome(runner: unknown, request: string): RunnerSuspendR
528558
return undefined;
529559
}
530560

561+
/** Polls a resume's Runner, read by `readRunner`, until it starts, fails or times out. */
562+
export async function awaitRestoreOf(
563+
readRunner: () => Promise<unknown>,
564+
opts: { pollMs: number; timeoutMs: number; onReadError: (err: unknown) => void }
565+
): Promise<RunnerRestoreResult> {
566+
const deadline = Date.now() + opts.timeoutMs;
567+
while (Date.now() < deadline) {
568+
await sleep(opts.pollMs);
569+
let runner: unknown;
570+
try {
571+
runner = await readRunner();
572+
} catch (err: unknown) {
573+
const code = statusCodeOf(err);
574+
if (code === 404) {
575+
return { ok: false, error: "the Runner no longer exists" };
576+
}
577+
if (code !== undefined && TERMINAL_READ_CODES.has(code)) {
578+
return { ok: false, error: `Runner read failed: ${messageOf(err)}` };
579+
}
580+
opts.onReadError(err);
581+
continue;
582+
}
583+
const outcome = restoreOutcome(runner);
584+
if (outcome) {
585+
return outcome;
586+
}
587+
}
588+
return { ok: false, error: `the Runner did not start within ${opts.timeoutMs}ms` };
589+
}
590+
591+
/**
592+
* Whether a resume started, or undefined while the operator is still restoring
593+
* it. A Runner that ended is a success only if it ran to completion.
594+
*/
595+
export function restoreOutcome(runner: unknown): RunnerRestoreResult | undefined {
596+
const status = (runner as RunnerSuspendStatus | null)?.status;
597+
switch (status?.phase) {
598+
case "Running":
599+
case "Suspending":
600+
case "Succeeded":
601+
return { ok: true };
602+
case "Failed": {
603+
const condition = status.conditions?.find((c) => c.type === "Failed");
604+
return {
605+
ok: false,
606+
error: condition
607+
? `${condition.reason}: ${condition.message}`
608+
: "the Runner failed with no reason recorded",
609+
};
610+
}
611+
default:
612+
return undefined;
613+
}
614+
}
615+
531616
/**
532617
* The Runner's latest suspend, when it has an outcome the platform has not
533618
* been marked as accepting. A request made without the run annotation, by an

0 commit comments

Comments
 (0)