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 ( /unknown variant [`']?image_url[`']?.*expected [`']?text/i.test(normalized) || /failed to deserialize.*image_url/i.test(normalized) ) { return 'SESSION_VISUAL_CONTEXT_UNSUPPORTED'; } 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); if (error.code === 'SESSION_VISUAL_CONTEXT_UNSUPPORTED') { error.retryable = false; } 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, requireRequestActivityBeforeUnscopedFinish = false, emptyFinishGraceMs = 20_000, } = {}, ) { 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 sawRequestScopedActivity = false; let buffer = ''; const deadline = Date.now() + Math.max(1, Number(timeoutMs) || 1); let timeoutHandle = null; const armTimeout = (delayMs, { code, message }) => { if (timeoutHandle) clearTimeout(timeoutHandle); timeoutHandle = setTimeout(() => { const error = new Error(message); error.code = code; error.retryable = false; reader.destroy(error); }, Math.max(1, delayMs)); }; const armOriginalDeadline = () => { armTimeout(Math.max(1, deadline - Date.now()), { code: 'SESSION_REPLY_TIMEOUT', message: `session reply timed out after ${timeoutMs}ms`, }); }; armOriginalDeadline(); try { 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; const routingId = event?.chat_request_id ?? event?.request_id ?? null; if (requestId && routingId === requestId) { sawRequestScopedActivity = true; } 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 === 'Message') { armOriginalDeadline(); } if (event.type !== 'Finish') continue; const tokenState = event.token_state ?? null; const tokenActivity = [ tokenState?.inputTokens, tokenState?.outputTokens, tokenState?.totalTokens, tokenState?.accumulatedInputTokens, tokenState?.accumulatedOutputTokens, tokenState?.accumulatedTotalTokens, ].some((value) => Number(value) > 0); const toolActivity = toolEvidence.snapshot().calls.length > 0; if ( requireRequestActivityBeforeUnscopedFinish && requestId && ( (!routingId && !sawRequestScopedActivity) || (!tokenActivity && !toolActivity) ) ) { armTimeout( Math.min( Math.max(1, deadline - Date.now()), Math.max(1, Number(emptyFinishGraceMs) || 1), ), { code: 'SESSION_EMPTY_FINISH', message: 'session emitted an empty Finish without model or tool activity', }, ); continue; } if (pendingProviderError) throw pendingProviderError; return { finishEvent: event, tokenState, toolEvidence: toolEvidence.snapshot(), }; } } } finally { if (timeoutHandle) clearTimeout(timeoutHandle); } const err = new Error('session event stream ended before Finish'); err.code = 'SESSION_REPLY_INCOMPLETE'; err.retryable = false; throw err; }