mirror of
https://github.com/tiennm99/ccs.git
synced 2026-10-11 03:13:12 +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 * as path from 'path';
|
||||||
import { BAR_AUTH_TOKEN_HEADER, getOrCreateBarAuthToken } from '../../utils/bar-auth-token';
|
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 {
|
export interface DashboardInfo {
|
||||||
port: number;
|
port: number;
|
||||||
baseUrl: string;
|
baseUrl: string;
|
||||||
@@ -65,9 +68,12 @@ export async function defaultFindRunningServer(ccsDir: string): Promise<Dashboar
|
|||||||
return new Promise((resolve) => {
|
return new Promise((resolve) => {
|
||||||
let rawResponse = '';
|
let rawResponse = '';
|
||||||
let settled = false;
|
let settled = false;
|
||||||
|
const absoluteDeadline = setTimeout(() => finish(), PROBE_TIMEOUT_MS);
|
||||||
|
absoluteDeadline.unref?.();
|
||||||
const finish = (statusCode = 0, headerSection = '') => {
|
const finish = (statusCode = 0, headerSection = '') => {
|
||||||
if (settled) return;
|
if (settled) return;
|
||||||
settled = true;
|
settled = true;
|
||||||
|
clearTimeout(absoluteDeadline);
|
||||||
// Tear down the socket the moment we have enough to decide. The summary
|
// 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
|
// endpoint only needs the status code for liveness, so a non-CCS
|
||||||
// loopback service that streams forever cannot block discovery from
|
// 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`
|
`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) => {
|
socket.on('data', (chunk) => {
|
||||||
rawResponse += chunk.toString('utf8');
|
rawResponse += chunk.toString('utf8');
|
||||||
|
if (rawResponse.length > MAX_PROBE_RESPONSE_BYTES) {
|
||||||
|
finish();
|
||||||
|
return;
|
||||||
|
}
|
||||||
const statusMatch = rawResponse.match(/^HTTP\/\d(?:\.\d)?\s+(\d{3})/);
|
const statusMatch = rawResponse.match(/^HTTP\/\d(?:\.\d)?\s+(\d{3})/);
|
||||||
if (statusMatch) {
|
if (statusMatch) {
|
||||||
const code = Number(statusMatch[1]);
|
const code = Number(statusMatch[1]);
|
||||||
|
|||||||
@@ -31,6 +31,9 @@ import {
|
|||||||
} from './bar-server-probe';
|
} from './bar-server-probe';
|
||||||
import type { DashboardInfo as _DashboardInfo } 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.
|
// Re-exports — backward compat for tests that import from this module.
|
||||||
// resolveBarPort + defaultFindRunningServer are canonical in bar-server-probe.ts;
|
// 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 {
|
function isAuthRequiredStatus(statusCode: number): boolean {
|
||||||
return statusCode === 401 || statusCode === 403;
|
return statusCode === 401 || statusCode === 403;
|
||||||
}
|
}
|
||||||
@@ -145,9 +155,12 @@ export async function defaultWaitForServerLive(baseUrl: string): Promise<void> {
|
|||||||
return new Promise((resolve) => {
|
return new Promise((resolve) => {
|
||||||
let rawResponse = '';
|
let rawResponse = '';
|
||||||
let settled = false;
|
let settled = false;
|
||||||
|
const absoluteDeadline = setTimeout(() => finish(), BAR_PROBE_TIMEOUT_MS);
|
||||||
|
absoluteDeadline.unref?.();
|
||||||
const finish = (statusCode: number | null = null, headerSection = '') => {
|
const finish = (statusCode: number | null = null, headerSection = '') => {
|
||||||
if (settled) return;
|
if (settled) return;
|
||||||
settled = true;
|
settled = true;
|
||||||
|
clearTimeout(absoluteDeadline);
|
||||||
socket.destroy();
|
socket.destroy();
|
||||||
if (statusCode === 200) {
|
if (statusCode === 200) {
|
||||||
const echoMatch = headerSection.match(
|
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) => {
|
socket.on('data', (chunk) => {
|
||||||
rawResponse += chunk.toString('utf8');
|
rawResponse += chunk.toString('utf8');
|
||||||
|
if (rawResponse.length > MAX_BAR_PROBE_RESPONSE_BYTES) {
|
||||||
|
finish();
|
||||||
|
return;
|
||||||
|
}
|
||||||
const statusMatch = rawResponse.match(/^HTTP\/\d(?:\.\d)?\s+(\d{3})/);
|
const statusMatch = rawResponse.match(/^HTTP\/\d(?:\.\d)?\s+(\d{3})/);
|
||||||
if (statusMatch) {
|
if (statusMatch) {
|
||||||
const code = Number(statusMatch[1]);
|
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));
|
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 {
|
function defaultWriteLaunchDescriptor(jsonPath: string, descriptor: LaunchJson): void {
|
||||||
|
|||||||
@@ -3162,12 +3162,7 @@ describe('defaultFindRunningServer: socket-level 401/403 classifies authRequired
|
|||||||
setImmediate(() => {
|
setImmediate(() => {
|
||||||
onConnect();
|
onConnect();
|
||||||
for (const cb of listeners.data ?? []) {
|
for (const cb of listeners.data ?? []) {
|
||||||
cb(
|
cb(Buffer.from(`HTTP/1.1 200 OK\r\nx-ccs-bar-token: ${token}\r\n\r\n`, 'utf8'));
|
||||||
Buffer.from(
|
|
||||||
`HTTP/1.1 200 OK\r\nx-ccs-bar-token: ${token}\r\n\r\n`,
|
|
||||||
'utf8'
|
|
||||||
)
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
return socket;
|
return socket;
|
||||||
@@ -3387,7 +3382,9 @@ describe('defaultFindRunningServer: streaming lower-priority probes', () => {
|
|||||||
// exactly what the production CCS Bar server does, and is the property
|
// exactly what the production CCS Bar server does, and is the property
|
||||||
// that prevents a rogue reflector from passing the check.
|
// that prevents a rogue reflector from passing the check.
|
||||||
for (const cb of data) {
|
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
|
// 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
|
// GH-1588 — `--await-quit` waits for a running app to exit, then swaps + relaunches
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|||||||
Reference in new issue
Block a user