| 1 | import { spawn as nodeSpawn, type ChildProcess, type SpawnOptions } from "node:child_process"; |
| 2 | import { randomUUID } from "node:crypto"; |
| 3 | import type { EventFrame, ServiceState } from "../shared/ipc.js"; |
| 4 | import { eventFrame } from "../shared/eventStream.js"; |
| 5 | import type { HelloResult } from "./handshake.js"; |
| 6 | import { errorText, type Logger } from "./log.js"; |
| 7 | import { RestartBudget } from "./restartBudget.js"; |
| 8 | import { RpcClient } from "./rpc.js"; |
| 9 | |
| 10 | export const LIFECYCLE_TIMEOUT_MS = 10_000; |
| 11 | export const EXIT_GRACE_MS = 5_000; |
| 12 | |
| 13 | export interface ShutdownTiming { |
| 14 | replyMs: number; // silence on desktop/shutdown before its status is followed |
| 15 | pollMs: number; |
| 16 | stallMs: number; // longest a live shutdown may keep one phase and outcome |
| 17 | ceilingMs: number; // hard bound on the whole wait, however it advances |
| 18 | } |
| 19 | |
| 20 | export const SHUTDOWN_TIMING: ShutdownTiming = { |
| 21 | replyMs: LIFECYCLE_TIMEOUT_MS, |
| 22 | pollMs: 250, |
| 23 | stallMs: 30_000, |
| 24 | ceilingMs: 60_000, |
| 25 | }; |
| 26 | |
| 27 | // code is the failure's identity; the message is only its display form. |
| 28 | export class ShutdownError extends Error { |
| 29 | constructor( |
| 30 | readonly code: string, |
| 31 | detail: string, |
| 32 | ) { |
| 33 | super(`${code}: ${detail}`); |
| 34 | this.name = "ShutdownError"; |
| 35 | } |
| 36 | } |
| 37 | |
| 38 | export interface ServiceHandlers { |
| 39 | hello(client: RpcClient): Promise<HelloResult>; |
| 40 | onRequest(method: string, params: unknown): Promise<unknown>; |
| 41 | onEvent(frame: EventFrame): void; |
| 42 | onState(state: ServiceState): void; |
| 43 | onReady(hello: HelloResult, restarted: boolean): void | Promise<void>; |
| 44 | onFailed(error: unknown): void; |
| 45 | } |
| 46 | |
| 47 | export type SpawnFn = (command: string, args: string[], options: SpawnOptions) => ChildProcess; |
| 48 | |
| 49 | export interface ServiceOptions { |
| 50 | binary: string; |
| 51 | args: string[]; |
| 52 | env: NodeJS.ProcessEnv; |
| 53 | onStderr(chunk: Buffer): void; |
| 54 | log: Logger; |
| 55 | spawn?: SpawnFn; |
| 56 | budget?: RestartBudget; |
| 57 | now?: () => number; |
| 58 | exitGraceMs?: number; |
| 59 | shutdownTiming?: Partial<ShutdownTiming>; |
| 60 | } |
| 61 | |
| 62 | interface Session { |
| 63 | child: ChildProcess; |
| 64 | client: RpcClient; |
| 65 | generation: string; |
| 66 | alive: boolean; |
| 67 | ready: boolean; |
| 68 | eventSeq: number; |
| 69 | expectExit: boolean; |
| 70 | exited: Promise<void>; |
| 71 | } |
| 72 | |
| 73 | type ShutdownResult = { |
| 74 | requestId: string; |
| 75 | reason: string; |
| 76 | phase: "idle" | "preparing" | "saving" | "closing" | "completed"; |
| 77 | outcome: "not_started" | "in_progress" | "success" | "failed"; |
| 78 | completed: boolean; |
| 79 | retryable: boolean; |
| 80 | errorCode?: string; |
| 81 | error?: string; |
| 82 | }; |
| 83 | |
| 84 | export type ShutdownPhase = ShutdownResult["phase"]; |
| 85 | |
| 86 | function shutdownResult(value: unknown): ShutdownResult { |
| 87 | if (!value || typeof value !== "object") throw new Error("desktop shutdown returned an invalid result"); |
| 88 | const result = value as Partial<ShutdownResult>; |
| 89 | if ( |
| 90 | typeof result.requestId !== "string" || |
| 91 | typeof result.phase !== "string" || |
| 92 | typeof result.outcome !== "string" || |
| 93 | typeof result.completed !== "boolean" || |
| 94 | typeof result.retryable !== "boolean" |
| 95 | ) { |
| 96 | throw new Error("desktop shutdown returned an invalid result"); |
| 97 | } |
| 98 | return result as ShutdownResult; |
| 99 | } |
| 100 | |
| 101 | function describeExit(code: number | null, signal: NodeJS.Signals | null): string { |
| 102 | return signal ? `signal ${signal}` : `code ${code ?? "unknown"}`; |
| 103 | } |
| 104 | |
| 105 | function delay(ms: number): Promise<void> { |
| 106 | return new Promise((resolve) => setTimeout(resolve, ms)); |
| 107 | } |
| 108 | |
| 109 | async function within<T>(promise: Promise<T>, ms: number): Promise<T | null> { |
| 110 | let timer: NodeJS.Timeout | undefined; |
| 111 | try { |
| 112 | return await Promise.race([promise, new Promise<null>((resolve) => (timer = setTimeout(() => resolve(null), ms)))]); |
| 113 | } finally { |
| 114 | clearTimeout(timer); |
| 115 | } |
| 116 | } |
| 117 | |
| 118 | export class ServiceSupervisor { |
| 119 | private session: Session | null = null; |
| 120 | private state: ServiceState = { phase: "starting", generation: "" }; |
| 121 | private hello: HelloResult | null = null; |
| 122 | private launching: Promise<HelloResult> | null = null; |
| 123 | private stopping = false; |
| 124 | private shutdownPending: Promise<void> | null = null; |
| 125 | private shutdownRequestId = ""; |
| 126 | private revision = 0; |
| 127 | private restarting: Promise<HelloResult> | null = null; |
| 128 | private readonly budget: RestartBudget; |
| 129 | private readonly spawnFn: SpawnFn; |
| 130 | private readonly now: () => number; |
| 131 | private readonly exitGraceMs: number; |
| 132 | private readonly shutdownTiming: ShutdownTiming; |
| 133 | |
| 134 | constructor( |
| 135 | private readonly options: ServiceOptions, |
| 136 | private readonly handlers: ServiceHandlers, |
| 137 | ) { |
| 138 | this.budget = options.budget ?? new RestartBudget(); |
| 139 | this.spawnFn = options.spawn ?? ((command, args, spawnOptions) => nodeSpawn(command, args, spawnOptions)); |
| 140 | this.now = options.now ?? (() => Date.now()); |
| 141 | this.exitGraceMs = options.exitGraceMs ?? EXIT_GRACE_MS; |
| 142 | this.shutdownTiming = { ...SHUTDOWN_TIMING, ...options.shutdownTiming }; |
| 143 | } |
| 144 | |
| 145 | get current(): ServiceState { |
| 146 | return this.state; |
| 147 | } |
| 148 | |
| 149 | get helloResult(): HelloResult | null { |
| 150 | return this.hello; |
| 151 | } |
| 152 | |
| 153 | get generation(): string { |
| 154 | return this.session?.alive && this.session.ready ? this.session.generation : ""; |
| 155 | } |
| 156 | |
| 157 | get ready(): boolean { |
| 158 | return !this.stopping && this.state.phase === "ready" && this.session?.alive === true; |
| 159 | } |
| 160 | |
| 161 | get shutdownRequestIdentity(): string { |
| 162 | return this.shutdownRequestId; |
| 163 | } |
| 164 | |
| 165 | start(): Promise<HelloResult> { |
| 166 | return this.begin(false); |
| 167 | } |
| 168 | |
| 169 | async restart(): Promise<HelloResult> { |
| 170 | if (this.stopping) throw new Error("desktop service is shutting down"); |
| 171 | if (this.launching) return this.launching; |
| 172 | if (!this.restarting) { |
| 173 | this.restarting = (async () => { |
| 174 | const old = this.session; |
| 175 | if (old?.alive) await this.terminate(old); |
| 176 | return this.begin(true); |
| 177 | })().finally(() => { |
| 178 | this.restarting = null; |
| 179 | }); |
| 180 | } |
| 181 | return this.restarting; |
| 182 | } |
| 183 | |
| 184 | async invoke(method: string, args: unknown[]): Promise<unknown> { |
| 185 | return this.live().client.request("desktop/invoke", { method, args }); |
| 186 | } |
| 187 | |
| 188 | async request(method: string, params: unknown, timeoutMs = LIFECYCLE_TIMEOUT_MS): Promise<unknown> { |
| 189 | return this.live().client.request(method, params, timeoutMs); |
| 190 | } |
| 191 | |
| 192 | async hostEvent(name: string, payload: unknown): Promise<void> { |
| 193 | try { |
| 194 | await this.request("desktop/hostEvent", { name, payload }); |
| 195 | } catch (error) { |
| 196 | this.options.log.warn(`hostEvent ${name} not delivered: ${errorText(error)}`); |
| 197 | } |
| 198 | } |
| 199 | |
| 200 | shutdown( |
| 201 | reason: "user_quit" | "update_restart" | "system_signal" = "user_quit", |
| 202 | onProgress?: (phase: ShutdownPhase) => void, |
| 203 | ): Promise<void> { |
| 204 | if (!this.stopping) { |
| 205 | this.stopping = true; |
| 206 | this.revision++; |
| 207 | this.setState({ phase: "stopping", generation: this.session?.generation ?? "" }); |
| 208 | } |
| 209 | if (!this.shutdownPending) { |
| 210 | this.shutdownPending = this.finishShutdown(reason, onProgress) |
| 211 | .catch((error) => { |
| 212 | if (this.state.phase !== "exited") { |
| 213 | this.setState({ phase: "stopping", generation: this.session?.generation ?? "", error: errorText(error) }); |
| 214 | } |
| 215 | throw error; |
| 216 | }) |
| 217 | .finally(() => { |
| 218 | this.shutdownPending = null; |
| 219 | }); |
| 220 | } |
| 221 | return this.shutdownPending; |
| 222 | } |
| 223 | |
| 224 | private async finishShutdown( |
| 225 | reason: "user_quit" | "update_restart" | "system_signal", |
| 226 | onProgress?: (phase: ShutdownPhase) => void, |
| 227 | ): Promise<void> { |
| 228 | // A restart may already own termination of the current generation. Let that |
| 229 | // operation settle before selecting the session to shut down, otherwise the |
| 230 | // shutdown RPC races a closing stdin and turns an intentional quit into an |
| 231 | // indeterminate transport error. |
| 232 | if (this.restarting) await this.restarting.catch(() => undefined); |
| 233 | const session = this.session; |
| 234 | if (!session?.alive) { |
| 235 | this.setState({ phase: "exited", generation: "" }); |
| 236 | return; |
| 237 | } |
| 238 | session.expectExit = true; |
| 239 | if (!this.shutdownRequestId) this.shutdownRequestId = randomUUID(); |
| 240 | const result = await this.awaitShutdown(session, reason, onProgress); |
| 241 | if (result.outcome === "failed") { |
| 242 | throw new ShutdownError(result.errorCode ?? "shutdown_failed", result.error ?? "desktop shutdown failed"); |
| 243 | } |
| 244 | if (!result.completed || result.outcome !== "success") { |
| 245 | throw new ShutdownError(result.errorCode ?? "shutdown_incomplete", result.error ?? `desktop shutdown stopped in ${result.phase}`); |
| 246 | } |
| 247 | // The service closes itself only after the completed result has reached the |
| 248 | // shell. stdin/kill are now a post-completion process-exit fallback. |
| 249 | await this.terminate(session); |
| 250 | this.setState({ phase: "exited", generation: "" }); |
| 251 | } |
| 252 | |
| 253 | // The shutdown reply stays awaited for the whole wait: a completed result that |
| 254 | // arrives after replyMs is the service's verdict, not an orphan. Status polls |
| 255 | // only decide whether the shutdown is still live and advancing meanwhile. |
| 256 | private async awaitShutdown( |
| 257 | session: Session, |
| 258 | reason: "user_quit" | "update_restart" | "system_signal", |
| 259 | onProgress?: (phase: ShutdownPhase) => void, |
| 260 | ): Promise<ShutdownResult> { |
| 261 | const timing = this.shutdownTiming; |
| 262 | const requestId = this.shutdownRequestId; |
| 263 | const startedAt = this.now(); |
| 264 | const reply = session.client |
| 265 | .request("desktop/shutdown", { requestId, reason }, timing.ceilingMs) |
| 266 | .then((value) => ({ result: shutdownResult(value) })) |
| 267 | .catch((error: unknown) => ({ error })); |
| 268 | const first = await within(reply, timing.replyMs); |
| 269 | const unanswered = new Promise<never>(() => undefined); |
| 270 | const answer = first ? unanswered : reply.then((settled) => ("result" in settled ? settled.result : unanswered)); |
| 271 | let result = null as ShutdownResult | null; |
| 272 | let advancedAt = startedAt; |
| 273 | const observe = (next: ShutdownResult) => { |
| 274 | if (!result || next.phase !== result.phase || next.outcome !== result.outcome) { |
| 275 | advancedAt = this.now(); |
| 276 | onProgress?.(next.phase); |
| 277 | } |
| 278 | result = next; |
| 279 | }; |
| 280 | if (first && "result" in first) observe(first.result); |
| 281 | else { |
| 282 | const why = first ? errorText(first.error) : `no reply after ${timing.replyMs} ms`; |
| 283 | this.options.log.warn(`desktop/shutdown result unknown: ${why}; following status`); |
| 284 | } |
| 285 | for (;;) { |
| 286 | if (result && (result.completed || result.outcome !== "in_progress")) return result; |
| 287 | const phase = result?.phase ?? "unknown"; |
| 288 | if (this.now() - startedAt >= timing.ceilingMs) { |
| 289 | throw new ShutdownError("shutdown_incomplete", `desktop shutdown did not complete within ${timing.ceilingMs} ms; last phase ${phase}`); |
| 290 | } |
| 291 | if (this.now() - advancedAt >= timing.stallMs) { |
| 292 | throw new ShutdownError("shutdown_incomplete", `desktop shutdown stopped in ${phase}`); |
| 293 | } |
| 294 | const late = await within(answer, timing.pollMs); |
| 295 | if (late) { |
| 296 | observe(late); |
| 297 | continue; |
| 298 | } |
| 299 | if (!session.alive) { |
| 300 | throw new ShutdownError("service_exited", `desktop service exited during shutdown in phase ${phase}`); |
| 301 | } |
| 302 | try { |
| 303 | observe(shutdownResult(await session.client.request("desktop/shutdownStatus", { requestId }, timing.replyMs))); |
| 304 | } catch (error) { |
| 305 | if (session.alive) this.options.log.warn(`desktop/shutdownStatus unanswered: ${errorText(error)}`); |
| 306 | } |
| 307 | } |
| 308 | } |
| 309 | |
| 310 | private live(): Session { |
| 311 | const session = this.session; |
| 312 | if (this.stopping) throw new Error("desktop service is shutting down"); |
| 313 | if (!session?.alive || !session.ready) throw new Error(`desktop service is not running (${this.state.phase})`); |
| 314 | return session; |
| 315 | } |
| 316 | |
| 317 | private begin(restarted: boolean): Promise<HelloResult> { |
| 318 | if (this.stopping) return Promise.reject(new Error("desktop service is shutting down")); |
| 319 | if (this.launching) return this.launching; |
| 320 | this.launching = this.launch(restarted).finally(() => { |
| 321 | this.launching = null; |
| 322 | }); |
| 323 | return this.launching; |
| 324 | } |
| 325 | |
| 326 | private async launch(restarted: boolean): Promise<HelloResult> { |
| 327 | const revision = ++this.revision; |
| 328 | this.setState({ phase: restarted ? "restarting" : "starting", generation: "" }); |
| 329 | let session: Session | null = null; |
| 330 | try { |
| 331 | session = this.spawnSession(); |
| 332 | const hello = await this.handlers.hello(session.client); |
| 333 | if (this.stopping || revision !== this.revision) throw new Error("desktop startup cancelled"); |
| 334 | session.generation = hello.runtimeGeneration; |
| 335 | await session.client.request("desktop/start", {}, LIFECYCLE_TIMEOUT_MS); |
| 336 | if (!session.alive || this.stopping || revision !== this.revision) throw new Error("desktop service exited during startup"); |
| 337 | session.ready = true; |
| 338 | this.hello = hello; |
| 339 | this.setState({ phase: "ready", generation: hello.runtimeGeneration }); |
| 340 | await this.handlers.onReady(hello, restarted); |
| 341 | return hello; |
| 342 | } catch (error) { |
| 343 | if ((!session || this.session === session) && !this.stopping && revision === this.revision) { |
| 344 | this.setState({ phase: "failed", generation: "", error: errorText(error) }); |
| 345 | this.handlers.onFailed(error); |
| 346 | if (session) void this.terminate(session).catch((error) => this.options.log.error(`service termination failed: ${errorText(error)}`)); |
| 347 | } |
| 348 | throw error; |
| 349 | } |
| 350 | } |
| 351 | |
| 352 | private spawnSession(): Session { |
| 353 | const { binary, args, env, log } = this.options; |
| 354 | const child = this.spawnFn(binary, args, { |
| 355 | stdio: ["pipe", "pipe", "pipe"], |
| 356 | windowsHide: true, |
| 357 | env, |
| 358 | }); |
| 359 | const session: Session = { |
| 360 | child, |
| 361 | client: null as unknown as RpcClient, |
| 362 | generation: "", |
| 363 | alive: true, |
| 364 | ready: false, |
| 365 | eventSeq: 0, |
| 366 | expectExit: false, |
| 367 | exited: Promise.resolve(), |
| 368 | }; |
| 369 | session.client = new RpcClient( |
| 370 | { |
| 371 | write: (line) => { |
| 372 | if (!child.stdin || child.stdin.destroyed) throw new Error("desktop service stdin is closed"); |
| 373 | child.stdin.write(line); |
| 374 | }, |
| 375 | }, |
| 376 | { |
| 377 | onRequest: (method, params) => this.handlers.onRequest(method, params), |
| 378 | onNotification: (method, params) => this.onNotification(session, method, params), |
| 379 | onProtocolError: (kind, line) => log.warn(`service stdout ${kind}: ${line.slice(0, 200)}`), |
| 380 | }, |
| 381 | ); |
| 382 | session.exited = new Promise<void>((resolve) => { |
| 383 | const finish = (error: Error, code: number | null, signal: NodeJS.Signals | null) => { |
| 384 | if (!session.alive) return; |
| 385 | session.alive = false; |
| 386 | session.client.close(error); |
| 387 | resolve(); |
| 388 | this.onExit(session, code, signal); |
| 389 | }; |
| 390 | child.once("exit", (code, signal) => finish(new Error(`desktop service exited (${describeExit(code, signal)})`), code, signal)); |
| 391 | child.once("error", (error) => finish(new Error(`desktop service failed to start: ${errorText(error)}`), null, null)); |
| 392 | }); |
| 393 | child.stdout?.on("data", (chunk: Buffer) => { |
| 394 | try { |
| 395 | session.client.feed(chunk); |
| 396 | } catch (error) { |
| 397 | log.error(`service stream unusable: ${errorText(error)}`); |
| 398 | child.kill(); |
| 399 | } |
| 400 | }); |
| 401 | child.stderr?.on("data", (chunk: Buffer) => this.options.onStderr(chunk)); |
| 402 | child.stdin?.on("error", (error) => log.warn(`service stdin: ${errorText(error)}`)); |
| 403 | this.session = session; |
| 404 | return session; |
| 405 | } |
| 406 | |
| 407 | private onNotification(session: Session, method: string, params: unknown): void { |
| 408 | if (method !== "desktop/event") { |
| 409 | this.options.log.warn(`unknown service notification ${method}`); |
| 410 | return; |
| 411 | } |
| 412 | const frame = eventFrame(params); |
| 413 | if (!frame) { |
| 414 | this.options.log.warn("malformed desktop/event frame dropped"); |
| 415 | return; |
| 416 | } |
| 417 | if (this.session !== session || !session.alive || frame.generation !== session.generation) { |
| 418 | this.options.log.warn(`event ${frame.name} from dead generation ${frame.generation} dropped`); |
| 419 | return; |
| 420 | } |
| 421 | if (frame.seq <= session.eventSeq) return; |
| 422 | if (frame.seq > session.eventSeq + 1) this.options.log.warn(`desktop event gap: ${session.eventSeq} -> ${frame.seq}`); |
| 423 | session.eventSeq = frame.seq; |
| 424 | this.handlers.onEvent(frame); |
| 425 | } |
| 426 | |
| 427 | private onExit(session: Session, code: number | null, signal: NodeJS.Signals | null): void { |
| 428 | if (this.session !== session || !session.ready) return; |
| 429 | const reason = describeExit(code, signal); |
| 430 | if (session.expectExit || this.stopping) { |
| 431 | this.setState({ phase: "exited", generation: "" }); |
| 432 | return; |
| 433 | } |
| 434 | this.options.log.error(`desktop service exited unexpectedly (${reason})`); |
| 435 | if (this.budget.allow(this.now())) { |
| 436 | void this.begin(true).catch(() => undefined); |
| 437 | return; |
| 438 | } |
| 439 | const error = new Error(`desktop service exited (${reason}); automatic restarts exhausted`); |
| 440 | this.setState({ phase: "failed", generation: "", error: error.message }); |
| 441 | this.handlers.onFailed(error); |
| 442 | } |
| 443 | |
| 444 | private async terminate(session: Session): Promise<void> { |
| 445 | session.expectExit = true; |
| 446 | try { |
| 447 | session.child.stdin?.end(); |
| 448 | } catch { |
| 449 | // Already closed. |
| 450 | } |
| 451 | if (!session.alive) return; |
| 452 | await Promise.race([session.exited, delay(this.exitGraceMs)]); |
| 453 | if (!session.alive) return; |
| 454 | this.options.log.warn("desktop service did not exit after stdin close; killing it"); |
| 455 | session.child.kill("SIGKILL"); |
| 456 | await Promise.race([session.exited, delay(1000)]); |
| 457 | if (session.alive) throw new Error("desktop service did not terminate; shell exit withheld"); |
| 458 | } |
| 459 | |
| 460 | private setState(state: ServiceState): void { |
| 461 | this.state = state; |
| 462 | // A destroyed renderer or failing observer must not turn confirmed service |
| 463 | // exit into a failed shutdown, nor skip the shell's remaining cleanup. |
| 464 | try { |
| 465 | this.handlers.onState(state); |
| 466 | } catch (error) { |
| 467 | this.options.log.warn(`service state observer failed: ${errorText(error)}`); |
| 468 | } |
| 469 | } |
| 470 | } |
| 471 |