| 1 | import assert from "node:assert/strict"; |
| 2 | import type { ChildProcess } from "node:child_process"; |
| 3 | import { EventEmitter } from "node:events"; |
| 4 | import { PassThrough } from "node:stream"; |
| 5 | import { test } from "node:test"; |
| 6 | import type { ServiceState } from "../shared/ipc.js"; |
| 7 | import { validateHelloResult, type HelloResult } from "./handshake.js"; |
| 8 | import { RestartBudget } from "./restartBudget.js"; |
| 9 | import { ServiceSupervisor, ShutdownError, type ShutdownTiming } from "./service.js"; |
| 10 | |
| 11 | const silent = { info() {}, warn() {}, error() {} }; |
| 12 | const tick = async (times = 4) => { |
| 13 | for (let i = 0; i < times; i++) await new Promise((resolve) => setImmediate(resolve)); |
| 14 | }; |
| 15 | |
| 16 | class FakeChild extends EventEmitter { |
| 17 | stdin = new PassThrough(); |
| 18 | stdout = new PassThrough(); |
| 19 | stderr = new PassThrough(); |
| 20 | alive = true; |
| 21 | requests: Array<{ id: number; method: string; params: unknown }> = []; |
| 22 | status: Record<string, unknown> | null = null; |
| 23 | heldShutdown: number | null = null; |
| 24 | private buffered = ""; |
| 25 | |
| 26 | constructor(readonly generation: string, readonly behaviour: { helloError?: { code: number; message: string }; exitOnStdinEnd?: boolean; holdShutdown?: boolean; |
| 27 | shutdownResults?: Array<Record<string, unknown>>; },) { |
| 28 | super(); |
| 29 | this.stdin.on("data", (chunk: Buffer) => { |
| 30 | this.buffered += chunk.toString("utf8"); |
| 31 | let index: number; |
| 32 | while ((index = this.buffered.indexOf("\n")) >= 0) { |
| 33 | const line = this.buffered.slice(0, index); |
| 34 | this.buffered = this.buffered.slice(index + 1); |
| 35 | this.handle(JSON.parse(line) as { id?: number; method?: string; params?: unknown; },); |
| 36 | } |
| 37 | }); |
| 38 | this.stdin.on("end", () => { |
| 39 | if (this.behaviour.exitOnStdinEnd !== false) this.exit(0, null); |
| 40 | }); |
| 41 | } |
| 42 | |
| 43 | kill(): boolean { |
| 44 | this.exit(null, "SIGKILL"); |
| 45 | return true; |
| 46 | } |
| 47 | |
| 48 | exit(code: number | null, signal: NodeJS.Signals | null): void { |
| 49 | if (!this.alive) return; |
| 50 | this.alive = false; |
| 51 | this.emit("exit", code, signal); |
| 52 | } |
| 53 | |
| 54 | send(frame: Record<string, unknown>): void { |
| 55 | this.stdout.write(JSON.stringify({ jsonrpc: "2.0", ...frame }) + "\n"); |
| 56 | } |
| 57 | |
| 58 | event(name: string, generation = this.generation): void { |
| 59 | this.send({ method: "desktop/event", params: { seq: 1, generation, name, args: [{ ok: true }] }, }); |
| 60 | } |
| 61 | |
| 62 | replyShutdown(result: Record<string, unknown>): void { |
| 63 | if (this.heldShutdown === null) throw new Error("no shutdown request is held"); |
| 64 | this.send({ id: this.heldShutdown, result: { requestId: "", reason: "user_quit", retryable: false, ...result } }); |
| 65 | this.heldShutdown = null; |
| 66 | } |
| 67 | |
| 68 | private handle(frame: { id?: number; method?: string; params?: unknown }): void { |
| 69 | if (typeof frame.id !== "number" || typeof frame.method !== "string") return; |
| 70 | this.requests.push({ id: frame.id, method: frame.method, params: frame.params, }); |
| 71 | if (frame.method === "desktop/shutdown" && this.behaviour.holdShutdown) { |
| 72 | this.heldShutdown = frame.id; |
| 73 | return; |
| 74 | } |
| 75 | if (frame.method === "desktop/shutdownStatus" && this.status) { |
| 76 | const params = frame.params as { requestId?: string }; |
| 77 | this.send({ id: frame.id, result: { requestId: params.requestId ?? "", reason: "user_quit", retryable: false, ...this.status } }); |
| 78 | return; |
| 79 | } |
| 80 | if (frame.method === "desktop/hello") { |
| 81 | if (this.behaviour.helloError) { |
| 82 | this.send({ id: frame.id, error: this.behaviour.helloError }); |
| 83 | return; |
| 84 | } |
| 85 | this.send({ |
| 86 | id: frame.id, |
| 87 | result: { |
| 88 | protocolVersion: 11, |
| 89 | contractDigest: "sha256:abc", |
| 90 | service: { version: "dev", channel: "dev", commit: "dev", pid: 1 }, |
| 91 | runtimeGeneration: this.generation, |
| 92 | runId: `run-${this.generation}`, |
| 93 | incidentId: `incident-${this.generation}`, |
| 94 | diagnosticsEnabled: true, |
| 95 | resources: { origin: "http://127.0.0.1:1", token: "t" }, |
| 96 | window: { width: 1000, height: 700, minWidth: 760, minHeight: 480, frameless: false, zoomFactor: 1, }, |
| 97 | }, |
| 98 | }); |
| 99 | return; |
| 100 | } |
| 101 | if (frame.method === "desktop/invoke") { |
| 102 | const params = frame.params as { method: string; args: unknown[] }; |
| 103 | if (params.method === "Fail") this.send({ id: frame.id, error: { code: -32000, message: "workspace not found", data: { method: "Fail" }, }, }); |
| 104 | else this.send({ id: frame.id, result: { method: params.method, args: params.args }, }); |
| 105 | return; |
| 106 | } |
| 107 | if (frame.method === "desktop/shutdown" || frame.method === "desktop/shutdownStatus") { |
| 108 | const params = frame.params as { requestId?: string; reason?: string }; |
| 109 | const configured = |
| 110 | this.behaviour.shutdownResults?.shift(); |
| 111 | this.send({ id: frame.id, result: { |
| 112 | requestId: params.requestId ?? "", |
| 113 | reason: params.reason ?? "user_quit", |
| 114 | phase: "completed", |
| 115 | outcome: "success", |
| 116 | completed: true, |
| 117 | retryable: false, |
| 118 | updatedAt: new Date().toISOString(), |
| 119 | ...configured,}, }); |
| 120 | return; |
| 121 | } |
| 122 | this.send({ id: frame.id, result: { |
| 123 | } }); |
| 124 | } |
| 125 | } |
| 126 | |
| 127 | function harness(options: { children?: FakeChild[]; budget?: RestartBudget; onState?(state: ServiceState): void; shutdownTiming?: Partial<ShutdownTiming>; } = {},) { |
| 128 | const spawned: FakeChild[] = []; |
| 129 | const states: ServiceState[] = []; |
| 130 | const events: string[] = []; |
| 131 | const ready: Array<{ generation: string; restarted: boolean }> = []; |
| 132 | const failures: string[] = []; |
| 133 | let index = 0; |
| 134 | const supervisor = new ServiceSupervisor( |
| 135 | { |
| 136 | binary: "fake", |
| 137 | args: ["--host-rpc"], |
| 138 | env: {}, |
| 139 | onStderr: () => undefined, |
| 140 | log: silent, |
| 141 | budget: options.budget ?? new RestartBudget(), |
| 142 | exitGraceMs: 10, |
| 143 | shutdownTiming: options.shutdownTiming, |
| 144 | spawn: () => { |
| 145 | const child = options.children?.[index++] ?? new FakeChild(`g-${spawned.length + 1}`, {}); |
| 146 | spawned.push(child); |
| 147 | return child as unknown as ChildProcess; |
| 148 | }, |
| 149 | }, |
| 150 | { |
| 151 | hello: async (client) => validateHelloResult(await client.request("desktop/hello", {}, 1000)), |
| 152 | onRequest: async () => ({}), |
| 153 | onEvent: (frame) => events.push(`${frame.generation}:${frame.name}`), |
| 154 | onState: (state) => { states.push(state); options.onState?.(state); }, |
| 155 | onReady: (hello: HelloResult, restarted) => { ready.push({ generation: hello.runtimeGeneration, restarted }); }, |
| 156 | onFailed: (error) => failures.push(error instanceof Error ? error.message : String(error)), |
| 157 | }, |
| 158 | ); |
| 159 | return { supervisor, spawned, states, events, ready, failures }; |
| 160 | } |
| 161 | |
| 162 | test("start runs hello then desktop/start and exposes the generation", async () => { |
| 163 | const h = harness(); |
| 164 | const hello = await h.supervisor.start(); |
| 165 | assert.equal(hello.runtimeGeneration, "g-1"); |
| 166 | assert.deepEqual(h.spawned[0]?.requests.map((r) => r.method), ["desktop/hello", "desktop/start"],); |
| 167 | assert.equal(h.supervisor.ready, true); |
| 168 | assert.equal(h.supervisor.generation, "g-1"); |
| 169 | assert.deepEqual(h.states.map((s) => s.phase), ["starting", "ready"],); |
| 170 | assert.deepEqual(h.ready, [{ generation: "g-1", restarted: false }]); |
| 171 | assert.deepEqual(await h.supervisor.invoke("OpenProjectTab", ["/p", true]), { method: "OpenProjectTab", args: ["/p", true], }); |
| 172 | await assert.rejects(h.supervisor.invoke("Fail", []), /workspace not found/); |
| 173 | }); |
| 174 | |
| 175 | test("events from the live generation are forwarded and stale ones dropped", async () => { |
| 176 | const h = harness(); |
| 177 | await h.supervisor.start(); |
| 178 | h.spawned[0]?.event("agent:event"); |
| 179 | h.spawned[0]?.event("agent:event", "g-old"); |
| 180 | h.spawned[0]?.event("duplicate"); |
| 181 | h.spawned[0]?.send({ method: "desktop/event", params: { seq: 3, generation: "g-1", name: "after-gap", args: [] }, }); |
| 182 | h.spawned[0]?.send({ method: "desktop/event", params: { seq: 2, generation: "g-1", name: "late", args: [] }, }); |
| 183 | await tick(); |
| 184 | assert.deepEqual(h.events, ["g-1:agent:event", "g-1:after-gap"]); |
| 185 | }); |
| 186 | |
| 187 | test("a handshake error fails the service without a restart and terminates the process", async () => { |
| 188 | const child = new FakeChild("g-1", { helloError: { code: -32003, message: "digest differs" }, }); |
| 189 | const h = harness({ children: [child] }); |
| 190 | await assert.rejects(h.supervisor.start(), /digest differs/); |
| 191 | await tick(); |
| 192 | assert.equal(h.supervisor.current.phase, "failed"); |
| 193 | assert.deepEqual(h.failures, ["digest differs"]); |
| 194 | assert.equal(child.alive, false, "stdin close makes the fake exit"); |
| 195 | assert.equal(h.spawned.length, 1); |
| 196 | }); |
| 197 | |
| 198 | test("an unexpected exit restarts automatically until the budget is exhausted", async () => { |
| 199 | const budget = new RestartBudget(2, 60_000); |
| 200 | const h = harness({ budget }); |
| 201 | await h.supervisor.start(); |
| 202 | h.spawned[0]?.exit(1, null); |
| 203 | await tick(8); |
| 204 | assert.equal(h.spawned.length, 2); |
| 205 | assert.equal(h.supervisor.generation, "g-2"); |
| 206 | assert.deepEqual(h.ready.map((r) => r.restarted), [false, true],); |
| 207 | h.spawned[1]?.exit(1, null); |
| 208 | await tick(8); |
| 209 | assert.equal(h.spawned.length, 3); |
| 210 | h.spawned[2]?.exit(1, null); |
| 211 | await tick(8); |
| 212 | assert.equal(h.spawned.length, 3, "no fourth spawn once the budget is spent"); |
| 213 | assert.equal(h.supervisor.current.phase, "failed"); |
| 214 | assert.match(h.supervisor.current.error ?? "", /automatic restarts exhausted/); |
| 215 | await assert.rejects(h.supervisor.invoke("X", []), /not running/); |
| 216 | await h.supervisor.restart(); |
| 217 | assert.equal(h.spawned.length, 4, "a manual restart is always allowed"); |
| 218 | assert.equal(h.supervisor.generation, "g-4"); |
| 219 | assert.deepEqual(h.states.map((s) => s.phase), ["starting", "ready", "restarting", "ready", "restarting", "ready", "failed", "restarting", "ready"],); |
| 220 | }); |
| 221 | |
| 222 | test("shutdown sends desktop/shutdown, closes stdin and waits for the exit", async () => { |
| 223 | const h = harness(); |
| 224 | await h.supervisor.start(); |
| 225 | const shutdown = h.supervisor.shutdown(); |
| 226 | assert.equal(h.supervisor.current.phase, "stopping"); |
| 227 | assert.equal(h.supervisor.ready, false); |
| 228 | await assert.rejects(h.supervisor.invoke("CloseMainWindow", []), /shutting down/); |
| 229 | await shutdown; |
| 230 | const child = h.spawned[0] as FakeChild; |
| 231 | assert.deepEqual(child.requests.map((r) => r.method), ["desktop/hello", "desktop/start", "desktop/shutdown"],); |
| 232 | assert.equal(child.alive, false); |
| 233 | assert.equal(h.supervisor.current.phase, "exited"); |
| 234 | assert.equal(h.spawned.length, 1, "a deliberate exit never restarts"); |
| 235 | assert.deepEqual(h.states.map((state) => state.phase), ["starting", "ready", "stopping", "exited", "exited"]); |
| 236 | }); |
| 237 | |
| 238 | test("a service that ignores stdin close is killed after the grace period", async () => { |
| 239 | const child = new FakeChild("g-1", { exitOnStdinEnd: false }); |
| 240 | const h = harness({ children: [child] }); |
| 241 | await h.supervisor.start(); |
| 242 | await h.supervisor.shutdown(); |
| 243 | assert.equal(child.alive, false); |
| 244 | assert.equal(h.supervisor.current.phase, "exited"); |
| 245 | }); |
| 246 | |
| 247 | test("a retryable save failure keeps the service alive and retry uses the same request identity", async () => { |
| 248 | const child = new FakeChild("g-1", { |
| 249 | shutdownResults: [ |
| 250 | { |
| 251 | phase: "saving", |
| 252 | outcome: "failed", |
| 253 | completed: false, |
| 254 | retryable: true, |
| 255 | errorCode: "session_save_failed", |
| 256 | error: "disk full", |
| 257 | }, |
| 258 | { |
| 259 | phase: "completed", |
| 260 | outcome: "success", |
| 261 | completed: true, |
| 262 | retryable: false, |
| 263 | }, |
| 264 | ], |
| 265 | }); |
| 266 | const h = harness({ children: [child] }); |
| 267 | await h.supervisor.start(); |
| 268 | await assert.rejects(h.supervisor.shutdown(), /session_save_failed/); |
| 269 | assert.equal(child.alive, true, "failed durable save must not close stdin or kill the service"); |
| 270 | await h.supervisor.shutdown(); |
| 271 | const shutdowns = child.requests.filter((request) => request.method === "desktop/shutdown"); |
| 272 | assert.equal(shutdowns.length, 2); |
| 273 | assert.equal((shutdowns[0].params as { requestId: string }).requestId, (shutdowns[1].params as { requestId: string }).requestId); |
| 274 | assert.equal(child.alive, false); |
| 275 | }); |
| 276 | |
| 277 | test("shutdown status polling waits through in-progress cleanup", async () => { |
| 278 | const child = new FakeChild("g-1", { |
| 279 | shutdownResults: [ |
| 280 | { |
| 281 | phase: "saving", |
| 282 | outcome: "in_progress", |
| 283 | completed: false, |
| 284 | retryable: false, |
| 285 | }, |
| 286 | { |
| 287 | phase: "closing", |
| 288 | outcome: "in_progress", |
| 289 | completed: false, |
| 290 | retryable: false, |
| 291 | }, |
| 292 | { |
| 293 | phase: "completed", |
| 294 | outcome: "success", |
| 295 | completed: true, |
| 296 | retryable: false, |
| 297 | }, |
| 298 | ], |
| 299 | }); |
| 300 | const h = harness({ children: [child] }); |
| 301 | const phases: string[] = []; |
| 302 | await h.supervisor.start(); |
| 303 | await h.supervisor.shutdown("user_quit", (phase) => phases.push(phase)); |
| 304 | assert.deepEqual(phases, ["saving", "closing", "completed"]); |
| 305 | assert.deepEqual( |
| 306 | child.requests.slice(2).map((request) => request.method), |
| 307 | ["desktop/shutdown", "desktop/shutdownStatus", "desktop/shutdownStatus"], |
| 308 | ); |
| 309 | }); |
| 310 | |
| 311 | const inProgress = (phase: string) => ({ phase, outcome: "in_progress", completed: false }); |
| 312 | const completed = { phase: "completed", outcome: "success", completed: true }; |
| 313 | const fastShutdown: ShutdownTiming = { replyMs: 40, pollMs: 5, stallMs: 150, ceilingMs: 2_000 }; |
| 314 | const wait = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); |
| 315 | |
| 316 | test("a shutdown reply that outlives the reply window is still a success when it arrives", async () => { |
| 317 | const child = new FakeChild("g-1", { holdShutdown: true }); |
| 318 | child.status = inProgress("saving"); |
| 319 | const h = harness({ children: [child], shutdownTiming: fastShutdown }); |
| 320 | await h.supervisor.start(); |
| 321 | const phases: string[] = []; |
| 322 | const shutdown = h.supervisor.shutdown("user_quit", (phase) => phases.push(phase)); |
| 323 | await wait(60); |
| 324 | child.status = inProgress("closing"); |
| 325 | await wait(80); |
| 326 | child.replyShutdown(completed); |
| 327 | child.exit(0, null); |
| 328 | await shutdown; |
| 329 | assert.equal(h.supervisor.current.phase, "exited"); |
| 330 | assert.equal(h.failures.length, 0); |
| 331 | assert.ok(phases.includes("closing"), `phases: ${phases.join(",")}`); |
| 332 | assert.equal(phases.at(-1), "completed"); |
| 333 | }); |
| 334 | |
| 335 | test("a late completed status ends a shutdown whose reply never arrives", async () => { |
| 336 | const child = new FakeChild("g-1", { holdShutdown: true }); |
| 337 | child.status = inProgress("saving"); |
| 338 | const h = harness({ children: [child], shutdownTiming: fastShutdown }); |
| 339 | await h.supervisor.start(); |
| 340 | const shutdown = h.supervisor.shutdown(); |
| 341 | await wait(70); |
| 342 | child.status = inProgress("closing"); |
| 343 | await wait(70); |
| 344 | child.status = completed; |
| 345 | await shutdown; |
| 346 | assert.equal(child.alive, false); |
| 347 | assert.equal(h.supervisor.current.phase, "exited"); |
| 348 | }); |
| 349 | |
| 350 | test("a shutdown that stops advancing is reported incomplete and the service is kept", async () => { |
| 351 | const child = new FakeChild("g-1", { holdShutdown: true }); |
| 352 | child.status = inProgress("closing"); |
| 353 | const h = harness({ children: [child], shutdownTiming: fastShutdown }); |
| 354 | await h.supervisor.start(); |
| 355 | const started = Date.now(); |
| 356 | await assert.rejects(h.supervisor.shutdown(), (error: unknown) => { |
| 357 | assert.ok(error instanceof ShutdownError); |
| 358 | assert.equal(error.code, "shutdown_incomplete"); |
| 359 | assert.match(error.message, /closing/); |
| 360 | return true; |
| 361 | }); |
| 362 | assert.ok(Date.now() - started < fastShutdown.ceilingMs, "a stall is reported before the ceiling"); |
| 363 | assert.equal(child.alive, true, "an incomplete shutdown leaves the service for a retry"); |
| 364 | }); |
| 365 | |
| 366 | test("a service that exits without reporting success fails the shutdown", async () => { |
| 367 | const child = new FakeChild("g-1", { holdShutdown: true }); |
| 368 | child.status = inProgress("saving"); |
| 369 | const h = harness({ children: [child], shutdownTiming: fastShutdown }); |
| 370 | await h.supervisor.start(); |
| 371 | const shutdown = h.supervisor.shutdown(); |
| 372 | await wait(60); |
| 373 | child.exit(1, null); |
| 374 | await assert.rejects(shutdown, (error: unknown) => error instanceof ShutdownError && error.code === "service_exited"); |
| 375 | }); |
| 376 | |
| 377 | test("a shutdown that keeps advancing still ends at the ceiling", async () => { |
| 378 | const child = new FakeChild("g-1", { holdShutdown: true }); |
| 379 | const phases = ["saving", "closing"]; |
| 380 | let turn = 0; |
| 381 | const flip = setInterval(() => { child.status = inProgress(phases[turn++ % 2]); }, 10); |
| 382 | child.status = inProgress("saving"); |
| 383 | const h = harness({ children: [child], shutdownTiming: { ...fastShutdown, ceilingMs: 300 } }); |
| 384 | await h.supervisor.start(); |
| 385 | try { |
| 386 | await assert.rejects(h.supervisor.shutdown(), (error: unknown) => error instanceof ShutdownError && error.code === "shutdown_incomplete"); |
| 387 | } finally { |
| 388 | clearInterval(flip); |
| 389 | } |
| 390 | }); |
| 391 | |
| 392 | test("shutdown fences a hello completion queued before shutdown", async () => { |
| 393 | const h = harness(); |
| 394 | const starting = h.supervisor.start(); |
| 395 | const rejected = assert.rejects(starting, /cancelled|exited/); |
| 396 | await h.supervisor.shutdown(); |
| 397 | await rejected; |
| 398 | assert.deepEqual(h.ready, []); |
| 399 | assert.deepEqual(h.failures, []); |
| 400 | assert.equal(h.spawned.length, 1); |
| 401 | await assert.rejects(h.supervisor.start(), /shutting down/); |
| 402 | }); |
| 403 | |
| 404 | test("a destroyed renderer throwing during exit notification cannot prevent shutdown", async () => { |
| 405 | const h = harness({ onState: (state) => { if (state.phase === "exited") throw new Error("Object has been destroyed"); }, }); |
| 406 | await h.supervisor.start(); |
| 407 | await h.supervisor.shutdown(); |
| 408 | assert.equal(h.supervisor.current.phase, "exited"); |
| 409 | assert.equal(h.spawned[0].alive, false); |
| 410 | }); |
| 411 | |
| 412 | test("concurrent restarts share one replacement and shutdown prevents its spawn", async () => { |
| 413 | const h = harness(); |
| 414 | await h.supervisor.start(); |
| 415 | const a = h.supervisor.restart(); |
| 416 | const b = h.supervisor.restart(); |
| 417 | const rejected = Promise.all([assert.rejects(a, /shutting down/), assert.rejects(b, /shutting down/)]); |
| 418 | await h.supervisor.shutdown(); |
| 419 | await rejected; |
| 420 | assert.equal(h.spawned.length, 1); |
| 421 | }); |
| 422 |