/** * Claude Code SDK → Unified Type Converters * * Converts SDK-specific message types to EngineOutput (unified types). * Uses proper SDK types (SDKAssistantMessage, SDKUserMessage, etc.) * instead of opaque Record casts. * * Mirrors the pattern in opencode/message-converter.ts. */ import type { SDKMessage, SDKAssistantMessage, SDKUserMessage, SDKPartialAssistantMessage, SDKResultSuccess, SDKResultError, SDKSystemMessage, SDKCompactBoundaryMessage, SDKRateLimitEvent, SDKModelRefusalFallbackMessage, SDKModelRefusalNoFallbackMessage, } from '@anthropic-ai/claude-agent-sdk'; // Infer block/usage types from the agent SDK so we stay in sync with whichever // @anthropic-ai/sdk version claude-agent-sdk bundles internally. type BetaContentBlock = SDKAssistantMessage['message']['content'][number]; type BetaUsage = SDKAssistantMessage['message']['usage']; import type { EngineOutput, UserMessage, AssistantMessage, ReasoningMessage, CompactBoundaryMessage, UserContentBlock, AssistantContentBlock, ToolUseBlock, TextDeltaEvent, StreamLifecycleEvent, SuccessResultEvent, ErrorResultEvent, SystemInitEvent, RateLimitEvent, NotificationEvent, TokenUsage, StopReason, } from '$shared/types/unified'; import { toCanonicalToolName } from '$shared/types/unified'; // ============================================================ // Helper Mappers // ============================================================ /** Map SDK stop_reason to unified StopReason */ function mapStopReason(sdkStop: string | null | undefined): StopReason | null { switch (sdkStop) { case 'end_turn': return 'end_turn'; case 'tool_use': return 'tool_use'; case 'max_tokens': return 'max_tokens'; case 'interrupted': return 'interrupted'; default: return sdkStop ? 'end_turn' : null; } } /** Map SDK BetaUsage to unified TokenUsage */ function mapUsage(sdkUsage: BetaUsage | null | undefined): TokenUsage | null { if (!sdkUsage) return null; return { inputTokens: sdkUsage.input_tokens || 0, outputTokens: sdkUsage.output_tokens || 0, cacheCreationInputTokens: sdkUsage.cache_creation_input_tokens || 0, cacheReadInputTokens: sdkUsage.cache_read_input_tokens || 0, }; } /** Map raw usage from result messages (NonNullableUsage shape) */ function mapRawUsage(raw: Record | null | undefined): TokenUsage { if (!raw) return { inputTokens: 0, outputTokens: 0, cacheCreationInputTokens: 0, cacheReadInputTokens: 0 }; return { inputTokens: raw.input_tokens || 0, outputTokens: raw.output_tokens || 0, cacheCreationInputTokens: raw.cache_creation_input_tokens || 0, cacheReadInputTokens: raw.cache_read_input_tokens || 0, }; } // ============================================================ // Tool Input Converter (snake_case → camelCase) // ============================================================ /** Grep option mapping for dash-prefixed parameters → camelCase */ const GREP_OPTION_MAP: Record = { '-i': 'caseInsensitive', '-n': 'lineNumbers', '-A': 'afterContext', '-B': 'beforeContext', '-C': 'context', }; /** Convert snake_case string to camelCase */ function snakeToCamel(str: string): string { return str.replace(/_([a-z])/g, (_, c: string) => c.toUpperCase()); } /** * Convert raw SDK tool input to unified camelCase format. * Handles standard snake_case → camelCase and Grep's dash-prefixed options. * Already-camelCase properties (e.g. multiline, activeForm) pass through unchanged. */ function convertToolInput(toolName: string, raw: Record): Record { const converted: Record = {}; for (const [key, value] of Object.entries(raw)) { if (toolName === 'Grep' && key in GREP_OPTION_MAP) { converted[GREP_OPTION_MAP[key]] = value; } else { converted[snakeToCamel(key)] = value; } } return converted; } // ============================================================ // Content Block Converters // ============================================================ /** Convert SDK BetaContentBlock[] to unified AssistantContentBlock[] (excluding thinking) */ export function mapAssistantContent(sdkContent: BetaContentBlock[]): AssistantContentBlock[] { const blocks: AssistantContentBlock[] = []; for (const block of sdkContent) { switch (block.type) { case 'text': blocks.push({ type: 'text', text: block.text }); break; case 'tool_use': { const rawName = block.name === 'Task' ? 'Agent' : block.name; blocks.push({ type: 'tool_use', id: block.id, name: toCanonicalToolName(rawName), input: convertToolInput(block.name, block.input as Record), result: null, subActivities: [], skillPrompt: null, interrupted: false, } as ToolUseBlock); break; } } } return blocks; } /** Convert SDK user content to unified UserContentBlock[] */ export function mapUserContent(sdkContent: string | Array>): UserContentBlock[] { if (typeof sdkContent === 'string') { return [{ type: 'text', text: sdkContent }]; } const blocks: UserContentBlock[] = []; for (const block of sdkContent) { switch (block.type) { case 'text': blocks.push({ type: 'text', text: block.text as string }); break; case 'tool_result': blocks.push({ type: 'tool_result', toolUseId: block.tool_use_id as string, content: (block.content as string) || '', isError: !!(block.is_error), }); break; case 'image': { const source = block.source as Record | undefined; blocks.push({ type: 'image', mediaType: (source?.media_type as string) || '', data: (source?.data as string) || '', }); break; } case 'document': { const source = block.source as Record | undefined; blocks.push({ type: 'document', mediaType: (source?.media_type as string) || '', data: (source?.data as string) || '', title: (block.title as string) || null, }); break; } } } if (blocks.length === 0) blocks.push({ type: 'text', text: '' }); return blocks; } // ============================================================ // Message Converters // ============================================================ /** Convert SDKAssistantMessage → ReasoningMessage (if thinking) + AssistantMessage */ export function convertAssistantMessage(msg: SDKAssistantMessage): EngineOutput[] { const sessionId = msg.session_id; const betaMessage = msg.message; const content = betaMessage.content; const outputs: EngineOutput[] = []; // Extract thinking blocks → ReasoningMessage. Skip emission when all // thinking blocks are empty (e.g. when adaptive thinking is in 'omitted' // mode and only signatures arrive); otherwise we'd persist a blank // reasoning entry that confuses the UI. const thinkingBlocks = content.filter(b => b.type === 'thinking'); if (thinkingBlocks.length > 0) { const reasoningText = thinkingBlocks .map(b => b.type === 'thinking' ? b.thinking : '') .join('\n') .trim(); if (reasoningText) { const reasoning: ReasoningMessage = { type: 'reasoning', createdAt: new Date().toISOString(), messageId: crypto.randomUUID(), sessionId, parent: { messageId: null, sessionId: null, toolUseId: null }, engine: { type: 'claude-code', provider: 'anthropic', model: { id: betaMessage.model || '', name: '' }, account: { id: 0, name: '' } }, text: reasoningText, }; outputs.push(reasoning); } } // Build AssistantMessage with non-thinking content const otherBlocks = content.filter(b => b.type !== 'thinking' && b.type !== 'redacted_thinking'); const assistantContent = mapAssistantContent( otherBlocks.length > 0 ? otherBlocks : [{ type: 'text', text: '', citations: null } as BetaContentBlock] ); const assistant: AssistantMessage = { type: 'assistant', createdAt: new Date().toISOString(), messageId: msg.uuid || crypto.randomUUID(), sessionId, parent: { messageId: null, sessionId: null, toolUseId: msg.parent_tool_use_id || null }, engine: { type: 'claude-code', provider: 'anthropic', model: { id: betaMessage.model || '', name: '' }, account: { id: 0, name: '' } }, content: assistantContent, stopReason: mapStopReason(betaMessage.stop_reason), usage: mapUsage(betaMessage.usage), }; outputs.push(assistant); return outputs; } /** Convert SDKUserMessage → UserMessage */ export function convertUserMessage(msg: SDKUserMessage): UserMessage { const sdkContent = msg.message.content; const content = mapUserContent(sdkContent as string | Array>); return { type: 'user', createdAt: new Date().toISOString(), messageId: msg.uuid || crypto.randomUUID(), sessionId: msg.session_id || null, parent: { messageId: null, sessionId: null, toolUseId: msg.parent_tool_use_id || null }, engine: { type: 'claude-code', provider: 'anthropic', model: { id: '', name: '' }, account: { id: 0, name: '' } }, sender: { id: '', name: '' }, content, synthetic: true, // SDK-generated tool_result messages }; } /** * Stateful per-stream tracking. The SDK emits content_block_stop without a * content_block payload, so we must remember which block indices belong to * thinking blocks to emit the correct `reasoning` flag on stop. A Map per * converter instance keeps state isolated between concurrent queries. */ export interface StreamConverterState { /** index → true when the block is a thinking block */ reasoningBlocks: Map; /** Workflow background task id → originating Workflow tool_use id. */ workflowParents: Map; } export function createStreamConverterState(): StreamConverterState { return { reasoningBlocks: new Map(), workflowParents: new Map(), }; } /** Convert SDKPartialAssistantMessage (stream_event) → TextDeltaEvent | StreamLifecycleEvent */ export function convertStreamEvent( msg: SDKPartialAssistantMessage, state: StreamConverterState, ): EngineOutput[] { const sessionId = msg.session_id; const event = msg.event; const outputs: EngineOutput[] = []; switch (event.type) { case 'message_start': case 'message_stop': { const lifecycle: StreamLifecycleEvent = { type: 'stream_event', event: event.type === 'message_start' ? 'start' : 'stop', sessionId, reasoning: false, }; outputs.push(lifecycle); break; } case 'content_block_start': { const blockType = event.content_block.type; const isReasoning = blockType === 'thinking' || blockType === 'redacted_thinking'; state.reasoningBlocks.set(event.index, isReasoning); if (isReasoning) { outputs.push({ type: 'stream_event', event: 'start', sessionId, reasoning: true } as StreamLifecycleEvent); } else if (blockType === 'text') { outputs.push({ type: 'stream_event', event: 'start', sessionId, reasoning: false } as StreamLifecycleEvent); } break; } case 'content_block_delta': { const delta = event.delta; if (delta.type === 'thinking_delta') { const text: TextDeltaEvent = { type: 'stream_event', event: 'delta', sessionId, text: delta.thinking || '', reasoning: true, }; outputs.push(text); } else if (delta.type === 'text_delta') { const text: TextDeltaEvent = { type: 'stream_event', event: 'delta', sessionId, text: delta.text || '', reasoning: false, }; outputs.push(text); } break; } case 'content_block_stop': { const isReasoning = state.reasoningBlocks.get(event.index) === true; state.reasoningBlocks.delete(event.index); outputs.push({ type: 'stream_event', event: 'stop', sessionId, reasoning: isReasoning } as StreamLifecycleEvent); break; } } return outputs; } /** Convert SDKResultSuccess → SuccessResultEvent */ export function convertResultSuccess(msg: SDKResultSuccess): SuccessResultEvent { return { type: 'result', subtype: 'success', sessionId: msg.session_id, numTurns: msg.num_turns || 0, totalCostUsd: msg.total_cost_usd || 0, usage: mapRawUsage(msg.usage as unknown as Record), stopReason: msg.stop_reason || null, }; } /** Convert SDKResultError → ErrorResultEvent */ export function convertResultError(msg: SDKResultError): ErrorResultEvent { return { type: 'result', subtype: msg.subtype === 'error_max_turns' ? 'error_max_turns' : msg.subtype === 'error_max_budget_usd' ? 'error_max_budget' : 'error_during_execution', sessionId: msg.session_id, errors: msg.errors || [], }; } /** Convert SDKSystemMessage (init) → SystemInitEvent */ export function convertSystemInit(msg: SDKSystemMessage): SystemInitEvent { return { type: 'system_init', sessionId: msg.session_id, model: msg.model || '', engine: 'claude-code', tools: msg.tools || [], mcpServers: (msg.mcp_servers || []).map(s => ({ name: s.name || '', status: (s.status as 'connected' | 'disconnected' | 'error') || 'disconnected', })), }; } function workflowActivityMessage( msg: { session_id: string; uuid?: string }, parentToolUseId: string, text: string, ): AssistantMessage { return { type: 'assistant', createdAt: new Date().toISOString(), messageId: msg.uuid || crypto.randomUUID(), sessionId: msg.session_id, parent: { messageId: null, sessionId: null, toolUseId: parentToolUseId }, engine: { type: 'claude-code', provider: 'anthropic', model: { id: '', name: '' }, account: { id: 0, name: '' } }, content: [{ type: 'text', text }], stopReason: null, usage: null, }; } /** Convert SDKCompactBoundaryMessage → CompactBoundaryMessage */ export function convertCompactBoundary(msg: SDKCompactBoundaryMessage): CompactBoundaryMessage { return { type: 'compact_boundary', createdAt: new Date().toISOString(), messageId: msg.uuid || crypto.randomUUID(), sessionId: msg.session_id, parent: { messageId: null, sessionId: null, toolUseId: null }, engine: { type: 'claude-code', provider: 'anthropic', model: { id: '', name: '' }, account: { id: 0, name: '' } }, trigger: msg.compact_metadata.trigger, preTokens: msg.compact_metadata.pre_tokens || 0, }; } /** * Convert a model-refusal banner → NotificationEvent. * * Both subtypes arrived in @anthropic-ai/claude-agent-sdk 0.3.229 and were * previously dropped by the `default` arm, which is the worst outcome for the * two cases users actually notice: a turn answered by a *different* model than * the one they picked, and a turn that produced nothing at all. Neither is an * error the stream should abort on, so they surface as toasts. * * NOT handled: `retracted_message_uuids`, the fallback's list of already- * delivered messages the CLI has retracted. Honouring it means deleting * persisted messages by id across adapter → stream-manager → WS → UI, so the * refused partial currently stays in the transcript above the fallback answer. * See docs/lessons-learned.md §10.25. */ function convertModelRefusalFallback(msg: SDKModelRefusalFallbackMessage): NotificationEvent { // 'local' means only this response came from the fallback (a subagent or a // side question); 'session' means the session's model itself was swapped. const scope = msg.scope ?? 'session'; const swapped = scope === 'session' ? `The session model is now ${msg.fallback_model}.` : `Only this response came from ${msg.fallback_model}; the session model is unchanged.`; return { type: 'notification', sessionId: msg.session_id, level: 'warning', title: 'Model refused — answered by a fallback', message: `${msg.original_model} declined this request${msg.api_refusal_category ? ` (${msg.api_refusal_category})` : ''}. ${swapped}`, }; } /** Convert SDKModelRefusalNoFallbackMessage → NotificationEvent (no retry ran). */ function convertModelRefusalNoFallback(msg: SDKModelRefusalNoFallbackMessage): NotificationEvent { return { type: 'notification', sessionId: msg.session_id, level: 'error', title: 'Model refused this request', // `content` is the CLI's own human-readable banner; explanation is // unstable prose meant for display only. message: msg.api_refusal_explanation || msg.content, }; } /** Convert SDKRateLimitEvent → RateLimitEvent (returns null for non-actionable 'allowed' status) */ export function convertRateLimit(msg: SDKRateLimitEvent): RateLimitEvent | null { const info = msg.rate_limit_info; // Only emit for actual warnings/rejections — 'allowed' is just informational if (info.status === 'allowed') return null; return { type: 'rate_limit', sessionId: msg.session_id, status: info.status === 'rejected' ? 'rejected' : 'allowed_warning', utilization: info.utilization || 0, resetsAt: info.resetsAt || null, rateLimitType: info.rateLimitType || null, }; } // ============================================================ // Top-Level Converter // ============================================================ /** * Create a stateful per-stream SDK message converter. * * Returns a generator function that converts a single SDKMessage into zero or * more EngineOutput events. The closure tracks per-block reasoning state so * `content_block_stop` events can carry the correct `reasoning` flag (the SDK * does not include the block payload on stop). * * Use one converter per query; do not share across concurrent streams. */ export function createSdkMessageConverter(): (msg: SDKMessage) => Generator { const state = createStreamConverterState(); return function* convertSdkMessage(msg: SDKMessage): Generator { yield* dispatchSdkMessage(msg, state); }; } /** * Stateless dispatcher kept for backwards compatibility. Allocates a fresh * state on every call, so block-stop reasoning tracking only works within a * single SDK message. Prefer `createSdkMessageConverter()` for streaming. */ export function* convertSdkMessage(msg: SDKMessage): Generator { yield* dispatchSdkMessage(msg, createStreamConverterState()); } function* dispatchSdkMessage(msg: SDKMessage, state: StreamConverterState): Generator { switch (msg.type) { case 'assistant': for (const output of convertAssistantMessage(msg as SDKAssistantMessage)) { yield output; } break; case 'user': // SDKUserMessage and SDKUserMessageReplay both have type: 'user' yield convertUserMessage(msg as SDKUserMessage); break; case 'stream_event': for (const output of convertStreamEvent(msg as SDKPartialAssistantMessage, state)) { yield output; } break; case 'result': if (msg.subtype === 'success') { yield convertResultSuccess(msg as SDKResultSuccess); } else { yield convertResultError(msg as SDKResultError); } break; case 'system': { // Read subtype as a plain string: the SDKMessage `system` arm covers // more subtypes (task_*, etc.) than the SDKSystemMessage | // SDKCompactBoundaryMessage union, so narrowing on that union would // collapse the task branch to `never`. const subtype = (msg as { subtype?: string }).subtype; if (subtype === 'init') { yield convertSystemInit(msg as SDKSystemMessage); } else if (subtype === 'compact_boundary') { yield convertCompactBoundary(msg as SDKCompactBoundaryMessage); } else if (subtype === 'model_refusal_fallback') { yield convertModelRefusalFallback(msg as SDKModelRefusalFallbackMessage); } else if (subtype === 'model_refusal_no_fallback') { yield convertModelRefusalNoFallback(msg as SDKModelRefusalNoFallbackMessage); } else if (subtype === 'task_started') { const task = msg as unknown as { session_id: string; uuid?: string; task_id: string; tool_use_id?: string; task_type?: string; workflow_name?: string; description?: string }; if (task.task_type === 'local_workflow' && task.tool_use_id) { state.workflowParents.set(task.task_id, task.tool_use_id); const detail = `Workflow ${task.workflow_name || task.description || task.task_id} started.`; yield workflowActivityMessage(task, task.tool_use_id, detail); } } else if (subtype === 'task_notification') { // Only workflow tasks surface anything, and only inline in the // conversation. Backgrounded Agent/Bash tasks stay silent: their // outcome already lands in the transcript, so a terminal toast // ("Background task failed") is pure noise — a backgrounded // `grep` with no match is a non-zero exit, not a problem. const task = msg as unknown as { session_id: string; uuid?: string; task_id: string; tool_use_id?: string; status: string; summary?: string }; const parent = task.tool_use_id || state.workflowParents.get(task.task_id); if (parent && state.workflowParents.has(task.task_id)) { yield workflowActivityMessage(task, parent, `Workflow ${task.status}.`); state.workflowParents.delete(task.task_id); } } // task_progress is a repeated heartbeat; real child activity is replayed // from Workflow-owned transcripts. task_updated has no tool_use_id. break; } case 'rate_limit_event': { const rlEvent = convertRateLimit(msg as SDKRateLimitEvent); if (rlEvent) yield rlEvent; break; } // Transient SDK message types — skip default: break; } } // ============================================================ // Prompt Converter (Unified → SDK) // ============================================================ /** * Convert a UserMessage prompt to SDKUserMessage format for the Claude SDK. */ export function toSdkUserMessage(msg: UserMessage): SDKUserMessage { const sdkContent: Array> = []; for (const block of msg.content) { switch (block.type) { case 'text': sdkContent.push({ type: 'text', text: block.text }); break; case 'image': sdkContent.push({ type: 'image', source: { type: 'base64', media_type: block.mediaType, data: block.data }, }); break; case 'document': sdkContent.push({ type: 'document', source: { type: 'base64', media_type: block.mediaType, data: block.data }, title: block.title, }); break; case 'tool_result': sdkContent.push({ type: 'tool_result', tool_use_id: block.toolUseId, content: block.content, is_error: block.isError, }); break; } } // SDK accepts string for simple text-only prompts const content = sdkContent.length === 1 && sdkContent[0].type === 'text' ? sdkContent[0].text as string : sdkContent; return { type: 'user', uuid: msg.messageId, session_id: msg.sessionId, parent_tool_use_id: msg.parent.toolUseId, message: { role: 'user', content }, } as unknown as SDKUserMessage; }