kimi-code/packages/server/test/ws-handshake.e2e.test.ts
Haozhe 60dfb68a2d
feat(server): add bearer-token auth and safe host exposure (#1006)
* test(server): add API surface snapshot guardrail

Boot startServer on port 0 and snapshot the documented v1 route table derived from /openapi.json paths, plus the reachability of doc/meta endpoints (/healthz, /openapi.json, /asyncapi.json, /). Gives later auth/--host phases an intentional diff when routes change. M0 makes no production behavior change.

* test(server): add e2e server harness with token support

Add test/helpers/serverHarness.ts: boot() wraps startServer with an isolated lock + home dir and returns a handle (server, address, baseUrl, wsUrl, token, close) plus authedFetch/authedWs that carry Authorization: Bearer <token> (and the kimi-code.bearer.<token> WS subprotocol). serviceOverrides is the generic DI seam later phases use to inject a fixed-token auth service; IAuthTokenService is not referenced yet. closeAll() tears down every booted server and socket. M0 makes no production behavior change; typecheck-only gate.

* feat(server): add privateFiles 0600 atomic write/read utility

* feat(server): add per-start tokenStore

* feat(server): add env-based bcrypt password hash utility

* feat(server): add IAuthTokenService DI seam

* feat(server): add global onRequest auth hook with bypass + redaction

* fix(server): stop reflecting Host header in /asyncapi.json

* feat(server): add WS bearer subprotocol constant and parser

* feat(server): enforce bearer token auth on WS upgrade

* feat(server): add Host header allowlist middleware

* feat(server): add Origin/CORS middleware

* feat(server): wire Host/Origin checks into HTTP and WS

* feat(server): wire token auth, Host/Origin, and WS auth into start.ts

* fix(server): create lock file with 0600 permissions

* fix(server): suppress debug routes on non-loopback binds

* feat(kimi-code): read server token and send Authorization on CLI calls

* feat(kimi-code): inject server token into /web URL fragment

* feat(server): add bindClassify for loopback/lan/public classification

* feat(kimi-code): register --host flag and pass it through the daemon

* feat(server): require password and TLS opt-out on non-loopback binds

* feat(server): rate-limit repeated auth failures on non-loopback binds

* feat(server): disable shutdown and terminals on public binds by default

* feat(server): add security response headers on non-loopback binds

* test(server): cover LAN/public host-exposure hardening end to end

* docs(server): add deployment security and threat-model guide

* changeset: minor kimi-code for server auth and host exposure

* feat(kimi-web): add server bearer-token auth support

* fix: repair CI for server auth and host exposure

- Replace native @node-rs/bcrypt with pure-JS bcryptjs so the ESM CLI
  bundle and the SEA native bundle both build without native-addon
  require issues (node-rs/bcrypt broke the ESM smoke and the SEA
  check-bundle allowlist).
- Remove dead cleanup references (stopSpinner, authLogoBlinkTimer) in
  apps/kimi-web App.vue that failed vue-tsc.
- Fix lint: drop empty spread fallbacks in the e2e auth-header merge,
  void the intentionally-async WS upgrade listener, add missing
  assertions to satisfy jest/expect-expect, and convert a ternary
  statement to if/else.
- Send the bearer token in the snapshot perf/smoke tests so they pass
  under the new global auth hook.
- Refresh the pnpmDeps hash in flake.nix for the updated lockfile.

* feat(server): persist bearer token and add rotate-token command

- persist the server bearer token in <home>/server.token (0600) and reuse it across restarts instead of per-start server-<pid>.token
- add `kimi server rotate-token` to regenerate the token; the token store reloads on mtime/inode change so rotation applies without restart
- print the token and Vite-style Local/Network URLs in the startup banner
- allow non-loopback binds with bearer-token-only auth (password now optional) and update SECURITY.md
- surface daemon boot failures immediately with the exit reason and log tail instead of waiting for the spawn timeout

* feat(server): print full token URLs and re-print links after rotate

- Drop the ready-panel border so token URLs print in full for copying; keep the Kimi sprite beside the title.
- Re-print Local/Network access links after `server rotate-token` (host/port from the lock).
- Extract shared access-URL helpers into access-urls.ts.
- Unify link and token colors between the banner and rotate-token.

* feat(server): dim URL #token= fragment and de-highlight token

- Render the `#token=…` fragment in a dim gray so the host/port stands out in the banner and rotate-token links.
- De-highlight the standalone token; set it off with surrounding whitespace instead of color.
- Add splitTokenFragment helper.

* refactor(cli): polish server ready banner and rotate-token output

- move version onto the ready banner title line; drop the separate
  Ready:/Version: rows and the startup-time metric
- reorder rotate-token output so the new token sits between the
  invalidation note and the access links
- update server CLI tests for the new layout

* feat(server): warn on reuse and refine ready banner

- Warn when `server run` reuses an already-running daemon (its options are not applied) and show the running server's actual URLs.
- Show a `Network: off  use --host 0.0.0.0 to enable` hint on loopback binds.
- Move the version onto the title line and drop the startup-time metric.

* fix(web): relabel auth dialog to token and cover full page

- Relabel the server auth dialog from "password" to "token"; the server accepts the bearer token, with the password only as a fallback.
- Make the auth dialog overlay fully opaque so it covers the whole page instead of revealing the login page underneath.

* fix: resolve CI failures on web auth PR

- Replace chalk.yellow named color with chalk.hex(darkColors.warning)
  in the server reuse notice to satisfy the chalk named color guard.
- Update pnpmDeps hash in flake.nix to match the regenerated
  pnpm-lock.yaml so the Nix build succeeds.
- Retry rmSync in ws-broadcast e2e teardown to ride out EBUSY /
  ENOTEMPTY races while the server flushes files after close().

* test(server): update API surface snapshot for warnings route

The feat/web-auth branch adds GET /api/v1/sessions/{session_id}/warnings
(packages/server/src/routes/sessions.ts), so the API surface guardrail
snapshot needs to record the new documented v1 route.
2026-06-25 17:57:56 +08:00

390 lines
12 KiB
TypeScript

/**
* WS handshake + heartbeat e2e (W5.1 / P0.15).
*
* Boots `startServer` on port 0 with a tmpdir lock, then connects a real
* `ws` client. Validates the full WS.md §1 lifecycle:
*
* 1. Server immediately sends `server_hello`.
* 2. Client sends `client_hello` → server acks (code=0).
* 3. Server pings every `pingIntervalMs` (overridden to 60ms here so the
* test finishes in <1s instead of waiting 30s).
* 4. Client `pong` resets the pong timer — connection stays alive.
* 5. Daemon close → WS code 1001 (going away).
*
* Uses `wsGatewayOptions.pingIntervalMs` + `.pongTimeoutMs` to shrink timers.
* The connection's `WsConnection` reads those through the gateway.
*/
import { mkdtempSync, rmSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import type { TelemetryClient, TelemetryProperties } from '@moonshot-ai/agent-core';
import { pino } from 'pino';
import { afterEach, beforeEach, describe, expect, it } from 'vitest';
import { WebSocket } from 'ws';
import { IConnectionRegistry, IWSGateway, startServer, type RunningServer } from '../src';
import { fixedTokenAuth } from './helpers/serverHarness';
import { rawDataToString } from '../src/ws/rawData';
interface TelemetryRecord {
readonly event: string;
readonly properties?: TelemetryProperties;
}
function recordingTelemetry(records: TelemetryRecord[]): TelemetryClient {
const client: TelemetryClient = {
track: (event, properties) => {
records.push({ event, properties });
},
withContext: () => client,
setContext: () => {},
};
return client;
}
async function waitForEvent(
records: readonly TelemetryRecord[],
event: string,
timeoutMs: number,
): Promise<TelemetryRecord> {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
const found = records.find((r) => r.event === event);
if (found !== undefined) return found;
await new Promise((res) => {
setTimeout(res, 10);
});
}
throw new Error(`event ${event} not observed; got ${JSON.stringify(records)}`);
}
let tmpDir: string;
let lockPath: string;
let bridgeHome: string;
const running: RunningServer[] = [];
beforeEach(() => {
tmpDir = mkdtempSync(join(tmpdir(), 'kimi-server-ws-handshake-'));
lockPath = join(tmpDir, 'lock');
bridgeHome = mkdtempSync(join(tmpdir(), 'kimi-server-ws-handshake-home-'));
});
afterEach(async () => {
for (const r of running.splice(0)) {
try {
await r.close();
} catch {
// ignore
}
}
rmSync(tmpDir, { recursive: true, force: true });
rmSync(bridgeHome, { recursive: true, force: true });
});
async function spawn(): Promise<RunningServer> {
const r = await startServer({
serviceOverrides: [fixedTokenAuth()],
host: '127.0.0.1',
port: 0,
lockPath,
logger: pino({ level: 'silent' }),
coreProcessOptions: { homeDir: bridgeHome },
wsGatewayOptions: { pingIntervalMs: 60, pongTimeoutMs: 200 },
});
running.push(r);
return r;
}
function wsUrl(http: string): string {
return http.replace(/^http:\/\//, 'ws://') + '/api/v1/ws';
}
interface WsFrame {
type: string;
payload?: unknown;
id?: string;
code?: number;
[k: string]: unknown;
}
/**
* Wraps a `WebSocket` with a message queue. Frames received between socket
* open and the first `await receive(...)` call go into `queue` so the test
* doesn't race the server's first push (which can land within the same tick
* as the upgrade callback). Without this, fast tests miss `server_hello`.
*/
interface Conn {
ws: WebSocket;
queue: WsFrame[];
waiters: Array<(frame: WsFrame) => void>;
closed: Promise<{ code: number; reason: string }>;
}
function openConn(url: string): Promise<Conn> {
return new Promise((resolve, reject) => {
const ws = new WebSocket(url, ['kimi-code.bearer.test-token']);
const queue: WsFrame[] = [];
const waiters: Array<(frame: WsFrame) => void> = [];
let closedResolve: (v: { code: number; reason: string }) => void;
const closed = new Promise<{ code: number; reason: string }>((res) => {
closedResolve = res;
});
ws.on('message', (data) => {
let parsed: WsFrame;
try {
parsed = JSON.parse(rawDataToString(data)) as WsFrame;
} catch {
return;
}
if (waiters.length > 0) {
const w = waiters.shift();
w?.(parsed);
} else {
queue.push(parsed);
}
});
ws.on('close', (code, reason) => {
closedResolve({ code, reason: String(reason) });
});
ws.once('open', () => resolve({ ws, queue, waiters, closed }));
ws.once('error', (err) => reject(err));
});
}
/** Pop the next frame (queued or yet-to-arrive). Rejects on timeout. */
function receive(conn: Conn, timeoutMs: number): Promise<WsFrame> {
return new Promise((resolve, reject) => {
if (conn.queue.length > 0) {
resolve(conn.queue.shift()!);
return;
}
const t = setTimeout(() => {
const idx = conn.waiters.indexOf(waiter);
if (idx >= 0) conn.waiters.splice(idx, 1);
reject(new Error(`no message in ${timeoutMs}ms`));
}, timeoutMs);
const waiter = (frame: WsFrame): void => {
clearTimeout(t);
resolve(frame);
};
conn.waiters.push(waiter);
});
}
/** Drain until a frame with the given `type` arrives. */
async function receiveType(conn: Conn, type: string, timeoutMs: number): Promise<WsFrame> {
const deadline = Date.now() + timeoutMs;
for (;;) {
const remaining = deadline - Date.now();
if (remaining <= 0) throw new Error(`no message of type ${type} within ${timeoutMs}ms`);
const frame = await receive(conn, remaining);
if (frame.type === type) return frame;
}
}
describe('WS gateway handshake + heartbeat (W5.1)', () => {
it('sends server_hello on connect and acks client_hello', async () => {
const r = await spawn();
const conn = await openConn(wsUrl(r.address));
const hello = await receiveType(conn, 'server_hello', 1000);
const helloPayload = hello.payload as { heartbeat_ms: number; max_event_buffer_size: number };
expect(helloPayload.heartbeat_ms).toBe(60);
expect(helloPayload.max_event_buffer_size).toBe(1000);
conn.ws.send(
JSON.stringify({
type: 'client_hello',
id: 'cli_test_1',
payload: { client_id: 'cli_1', subscriptions: [] },
}),
);
const ack = await receiveType(conn, 'ack', 1000);
expect(ack.id).toBe('cli_test_1');
expect(ack.code).toBe(0);
// The registry should now have exactly one attached connection.
r.services.invokeFunction((a) => {
expect(a.get(IConnectionRegistry).size()).toBe(1);
expect(a.get(IWSGateway).size).toBe(1);
});
conn.ws.close();
await conn.closed;
});
it('sends ping after pingIntervalMs', async () => {
const r = await spawn();
const conn = await openConn(wsUrl(r.address));
await receiveType(conn, 'server_hello', 1000);
// First ping should land within ~pingIntervalMs (=60ms). Give 5x slack.
const ping = await receiveType(conn, 'ping', 600);
const pingPayload = ping.payload as { nonce: string };
expect(pingPayload.nonce).toMatch(/^[0-9A-Z]{26}$/);
conn.ws.close();
await conn.closed;
});
it('pong from client keeps the connection alive past pongTimeout', async () => {
const r = await spawn();
const conn = await openConn(wsUrl(r.address));
await receiveType(conn, 'server_hello', 1000);
const ping = await receiveType(conn, 'ping', 600);
const nonce = (ping.payload as { nonce: string }).nonce;
conn.ws.send(JSON.stringify({ type: 'pong', payload: { nonce } }));
// Wait > pongTimeoutMs (200ms) — without the pong reset above we'd be
// terminated. Assert by observing the next ping arrives normally.
const nextPing = await receiveType(conn, 'ping', 600);
expect(nextPing.type).toBe('ping');
expect(conn.ws.readyState).toBe(WebSocket.OPEN);
conn.ws.close();
await conn.closed;
});
it('server close sends WS close code 1001 to attached clients', async () => {
const r = await spawn();
const conn = await openConn(wsUrl(r.address));
await receiveType(conn, 'server_hello', 1000);
await r.close();
const { code } = await conn.closed;
expect(code).toBe(1001);
});
it('non-/api/v1/ws upgrade requests are rejected', async () => {
const r = await spawn();
const badUrl = r.address.replace(/^http:\/\//, 'ws://') + '/api/v1/other';
await expect(openConn(badUrl)).rejects.toBeInstanceOf(Error);
});
});
async function waitForCount(
counts: readonly number[],
target: number,
timeoutMs: number,
): Promise<void> {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
if (counts.includes(target)) return;
await new Promise((res) => {
setTimeout(res, 10);
});
}
throw new Error(`count ${String(target)} not observed; got ${JSON.stringify(counts)}`);
}
describe('WS gateway connection-count observer', () => {
it('reports size 1 after a connect and size 0 after the last disconnect', async () => {
const counts: number[] = [];
const r = await startServer({
serviceOverrides: [fixedTokenAuth()],
host: '127.0.0.1',
port: 0,
lockPath,
logger: pino({ level: 'silent' }),
coreProcessOptions: { homeDir: bridgeHome },
wsGatewayOptions: {
pingIntervalMs: 60,
pongTimeoutMs: 200,
onConnectionCountChange: (size) => {
counts.push(size);
},
},
});
running.push(r);
const conn = await openConn(wsUrl(r.address));
await receiveType(conn, 'server_hello', 1000);
await waitForCount(counts, 1, 1000);
conn.ws.close();
await conn.closed;
await waitForCount(counts, 0, 1000);
expect(counts.at(-1)).toBe(0);
});
it('tracks multiple concurrent connections independently', async () => {
const counts: number[] = [];
const r = await startServer({
serviceOverrides: [fixedTokenAuth()],
host: '127.0.0.1',
port: 0,
lockPath,
logger: pino({ level: 'silent' }),
coreProcessOptions: { homeDir: bridgeHome },
wsGatewayOptions: {
pingIntervalMs: 60,
pongTimeoutMs: 200,
onConnectionCountChange: (size) => {
counts.push(size);
},
},
});
running.push(r);
const a = await openConn(wsUrl(r.address));
await receiveType(a, 'server_hello', 1000);
const b = await openConn(wsUrl(r.address));
await receiveType(b, 'server_hello', 1000);
await waitForCount(counts, 2, 1000);
a.ws.close();
await a.closed;
await waitForCount(counts, 1, 1000);
b.ws.close();
await b.closed;
await waitForCount(counts, 0, 1000);
expect(counts).toEqual([1, 2, 1, 0]);
});
});
describe('WS gateway telemetry (ws_connected / ws_disconnected)', () => {
it('emits ws_connected on connect and ws_disconnected on close', async () => {
const records: TelemetryRecord[] = [];
const r = await startServer({
serviceOverrides: [fixedTokenAuth()],
host: '127.0.0.1',
port: 0,
lockPath,
logger: pino({ level: 'silent' }),
coreProcessOptions: { homeDir: bridgeHome },
wsGatewayOptions: {
pingIntervalMs: 60,
pongTimeoutMs: 200,
telemetry: recordingTelemetry(records),
},
});
running.push(r);
const conn = await openConn(wsUrl(r.address));
await receiveType(conn, 'server_hello', 1000);
const connected = await waitForEvent(records, 'ws_connected', 1000);
expect(connected.properties).toMatchObject({
connection_id: expect.any(String),
connection_count: 1,
});
conn.ws.close();
await conn.closed;
const disconnected = await waitForEvent(records, 'ws_disconnected', 1000);
expect(disconnected.properties).toMatchObject({
connection_id: expect.any(String),
connection_count: 0,
duration_ms: expect.any(Number),
});
});
});