import Anthropic from '@anthropic-ai/sdk'; import { z } from 'zod'; import { estimateUsd, addUsage, ZERO_USAGE } from '../pricing.js'; import { recordCost } from '../costs.js'; import { AiNotConfiguredError, AiRefusalError, type ChatMessage, type Completion, type CompletionRequest, type ContentPart, type EmbedRequest, type EmbedResult, type ExtractRequest, type ExtractResult, type ModelProvider, type Role, type StreamDelta, type ToolAgentRequest, type Usage } from '../types.js'; export interface AnthropicProviderOptions { apiKey: string; /** model per role; defaults below */ models?: Partial>; /** stable system-prompt prefix cached across calls */ defaultTimeoutMs?: number; } /** * Default role → model mapping. Fast/cheap for bulk pipeline work, strong for research and vision. * All ids are the exact strings from the current model table (no date suffixes). */ export const ANTHROPIC_DEFAULT_MODELS: Record = { normalize: 'claude-haiku-4-5', classify: 'claude-haiku-4-5', summarize: 'claude-haiku-4-5', resolve: 'claude-sonnet-5', research: 'claude-opus-5', vision: 'claude-opus-5', embed: '', }; function usageOf(u: Anthropic.Usage | Anthropic.MessageDeltaUsage | null | undefined): Usage { return { inputTokens: u?.input_tokens ?? 0, outputTokens: u?.output_tokens ?? 0, cacheReadTokens: u?.cache_read_input_tokens ?? 0, cacheWriteTokens: u?.cache_creation_input_tokens ?? 0, }; } function toBlocks(content: string | ContentPart[]): string | Anthropic.ContentBlockParam[] { if (typeof content === 'string') return content; return content.map((p): Anthropic.ContentBlockParam => { if (p.type === 'text') return { type: 'text', text: p.text }; const img = p.image; if (img.url) return { type: 'image', source: { type: 'url', url: img.url } }; return { type: 'image', source: { type: 'base64', media_type: img.mediaType, data: img.data ?? '' } }; }); } function toMessages(messages: ChatMessage[]): Anthropic.MessageParam[] { return messages.map((m) => ({ role: m.role, content: toBlocks(m.content) })); } /** Haiku 4.5 still uses budgeted thinking; current-generation models use adaptive thinking. */ function thinkingFor(model: string, effort: string | undefined): Anthropic.ThinkingConfigParam | undefined { if (model.startsWith('claude-haiku')) return undefined; if (effort === 'low') return { type: 'adaptive' }; return { type: 'adaptive' }; } function outputConfig(model: string, effort: string | undefined, format?: Anthropic.JSONOutputFormat): Anthropic.OutputConfig | undefined { const cfg: Anthropic.OutputConfig = {}; if (effort && !model.startsWith('claude-haiku')) cfg.effort = effort as Anthropic.OutputConfig['effort']; if (format) cfg.format = format; return Object.keys(cfg).length ? cfg : undefined; } export class AnthropicProvider implements ModelProvider { readonly id = 'anthropic'; private readonly client: Anthropic; private readonly models: Record; constructor(opts: AnthropicProviderOptions) { if (!opts.apiKey) throw new AiNotConfiguredError('ANTHROPIC_API_KEY missing'); this.client = new Anthropic({ apiKey: opts.apiKey, timeout: opts.defaultTimeoutMs ?? 10 * 60_000, maxRetries: 2 }); this.models = { ...ANTHROPIC_DEFAULT_MODELS, ...(opts.models ?? {}) }; } supports(role: Role): boolean { return role !== 'embed'; } modelFor(role: Role): string { return this.models[role] || ANTHROPIC_DEFAULT_MODELS[role]; } private async account(role: Role, model: string, usage: Usage, cost?: CompletionRequest['cost']): Promise { const usdEst = estimateUsd(model, usage); await recordCost({ provider: this.id, model, role, usage, usdEst, context: cost }); return usdEst; } async complete(role: Role, req: CompletionRequest): Promise { const model = this.modelFor(role); const system = req.json ? `${req.system ?? ''}\nRespond with a single JSON object and nothing else.`.trim() : req.system; const res = await this.client.messages.create({ model, max_tokens: req.maxTokens ?? 4096, ...(system ? { system: [{ type: 'text', text: system, cache_control: { type: 'ephemeral' } }] } : {}), messages: toMessages(req.messages), ...(thinkingFor(model, req.effort) ? { thinking: thinkingFor(model, req.effort) } : {}), ...(outputConfig(model, req.effort) ? { output_config: outputConfig(model, req.effort) } : {}), ...(req.stopSequences ? { stop_sequences: req.stopSequences } : {}), }); const usage = usageOf(res.usage); const usdEst = await this.account(role, model, usage, req.cost); if (res.stop_reason === 'refusal') { const d = res.stop_details; return { text: '', model, provider: this.id, usage, usdEst, stopReason: 'refusal', refusal: { category: d?.category ?? null, explanation: d?.explanation ?? null } }; } const text = res.content .filter((b): b is Anthropic.TextBlock => b.type === 'text') .map((b) => b.text) .join(''); return { text, model, provider: this.id, usage, usdEst, stopReason: res.stop_reason, refusal: null }; } async *stream(role: Role, req: CompletionRequest): AsyncIterable { yield* this.runTools(role, { ...req, tools: [], execute: async () => null }); } async *runTools(role: Role, req: ToolAgentRequest): AsyncIterable { const model = this.modelFor(role); const messages = toMessages(req.messages); const tools: Anthropic.Tool[] = req.tools.map((t) => ({ name: t.name, description: t.description, input_schema: t.inputSchema as Anthropic.Tool.InputSchema, })); const maxIter = req.maxIterations ?? 8; let total: Usage = ZERO_USAGE; let iterations = 0; let stopReason: string | null = null; try { while (iterations < maxIter) { iterations++; const stream = this.client.messages.stream( { model, max_tokens: req.maxTokens ?? 8192, ...(req.system ? { system: [{ type: 'text', text: req.system, cache_control: { type: 'ephemeral' } }] } : {}), messages, ...(tools.length ? { tools } : {}), ...(thinkingFor(model, req.effort) ? { thinking: { ...thinkingFor(model, req.effort)!, display: 'summarized' } as Anthropic.ThinkingConfigParam } : {}), ...(outputConfig(model, req.effort) ? { output_config: outputConfig(model, req.effort) } : {}), }, { signal: req.signal }, ); const queue: StreamDelta[] = []; let resolveWake: (() => void) | null = null; const wake = () => { resolveWake?.(); resolveWake = null; }; stream.on('text', (delta) => { queue.push({ type: 'text', text: delta }); wake(); }); stream.on('thinking', (delta) => { queue.push({ type: 'thinking', text: delta }); wake(); }); let finished = false; let failure: unknown = null; const finalP = stream .finalMessage() .catch((err) => { failure = err; return null; }) .finally(() => { finished = true; wake(); }); while (!finished || queue.length) { if (queue.length) { yield queue.shift()!; continue; } await new Promise((r) => { resolveWake = r; }); } const message = await finalP; if (failure || !message) throw failure ?? new Error('stream ended without a message'); const u = usageOf(message.usage); total = addUsage(total, u); yield { type: 'usage', usage: u, usdEst: estimateUsd(model, u), model }; stopReason = message.stop_reason; if (message.stop_reason === 'refusal') { const d = message.stop_details; throw new AiRefusalError(d?.category ?? null, d?.explanation ?? null); } if (message.stop_reason === 'pause_turn') { messages.push({ role: 'assistant', content: message.content }); continue; } const toolUses = message.content.filter((b): b is Anthropic.ToolUseBlock => b.type === 'tool_use'); if (message.stop_reason !== 'tool_use' || toolUses.length === 0) break; messages.push({ role: 'assistant', content: message.content }); const results: Anthropic.ToolResultBlockParam[] = []; for (const tu of toolUses) { yield { type: 'tool_call', id: tu.id, name: tu.name, input: tu.input }; const started = Date.now(); let output: unknown; let isError = false; try { output = await req.execute(tu.name, tu.input); } catch (err) { isError = true; output = { error: err instanceof Error ? err.message : String(err) }; } yield { type: 'tool_result', id: tu.id, name: tu.name, output, isError, durationMs: Date.now() - started }; results.push({ type: 'tool_result', tool_use_id: tu.id, content: typeof output === 'string' ? output : JSON.stringify(output ?? null), is_error: isError }); } messages.push({ role: 'user', content: results }); } } finally { await this.account(role, model, total, req.cost); } yield { type: 'done', stopReason }; } async extract(role: Role, req: ExtractRequest): Promise> { const model = this.modelFor(role); const schema = z.toJSONSchema(req.schema as z.ZodType, { target: 'draft-2020-12', io: 'output' }) as Record; const format: Anthropic.JSONOutputFormat = { type: 'json_schema', schema: stripUnsupported(schema) }; const content = typeof req.input === 'string' ? `${req.prompt}\n\n${req.input}` : [{ type: 'text', text: req.prompt } as ContentPart, ...req.input]; const res = await this.client.messages.create({ model, max_tokens: req.maxTokens ?? 4096, ...(req.system ? { system: [{ type: 'text', text: req.system, cache_control: { type: 'ephemeral' } }] } : {}), messages: toMessages([{ role: 'user', content }]), ...(thinkingFor(model, req.effort) ? { thinking: thinkingFor(model, req.effort) } : {}), output_config: outputConfig(model, req.effort, format)!, }); const usage = usageOf(res.usage); const usdEst = await this.account(role, model, usage, req.cost); if (res.stop_reason === 'refusal') throw new AiRefusalError(res.stop_details?.category ?? null, res.stop_details?.explanation ?? null); const text = res.content .filter((b): b is Anthropic.TextBlock => b.type === 'text') .map((b) => b.text) .join(''); const parsed = req.schema.safeParse(JSON.parse(text)); if (!parsed.success) throw new Error(`extract: model output did not match schema: ${parsed.error.issues.map((i) => `${i.path.join('.')}: ${i.message}`).join('; ')}`); const data = parsed.data as T & { confidence?: number }; const confidence = typeof data.confidence === 'number' ? Math.max(0, Math.min(1, data.confidence)) : 1; return { data: parsed.data, confidence, model, provider: this.id, usage, usdEst }; } vision(req: ExtractRequest): Promise> { return this.extract('vision', req); } async embed(_req: EmbedRequest): Promise { throw new AiNotConfiguredError('Anthropic has no embeddings endpoint; configure OPENAI_API_KEY or a local embedding model'); } } /** Structured-output schemas reject a few JSON-schema keywords; strip them recursively. */ function stripUnsupported(schema: Record): Record { const drop = new Set(['$schema', 'default', 'minLength', 'maxLength', 'minimum', 'maximum', 'exclusiveMinimum', 'exclusiveMaximum', 'pattern', 'format', 'minItems', 'maxItems', 'multipleOf']); const walk = (node: unknown): unknown => { if (Array.isArray(node)) return node.map(walk); if (node && typeof node === 'object') { const out: Record = {}; for (const [k, v] of Object.entries(node as Record)) { if (drop.has(k)) continue; out[k] = walk(v); } if (out.type === 'object' && out.properties && out.additionalProperties === undefined) out.additionalProperties = false; return out; } return node; }; return walk(schema) as Record; }