agenthub/tests/realtime-regression.test.ts
chahinebrini 9551c90597 fix(realtime): discovery self-heal, polling fallback, message ack semantics, CRLF SSE parser
- 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)
2026-07-23 16:53:33 +02:00

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);
});
});