Skip to content

Commit e39f445

Browse files
1stvampTrigger.dev RepoOps
authored andcommitted
feat(supervisor): stop gracefully on SIGTERM, send the snapshot route with restore reports, keep the failed-pod informer alive
The supervisor now shuts down gracefully on SIGTERM, and its Kubernetes pod watcher recovers from API errors instead of stopping. Mono-RevId: a87068c3ab16ad9d6955fcdfd69e33c8251e9e75
1 parent af0453e commit e39f445

18 files changed

Lines changed: 1197 additions & 278 deletions

‎apps/supervisor/package.json‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@
2929
},
3030
"devDependencies": {
3131
"@internal/testcontainers": "workspace:*",
32-
"@types/dockerode": "^3.3.33"
32+
"@types/dockerode": "^3.3.33",
33+
"socket.io-client": "4.7.5"
3334
}
3435
}

‎apps/supervisor/src/clients/kubernetes.ts‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import * as k8s from "@kubernetes/client-node";
2-
import type { Informer, KubernetesObject, ListPromise, ObjectCache } from "@kubernetes/client-node";
2+
import type { KubernetesObject, ListPromise } from "@kubernetes/client-node";
33
import { assertExhaustive } from "@trigger.dev/core/utils";
44
import { SimpleStructuredLogger } from "@trigger.dev/core/v3/utils/structuredLogger";
55

@@ -15,15 +15,15 @@ export function createK8sApi() {
1515
listPromiseFn: ListPromise<T>,
1616
labelSelector?: string,
1717
fieldSelector?: string
18-
): Informer<T> & ObjectCache<T> {
18+
): k8s.ListWatch<T> {
1919
// The client's informer is a ListWatch, which is also the cache it keeps.
2020
return k8s.makeInformer(
2121
kubeConfig,
2222
path,
2323
listPromiseFn,
2424
labelSelector,
2525
fieldSelector
26-
) as Informer<T> & ObjectCache<T>;
26+
) as k8s.ListWatch<T>;
2727
}
2828

