import assert from 'node:assert/strict'; import test from 'node:test'; import { extractRunFromStreamEvent, formatRunStreamSseChunk, isRunStreamReplayEnabled, isTerminalRunStatus, projectRunStreamEvents, shouldEmitRunUpdateForStreamEvent, } from './agent-run-stream.mjs'; test('formatRunStreamSseChunk includes SSE id for replay cursor', () => { const chunk = formatRunStreamSseChunk({ id: 'evt-1', event: 'run', data: { run: { id: 'run-1', status: 'running' } }, }); assert.match(chunk, /^id: evt-1\n/); assert.match(chunk, /event: run\n/); assert.match(chunk, /"status":"running"/); }); test('projectRunStreamEvents prefers run_snapshot payloads', () => { const { emittedRuns, latestRun } = projectRunStreamEvents([ { id: 'e1', eventType: 'queued', data: null, createdAt: 1 }, { id: 'e2', eventType: 'run_snapshot', data: { run: { id: 'run-1', status: 'running', sessionId: 's1' } }, createdAt: 2, }, { id: 'e3', eventType: 'succeeded', data: null, createdAt: 3 }, ]); assert.equal(emittedRuns.length, 2); assert.equal(latestRun?.sessionId, 's1'); assert.equal(extractRunFromStreamEvent(emittedRuns[0] && { eventType: 'run_snapshot', data: { run: emittedRuns[0].run }, })?.sessionId, 's1'); }); test('shouldEmitRunUpdateForStreamEvent covers status and snapshot events', () => { assert.equal(shouldEmitRunUpdateForStreamEvent('run_snapshot'), true); assert.equal(shouldEmitRunUpdateForStreamEvent('worker_heartbeat'), false); assert.equal(isTerminalRunStatus('succeeded'), true); assert.equal(isRunStreamReplayEnabled({ MEMIND_RUN_STREAM_REPLAY: '0' }), false); assert.equal(isRunStreamReplayEnabled({ MEMIND_RUN_STREAM_REPLAY: '1' }), true); });