import { Readable } from 'node:stream'; function parseSseDataLines(frame) { let data = ''; for (const line of String(frame ?? '').split('\n')) { if (line.startsWith('data:')) data += line.slice(5).trim(); } return data; } export function parseSessionStreamEvent(frame) { const data = parseSseDataLines(frame); if (!data) return null; try { return JSON.parse(data); } catch { return null; } } export function eventMatchesRequest(event, requestId) { if (!requestId) return true; const routingId = event?.chat_request_id ?? event?.request_id ?? null; return !routingId || routingId === requestId; } function messageContentItems(event) { const candidates = [ event?.message?.content, event?.data?.message?.content, event?.content, ]; return candidates.find((content) => Array.isArray(content)) ?? []; } function parseGenerateImageResult(toolResult) { const content = toolResult?.value?.content; if (!Array.isArray(content)) return null; for (const item of content) { if (item?.type !== 'text' || !String(item.text ?? '').trim()) continue; try { const result = JSON.parse(String(item.text)); const mimeType = String(result?.source?.mimeType ?? result?.asset?.mimeType ?? '').toLowerCase(); if ( result?.ok === true && String(result?.jobId ?? '').trim() && /^image\/(?:png|jpeg|webp)$/.test(mimeType) ) { return { jobId: String(result.jobId), mimeType, assetId: String(result?.asset?.id ?? '').trim() || null, publicUrl: String(result?.asset?.publicUrl ?? '').trim() || null, workspaceRelativePath: String(result?.asset?.workspaceRelativePath ?? '').trim() || null, }; } } catch { // Ignore non-JSON tool text; it cannot prove a successful image_make delivery. } } return null; } function createToolEvidenceCollector() { const requestNames = new Map(); const calls = new Set(); const successfulCalls = new Set(); let generateImage = { called: false, succeeded: false }; return { observe(event) { for (const item of messageContentItems(event)) { if (item?.type === 'toolRequest') { const name = String(item?.toolCall?.value?.name ?? '').trim(); if (!name) continue; calls.add(name); if (item.id) requestNames.set(String(item.id), name); if (name.endsWith('generate_image')) generateImage = { called: true, succeeded: false }; continue; } if (item?.type !== 'toolResponse') continue; const name = requestNames.get(String(item.id ?? '')) ?? ''; const result = item.toolResult ?? {}; const succeeded = result.status === 'success' && result?.value?.isError !== true; if (name && succeeded) successfulCalls.add(name); if (!name.endsWith('generate_image')) continue; const generated = succeeded ? parseGenerateImageResult(result) : null; generateImage = generated ? { called: true, succeeded: true, ...generated } : { called: true, succeeded: false }; } }, snapshot() { return { calls: [...calls], successfulCalls: [...successfulCalls], generateImage, }; }, }; } export async function consumeSessionEventsUntilFinish( body, { requestId = null, timeoutMs = 15 * 60 * 1000, onEvent = null, } = {}, ) { if (!body) { const err = new Error('session event stream unavailable'); err.code = 'SESSION_EVENT_STREAM_UNAVAILABLE'; throw err; } const reader = Readable.fromWeb(body); const decoder = new TextDecoder(); const toolEvidence = createToolEvidenceCollector(); let buffer = ''; const deadline = Date.now() + Math.max(1, Number(timeoutMs) || 1); for await (const chunk of reader) { if (Date.now() > deadline) { const err = new Error(`session reply timed out after ${timeoutMs}ms`); err.code = 'SESSION_REPLY_TIMEOUT'; err.retryable = false; throw err; } buffer += decoder.decode(chunk, { stream: true }); const frames = buffer.split('\n\n'); buffer = frames.pop() ?? ''; for (const frame of frames) { const trimmed = frame.trim(); if (!trimmed || trimmed.startsWith(':')) continue; const event = parseSessionStreamEvent(trimmed); if (!event) continue; if (!eventMatchesRequest(event, requestId)) continue; toolEvidence.observe(event); onEvent?.(event); if (event.type === 'Error') { const err = new Error(String(event.error ?? 'session reply failed')); err.code = 'SESSION_REPLY_ERROR'; err.retryable = false; throw err; } if (event.type === 'Finish') { return { finishEvent: event, tokenState: event.token_state ?? null, toolEvidence: toolEvidence.snapshot(), }; } } } const err = new Error('session event stream ended before Finish'); err.code = 'SESSION_REPLY_INCOMPLETE'; err.retryable = false; throw err; }