- probe configured server URL (~1s) before remote commands; on failure re-discover via LAN discovery, persist and use the resolved URL - poll tryClaim/trySurfaceMessages every ~4s during the SSE wait in CLI work and MCP waitForTask — a lost SSE frame costs seconds, never sleep - work/MCP no longer auto-mark messages read: surfacing sets delivered, explicit ack/read required; loops still wake only on unread (no spin) - parseSSEBuffer accepts CRLF frame separators and joins multi data: lines - bump version to 0.10.1 (package.json, CLI --version, MCP server)
211 lines
8.3 KiB
TypeScript
211 lines
8.3 KiB
TypeScript
/**
|
|
* Regression tests for the realtime reliability patch (TSK-0224):
|
|
* (a) polling fallback claims a task even when its SSE event is lost;
|
|
* (b) parseSSEBuffer accepts CRLF frames and joins multiple data: lines;
|
|
* (c) delivered-but-unacked messages stay visible in the inbox;
|
|
* (d) discovery self-heal falls back when the configured URL is dead.
|
|
*/
|
|
import { describe, it, expect, beforeEach, afterEach } from 'vitest';
|
|
import { mkdtempSync, rmSync } from 'fs';
|
|
import { tmpdir } from 'os';
|
|
import { join } from 'path';
|
|
import { init } from '../src/cli/commands/init.js';
|
|
import { workAgent } from '../src/cli/commands/work.js';
|
|
import { parseSSEBuffer } from '../src/cli/commands/watch.js';
|
|
import { startServer } from '../src/server/index.js';
|
|
import { createMessage, listInbox, getMessage } from '../src/core/services/messageService.js';
|
|
import { probeServer, resolveReachableServerUrl } from '../src/discovery.js';
|
|
import type { Task } from '../src/core/schema.js';
|
|
|
|
// ─── (a) polling fallback: claim despite a lost SSE event ────────────────────
|
|
|
|
describe('work — polling fallback (lost SSE event)', () => {
|
|
let cwd: string;
|
|
let server: Awaited<ReturnType<typeof startServer>>;
|
|
const realFetch = globalThis.fetch;
|
|
|
|
beforeEach(async () => {
|
|
cwd = mkdtempSync(join(tmpdir(), 'ah-poll-'));
|
|
init(cwd, { projectName: 'poll-test', yes: true });
|
|
server = await startServer(cwd, { host: '127.0.0.1', port: 0 });
|
|
|
|
// Suppress the SSE channel: /events returns an open stream that never
|
|
// delivers a frame. Every other request passes through to the real server.
|
|
globalThis.fetch = (async (input: RequestInfo | URL, init?: RequestInit) => {
|
|
const url = typeof input === 'string' ? input : input instanceof URL ? input.toString() : input.url;
|
|
if (url.endsWith('/events')) {
|
|
const stream = new ReadableStream<Uint8Array>({
|
|
start(controller) {
|
|
init?.signal?.addEventListener('abort', () => {
|
|
try {
|
|
controller.close();
|
|
} catch {
|
|
/* already closed */
|
|
}
|
|
});
|
|
// Never enqueue: the SSE frame is "lost".
|
|
},
|
|
});
|
|
return new Response(stream, { status: 200, headers: { 'Content-Type': 'text/event-stream' } });
|
|
}
|
|
return realFetch(input, init);
|
|
}) as typeof fetch;
|
|
});
|
|
|
|
afterEach(async () => {
|
|
globalThis.fetch = realFetch;
|
|
await server.app.close().catch(() => undefined);
|
|
rmSync(cwd, { recursive: true, force: true });
|
|
});
|
|
|
|
it('claims a task via the polling fallback when its SSE event never arrives', async () => {
|
|
const workDone = workAgent({
|
|
serverUrl: server.url,
|
|
projectCwd: cwd,
|
|
agent: 'kimi',
|
|
role: 'implementer',
|
|
timeoutSec: 8,
|
|
pollIntervalMs: 100,
|
|
});
|
|
|
|
// Delegate a task after the wait started. The SSE event is swallowed by the
|
|
// stub above, so only the polling fallback can notice it.
|
|
await new Promise((r) => setTimeout(r, 300));
|
|
await realFetch(`${server.url}/tasks`, {
|
|
method: 'POST',
|
|
headers: { 'Content-Type': 'application/json' },
|
|
body: JSON.stringify({ title: 'kimi: claimed via polling', role: 'implementer' }),
|
|
});
|
|
|
|
await workDone; // must resolve via the poll, long before the 8s timeout
|
|
|
|
const tasks = (await realFetch(`${server.url}/tasks`).then((r) => r.json())) as Task[];
|
|
const mine = tasks.find((t) => (t.title ?? '').startsWith('kimi:'));
|
|
expect(mine).toBeDefined();
|
|
expect(mine?.status).toBe('in_progress');
|
|
expect(mine?.assignedTo).toBe('kimi');
|
|
}, 10000);
|
|
});
|
|
|
|
// ─── (b) CRLF-tolerant SSE parser ────────────────────────────────────────────
|
|
|
|
describe('parseSSEBuffer — CRLF + multi-line data', () => {
|
|
it('parses CRLF-framed events', () => {
|
|
const buf = 'data: {"type":"task","action":"created","id":"TSK-0001"}\r\n\r\n';
|
|
const { events, remaining } = parseSSEBuffer(buf);
|
|
expect(events).toHaveLength(1);
|
|
expect(events[0].id).toBe('TSK-0001');
|
|
expect(remaining).toBe('');
|
|
});
|
|
|
|
it('parses mixed LF/CRLF frames and CRLF keepalives', () => {
|
|
const buf =
|
|
':\r\n\r\n' +
|
|
'data: {"type":"memory","action":"created","id":"MEM-001"}\n\n' +
|
|
'data: {"type":"decision","action":"created","id":"DEC-001"}\r\n\r\n';
|
|
const { events, remaining } = parseSSEBuffer(buf);
|
|
expect(events).toHaveLength(2);
|
|
expect(events[0].id).toBe('MEM-001');
|
|
expect(events[1].id).toBe('DEC-001');
|
|
expect(remaining).toBe('');
|
|
});
|
|
|
|
it('joins multiple data: lines of one event with a newline', () => {
|
|
const buf = 'data: {"type":"task",\r\ndata: "id":"TSK-0002"}\r\n\r\n';
|
|
const { events } = parseSSEBuffer(buf);
|
|
expect(events).toHaveLength(1);
|
|
expect(events[0].type).toBe('task');
|
|
expect(events[0].id).toBe('TSK-0002');
|
|
});
|
|
|
|
it('keeps an incomplete CRLF tail in `remaining`', () => {
|
|
const buf = 'data: {"type":"task","id":"TSK-0001"}\r\n\r\ndata: {"type":"deci';
|
|
const { events, remaining } = parseSSEBuffer(buf);
|
|
expect(events).toHaveLength(1);
|
|
expect(remaining).toBe('data: {"type":"deci');
|
|
});
|
|
});
|
|
|
|
// ─── (c) delivered-unacked messages stay visible ─────────────────────────────
|
|
|
|
describe('message semantics — delivered stays visible until explicit ack/read', () => {
|
|
let cwd: string;
|
|
|
|
beforeEach(() => {
|
|
cwd = mkdtempSync(join(tmpdir(), 'ah-msg-vis-'));
|
|
init(cwd, { projectName: 'msg-vis', yes: true });
|
|
});
|
|
|
|
afterEach(() => {
|
|
rmSync(cwd, { recursive: true, force: true });
|
|
});
|
|
|
|
it('a surfaced (delivered) message stays in the default inbox and never re-wakes', () => {
|
|
const msg = createMessage(cwd, { from: 'claude', to: 'kimi', text: 'look at this' });
|
|
|
|
// The work-loop drain path: fetch unread. First call surfaces it and flips
|
|
// it to delivered; nothing marks it read.
|
|
const first = listInbox(cwd, 'kimi', { unreadOnly: true });
|
|
expect(first).toHaveLength(1);
|
|
expect(first[0].status).toBe('delivered');
|
|
expect(getMessage(cwd, msg.id).message.status).toBe('delivered');
|
|
|
|
// No spin: the unread-filtered wake check is empty on the next pass…
|
|
expect(listInbox(cwd, 'kimi', { unreadOnly: true })).toHaveLength(0);
|
|
|
|
// …but the info stays visible in the default inbox until explicit ack/read.
|
|
const inbox = listInbox(cwd, 'kimi');
|
|
expect(inbox.map((m) => m.id)).toContain(msg.id);
|
|
expect(inbox.find((m) => m.id === msg.id)?.status).toBe('delivered');
|
|
});
|
|
});
|
|
|
|
// ─── (d) discovery self-heal ─────────────────────────────────────────────────
|
|
|
|
describe('discovery self-heal — stale configured URL', () => {
|
|
let cwd: string;
|
|
let server: Awaited<ReturnType<typeof startServer>>;
|
|
|
|
beforeEach(async () => {
|
|
cwd = mkdtempSync(join(tmpdir(), 'ah-heal-'));
|
|
init(cwd, { projectName: 'heal-test', yes: true });
|
|
server = await startServer(cwd, { host: '127.0.0.1', port: 0 });
|
|
});
|
|
|
|
afterEach(async () => {
|
|
await server.app.close().catch(() => undefined);
|
|
rmSync(cwd, { recursive: true, force: true });
|
|
});
|
|
|
|
it('falls back to the discovered URL when the configured one is dead', async () => {
|
|
const discovered = await resolveReachableServerUrl('http://127.0.0.1:1', {
|
|
discover: async () => server.url,
|
|
});
|
|
expect(discovered).toBe(server.url);
|
|
});
|
|
|
|
it('keeps the configured URL when it answers (discovery not consulted)', async () => {
|
|
let discoverCalled = false;
|
|
const resolved = await resolveReachableServerUrl(server.url, {
|
|
discover: async () => {
|
|
discoverCalled = true;
|
|
return 'http://example.invalid:9';
|
|
},
|
|
});
|
|
expect(resolved).toBe(server.url);
|
|
expect(discoverCalled).toBe(false);
|
|
});
|
|
|
|
it('returns undefined when the configured URL is dead and nothing is discovered', async () => {
|
|
const resolved = await resolveReachableServerUrl('http://127.0.0.1:1', {
|
|
discover: async () => undefined,
|
|
});
|
|
expect(resolved).toBeUndefined();
|
|
});
|
|
|
|
it('probeServer distinguishes a live server from a dead one', async () => {
|
|
expect(await probeServer(server.url)).toBe(true);
|
|
expect(await probeServer('http://127.0.0.1:1', 300)).toBe(false);
|
|
});
|
|
});
|