Download packages/kap-server/test/instanceRegistry.test.ts from SaylorTwift/kimi-code: direct link, hf CLI and curl.
- Browser
- Download file 11.6 kB
-
https://huggingface.co/SaylorTwift/kimi-code/resolve/main/packages/kap-server/test/instanceRegistry.test.ts
- Command line
-
hf download hf://SaylorTwift/kimi-code/packages/kap-server/test/instanceRegistry.test.ts
-
curl -L -o instanceRegistry.test.ts https://huggingface.co/SaylorTwift/kimi-code/resolve/main/packages/kap-server/test/instanceRegistry.test.ts
11.6 kB
| 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); | |
| }); | |
| }); | |