From 396e01ab6fbaddf74949a1777dc34a0ebbdc9677 Mon Sep 17 00:00:00 2001 From: "Kai (Tam Nhu) Tran" <61256810+kaitranntt@users.noreply.github.com> Date: Tue, 30 Jun 2026 12:53:35 -0400 Subject: [PATCH] fix: bound bar raw socket probes (#1618) --- src/commands/bar/bar-server-probe.ts | 12 +- src/commands/bar/launch-subcommand.ts | 21 +++- tests/unit/commands/bar-command.test.ts | 155 ++++++++++++++++++++++-- 3 files changed, 178 insertions(+), 10 deletions(-) diff --git a/src/commands/bar/bar-server-probe.ts b/src/commands/bar/bar-server-probe.ts index 2d489d9d..c2a4cf42 100644 --- a/src/commands/bar/bar-server-probe.ts +++ b/src/commands/bar/bar-server-probe.ts @@ -10,6 +10,9 @@ import * as fs from 'fs'; import * as path from 'path'; import { BAR_AUTH_TOKEN_HEADER, getOrCreateBarAuthToken } from '../../utils/bar-auth-token'; +const PROBE_TIMEOUT_MS = 1500; +const MAX_PROBE_RESPONSE_BYTES = 8192; + export interface DashboardInfo { port: number; baseUrl: string; @@ -65,9 +68,12 @@ export async function defaultFindRunningServer(ccsDir: string): Promise { let rawResponse = ''; let settled = false; + const absoluteDeadline = setTimeout(() => finish(), PROBE_TIMEOUT_MS); + absoluteDeadline.unref?.(); const finish = (statusCode = 0, headerSection = '') => { if (settled) return; settled = true; + clearTimeout(absoluteDeadline); // Tear down the socket the moment we have enough to decide. The summary // endpoint only needs the status code for liveness, so a non-CCS // loopback service that streams forever cannot block discovery from @@ -98,9 +104,13 @@ export async function defaultFindRunningServer(ccsDir: string): Promise finish()); + socket.setTimeout(PROBE_TIMEOUT_MS, () => finish()); socket.on('data', (chunk) => { rawResponse += chunk.toString('utf8'); + if (rawResponse.length > MAX_PROBE_RESPONSE_BYTES) { + finish(); + return; + } const statusMatch = rawResponse.match(/^HTTP\/\d(?:\.\d)?\s+(\d{3})/); if (statusMatch) { const code = Number(statusMatch[1]); diff --git a/src/commands/bar/launch-subcommand.ts b/src/commands/bar/launch-subcommand.ts index 4b22a182..7a0e3248 100644 --- a/src/commands/bar/launch-subcommand.ts +++ b/src/commands/bar/launch-subcommand.ts @@ -31,6 +31,9 @@ import { } from './bar-server-probe'; import type { DashboardInfo as _DashboardInfo } from './bar-server-probe'; +const BAR_PROBE_TIMEOUT_MS = 1500; +const MAX_BAR_PROBE_RESPONSE_BYTES = 8192; + // --------------------------------------------------------------------------- // Re-exports — backward compat for tests that import from this module. // resolveBarPort + defaultFindRunningServer are canonical in bar-server-probe.ts; @@ -129,6 +132,13 @@ export class BarServerAuthRequiredError extends Error { } } +export class BarServerTimeoutError extends Error { + constructor(baseUrl: string, timeoutSeconds: number) { + super(`CCS Bar server did not become live at ${baseUrl} within ${timeoutSeconds}s`); + this.name = 'BarServerTimeoutError'; + } +} + function isAuthRequiredStatus(statusCode: number): boolean { return statusCode === 401 || statusCode === 403; } @@ -145,9 +155,12 @@ export async function defaultWaitForServerLive(baseUrl: string): Promise { return new Promise((resolve) => { let rawResponse = ''; let settled = false; + const absoluteDeadline = setTimeout(() => finish(), BAR_PROBE_TIMEOUT_MS); + absoluteDeadline.unref?.(); const finish = (statusCode: number | null = null, headerSection = '') => { if (settled) return; settled = true; + clearTimeout(absoluteDeadline); socket.destroy(); if (statusCode === 200) { const echoMatch = headerSection.match( @@ -169,9 +182,13 @@ export async function defaultWaitForServerLive(baseUrl: string): Promise { ); } ); - socket.setTimeout(1500, () => finish()); + socket.setTimeout(BAR_PROBE_TIMEOUT_MS, () => finish()); socket.on('data', (chunk) => { rawResponse += chunk.toString('utf8'); + if (rawResponse.length > MAX_BAR_PROBE_RESPONSE_BYTES) { + finish(); + return; + } const statusMatch = rawResponse.match(/^HTTP\/\d(?:\.\d)?\s+(\d{3})/); if (statusMatch) { const code = Number(statusMatch[1]); @@ -204,7 +221,7 @@ export async function defaultWaitForServerLive(baseUrl: string): Promise { await new Promise((resolve) => setTimeout(resolve, INTERVAL_MS)); } - throw new Error(`CCS Bar server did not become live at ${baseUrl} within ${TIMEOUT_MS / 1000}s`); + throw new BarServerTimeoutError(baseUrl, TIMEOUT_MS / 1000); } function defaultWriteLaunchDescriptor(jsonPath: string, descriptor: LaunchJson): void { diff --git a/tests/unit/commands/bar-command.test.ts b/tests/unit/commands/bar-command.test.ts index f94de371..fbc5e2cc 100644 --- a/tests/unit/commands/bar-command.test.ts +++ b/tests/unit/commands/bar-command.test.ts @@ -3162,12 +3162,7 @@ describe('defaultFindRunningServer: socket-level 401/403 classifies authRequired setImmediate(() => { onConnect(); for (const cb of listeners.data ?? []) { - cb( - Buffer.from( - `HTTP/1.1 200 OK\r\nx-ccs-bar-token: ${token}\r\n\r\n`, - 'utf8' - ) - ); + cb(Buffer.from(`HTTP/1.1 200 OK\r\nx-ccs-bar-token: ${token}\r\n\r\n`, 'utf8')); } }); return socket; @@ -3387,7 +3382,9 @@ describe('defaultFindRunningServer: streaming lower-priority probes', () => { // exactly what the production CCS Bar server does, and is the property // that prevents a rogue reflector from passing the check. for (const cb of data) { - cb(Buffer.from(`HTTP/1.1 200 OK\r\nx-ccs-bar-token: ${expectedToken}\r\n\r\n`, 'utf8')); + cb( + Buffer.from(`HTTP/1.1 200 OK\r\nx-ccs-bar-token: ${expectedToken}\r\n\r\n`, 'utf8') + ); } } // Port 3000 never emits a status line: simulate an endlessly @@ -3424,6 +3421,150 @@ describe('defaultFindRunningServer: streaming lower-priority probes', () => { }); }); +describe('bar raw socket probes: absolute deadline for malformed streaming peers', () => { + it('continues past a higher-priority trickling non-HTTP response and returns a lower-priority hit', async () => { + const ccsDir = path.join(tempHome, '.ccs'); + fs.mkdirSync(ccsDir, { recursive: true }); + fs.writeFileSync( + path.join(ccsDir, 'bar.json'), + JSON.stringify({ port: 41236, baseUrl: 'http://127.0.0.1:41236', authMode: 'loopback' }) + ); + + const { getOrCreateBarAuthToken: getToken } = await import( + `../../../src/utils/bar-auth-token?test=${Date.now()}-deadline-find` + ); + const expectedToken = getToken(ccsDir); + + mock.module('net', () => ({ + connect: (opts: { host: string; port: number }, onConnect: () => void): unknown => { + const listeners: Record void>> = {}; + let interval: ReturnType | undefined; + const socket = { + on(event: string, cb: (arg?: unknown) => void) { + (listeners[event] ??= []).push(cb); + return socket; + }, + setTimeout() { + return socket; + }, + write() { + return true; + }, + destroy() { + if (interval) clearInterval(interval); + return socket; + }, + }; + + setImmediate(() => { + onConnect(); + if (opts.port === 41236) { + interval = setInterval(() => { + for (const cb of listeners.data ?? []) cb(Buffer.from('x', 'utf8')); + }, 25); + interval.unref?.(); + return; + } + if (opts.port === 3000) { + for (const cb of listeners.data ?? []) { + cb( + Buffer.from(`HTTP/1.1 200 OK\r\nx-ccs-bar-token: ${expectedToken}\r\n\r\n`, 'utf8') + ); + } + } + }); + + return socket; + }, + })); + + moduleSeq++; + const { defaultFindRunningServer } = (await import( + `../../../src/commands/bar/bar-server-probe?test=${Date.now()}-${moduleSeq}` + )) as { + defaultFindRunningServer: ( + ccsDir: string + ) => Promise<{ port: number; baseUrl: string; authRequired?: boolean } | null>; + }; + + const result = await Promise.race([ + defaultFindRunningServer(ccsDir), + new Promise<'timeout'>((resolve) => setTimeout(() => resolve('timeout'), 2500)), + ]); + + expect(result).toEqual({ + port: 3000, + baseUrl: 'http://127.0.0.1:3000', + authRequired: false, + }); + }); + + it('does not let a single launch wait probe outlive the outer deadline', async () => { + const ccsDir = path.join(tempHome, '.ccs'); + fs.mkdirSync(ccsDir, { recursive: true }); + + const realDateNow = Date.now; + let dateCall = 0; + Date.now = () => { + dateCall++; + return dateCall < 4 ? 0 : 10_001; + }; + + mock.module('net', () => ({ + connect: (_opts: { host: string; port: number }, onConnect: () => void): unknown => { + const listeners: Record void>> = {}; + let interval: ReturnType | undefined; + const socket = { + on(event: string, cb: (arg?: unknown) => void) { + (listeners[event] ??= []).push(cb); + return socket; + }, + setTimeout() { + return socket; + }, + write() { + return true; + }, + destroy() { + if (interval) clearInterval(interval); + return socket; + }, + }; + + setImmediate(() => { + onConnect(); + interval = setInterval(() => { + for (const cb of listeners.data ?? []) cb(Buffer.from('x', 'utf8')); + }, 25); + interval.unref?.(); + }); + + return socket; + }, + })); + + try { + moduleSeq++; + const { defaultWaitForServerLive } = (await import( + `../../../src/commands/bar/launch-subcommand?test=${Date.now()}-${moduleSeq}` + )) as { + defaultWaitForServerLive: (baseUrl: string) => Promise; + }; + + const result = await Promise.race([ + defaultWaitForServerLive('http://127.0.0.1:9996') + .then(() => 'resolved' as const) + .catch(() => 'rejected' as const), + new Promise<'timeout'>((resolve) => setTimeout(() => resolve('timeout'), 2500)), + ]); + + expect(result).toBe('rejected'); + } finally { + Date.now = realDateNow; + } + }); +}); + // --------------------------------------------------------------------------- // GH-1588 — `--await-quit` waits for a running app to exit, then swaps + relaunches // ---------------------------------------------------------------------------