mirror of
https://github.com/MoonshotAI/kimi-code.git
synced 2026-08-22 07:05:41 +00:00
Some checks are pending
CI / build (push) Waiting to run
CI / test (1) (push) Waiting to run
CI / test (2) (push) Waiting to run
CI / test (3) (push) Waiting to run
CI / test (4) (push) Waiting to run
CI / test (5) (push) Waiting to run
CI / test-pi-tui (push) Waiting to run
CI / test-vscode-legacy (push) Waiting to run
CI / test-windows (push) Waiting to run
CI / lint (push) Waiting to run
CI / typecheck (push) Waiting to run
Nix Build / Check flake.nix workspace sync (push) Waiting to run
Nix Build / nix build .#kimi-code (push) Blocked by required conditions
Release / Release (push) Waiting to run
Release / Deploy docs (push) Blocked by required conditions
Release / Native release artifact (push) Blocked by required conditions
Release / Publish native release assets (push) Blocked by required conditions
318 lines
11 KiB
TypeScript
318 lines
11 KiB
TypeScript
import { existsSync, mkdirSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from 'node:fs';
|
|
import { tmpdir } from 'node:os';
|
|
import { join } from 'node:path';
|
|
|
|
import { afterEach, beforeEach, describe, expect, it } from 'vitest';
|
|
|
|
import {
|
|
createInstanceRegistry,
|
|
getLiveServerInstance,
|
|
listLiveServerInstances,
|
|
type ServerInstanceInfo,
|
|
} from '../src/instanceRegistry';
|
|
import { type RunningServer, startServer } from '../src/start';
|
|
import { TEST_HOST_IDENTITY } from './helpers/hostIdentity';
|
|
|
|
let tmpDir: string;
|
|
let instancesDir: string;
|
|
|
|
const DEAD_PID = 0x7fffffff;
|
|
|
|
beforeEach(() => {
|
|
tmpDir = mkdtempSync(join(tmpdir(), 'kimi-instance-registry-test-'));
|
|
instancesDir = join(tmpDir, 'instances');
|
|
});
|
|
|
|
afterEach(() => {
|
|
rmSync(tmpDir, { recursive: true, force: true });
|
|
});
|
|
|
|
interface DiskInstance {
|
|
server_id: string;
|
|
pid: number;
|
|
host: string;
|
|
port: number;
|
|
started_at: number;
|
|
heartbeat_at: number;
|
|
host_version?: string;
|
|
}
|
|
|
|
function writeInstance(serverId: string, fields: Partial<DiskInstance> & { pid: number }): void {
|
|
mkdirSync(instancesDir, { recursive: true });
|
|
const disk: DiskInstance = {
|
|
server_id: serverId,
|
|
pid: fields.pid,
|
|
host: fields.host ?? '127.0.0.1',
|
|
port: fields.port ?? 58627,
|
|
started_at: fields.started_at ?? 1000,
|
|
heartbeat_at: fields.heartbeat_at ?? 1000,
|
|
...(fields.host_version !== undefined ? { host_version: fields.host_version } : {}),
|
|
};
|
|
writeFileSync(join(instancesDir, `${serverId}.json`), JSON.stringify(disk));
|
|
}
|
|
|
|
function readInstance(serverId: string): DiskInstance {
|
|
return JSON.parse(readFileSync(join(instancesDir, `${serverId}.json`), 'utf8')) as DiskInstance;
|
|
}
|
|
|
|
function sleep(ms: number): Promise<void> {
|
|
return new Promise((resolvePromise) => setTimeout(resolvePromise, ms));
|
|
}
|
|
|
|
const baseInfo = {
|
|
pid: process.pid,
|
|
host: '127.0.0.1',
|
|
port: 58627,
|
|
startedAt: 1000,
|
|
};
|
|
|
|
describe('createInstanceRegistry — register / release', () => {
|
|
it('writes a <serverId>.json file and release removes it', async () => {
|
|
const registry = createInstanceRegistry({ instancesDir, now: () => 2000 });
|
|
const reg = await registry.register(baseInfo);
|
|
|
|
expect(typeof reg.serverId).toBe('string');
|
|
expect(reg.serverId.length).toBeGreaterThan(0);
|
|
const filePath = join(instancesDir, `${reg.serverId}.json`);
|
|
expect(existsSync(filePath)).toBe(true);
|
|
expect(readInstance(reg.serverId)).toEqual({
|
|
server_id: reg.serverId,
|
|
pid: process.pid,
|
|
host: '127.0.0.1',
|
|
port: 58627,
|
|
started_at: 1000,
|
|
heartbeat_at: 2000,
|
|
});
|
|
|
|
await reg.release();
|
|
expect(existsSync(filePath)).toBe(false);
|
|
});
|
|
|
|
it('records host_version when provided', async () => {
|
|
const registry = createInstanceRegistry({ instancesDir, now: () => 1 });
|
|
const reg = await registry.register({ ...baseInfo, serverVersion: '1.2.3' });
|
|
expect(readInstance(reg.serverId).host_version).toBe('1.2.3');
|
|
await reg.release();
|
|
});
|
|
|
|
it('assigns distinct serverIds to concurrent registrations', async () => {
|
|
const registry = createInstanceRegistry({ instancesDir, now: () => 1 });
|
|
const a = await registry.register(baseInfo);
|
|
const b = await registry.register(baseInfo);
|
|
expect(a.serverId).not.toBe(b.serverId);
|
|
await a.release();
|
|
await b.release();
|
|
});
|
|
|
|
it('release is idempotent', async () => {
|
|
const registry = createInstanceRegistry({ instancesDir, now: () => 1 });
|
|
const reg = await registry.register(baseInfo);
|
|
await reg.release();
|
|
await expect(reg.release()).resolves.toBeUndefined();
|
|
});
|
|
|
|
it('update is a no-op after release', async () => {
|
|
const registry = createInstanceRegistry({ instancesDir, now: () => 1 });
|
|
const reg = await registry.register(baseInfo);
|
|
const filePath = join(instancesDir, `${reg.serverId}.json`);
|
|
await reg.release();
|
|
await expect(reg.update({ port: 9999 })).resolves.toBeUndefined();
|
|
expect(existsSync(filePath)).toBe(false);
|
|
});
|
|
});
|
|
|
|
describe('createInstanceRegistry — stale sweep on register', () => {
|
|
it('removes dead-pid entries and keeps live ones when registering', async () => {
|
|
const registry = createInstanceRegistry({ instancesDir, now: () => 1 });
|
|
writeInstance('stale', { pid: DEAD_PID });
|
|
writeInstance('live-peer', { pid: process.pid, started_at: 500 });
|
|
|
|
const reg = await registry.register(baseInfo);
|
|
|
|
expect(existsSync(join(instancesDir, 'stale.json'))).toBe(false);
|
|
expect(existsSync(join(instancesDir, 'live-peer.json'))).toBe(true);
|
|
expect(existsSync(join(instancesDir, `${reg.serverId}.json`))).toBe(true);
|
|
await reg.release();
|
|
});
|
|
|
|
it('leaves unparseable entries alone (may be a live peer mid-write)', async () => {
|
|
const registry = createInstanceRegistry({ instancesDir, now: () => 1 });
|
|
mkdirSync(instancesDir, { recursive: true });
|
|
writeFileSync(join(instancesDir, 'garbage.json'), '{not valid');
|
|
const reg = await registry.register(baseInfo);
|
|
expect(existsSync(join(instancesDir, 'garbage.json'))).toBe(true);
|
|
await reg.release();
|
|
});
|
|
});
|
|
|
|
describe('createInstanceRegistry — listLive', () => {
|
|
it('returns live instances, drops dead ones, and sorts by startedAt', async () => {
|
|
const registry = createInstanceRegistry({ instancesDir });
|
|
writeInstance('dead', { pid: DEAD_PID, started_at: 1 });
|
|
writeInstance('older', { pid: process.pid, started_at: 100 });
|
|
writeInstance('newer', { pid: process.pid, started_at: 200 });
|
|
|
|
const live = await registry.listLive();
|
|
expect(live.map((i) => i.serverId)).toEqual(['older', 'newer']);
|
|
expect(existsSync(join(instancesDir, 'dead.json'))).toBe(false);
|
|
});
|
|
|
|
it('returns an empty array when the directory is missing', async () => {
|
|
const registry = createInstanceRegistry({ instancesDir });
|
|
await expect(registry.listLive()).resolves.toEqual([]);
|
|
});
|
|
});
|
|
|
|
describe('createInstanceRegistry — update', () => {
|
|
it('rewrites the port and refreshes heartbeatAt', async () => {
|
|
let t = 1000;
|
|
const registry = createInstanceRegistry({ instancesDir, now: () => t });
|
|
const reg = await registry.register(baseInfo);
|
|
expect(readInstance(reg.serverId).port).toBe(58627);
|
|
expect(readInstance(reg.serverId).heartbeat_at).toBe(1000);
|
|
|
|
t = 2000;
|
|
await reg.update({ port: 58628 });
|
|
|
|
const after = readInstance(reg.serverId);
|
|
expect(after.port).toBe(58628);
|
|
expect(after.heartbeat_at).toBe(2000);
|
|
expect(after.pid).toBe(process.pid);
|
|
expect(after.started_at).toBe(1000);
|
|
await reg.release();
|
|
});
|
|
|
|
it('refreshes heartbeat without changing port when patch is empty', async () => {
|
|
let t = 1000;
|
|
const registry = createInstanceRegistry({ instancesDir, now: () => t });
|
|
const reg = await registry.register(baseInfo);
|
|
t = 3000;
|
|
await reg.update({});
|
|
const after = readInstance(reg.serverId);
|
|
expect(after.port).toBe(58627);
|
|
expect(after.heartbeat_at).toBe(3000);
|
|
await reg.release();
|
|
});
|
|
});
|
|
|
|
describe('createInstanceRegistry — heartbeat', () => {
|
|
it('periodically rewrites heartbeatAt until released', async () => {
|
|
let tick = 0;
|
|
const registry = createInstanceRegistry({
|
|
instancesDir,
|
|
heartbeatIntervalMs: 1,
|
|
now: () => ++tick,
|
|
});
|
|
const reg = await registry.register(baseInfo);
|
|
const first = readInstance(reg.serverId).heartbeat_at;
|
|
|
|
await sleep(90);
|
|
const later = readInstance(reg.serverId).heartbeat_at;
|
|
expect(later).toBeGreaterThan(first);
|
|
|
|
await reg.release();
|
|
expect(existsSync(join(instancesDir, `${reg.serverId}.json`))).toBe(false);
|
|
await sleep(30);
|
|
expect(existsSync(join(instancesDir, `${reg.serverId}.json`))).toBe(false);
|
|
});
|
|
});
|
|
|
|
describe('convenience readers', () => {
|
|
it('listLiveServerInstances reads <homeDir>/server/instances', async () => {
|
|
const dir = join(tmpDir, 'server', 'instances');
|
|
const registry = createInstanceRegistry({ instancesDir: dir });
|
|
const reg = await registry.register({ ...baseInfo, startedAt: 100 });
|
|
|
|
const live = await listLiveServerInstances(tmpDir);
|
|
expect(live.map((i: ServerInstanceInfo) => i.serverId)).toEqual([reg.serverId]);
|
|
await reg.release();
|
|
});
|
|
|
|
it('getLiveServerInstance returns the longest-running live instance', async () => {
|
|
const dir = join(tmpDir, 'server', 'instances');
|
|
const registry = createInstanceRegistry({ instancesDir: dir });
|
|
const older = await registry.register({ ...baseInfo, startedAt: 100 });
|
|
const newer = await registry.register({ ...baseInfo, startedAt: 200 });
|
|
|
|
const first = await getLiveServerInstance(tmpDir);
|
|
expect(first?.serverId).toBe(older.serverId);
|
|
await older.release();
|
|
await newer.release();
|
|
});
|
|
|
|
it('getLiveServerInstance returns undefined when no live instance exists', async () => {
|
|
await expect(getLiveServerInstance(tmpDir)).resolves.toBeUndefined();
|
|
});
|
|
});
|
|
|
|
describe('startServer — instance registry wiring', () => {
|
|
let home: string | undefined;
|
|
const servers: RunningServer[] = [];
|
|
|
|
afterEach(async () => {
|
|
while (servers.length > 0) {
|
|
await servers.pop()!.close();
|
|
}
|
|
if (home !== undefined) {
|
|
rmSync(home, { recursive: true, force: true });
|
|
home = undefined;
|
|
}
|
|
});
|
|
|
|
it('lets two servers share one homeDir, each registering a distinct instance and port', async () => {
|
|
home = mkdtempSync(join(tmpdir(), 'kimi-server-multi-server-'));
|
|
const a = await startServer({ hostIdentity: TEST_HOST_IDENTITY, host: '127.0.0.1', port: 0, homeDir: home, logLevel: 'silent' });
|
|
servers.push(a);
|
|
const b = await startServer({ hostIdentity: TEST_HOST_IDENTITY, host: '127.0.0.1', port: 0, homeDir: home, logLevel: 'silent' });
|
|
servers.push(b);
|
|
|
|
expect(b.port).not.toBe(a.port);
|
|
|
|
const live = await listLiveServerInstances(home);
|
|
expect(live).toHaveLength(2);
|
|
expect(new Set(live.map((i) => i.serverId)).size).toBe(2);
|
|
expect(live.map((i) => i.port).sort((x, y) => x - y)).toEqual(
|
|
[a.port, b.port].sort((x, y) => x - y),
|
|
);
|
|
expect(existsSync(join(home, 'server', 'lock'))).toBe(false);
|
|
});
|
|
|
|
it('removes its instance file on close so peers no longer list it', async () => {
|
|
home = mkdtempSync(join(tmpdir(), 'kimi-server-multi-server-'));
|
|
const a = await startServer({ hostIdentity: TEST_HOST_IDENTITY, host: '127.0.0.1', port: 0, homeDir: home, logLevel: 'silent' });
|
|
servers.push(a);
|
|
const b = await startServer({ hostIdentity: TEST_HOST_IDENTITY, host: '127.0.0.1', port: 0, homeDir: home, logLevel: 'silent' });
|
|
servers.push(b);
|
|
expect(await listLiveServerInstances(home)).toHaveLength(2);
|
|
|
|
await a.close();
|
|
servers.splice(servers.indexOf(a), 1);
|
|
|
|
const live = await listLiveServerInstances(home);
|
|
expect(live).toHaveLength(1);
|
|
expect(live[0]?.port).toBe(b.port);
|
|
});
|
|
|
|
it('releases its registration on close so a fresh instance on the same home can start', async () => {
|
|
home = mkdtempSync(join(tmpdir(), 'kimi-server-multi-server-'));
|
|
const first = await startServer({
|
|
hostIdentity: TEST_HOST_IDENTITY,
|
|
host: '127.0.0.1',
|
|
port: 0,
|
|
homeDir: home,
|
|
logLevel: 'silent',
|
|
});
|
|
await first.close();
|
|
|
|
const restarted = await startServer({
|
|
hostIdentity: TEST_HOST_IDENTITY,
|
|
host: '127.0.0.1',
|
|
port: 0,
|
|
homeDir: home,
|
|
logLevel: 'silent',
|
|
});
|
|
servers.push(restarted);
|
|
expect(await listLiveServerInstances(home)).toHaveLength(1);
|
|
expect((await listLiveServerInstances(home))[0]?.port).toBe(restarted.port);
|
|
});
|
|
});
|