Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion test/activity-entries.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,9 @@ function abortMidBody(port, opened) {
}

// Poll `cond` until it holds or `ms` elapse; the caller asserts afterwards.
async function until(cond, ms = 5000) {
// The default is a watchdog against a condition that never comes, well above
// anything a loaded machine adds, not a bound on how fast it should.
async function until(cond, ms = 60_000) {
const deadline = Date.now() + ms;
while (!cond() && Date.now() < deadline) await new Promise(r => setTimeout(r, 10));
}
Expand Down
7 changes: 5 additions & 2 deletions test/client-disconnect.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,10 @@ async function fixture(t) {
return { am, url, upload, started, ended, forwarded: () => forwarded, upstreamClosed: () => upstreamClosed };
}

test('disconnect during a quota hold clears activity without waiting for the retry timer', { timeout: 4000 }, async t => {
// No per-test timeout on either: what must not happen is a wait for the
// 60 s retry timer (or the 2-minute idle watchdog), and the runner's own
// timeout is what catches that; a shorter bound here only measures the machine.
test('disconnect during a quota hold clears activity without waiting for the retry timer', async t => {
const f = await fixture(t);
// No account can serve: with holdSeconds set the proxy holds the connection
// and sleeps (60s here) before polling again. The sleep must end with the client.
Expand All @@ -64,7 +67,7 @@ test('disconnect during a quota hold clears activity without waiting for the ret
assert.equal(f.forwarded(), 0);
});

test('disconnect before upstream headers cancels upstream and closes activity promptly', { timeout: 4000 }, async t => {
test('disconnect before upstream headers cancels upstream and closes activity promptly', async t => {
const f = await fixture(t);
const client = f.upload({ 'x-test-stall': '1' }); client.req.end('{}');
await until(() => f.forwarded() === 1);
Expand Down
5 changes: 2 additions & 3 deletions test/codex-reset-credits.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -498,12 +498,11 @@ test('a token refresh that outlives the budget is left behind, not waited out',
usageFn: async () => ({}),
});

const started = Date.now();
const result = await redeemer.maybeRedeemForPool([am.accounts[0]]);
const waited = Date.now() - started;
assert.equal(result.redeemed, false);
// The reason is the budget: the attempt gave up on the refresh rather than
// waiting for it (which would be forever — finishRefresh is never called).
assert.match(result.reason, /budget/);
assert.ok(waited < 2000, `the refusal waited ${waited}ms on a refresh that never finished`);
assert.equal(calls.details, 0, 'nothing is read on a token the attempt never got');

// And it stays left behind: a refresh landing after the deadline is no longer
Expand Down
18 changes: 9 additions & 9 deletions test/config-lock-file.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -76,9 +76,9 @@ test('two concurrent atomicConfigUpdate calls both land', async () => {
test('a lock older than 10 s is broken even though its pid is alive', async () => {
await withConfigDir(async ({ dir, cfg, path, lockPath }) => {
await writeFile(lockPath, JSON.stringify({ pid: process.pid, at: Date.now() - 60_000 }));
const started = Date.now();
await cfg.saveConfig({ fresh: true });
assert.ok(Date.now() - started < 1000, 'a stale lock does not cost the 2 s wait');
const { warnings, restore } = captureBypassWarnings();
try { await cfg.saveConfig({ fresh: true }); } finally { restore(); }
assert.deepEqual(warnings, [], 'a stale lock is broken, not waited out and bypassed');
assert.deepEqual(await readJson(path), { fresh: true });
assert.deepEqual(await readdir(dir), ['teamclaude.json'], 'the stale lock is gone');
});
Expand All @@ -87,9 +87,9 @@ test('a lock older than 10 s is broken even though its pid is alive', async () =
test('a fresh lock whose pid is dead is broken', async () => {
await withConfigDir(async ({ dir, cfg, path, lockPath }) => {
await writeFile(lockPath, JSON.stringify({ pid: await deadPid(), at: Date.now() }));
const started = Date.now();
await cfg.saveConfig({ fresh: true });
assert.ok(Date.now() - started < 1000, 'a dead holder does not cost the 2 s wait');
const { warnings, restore } = captureBypassWarnings();
try { await cfg.saveConfig({ fresh: true }); } finally { restore(); }
assert.deepEqual(warnings, [], 'a dead holder\'s lock is broken, not waited out and bypassed');
assert.deepEqual(await readJson(path), { fresh: true });
assert.deepEqual(await readdir(dir), ['teamclaude.json']);
});
Expand Down Expand Up @@ -145,9 +145,9 @@ test('the lock is released after a success and after a throwing mutator', async
assert.equal((await readJson(path)).ok, true, 'the failed update wrote nothing');

// And the next writer is not held up by anything the failure left behind.
const started = Date.now();
await cfg.saveConfig({ after: true });
assert.ok(Date.now() - started < 1000);
const { warnings, restore } = captureBypassWarnings();
try { await cfg.saveConfig({ after: true }); } finally { restore(); }
assert.deepEqual(warnings, [], 'nothing left behind was waited out');
assert.deepEqual(await readJson(path), { after: true });
});
});
Expand Down
19 changes: 11 additions & 8 deletions test/mcp-tools.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -347,20 +347,23 @@ test('a tool that blows up reports a generic failure, never the exception', asyn
test('write tools run one at a time, so a removal cannot race a reload', async () => {
const order = [];
let releaseReload;
const reload = () => new Promise(resolve => { order.push('reload:start'); releaseReload = () => { order.push('reload:end'); resolve(0); }; });
// The reload announces that it was reached, so the test waits for that —
// not for a duration, and not against a deadline of its own: the runner's
// timeout bounds a write that never arrives.
let reachedReload;
const reached = new Promise(resolve => { reachedReload = resolve; });
const reload = () => new Promise(resolve => {
order.push('reload:start');
releaseReload = () => { order.push('reload:end'); resolve(0); };
reachedReload();
});
const persistAccounts = async () => { order.push('persist'); };
const { tools } = await fixture({ hooks: { reload, persistAccounts } });

const first = tools.call('set_threshold', { percent: 70 });
const second = tools.call('remove_account', { account: 'alice@example.com' });
try {
// Wait for the first call to reach its reload — not for a duration — before
// judging what the second has done meanwhile.
const deadline = Date.now() + 5000;
while (!releaseReload) {
if (Date.now() > deadline) throw new Error('the first write never reached its reload');
await new Promise(r => setTimeout(r, 5));
}
await reached;
assert.deepEqual(order, ['reload:start'], 'the removal must wait for the running write to finish');
} finally {
// Released whatever the verdict: the queue is shared by every write tool
Expand Down
5 changes: 3 additions & 2 deletions test/request-log-sweep.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -107,8 +107,9 @@ test('a server with a logDir sweeps it once at startup', { timeout: 20000 }, asy
});
await new Promise(r => proxy.listen(0, '127.0.0.1', r));
try {
const deadline = Date.now() + 5000;
while (readdirSync(dir).includes(expired) && Date.now() < deadline) {
// The startup sweep is asynchronous; wait for it to have happened, with
// the runner's timeout as the only bound.
while (readdirSync(dir).includes(expired)) {
await new Promise(r => setTimeout(r, 25));
}
const left = readdirSync(dir);
Expand Down
8 changes: 4 additions & 4 deletions test/routing-failover.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -243,17 +243,17 @@ test('a proxy that accepts the connection and never answers is a routing failure
const port = await listen(proxy);
try {
let first;
const started = Date.now();
const lines = await captureLogs(async () => { first = await post(port); });
assert.equal(first.type, 'response', `reset instead of failing over: ${JSON.stringify(first)}\n${lines.join('\n')}`);
assert.equal(first.status, 200, lines.join('\n'));
assert.deepEqual(upstream.hits.map(h => h.key), ['sk-direct'], 'served by the account whose path works');
// The tunnel's own budget gave up, not some caller's longer signal (which
// would have surfaced as a generic timeout and been retried as transient).
assert.ok(Date.now() - started < 10_000, `the forward waited ${Date.now() - started}ms on the wedged proxy`);

const failed = lines.filter(l => l.includes('Routing proxy failed for account "routed"'));
assert.equal(failed.length, 1, lines.join('\n'));
// The tunnel's own budget (TEAMCLAUDE_ROUTING_TIMEOUT_MS above) gave up,
// not some caller's longer signal, which would have surfaced as a generic
// timeout and been retried as transient. The failure names its budget.
assert.match(failed[0], /handshake timed out after 800ms/);
assert.match(failed[0], /SOCKS5 handshake timed out after 800ms/);
assert.equal(am.unavailableReason(am.accounts[0]), 'routing');
assert.ok(am.accounts[0].routingFailedUntil > Date.now(), 'the account sits out the cooldown');
Expand Down
9 changes: 6 additions & 3 deletions test/server-429.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -90,23 +90,26 @@ test('long upstream Retry-After is surfaced without sleeping in client request',
const proxyPort = await listen(proxy);

try {
// The bound is the 300 s retry window the proxy must NOT sleep through,
// not a measure of speed: a signal at a fraction of that window fails a
// proxy that absorbed it and nothing that merely ran on a busy machine.
const started = Date.now();
let res;
try {
res = await fetch(`http://127.0.0.1:${proxyPort}/v1/messages`, {
method: 'POST',
headers: { 'content-type': 'application/json' },
body: JSON.stringify({ model: 'x', messages: [] }),
signal: AbortSignal.timeout(2000),
signal: AbortSignal.timeout(60_000),
});
} catch (err) {
assert.fail(`request should return 429 promptly, got ${err.name}`);
assert.fail(`request should return 429 without waiting out the retry window, got ${err.name}`);
}

await res.text();
assert.equal(res.status, 429);
assert.equal(upstreamHits, 1, 'long Retry-After should not be retried inline');
assert.ok(Date.now() - started < 2000, 'request should not sleep for upstream retry window');
assert.ok(Date.now() - started < 150_000, 'request should not sleep for upstream retry window');
assert.equal(am.accounts[0].status, 'active', 'rate-limit 429 must not throttle/rotate the account');
assert.ok(am.accounts[0].pausedUntil > Date.now(), 'account should be paused so concurrent requests wait');
} finally {
Expand Down
3 changes: 2 additions & 1 deletion test/server-warmup-schedule.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,8 @@ function availablePort() {
function waitForOutput(child, pattern) {
return new Promise((resolve, reject) => {
let output = '';
const timer = setTimeout(() => reject(new Error(`server did not start:\n${output}`)), 10_000);
// A watchdog against a child that never gets there, not a bound on startup.
const timer = setTimeout(() => reject(new Error(`server did not start:\n${output}`)), 60_000);
const onData = chunk => {
output += chunk;
if (pattern.test(output)) {
Expand Down
3 changes: 2 additions & 1 deletion test/tui-remote.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,8 @@ const stripAnsi = s => s.replace(/\x1b\[[0-9;]*m/g, '');
// Every action here crosses a real socket, so wait for the thing to have
// happened rather than for a duration: a fixed sleep is a race that a loaded
// machine loses, and these tests run alongside the rest of the suite.
async function waitFor(predicate, what, timeoutMs = 5000) {
// The default is a watchdog, far above what a loaded machine adds.
async function waitFor(predicate, what, timeoutMs = 60_000) {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
if (predicate()) return;
Expand Down
5 changes: 4 additions & 1 deletion test/upgrade-response-head.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,10 @@ async function rawUpstream(response) {
// refusal, and after a 101 once the upstream side ends, so the read completes
// on 'close'. A bounded timeout turns a hang into a failure that shows what
// arrived.
async function handshake(upstreamPort, timeoutMs = 5000) {
// `timeoutMs` is a watchdog against a proxy that never closes the socket,
// set well above anything a loaded machine adds; the runner's timeout is the
// other bound.
async function handshake(upstreamPort, timeoutMs = 60_000) {
const proxy = http.createServer();
proxy.on('upgrade', (req, socket, head) => relayUpgrade(req, socket, head, `http://127.0.0.1:${upstreamPort}`, null, { log: () => {} }));
const port = await listen(proxy);
Expand Down
2 changes: 1 addition & 1 deletion test/upstream-overload.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ function post(port) {
// its queue gets a 503 with Retry-After, is never retried on another account
// (that would only add load to the same saturated origin), and its activity row
// is not attributed to an account it never reached.
test('a request past the upstream pool and queue gets a 503, with no account rotation or attribution', { timeout: 4000 }, async () => {
test('a request past the upstream pool and queue gets a 503, with no account rotation or attribution', async () => {
let reached = 0;
const held = [];
const upstream = http.createServer((req, res) => {
Expand Down
40 changes: 26 additions & 14 deletions test/upstream-pool.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -18,39 +18,48 @@ async function listen(handler) {
// HTTP/2 connection does under concurrent uploads.
test('concurrent requests each open their own connection and run in parallel', async () => {
let conns = 0;
const HEADER_DELAY = 300;
const N = 8;
// A barrier, not a delay: no request is answered until all N have arrived.
// Requests that serialize behind one connection can never all arrive, so
// the barrier never opens and the watchdog fails the test; requests that
// run in parallel all arrive, however slowly the machine gets them there.
/** @type {import('node:http').ServerResponse[]} */
const pending = [];
let serialized = false;
const answerAll = () => { for (const res of pending.splice(0)) { res.writeHead(200); res.end('ok'); } };
const watchdog = setTimeout(() => { serialized = true; answerAll(); }, 30_000);
const { server, port } = await listen((req, res) => {
setTimeout(() => { res.writeHead(200); res.end('ok'); }, HEADER_DELAY);
pending.push(res);
if (pending.length === N) { clearTimeout(watchdog); answerAll(); }
});
server.on('connection', () => { conns += 1; });

const N = 8;
const started = Date.now();
const bodies = await Promise.all(
Array.from({ length: N }, () =>
upstreamFetch(`http://127.0.0.1:${port}/`, { headersTimeoutMs: 5000 }).then((r) => r.text())),
upstreamFetch(`http://127.0.0.1:${port}/`, { headersTimeoutMs: 60_000 }).then((r) => r.text())),
);
const elapsed = Date.now() - started;

assert.deepEqual(bodies, Array(N).fill('ok'));
assert.equal(conns, N, `expected ${N} parallel connections, saw ${conns}`);
// Parallel: total ≈ one request's delay, NOT N × delay (serialization).
assert.ok(elapsed < HEADER_DELAY * 3, `expected parallel (~${HEADER_DELAY}ms), took ${elapsed}ms`);
assert.equal(serialized, false, `the ${N} requests did not all arrive while the first was still pending`);

server.close();
});

// Long-lived streams hold their pooled socket (and admission permit) until the
// body ends; with the default pool width a dozen of them must not wait on one
// another for headers.
test('twelve long-lived streams receive headers without waiting for another stream to end', { timeout: 4000 }, async t => {
test('twelve long-lived streams receive headers without waiting for another stream to end', async t => {
const { server, port } = await listen((req, res) => {
res.writeHead(200, { 'content-type': 'text/event-stream' });
res.write('event: ping\n\n');
});
t.after(() => { server.closeAllConnections(); server.close(); });
// The streams never end, so a fetch that waited on another's end would
// wait forever: the headers deadline only has to be shorter than that, and
// generous enough that a busy machine is not what trips it.
const results = await Promise.allSettled(Array.from({ length: 12 }, () =>
upstreamFetch(`http://127.0.0.1:${port}/`, { headersTimeoutMs: 500 })));
upstreamFetch(`http://127.0.0.1:${port}/`, { headersTimeoutMs: 60_000 })));
for (const r of results) if (r.status === 'fulfilled') await r.value.body.cancel();
assert.equal(results.filter(r => r.status === 'fulfilled').length, 12);
});
Expand All @@ -60,7 +69,7 @@ test('twelve long-lived streams receive headers without waiting for another stre
// signal aborts, is refused outright when the queue is full, and is admitted
// (with the headers deadline armed only then) when a stream ahead of it ends.
// None of the refused requests ever reach the server.
test('upstream queue is bounded, cancellable, separately timed, and recovers on stream release', { timeout: 5000 }, async t => {
test('upstream queue is bounded, cancellable, separately timed, and recovers on stream release', async t => {
const keys = ['TEAMCLAUDE_UPSTREAM_MAX_SOCKETS', 'TEAMCLAUDE_UPSTREAM_MAX_QUEUE'];
const saved = keys.map(k => process.env[k]);
process.env[keys[0]] = '1'; process.env[keys[1]] = '1';
Expand All @@ -81,9 +90,12 @@ test('upstream queue is bounded, cancellable, separately timed, and recovers on
await assert.rejects(limited(`${base}/overflow`), { code: 'TEAMCLAUDE_UPSTREAM_OVERLOADED' });
ac.abort(); await cancelled;
await assert.rejects(limited(`${base}/expired`, { queueTimeoutMs: 25 }), { code: 'TEAMCLAUDE_UPSTREAM_OVERLOADED' });
// Waiting 75ms must not consume the 50ms upstream headers deadline.
const next = limited(`${base}/next`, { queueTimeoutMs: 1000, headersTimeoutMs: 50 });
await delay(75);
// Queueing for longer than the headers deadline must not consume it: the
// deadline is armed at admission. Armed at enqueue instead, it would fire
// while the request is still queued, whatever the machine is doing; armed
// correctly, a loopback response has the whole 2 s to arrive.
const next = limited(`${base}/next`, { queueTimeoutMs: 60_000, headersTimeoutMs: 2_000 });
await delay(2_500);
await hold.body.cancel();
assert.equal(await (await next).text(), 'ok');
assert.deepEqual(reached, ['/hold', '/next']);
Expand Down
Loading
Loading