mirror of
https://github.com/tiennm99/ccs.git
synced 2026-10-06 08:13:49 +00:00
fix: bound bar raw socket probes (#1618)
This commit is contained in:
1 parent
2f1453e397
commit
396e01ab6f
3 files changed
+178
-10
No files matched your search
@@ -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<Dashboar
|
||||
return new Promise((resolve) => {
|
||||
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<Dashboar
|
||||
`GET ${parsed.pathname}${parsed.search} HTTP/1.1\r\nHost: ${parsed.host}\r\nConnection: close\r\n\r\n`
|
||||
);
|
||||
});
|
||||
socket.setTimeout(1500, () => 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]);
|
||||
|
||||
@@ -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<void> {
|
||||
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<void> {
|
||||
);
|
||||
}
|
||||
);
|
||||
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<void> {
|
||||
await new Promise<void>((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 {
|
||||
|
||||
@@ -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<string, Array<(arg?: unknown) => void>> = {};
|
||||
let interval: ReturnType<typeof setInterval> | 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<string, Array<(arg?: unknown) => void>> = {};
|
||||
let interval: ReturnType<typeof setInterval> | 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<void>;
|
||||
};
|
||||
|
||||
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
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
Reference in new issue
Block a user