import net from "node:net"; import { INITIAL_RETRY_MS, MAX_RETRY_MS, SPAWN_COOLDOWN_MS, } from "./constants.ts"; /** Minimal lease handle: an open connection is the acquire. */ export type LeaseSocket = { destroy(): void; on(event: "close", listener: () => void): void; on(event: "error", listener: (err: Error) => void): void; off?(event: "close", listener: () => void): void; off?(event: "error", listener: (err: Error) => void): void; }; export type SpawnDaemonArgs = { runtimeDir: string; scriptPath: string; onError?: (err: Error) => void; }; export type CapsBlinkClientDeps = { socketPath: string; runtimeDir: string; scriptPath: string; connect: (socketPath: string) => Promise; spawnDaemon: (args: SpawnDaemonArgs) => void; setTimeout: (fn: () => void, ms: number) => unknown; clearTimeout: (id: unknown) => void; random: () => number; now: () => number; warn: (message: string) => void; initialRetryMs?: number; maxRetryMs?: number; spawnCooldownMs?: number; }; /** * Default Unix-domain connect. Connection success is the lease; no protocol. */ export function defaultConnect(socketPath: string): Promise { return new Promise((resolve, reject) => { const socket = net.createConnection({ path: socketPath }); const onError = (err: Error) => { socket.off("connect", onConnect); socket.destroy(); reject(err); }; const onConnect = () => { socket.off("error", onError); // Prevent unhandled error events after lease is held; close drives retry. socket.on("error", () => {}); resolve(socket); }; socket.once("error", onError); socket.once("connect", onConnect); }); } /** * Explicit desiredBusy / socket / retry / generation state machine. * Handlers should call {@link CapsBlinkClient.setDesiredBusy} without awaiting. */ export class CapsBlinkClient { private desiredBusy = false; private socket: LeaseSocket | null = null; /** Generation of the in-flight connect, or null if none for the current lease cycle. */ private connectingGeneration: number | null = null; private generation = 0; private retryTimer: unknown = null; private retryAttempt = 0; private lastSpawnAt = Number.NEGATIVE_INFINITY; private busySince = 0; private warned = false; private readonly socketPath: string; private readonly runtimeDir: string; private readonly scriptPath: string; private readonly connectFn: CapsBlinkClientDeps["connect"]; private readonly spawnDaemonFn: CapsBlinkClientDeps["spawnDaemon"]; private readonly setTimeoutFn: CapsBlinkClientDeps["setTimeout"]; private readonly clearTimeoutFn: CapsBlinkClientDeps["clearTimeout"]; private readonly randomFn: CapsBlinkClientDeps["random"]; private readonly nowFn: CapsBlinkClientDeps["now"]; private readonly warnFn: CapsBlinkClientDeps["warn"]; private readonly initialRetryMs: number; private readonly maxRetryMs: number; private readonly spawnCooldownMs: number; constructor(deps: CapsBlinkClientDeps) { this.socketPath = deps.socketPath; this.runtimeDir = deps.runtimeDir; this.scriptPath = deps.scriptPath; this.connectFn = deps.connect; this.spawnDaemonFn = deps.spawnDaemon; this.setTimeoutFn = deps.setTimeout; this.clearTimeoutFn = deps.clearTimeout; this.randomFn = deps.random; this.nowFn = deps.now; this.warnFn = deps.warn; this.initialRetryMs = deps.initialRetryMs ?? INITIAL_RETRY_MS; this.maxRetryMs = deps.maxRetryMs ?? MAX_RETRY_MS; this.spawnCooldownMs = deps.spawnCooldownMs ?? SPAWN_COOLDOWN_MS; } /** Test/observation helpers. */ getDesiredBusy(): boolean { return this.desiredBusy; } getGeneration(): number { return this.generation; } hasSocket(): boolean { return this.socket !== null; } isConnecting(): boolean { return this.connectingGeneration !== null; } hasRetryTimer(): boolean { return this.retryTimer !== null; } /** * Synchronously update desired busy state and fire-and-forget reconcile. * Never awaits daemon startup. */ setDesiredBusy(desired: boolean): void { if (desired) { if ( this.desiredBusy && (this.socket || this.connectingGeneration === this.generation || this.retryTimer) ) { return; } if (!this.desiredBusy) { this.busySince = this.nowFn(); this.warned = false; } this.desiredBusy = true; this.reconcile(); return; } if ( !this.desiredBusy && !this.socket && this.connectingGeneration === null && !this.retryTimer ) { return; } this.desiredBusy = false; this.generation += 1; this.clearRetryTimer(); this.retryAttempt = 0; this.destroySocket(); // In-flight connect completions are generation-gated and will destroy. // connectingGeneration stays until that attempt settles (may be stale). } /** Idempotent shutdown helper for session_shutdown. */ dispose(): void { this.setDesiredBusy(false); } private reconcile(): void { if (!this.desiredBusy) { return; } if (this.socket) { return; } // Only block on a connect for *this* generation; a stale pending connect // must not prevent re-acquire after busy→false→true. if (this.connectingGeneration === this.generation) { return; } this.clearRetryTimer(); this.beginConnect(this.generation); } private beginConnect(generation: number): void { if (!this.desiredBusy || generation !== this.generation) { return; } if (this.socket || this.connectingGeneration === generation) { return; } this.connectingGeneration = generation; void this.connectFn(this.socketPath).then( (socket) => this.onConnectSuccess(generation, socket), (err: unknown) => this.onConnectFailure(generation, err), ); } private clearConnectingIf(generation: number): void { if (this.connectingGeneration === generation) { this.connectingGeneration = null; } } private onConnectSuccess(generation: number, socket: LeaseSocket): void { this.clearConnectingIf(generation); if (!this.desiredBusy || generation !== this.generation) { socket.destroy(); // Stale completion after re-acquire: ensure a current connect is running. this.reconcile(); return; } this.socket = socket; this.retryAttempt = 0; this.clearRetryTimer(); const onClose = () => { socket.off?.("close", onClose); this.onSocketClose(generation, socket); }; socket.on("close", onClose); // Swallow post-lease errors; close drives reconnection. socket.on("error", () => {}); } private onConnectFailure(generation: number, err: unknown): void { this.clearConnectingIf(generation); if (!this.desiredBusy || generation !== this.generation) { this.reconcile(); return; } // A missing socket is expected on first use: the daemon is on-demand. // Warn only if startup has remained unavailable for a full spawn cooldown. if (this.nowFn() - this.busySince >= this.spawnCooldownMs) { this.warnOnce( `pi-caps-blink: unavailable (${formatError(err)}); still retrying`, ); } this.maybeSpawnDaemon(); this.scheduleRetry(generation); } private onSocketClose(generation: number, socket: LeaseSocket): void { if (this.socket === socket) { this.socket = null; } // Retry daemon death only from close, only while same generation remains desired. if (!this.desiredBusy || generation !== this.generation) { return; } this.warnOnce("pi-caps-blink: lease socket closed; reconnecting"); this.maybeSpawnDaemon(); this.scheduleRetry(generation); } private scheduleRetry(generation: number): void { if (!this.desiredBusy || generation !== this.generation) { return; } if ( this.retryTimer !== null || this.connectingGeneration === generation || this.socket ) { return; } const base = Math.min( this.maxRetryMs, this.initialRetryMs * 2 ** this.retryAttempt, ); this.retryAttempt += 1; const jitter = this.randomFn() * base; const delay = Math.min(this.maxRetryMs, base * 0.5 + jitter * 0.5); this.retryTimer = this.setTimeoutFn(() => { this.retryTimer = null; if (!this.desiredBusy || generation !== this.generation) { return; } this.beginConnect(generation); }, delay); } private maybeSpawnDaemon(): void { const now = this.nowFn(); if (now - this.lastSpawnAt < this.spawnCooldownMs) { return; } this.lastSpawnAt = now; try { this.spawnDaemonFn({ runtimeDir: this.runtimeDir, scriptPath: this.scriptPath, onError: (err) => { this.warnOnce( `pi-caps-blink: daemon spawn error (${formatError(err)})`, ); }, }); } catch (err) { this.warnOnce(`pi-caps-blink: failed to spawn daemon (${formatError(err)})`); } } private destroySocket(): void { const socket = this.socket; this.socket = null; if (socket) { try { socket.destroy(); } catch { // ignore } } } private clearRetryTimer(): void { if (this.retryTimer !== null) { this.clearTimeoutFn(this.retryTimer); this.retryTimer = null; } } private warnOnce(message: string): void { if (this.warned) { return; } this.warned = true; this.warnFn(message); } } function formatError(err: unknown): string { if (err instanceof Error) { return err.message; } return String(err); }