Files
Codeman/test/tui/tui-client-sse.test.ts
T
Codeman maintainer e140132e45 feat: add the TUI's API, SSE and degraded-mode client
Everything the dashboard needs from outside the process, behind one typed
surface, so the app loop stays a loop. It is a client of the running server and
nothing else: rows come from the unified list, blocked states from the
approvals inbox, and answering goes through the endpoint that re-captures the
pane and refuses with a 409 when the dialog has already been answered in tmux.
That refusal is a typed result rather than an exception, because a human
beating you to a prompt is normal operation.

Discovery mirrors the daemon probe (`CODEMAN_API_URL`, else loopback on
`CODEMAN_PORT`, self-signed TLS accepted) and credentials come from where
`codeman attach` already reads them. An explicit port outranks the ambient
`CODEMAN_API_URL`, which every managed session exports: a caller that named a
port must not be redirected at whatever server owns its shell.

Input is single-line and `\r`-terminated at this layer, so no caller can strand
text on an unsubmitted composer, and each send is tagged for the server's
exactly-once path. The event stream defaults to a `?sessions=` filter that
matches nothing, which drops the terminal firehose while lifecycle, hook and
approval events still arrive. A silent-but-open stream is caught by a watchdog
rather than a socket error, since that failure mode reports nothing at all.

With no server answering, sessions are listed from tmux on the instance socket
(argv, never a shell string) and decorated from a read-only peek at state.json,
which keeps the "the server died, get me to my sessions" path alive.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-22 14:13:58 +02:00

207 lines
8.1 KiB
TypeScript

