/** * 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>; 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({ 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>; 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); }); });