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
181 changes: 178 additions & 3 deletions packages/cli/src/serve/run-qwen-serve.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4063,11 +4063,12 @@ describe('runQwenServe channel worker supervisor', () => {
);
const listenError = new Error('listen failed') as NodeJS.ErrnoException;
listenError.code = 'EADDRINUSE';
const fakeServer = createServer();
vi.spyOn(serverModule, 'createServeApp').mockReturnValue({
locals: {},
listen: vi.fn(() => {
setImmediate(() => fakeServer.emit('error', listenError));
return fakeServer;
const srv = createServer();
setImmediate(() => srv.emit('error', listenError));
return srv;
}),
} as unknown as express.Application);
const worker = makeWorker({
Expand Down Expand Up @@ -4103,6 +4104,180 @@ describe('runQwenServe channel worker supervisor', () => {
expect(pidfile.removeServeServiceInfo).toHaveBeenCalledWith(process.pid);
});

it('retries the next port on EADDRINUSE and succeeds', async () => {
Comment thread
wenshao marked this conversation as resolved.
tmpDir = fs.realpathSync(
fs.mkdtempSync(path.join(os.tmpdir(), 'qws-port-retry-')),
);
const portsAttempted: number[] = [];
const stderrWrites: string[] = [];
const stderrSpy = vi
.spyOn(process.stderr, 'write')
.mockImplementation((chunk) => {
stderrWrites.push(String(chunk));
return true;
});
vi.spyOn(serverModule, 'createServeApp').mockReturnValue({
locals: {},
listen: vi.fn((port, _host, cb) => {
portsAttempted.push(port);
const srv = createServer();
if (portsAttempted.length === 1) {
const err = new Error('address in use') as NodeJS.ErrnoException;
err.code = 'EADDRINUSE';
setImmediate(() => srv.emit('error', err));
} else {
srv.listen(0, '127.0.0.1', () => {
setImmediate(() => {
srv.emit('listening');
if (typeof cb === 'function') cb();
});
});
}
return srv;
}),
} as unknown as express.Application);

const handle = await runQwenServe(
{
port: 4170,
hostname: '127.0.0.1',
mode: 'http-bridge',
workspace: tmpDir,
serveWebShell: false,
},
{
bridge: makeFakeBridge(),
resolveOnListen: true,
},
);

try {
stderrSpy.mockRestore();
expect(portsAttempted).toEqual([4170, 4171]);
expect(handle.url).toMatch(/^http:\/\/127\.0\.0\.1:\d+$/);
expect(handle.url).not.toContain(':4170');
expect(
stderrWrites.some((w) =>
w.includes('port 4170 is in use, trying 4171'),
),
).toBe(true);
} finally {
await handle.close();
}
});

it('does not retry on non-EADDRINUSE listen errors', async () => {
Comment thread
wenshao marked this conversation as resolved.
tmpDir = fs.realpathSync(
fs.mkdtempSync(path.join(os.tmpdir(), 'qws-port-no-retry-')),
);
const portsAttempted: number[] = [];
const listenError = new Error('permission denied') as NodeJS.ErrnoException;
listenError.code = 'EACCES';
vi.spyOn(serverModule, 'createServeApp').mockReturnValue({
locals: {},
listen: vi.fn((port, _host, _cb) => {
portsAttempted.push(port);
const srv = createServer();
setImmediate(() => srv.emit('error', listenError));
return srv;
}),
} as unknown as express.Application);

await expect(
runQwenServe(
{
port: 4170,
hostname: '127.0.0.1',
mode: 'http-bridge',
workspace: tmpDir,
serveWebShell: false,
},
{ bridge: makeFakeBridge() },
),
).rejects.toBe(listenError);

expect(portsAttempted).toEqual([4170]);
});

it('rejects after exhausting all port retry attempts', async () => {
tmpDir = fs.realpathSync(
fs.mkdtempSync(path.join(os.tmpdir(), 'qws-port-exhaust-')),
);
const portsAttempted: number[] = [];
const stderrWrites: string[] = [];
const stderrSpy = vi
.spyOn(process.stderr, 'write')
.mockImplementation((chunk) => {
stderrWrites.push(String(chunk));
return true;
});
const listenError = new Error('address in use') as NodeJS.ErrnoException;
listenError.code = 'EADDRINUSE';
vi.spyOn(serverModule, 'createServeApp').mockReturnValue({
locals: {},
listen: vi.fn((port) => {
portsAttempted.push(port);
const srv = createServer();
setImmediate(() => srv.emit('error', listenError));
return srv;
}),
} as unknown as express.Application);

await expect(
runQwenServe(
{
port: 4170,
hostname: '127.0.0.1',
mode: 'http-bridge',
workspace: tmpDir,
serveWebShell: false,
},
{ bridge: makeFakeBridge() },
),
).rejects.toBe(listenError);

stderrSpy.mockRestore();
expect(portsAttempted).toEqual(
Array.from({ length: 10 }, (_, i) => 4170 + i),
);
expect(
stderrWrites.some((w) => w.includes('all ports 4170–4179 are in use')),
).toBe(true);
});

it('does not retry EADDRINUSE when port is 0 (ephemeral)', async () => {
tmpDir = fs.realpathSync(
fs.mkdtempSync(path.join(os.tmpdir(), 'qws-port0-no-retry-')),
);
const portsAttempted: number[] = [];
const listenError = new Error('address in use') as NodeJS.ErrnoException;
listenError.code = 'EADDRINUSE';
vi.spyOn(serverModule, 'createServeApp').mockReturnValue({
locals: {},
listen: vi.fn((port) => {
portsAttempted.push(port);
const srv = createServer();
setImmediate(() => srv.emit('error', listenError));
return srv;
}),
} as unknown as express.Application);

await expect(
runQwenServe(
{
port: 0,
hostname: '127.0.0.1',
mode: 'http-bridge',
workspace: tmpDir,
serveWebShell: false,
},
{ bridge: makeFakeBridge() },
),
).rejects.toBe(listenError);

expect(portsAttempted).toEqual([0]);
});

it('does not remove the channel pidfile reservation for handled uncaught exceptions', async () => {
tmpDir = fs.realpathSync(
fs.mkdtempSync(path.join(os.tmpdir(), 'qws-channel-worker-crash-')),
Expand Down
68 changes: 55 additions & 13 deletions packages/cli/src/serve/run-qwen-serve.ts
Original file line number Diff line number Diff line change
Expand Up @@ -183,6 +183,8 @@ function isNonNegativeIntegerMs(value: number): boolean {

const MAX_TIMEOUT_MS = 2_147_483_647;

const MAX_PORT_ATTEMPTS = 10;

function assertTimerDelayInRange(name: string, value: number): void {
if (value > MAX_TIMEOUT_MS) {
throw new TypeError(
Expand Down Expand Up @@ -3479,11 +3481,11 @@ export async function runQwenServe(
process.on('uncaughtExceptionMonitor', onUncaughtExceptionMonitor);

// Swap the boot-error listener for a runtime-error one
// before resolving. `server.once('error', reject)` at the
// bottom only catches errors BEFORE listening; post-listen
// errors (EMFILE after FD exhaustion, runtime errors on the
// listener) would be unhandled and crash the daemon. Use a
// persistent listener that logs to stderr instead.
// before resolving. `tryListen`'s `server.once('error', ...)`
// only catches errors BEFORE listening; post-listen errors
// (EMFILE after FD exhaustion, runtime errors on the listener)
// would be unhandled and crash the daemon. Use a persistent
// listener that logs to stderr instead.
server.removeAllListeners('error');
server.on('error', (err) => {
daemonLog.error('server error', err instanceof Error ? err : null);
Expand Down Expand Up @@ -3530,8 +3532,8 @@ export async function runQwenServe(
}
};
let server: Server;
let httpsServer: https.Server | undefined;
if (tlsOptions) {
let httpsServer: https.Server;
try {
httpsServer = https.createServer(tlsOptions, app);
} catch (err) {
Expand All @@ -3548,13 +3550,53 @@ export async function runQwenServe(
);
return;
}
server = httpsServer.listen(opts.port, listenHostname, onListening);
} else {
server = app.listen(opts.port, listenHostname, onListening);
}
server.once('error', (err) => {
removeCurrentServePidfile();
reject(err);
});

const tryListen = (attemptPort: number, attempt: number): void => {
try {
if (httpsServer) {
// server.listen(port, host, cb) registers `cb` as a one-time
// `listening` listener. On failed attempts (EADDRINUSE),
// `listening` never fires so the listener accumulates. Clear
// stale listeners before each retry.
httpsServer.removeAllListeners('listening');
server = httpsServer.listen(attemptPort, listenHostname, onListening);
} else {
server = app.listen(attemptPort, listenHostname, onListening);
}
} catch (err) {
// Synchronous listen failure (e.g. invalid address) — not
// recoverable via port bump.
removeCurrentServePidfile();
reject(err instanceof Error ? err : new Error(String(err)));
return;
}

server.once('error', (err: NodeJS.ErrnoException) => {
server.close();
const nextPort = attemptPort + 1;
if (
err.code === 'EADDRINUSE' &&
opts.port !== 0 &&
nextPort <= 65535 &&
attempt < MAX_PORT_ATTEMPTS - 1
) {
writeStderrLine(
`qwen serve: port ${attemptPort} is in use, trying ${nextPort}...`,
);
tryListen(nextPort, attempt + 1);
} else {
if (err.code === 'EADDRINUSE' && attempt > 0) {
writeStderrLine(
`qwen serve: all ports ${opts.port}–${attemptPort} are in use`,
);
}
removeCurrentServePidfile();
reject(err);
}
});
};

tryListen(opts.port, 0);
});
}
60 changes: 56 additions & 4 deletions scripts/daemon-dev.js
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
import { spawn } from 'node:child_process';
import crypto from 'node:crypto';
import http from 'node:http';
import net from 'node:net';
import { dirname, join, resolve } from 'node:path';
import { fileURLToPath, pathToFileURL } from 'node:url';
import { platform } from 'node:os';
Expand All @@ -24,6 +25,7 @@ const serveOptionNames = new Set([
'--max-connections',
'--require-auth',
'--event-ring-size',
'--compacted-replay-max-bytes',
]);

function readOption(name) {
Expand Down Expand Up @@ -86,6 +88,42 @@ function daemonUrl(hostname, port) {
return `http://${host}:${port}`;
}

const MAX_PORT_ATTEMPTS = 10;

function findAvailablePort(host, startPort) {
return new Promise((resolveFind, rejectFind) => {
let attempt = 0;
const tryNext = () => {
const port = startPort + attempt;
if (port > 65535 || attempt >= MAX_PORT_ATTEMPTS) {
rejectFind(
new Error(
`No available port found in range ${startPort}–${Math.min(startPort + MAX_PORT_ATTEMPTS - 1, 65535)}`,
),
);
return;
}
const probe = net.createServer();
probe.once('error', (err) => {
probe.close();
if (err.code === 'EADDRINUSE') {
console.log(
`[daemon-dev] port ${port} is in use, trying ${port + 1}...`,
);
attempt++;
tryNext();
Comment thread
wenshao marked this conversation as resolved.
} else {
rejectFind(err);
}
});
probe.listen(port, host, () => {
probe.close(() => resolveFind(port));
});
};
tryNext();
});
}

function spawnDevProcess(label, command, commandArgs, options) {
const child = spawn(command, commandArgs, {
stdio: 'inherit',
Expand Down Expand Up @@ -199,22 +237,36 @@ try {
process.exit(1);
}

const port = readOption('--port') || '4170';
if (port === '0') {
const hostname = readOption('--hostname') || '127.0.0.1';
const rawPort = readOption('--port') || '4170';
const startPort = parseInt(rawPort, 10);
if (!Number.isInteger(startPort) || startPort < 0 || startPort > 65535) {
console.error(
`daemon-dev: --port must be an integer 0–65535, got "${rawPort}".`,
);
process.exit(1);
}
if (startPort === 0) {
console.error(
'daemon-dev: --port 0 is not supported; the launcher needs a fixed port to poll for health.',
);
process.exit(1);
}

const hostname = readOption('--hostname') || '127.0.0.1';
const probeHostname =
hostname.startsWith('[') && hostname.endsWith(']')
? hostname.slice(1, -1)
: hostname;
const port = hasOption('--port')
? String(startPort)
: String(await findAvailablePort(probeHostname, startPort));
const token =
readOption('--token') ||
process.env.QWEN_SERVER_TOKEN ||
crypto.randomBytes(16).toString('hex');
const workspace = resolve(readOption('--workspace') || process.cwd());

const serveArgs = serveArgsFromLauncherArgs();
if (!hasOption('--port')) serveArgs.push('--port', port);
Comment thread
wenshao marked this conversation as resolved.
if (!hasOption('--workspace')) serveArgs.push('--workspace', workspace);

const tsxLoaderUrl = pathToFileURL(
Expand Down
Loading