/**
* @fileoverview Integration tests for the TUI's live-update stream.
*
* These run against a real loopback `text/event-stream` endpoint rather than a
* mocked socket, because the behaviours that matter here are all socket-level:
* a stream that ENDS, a stream that goes SILENT without erroring (the failure
* mode `EventSource` cannot see, which is why the server heartbeats), and a
* teardown that must leave no timer behind.
*/
import { describe, it, expect, beforeAll, afterAll, beforeEach, afterEach } from 'vitest';
import http from 'node:http';
import { TuiClient, type TuiApprovalEvent, type TuiSseStatusDetail } from '../../src/tui/tui-client.js';
const PORT = 3242;
const BASE_URL = `http://127.0.0.1:${PORT}`;
interface Connection {
url: string;
headers: http.IncomingHttpHeaders;
res: http.ServerResponse;
}
const connections: Connection[] = [];
/** Flipped by a test that wants every connect attempt to fail. */
let refuse = false;
let server: http.Server;
let client: TuiClient | null = null;
function frame(event: string, data: unknown): string {
return `event: ${event}\ndata: ${JSON.stringify(data)}\n\n`;
}
async function until(predicate: () => boolean, timeoutMs = 3000): Promise<void> {
const deadline = Date.now() + timeoutMs;
while (Date.now() < deadline) {
if (predicate()) return;
await new Promise((resolve) => setTimeout(resolve, 10));
}
throw new Error('condition not met before the deadline');
}
beforeAll(async () => {
server = http.createServer((req, res) => {
if (!req.url?.startsWith('/api/events')) {
res.writeHead(404).end();
return;
}
if (refuse) {
res.writeHead(503).end('busy');
return;
}
res.writeHead(200, { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache' });
// Node holds headers back until the first body write; the real server sends
// an `init` frame immediately, so flush to match it. Without this the
// client never sees a response and every test here waits forever.
res.flushHeaders();
connections.push({ url: req.url, headers: req.headers, res });
});
await new Promise<void>((resolve) => server.listen(PORT, '127.0.0.1', resolve));
});
afterAll(async () => {
await new Promise<void>((resolve) => server.close(() => resolve()));
});
beforeEach(() => {
connections.length = 0;
refuse = false;
});
afterEach(() => {
client?.close();
client = null;
for (const connection of connections) connection.res.end();
});
describe('subscribeEvents', () => {
it('routes each frame to the handler that owns it', async () => {
const resyncs: string[] = [];
const approvals: TuiApprovalEvent[] = [];
let planUsage: unknown = null;
let init: unknown = null;
client = new TuiClient({ baseUrl: BASE_URL, password: 's3cret' });
client.subscribeEvents({
onInit: (state) => {
init = state;
},
onResync: (event) => resyncs.push(event),
onApproval: (event) => approvals.push(event),
onPlanUsage: (usage) => {
planUsage = usage;
},
});
await until(() => connections.length === 1);
const { res } = connections[0];
res.write(frame('init', { version: '9.9.9', planUsage: { fiveHour: { usedPercentage: 5, resetAt: 1 } } }));
res.write(frame('session:created', { id: 'a' }));
// Split across writes on purpose: the parser must not need frame-aligned reads.
res.write('event: approval:pending\ndata: {"id":"a:1","sessionId":"a",');
res.write('"kind":"permission","createdAt":7}\n\n');
res.write(frame('session:terminal', { id: 'a', data: 'noise' }));
res.write(frame('sse:heartbeat', { t: 1 }));
res.write(frame('session:statusTelemetry', { sessionId: 'a', fiveHour: { usedPercentage: 41, resetAt: 2 } }));
await until(() => planUsage !== null);
expect(init).toEqual({ version: '9.9.9', planUsage: { fiveHour: { usedPercentage: 5, resetAt: 1 } } });
expect(approvals).toEqual([
{ kind: 'pending', item: { id: 'a:1', sessionId: 'a', kind: 'permission', createdAt: 7 } },
]);
expect(planUsage).toEqual({ sessionId: 'a', fiveHour: { usedPercentage: 41, resetAt: 2 } });
// The approval also regrouped a row, so it resyncs too. Terminal and
// heartbeat frames never do.
expect(resyncs).toEqual(['session:created', 'approval:pending']);
});
it('suppresses the terminal firehose by default and carries the auth header', async () => {
client = new TuiClient({ baseUrl: BASE_URL, password: 's3cret' });
client.subscribeEvents({});
await until(() => connections.length === 1);
expect(connections[0].url).toBe('/api/events?sessions=tui-no-terminal');
expect(connections[0].headers.authorization).toBe(`Basic ${Buffer.from('admin:s3cret').toString('base64')}`);
expect(connections[0].headers.accept).toBe('text/event-stream');
});
it('subscribes to the terminal stream of named sessions when asked', async () => {
client = new TuiClient({ baseUrl: BASE_URL });
client.subscribeEvents({}, { sessionIds: ['a', 'b'] });
await until(() => connections.length === 1);
expect(connections[0].url).toBe('/api/events?sessions=a%2Cb');
});
it('reconnects when the stream ends', async () => {
const statuses: Array<[string, TuiSseStatusDetail]> = [];
client = new TuiClient({ baseUrl: BASE_URL });
const stream = client.subscribeEvents(
{ onStatus: (status, detail) => statuses.push([status, detail]) },
{ baseBackoffMs: 10, maxBackoffMs: 20 }
);
await until(() => stream.status === 'connected');
connections[0].res.end();
await until(() => connections.length === 2 && stream.status === 'connected');
expect(statuses.map(([status]) => status)).toEqual(['connected', 'reconnecting', 'connected']);
expect(statuses[1][1].message).toBeTruthy();
});
it('reconnects when a live stream goes silent, which no socket error reports', async () => {
client = new TuiClient({ baseUrl: BASE_URL });
client.subscribeEvents({}, { staleTimeoutMs: 150, checkIntervalMs: 25, baseBackoffMs: 10, maxBackoffMs: 20 });
await until(() => connections.length === 1);
// The server holds the connection open and says nothing: exactly the case
// the watchdog exists for.
await until(() => connections.length === 2);
expect(connections).toHaveLength(2);
});
it('recommends polling once connecting keeps failing', async () => {
refuse = true;
const details: TuiSseStatusDetail[] = [];
client = new TuiClient({ baseUrl: BASE_URL });
const stream = client.subscribeEvents(
{ onStatus: (_status, detail) => details.push(detail) },
{ baseBackoffMs: 10, maxBackoffMs: 20, pollingAfterFailures: 2 }
);
await until(() => details.length >= 2);
expect(details[0]).toMatchObject({ attempt: 1, recommendPolling: false });
expect(details[1]).toMatchObject({ attempt: 2, recommendPolling: true });
expect(stream.recommendPolling).toBe(true);
expect(stream.status).toBe('reconnecting');
});
it('stops reconnecting after close, so the process can exit', async () => {
client = new TuiClient({ baseUrl: BASE_URL });
const stream = client.subscribeEvents({}, { baseBackoffMs: 10, maxBackoffMs: 20 });
await until(() => connections.length === 1);
connections[0].res.end();
stream.close();
const seen = connections.length;
await new Promise((resolve) => setTimeout(resolve, 120));
expect(connections.length).toBe(seen);
});
it('closes every stream the client opened', async () => {
client = new TuiClient({ baseUrl: BASE_URL });
client.subscribeEvents({});
client.subscribeEvents({});
await until(() => connections.length === 2);
client.close();
await until(() => connections.every((connection) => connection.res.socket === null || connection.res.destroyed));
const seen = connections.length;
await new Promise((resolve) => setTimeout(resolve, 120));
expect(connections.length).toBe(seen);
});
it('refuses to subscribe before the client knows where the server is', () => {
const disconnected = new TuiClient({ port: 3999 });
expect(() => disconnected.subscribeEvents({})).toThrow(/connect\(\)/);
});
});