2929
const api = {
Lines changed: 151 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,151 @@
1+
import { setTimeout as sleep } from "node:timers/promises";
2+
import type { KubernetesObject, ListPromise, ListWatch } from "@kubernetes/client-node";
3+
import type { SimpleStructuredLogger } from "@trigger.dev/core/v3/utils/structuredLogger";
4+
5+
const MAX_RECONNECT_BACKOFF_MS = 30_000;
6+
7+
type ReconnectingInformerOptions<T extends KubernetesObject> = {
8+
/** Used in log fields, to tell informers apart. */
9+
name: string;
10+
logger: SimpleStructuredLogger;
11+
reconnectIntervalMs: number;
12+
list: ListPromise<T>;
13+
/** Builds the informer around the list it is given, which wraps `list`. */
14+
makeInformer: (list: ListPromise<T>) => ListWatch<T>;
15+
/** Called for every error the informer raises, before any reconnect. */
16+
onError?: (err: unknown) => void;
17+
};
18+
19+
/**
20+
* An informer that reconnects with a capped backoff until a start succeeds, for
21+
* as long as it runs. The client's own error handling gives up after one failure.
22+
*/
23+
export class ReconnectingInformer<T extends KubernetesObject> {
24+
readonly informer: ListWatch<T>;
25+
private readonly name: string;
26+
private readonly logger: SimpleStructuredLogger;
27+
private readonly reconnectIntervalMs: number;
28+
private readonly listFn: ListPromise<T>;
29+
private readonly onErrorHook?: (err: unknown) => void;
30+
private running = false;
31+
private reconnecting = false;
32+
private erroredDuringReconnect = false;
33+
private starting = false;
34+
private ownListPending = false;
35+
36+
constructor(opts: ReconnectingInformerOptions<T>) {
37+
this.name = opts.name;
38+
this.logger = opts.logger;
39+
this.reconnectIntervalMs = opts.reconnectIntervalMs;
40+
this.listFn = opts.list;
41+
this.onErrorHook = opts.onError;
42+
this.informer = opts.makeInformer(() => this.list());
43+
this.informer.on("error", (err?: unknown) => void this.onError(err));
44+
}
45+
46+
get isRunning(): boolean {
47+
return this.running;
48+
}
49+
50+
/** Rejects when the first list fails, so the caller learns the informer never started. */
51+
async start() {
52+
if (this.running) {
53+
return;
54+
}
55+
this.running = true;
56+
await this.startInformer();
57+
}
58+
59+
async stop() {
60+
if (!this.running) {
61+
return;
62+
}
63+
this.running = false;
64+
await this.informer.stop();
65+
}
66+
67+
/**
68+
* The client relists on its own after a 410 and leaves that list's rejection
69+
* unhandled, so a failure there becomes a reconnect and the abandoned relist
70+
* never settles. An empty list instead would delete every cached object. Only
71+
* a start's own list rejects into the start: a 410 on the watch that start
72+
* opens relists inside it too, and nothing awaits that one.
73+
*/
74+
private async list() {
75+
const ownList = this.starting && this.ownListPending;
76+
this.ownListPending = false;
77+
try {
78+
return await this.listFn();
79+
} catch (err: unknown) {
80+
if (ownList) {
81+
throw err;
82+
}
83+
void this.onError(err);
84+
return new Promise<never>(() => {});
85+
}
86+
}
87+
88+
private async startInformer() {
89+
this.starting = true;
90+
// A start lists first only when it has no resourceVersion to watch from.
91+
this.ownListPending = !this.informer.latestResourceVersion();
92+
try {
93+
await this.informer.start();
94+
} finally {
95+
this.starting = false;
96+
}
97+
// A stop during the list still lets the client open its watch afterwards.
98+
if (!this.running) {
99+
await this.informer.stop();
100+
}
101+
}
102+
103+
/**
104+
* Retries until a start ends with no error raised during it. A watch that
105+
* fails to connect raises its error inside the start, which still resolves.
106+
*/
107+
private async onError(err: unknown) {
108+
if (!this.running) {
109+
return;
110+
}
111+
this.onErrorHook?.(err);
112+
if (this.reconnecting) {
113+
this.erroredDuringReconnect = true;
114+
return;
115+
}
116+
this.reconnecting = true;
117+
this.logger.error("Informer watch failed, reconnecting", {
118+
informer: this.name,
119+
error: messageOf(err),
120+
});
121+
let delayMs = this.reconnectIntervalMs;
122+
try {
123+
do {
124+
await sleep(delayMs);
125+
if (!this.running) {
126+
return;
127+
}
128+
this.erroredDuringReconnect = false;
129+
try {
130+
await this.startInformer();
131+
} catch (reconnectErr: unknown) {
132+
this.erroredDuringReconnect = true;
133+
this.logger.error("Informer reconnect failed", {
134+
informer: this.name,
135+
error: messageOf(reconnectErr),
136+
});
137+
}
138+
delayMs = Math.min(
139+
delayMs * 2,
140+
Math.max(this.reconnectIntervalMs, MAX_RECONNECT_BACKOFF_MS)
141+
);
142+
} while (this.running && this.erroredDuringReconnect);
143+
} finally {
144+
this.reconnecting = false;
145+
}
146+
}
147+
}
148+
149+
function messageOf(err: unknown): string {
150+
return err instanceof Error ? err.message : String(err);
151+
}

‎apps/supervisor/src/env.ts‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,10 @@ export const Env = z
1919
// Opt-in, dev-only: stream this process's logs over a local telnet/TCP socket on this port.
2020
SUPERVISOR_TELNET_LOGS_PORT: z.coerce.number().optional(),
2121

22+
// How long a SIGTERM or SIGINT waits for a clean stop before exiting anyway. Keep it below
23+
// the pod's terminationGracePeriodSeconds, so the exit is this process's and not a SIGKILL.
24+
SUPERVISOR_SHUTDOWN_TIMEOUT_MS: z.coerce.number().int().positive().default(8_000),
25+
2226
// Required settings
2327
TRIGGER_API_URL: z.string().url(),
2428
TRIGGER_WORKER_TOKEN: z.string().min(1), // accepts file:// path to read from a file

0 commit comments

Comments
 (0)