diff --git a/packages/stats/CHANGELOG.md b/packages/stats/CHANGELOG.md index ce319cddb..8a8a2009a 100644 --- a/packages/stats/CHANGELOG.md +++ b/packages/stats/CHANGELOG.md @@ -2,6 +2,10 @@ ## [Unreleased] +### Fixed + +- Reused a live stats dashboard on the requested port and reclaimed stale Bun, Node, or omp listeners instead of failing with `EADDRINUSE` ([#5970](https://github.com/can1357/oh-my-pi/issues/5970)). + ## [17.0.2] - 2026-07-17 ### Fixed diff --git a/packages/stats/src/port-conflict.ts b/packages/stats/src/port-conflict.ts new file mode 100644 index 000000000..db290a8cc --- /dev/null +++ b/packages/stats/src/port-conflict.ts @@ -0,0 +1,206 @@ +import type { Dirent } from "node:fs"; +import * as fs from "node:fs/promises"; +import * as path from "node:path"; +import { $which } from "@oh-my-pi/pi-utils"; +import { $ } from "bun"; + +const STATS_PROBE_TIMEOUT_MS = 500; +const PROCESS_EXIT_POLL_MS = 50; +const PROCESS_EXIT_POLLS = 10; +const RECLAIMABLE_IMAGES = new Set(["bun", "node", "omp"]); + +interface PortHolder { + pid: number; + image: string; +} + +async function probeStatsDashboard(port: number): Promise { + try { + const response = await fetch(`http://localhost:${port}/api/stats/models`, { + signal: AbortSignal.timeout(STATS_PROBE_TIMEOUT_MS), + }); + const isDashboard = response.status === 200; + await response.body?.cancel(); + return isDashboard; + } catch { + return false; + } +} + +async function findLinuxPortHolder(port: number): Promise { + const socketInodes = new Set(); + for (const tablePath of ["/proc/net/tcp", "/proc/net/tcp6"]) { + let table: string; + try { + table = await Bun.file(tablePath).text(); + } catch { + continue; + } + + for (const line of table.split("\n").slice(1)) { + const fields = line.trim().split(/\s+/); + const localAddress = fields[1]; + const state = fields[3]; + const inode = fields[9]; + if (!localAddress || state !== "0A" || !inode) continue; + const encodedPort = localAddress.slice(localAddress.lastIndexOf(":") + 1); + if (Number.parseInt(encodedPort, 16) === port) socketInodes.add(inode); + } + } + if (socketInodes.size === 0) return null; + + let processes: Dirent[]; + try { + processes = await fs.readdir("/proc", { withFileTypes: true }); + } catch { + return null; + } + + for (const entry of processes) { + if (!entry.isDirectory() || !/^\d+$/.test(entry.name)) continue; + const pid = Number.parseInt(entry.name, 10); + let descriptors: string[]; + try { + descriptors = await fs.readdir(`/proc/${pid}/fd`); + } catch { + continue; + } + + let ownsSocket = false; + for (const descriptor of descriptors) { + try { + const target = await fs.readlink(`/proc/${pid}/fd/${descriptor}`); + const match = /^socket:\[(\d+)]$/.exec(target); + if (match?.[1] && socketInodes.has(match[1])) { + ownsSocket = true; + break; + } + } catch {} + } + if (!ownsSocket) continue; + + try { + const executable = await fs.readlink(`/proc/${pid}/exe`); + return { pid, image: path.basename(executable) }; + } catch { + try { + const commandLine = await Bun.file(`/proc/${pid}/cmdline`).text(); + const executable = commandLine.split("\0", 1)[0]; + return { pid, image: executable ? path.basename(executable) : "unknown" }; + } catch { + return { pid, image: "unknown" }; + } + } + } + return null; +} + +async function findMacPortHolder(port: number): Promise { + const lsof = $which("lsof") ?? ((await Bun.file("/usr/sbin/lsof").exists()) ? "/usr/sbin/lsof" : null); + if (!lsof) return null; + + const selector = `-iTCP:${port}`; + const result = await $`${lsof} -nP ${selector} -sTCP:LISTEN -Fpc`.quiet().nothrow(); + if (result.exitCode !== 0) return null; + + let pid: number | null = null; + for (const line of result.text().split("\n")) { + if (line.startsWith("p")) { + const parsed = Number.parseInt(line.slice(1), 10); + pid = Number.isSafeInteger(parsed) ? parsed : null; + } else if (line.startsWith("c") && pid !== null) { + return { pid, image: line.slice(1) || "unknown" }; + } + } + return null; +} + +async function findWindowsPortHolder(port: number): Promise { + const netstat = $which("netstat"); + if (!netstat) return null; + + const result = await $`${netstat} -ano -p TCP`.quiet().nothrow(); + if (result.exitCode !== 0) return null; + + let pid: number | null = null; + for (const line of result.text().split("\n")) { + const fields = line.trim().split(/\s+/); + if (fields[0]?.toUpperCase() !== "TCP" || fields[3]?.toUpperCase() !== "LISTENING") continue; + const localAddress = fields[1]; + if (!localAddress || Number.parseInt(localAddress.slice(localAddress.lastIndexOf(":") + 1), 10) !== port) { + continue; + } + const parsed = Number.parseInt(fields[4] ?? "", 10); + if (Number.isSafeInteger(parsed)) { + pid = parsed; + break; + } + } + if (pid === null) return null; + + const tasklist = $which("tasklist"); + if (!tasklist) return { pid, image: "unknown" }; + const filter = `PID eq ${pid}`; + const task = await $`${tasklist} /FI ${filter} /FO CSV /NH`.quiet().nothrow(); + if (task.exitCode !== 0) return { pid, image: "unknown" }; + const imageMatch = /^"((?:[^"]|"")*)"/.exec(task.text().trim()); + return { pid, image: imageMatch?.[1]?.replaceAll('""', '"') || "unknown" }; +} + +async function findPortHolder(port: number): Promise { + if (process.platform === "linux") return findLinuxPortHolder(port); + if (process.platform === "darwin") return findMacPortHolder(port); + if (process.platform === "win32") return findWindowsPortHolder(port); + return null; +} + +async function terminatePortHolder(holder: PortHolder): Promise { + try { + process.kill(holder.pid, "SIGTERM"); + } catch (error) { + if (error instanceof Error && "code" in error && error.code === "ESRCH") return; + throw new Error(`Failed to stop ${holder.image} (PID ${holder.pid})`, { cause: error }); + } + + for (let attempt = 0; attempt < PROCESS_EXIT_POLLS; attempt++) { + await Bun.sleep(PROCESS_EXIT_POLL_MS); + try { + process.kill(holder.pid, 0); + } catch (error) { + if (error instanceof Error && "code" in error && error.code === "ESRCH") return; + throw new Error(`Failed to inspect ${holder.image} (PID ${holder.pid})`, { cause: error }); + } + } + + try { + process.kill(holder.pid, "SIGKILL"); + } catch (error) { + if (error instanceof Error && "code" in error && error.code === "ESRCH") return; + throw new Error(`Failed to kill ${holder.image} (PID ${holder.pid})`, { cause: error }); + } + await Bun.sleep(PROCESS_EXIT_POLL_MS); +} + +/** Reuse a live stats dashboard or reclaim the port from a stale omp runtime. */ +export async function recoverStatsPort(port: number): Promise<"retry" | "reuse"> { + if (await probeStatsDashboard(port)) return "reuse"; + + const holder = await findPortHolder(port); + if (!holder) { + throw new Error(`Port ${port} is in use, but the listening process could not be identified.`); + } + if (holder.pid === process.pid) { + throw new Error(`Port ${port} is held by the current process (${holder.image}, PID ${holder.pid}).`); + } + + const normalizedImage = holder.image + .toLowerCase() + .replace(/\.exe$/, "") + .replace(/ \(deleted\)$/, ""); + if (!RECLAIMABLE_IMAGES.has(normalizedImage)) { + throw new Error(`Port ${port} is in use by ${holder.image} (PID ${holder.pid}); refusing to stop it.`); + } + + await terminatePortHolder(holder); + return "retry"; +} diff --git a/packages/stats/src/server.ts b/packages/stats/src/server.ts index 607de3f88..6c77364cc 100644 --- a/packages/stats/src/server.ts +++ b/packages/stats/src/server.ts @@ -20,6 +20,7 @@ import { import { decodeEmbeddedClientArchive } from "./embedded-client"; import embeddedClientArchiveTxt from "./embedded-client.generated.txt"; import { getGainDashboardStats } from "./gain-aggregator"; +import { recoverStatsPort } from "./port-conflict"; const EMBEDDED_CLIENT_ARCHIVE = decodeEmbeddedClientArchive(embeddedClientArchiveTxt); @@ -293,12 +294,7 @@ async function handleStatic(requestPath: string): Promise { return new Response("Not Found", { status: 404 }); } -/** - * Start the HTTP server. - */ -export async function startServer(port = 3847): Promise<{ port: number; stop: () => void }> { - await ensureClientBuild(); - +function createDashboardServer(port: number) { const server = Bun.serve({ port, async fetch(req) { @@ -306,7 +302,7 @@ export async function startServer(port = 3847): Promise<{ port: number; stop: () const path = url.pathname; // CORS headers for local development - const corsHeaders = { + const corsHeaders: Record = { "Access-Control-Allow-Origin": "*", "Access-Control-Allow-Methods": "GET, POST, OPTIONS", "Access-Control-Allow-Headers": "Content-Type", @@ -327,8 +323,8 @@ export async function startServer(port = 3847): Promise<{ port: number; stop: () // Add CORS headers to all responses const headers = new Headers(response.headers); - for (const [key, value] of Object.entries(corsHeaders)) { - headers.set(key, value); + for (const key in corsHeaders) { + headers.set(key, corsHeaders[key]); } return new Response(response.body, { @@ -344,9 +340,39 @@ export async function startServer(port = 3847): Promise<{ port: number; stop: () } }, }); - - return { - port: server.port ?? port, - stop: () => server.stop(), - }; + return server; +} + +/** + * Start the HTTP server, reusing a live dashboard or reclaiming a stale omp listener. + */ +export async function startServer(port = 3847): Promise<{ port: number; stop: () => void }> { + await ensureClientBuild(); + + try { + const server = createDashboardServer(port); + return { + port: server.port ?? port, + stop: () => server.stop(), + }; + } catch (error) { + if (!(error instanceof Error && "code" in error && error.code === "EADDRINUSE")) throw error; + + const recovery = await recoverStatsPort(port); + if (recovery === "reuse") { + return { port, stop: () => {} }; + } + + try { + const server = createDashboardServer(port); + return { + port: server.port ?? port, + stop: () => server.stop(), + }; + } catch (retryError) { + throw new Error(`Failed to start stats dashboard on port ${port} after reclaiming it.`, { + cause: retryError, + }); + } + } } diff --git a/packages/stats/test/server-port-conflict.test.ts b/packages/stats/test/server-port-conflict.test.ts new file mode 100644 index 000000000..4abbd8408 --- /dev/null +++ b/packages/stats/test/server-port-conflict.test.ts @@ -0,0 +1,77 @@ +import { afterEach, describe, expect, it } from "bun:test"; +import type { Subprocess } from "bun"; +import { startServer } from "../src/server"; + +const holderProcesses: Array> = []; + +async function startBunHolder(status: number) { + const reservation = Bun.serve({ + hostname: "127.0.0.1", + port: 0, + fetch: () => new Response("reserved"), + }); + const port = reservation.port; + reservation.stop(true); + + const source = `Bun.serve({ hostname: "127.0.0.1", port: ${port}, fetch: () => new Response("holder", { status: ${status} }) }); process.stdout.write("ready"); await Promise.withResolvers().promise;`; + const child = Bun.spawn([process.execPath, "-e", source], { + stdin: "ignore", + stdout: "pipe", + stderr: "pipe", + }); + holderProcesses.push(child); + + const reader = child.stdout.getReader(); + const ready = await reader.read(); + reader.releaseLock(); + if (!ready.done && new TextDecoder().decode(ready.value) === "ready") { + return { child, port }; + } + + await child.exited; + const stderr = await new Response(child.stderr).text(); + throw new Error(`Holder failed to listen on port ${port}: ${stderr}`); +} + +afterEach(async () => { + for (const child of holderProcesses) { + child.kill(); + await child.exited; + } + holderProcesses.length = 0; +}); + +describe("startServer port conflicts", () => { + it("reuses a live stats dashboard without stopping it", async () => { + const existing = Bun.serve({ + hostname: "127.0.0.1", + port: 0, + fetch: request => + new URL(request.url).pathname === "/api/stats/models" ? Response.json([]) : new Response("dashboard"), + }); + + try { + const server = await startServer(existing.port); + expect(server.port).toBe(existing.port); + server.stop(); + + const response = await fetch(`http://127.0.0.1:${existing.port}/api/stats/models`); + expect(response.status).toBe(200); + await response.body?.cancel(); + } finally { + existing.stop(true); + } + }); + + it("reclaims an unresponsive Bun listener and starts the dashboard", async () => { + const holder = await startBunHolder(404); + const server = await startServer(holder.port); + + try { + expect(server.port).toBe(holder.port); + expect(await holder.child.exited).not.toBe(0); + } finally { + server.stop(); + } + }); +});