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)) ?? []; } export function classifySessionProviderErrorMessage(message) { const normalized = String(message ?? '').trim(); if (!normalized) return 'SESSION_REPLY_ERROR'; if (/tool_calls|tool_call_id|insufficient tool messages/i.test(normalized)) { return 'SESSION_TOOL_HISTORY_POISONED'; } // DeepSeek/Kimi thinking mode: assistant tool-call turns must replay reasoning_content. if (/reasoning_content/i.test(normalized)) { return 'SESSION_REASONING_CONTENT_POISONED'; } return 'SESSION_REPLY_ERROR'; } function providerReplyError(event) { if (event?.type !== 'Message') return null; const text = messageContentItems(event) .filter((item) => item?.type === 'text') .map((item) => String(item?.text ?? '')) .join('') .trim(); if (!/^Ran into this error:\s*/i.test(text)) return null; const normalized = text .replace(/^Ran into this error:\s*/i, '') .replace(/\n\nPlease retry if you think this is a transient or recoverable error\.?\s*$/i, '') .trim(); const error = new Error(normalized || 'session reply failed'); error.code = classifySessionProviderErrorMessage(normalized); return error; } 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 pendingProviderError = null; 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; pendingProviderError = providerReplyError(event) ?? pendingProviderError; toolEvidence.observe(event); onEvent?.(event); if (event.type === 'Error') { const err = new Error(String(event.error ?? 'session reply failed')); err.code = classifySessionProviderErrorMessage(err.message); err.retryable = false; throw err; } if (event.type === 'Finish') { if (pendingProviderError) throw pendingProviderError; 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; }