File
Blob: src/worker/lib/ai/openai-compat.ts
| 1 | import type { AiUsage } from "@/shared/ai"; |
| 2 | import { |
| 3 | AiBackendError, |
| 4 | AiMisconfiguredError, |
| 5 | type AiChatMessage, |
| 6 | type AiChatOptions, |
| 7 | type AiClient, |
| 8 | type AiFrame, |
| 9 | type AiSummarizeOptions, |
| 10 | type AiSummarizeResult, |
| 11 | } from "@/worker/lib/ai/types"; |
| 12 | import { parseProviderSseFrames, parseProviderTokenUsage } from "@/worker/lib/ai/sse-transform"; |
| 13 | |
| 14 | interface OpenAiCompatConfig { |
| 15 | endpoint: string; |
| 16 | apiKey: string; |
| 17 | chatModel: string; |
| 18 | summarizeModel: string; |
| 19 | } |
| 20 | |
| 21 | export function createOpenAiCompatClient(config: OpenAiCompatConfig): AiClient { |
| 22 | if (!config.endpoint) { |
| 23 | throw new AiMisconfiguredError("BLAND_AI_OPENAI_ENDPOINT is required when BLAND_AI_MODE=openai-compat"); |
| 24 | } |
| 25 | const baseUrl = config.endpoint.replace(/\/$/, ""); |
| 26 | |
| 27 | return { |
| 28 | async chat(messages: AiChatMessage[], opts?: AiChatOptions): Promise<AsyncIterable<AiFrame>> { |
| 29 | const response = await fetch(`${baseUrl}/chat/completions`, { |
| 30 | method: "POST", |
| 31 | headers: chatHeaders(config.apiKey), |
| 32 | signal: opts?.signal, |
| 33 | body: JSON.stringify({ |
| 34 | model: config.chatModel, |
| 35 | messages, |
| 36 | stream: true, |
| 37 | stream_options: { include_usage: true }, |
| 38 | // Reasoning-enabled servers (llama.cpp --reasoning on, OpenAI o-series) spend a |
| 39 | // chunk of this budget on the thinking trace before emitting content, so the |
| 40 | // default has to comfortably fit reasoning + answer. Callers can override. |
| 41 | max_tokens: opts?.maxTokens ?? 1024, |
| 42 | temperature: opts?.temperature ?? 0.7, |
| 43 | }), |
| 44 | }); |
| 45 | |
| 46 | if (!response.ok || !response.body) { |
| 47 | await drainAndDiscard(response); |
| 48 | throw new AiBackendError(`openai-compat chat failed: ${response.status}`, "ai_chat_failed"); |
| 49 | } |
| 50 | |
| 51 | return parseProviderSseFrames(response.body, { |
| 52 | extractChunkText: extractOpenAiChunk, |
| 53 | extractUsage: extractOpenAiUsage, |
| 54 | errorLabel: "openai-compat stream error", |
| 55 | }); |
| 56 | }, |
| 57 | |
| 58 | async summarize(text: string, opts?: AiSummarizeOptions): Promise<AiSummarizeResult> { |
| 59 | const response = await fetch(`${baseUrl}/chat/completions`, { |
| 60 | method: "POST", |
| 61 | headers: chatHeaders(config.apiKey), |
| 62 | body: JSON.stringify({ |
| 63 | model: config.summarizeModel, |
| 64 | stream: false, |
| 65 | // Same story as chat(): reasoning models need headroom for the trace before the |
| 66 | // summary is emitted. 768 comfortably fits a 256-token reasoning budget plus |
| 67 | // a 3–5 sentence summary. |
| 68 | max_tokens: opts?.maxTokens ?? 768, |
| 69 | temperature: 0.2, |
| 70 | messages: [ |
| 71 | { role: "system", content: "Summarize the user's document in 3–5 sentences. Be concise and faithful." }, |
| 72 | { role: "user", content: text }, |
| 73 | ], |
| 74 | }), |
| 75 | }); |
| 76 | |
| 77 | if (!response.ok) { |
| 78 | await drainAndDiscard(response); |
| 79 | throw new AiBackendError(`openai-compat summarize failed: ${response.status}`, "ai_summarize_failed"); |
| 80 | } |
| 81 | |
| 82 | const body = (await response.json()) as { |
| 83 | choices?: Array<{ message?: { content?: string } }>; |
| 84 | usage?: unknown; |
| 85 | }; |
| 86 | const content = body.choices?.[0]?.message?.content; |
| 87 | if (typeof content !== "string" || content.length === 0) { |
| 88 | throw new AiBackendError("OpenAI-compat summarize returned empty body", "ai_summarize_empty"); |
| 89 | } |
| 90 | const usage = parseProviderTokenUsage(body.usage); |
| 91 | return usage ? { summary: content.trim(), usage } : { summary: content.trim() }; |
| 92 | }, |
| 93 | }; |
| 94 | } |
| 95 | |
| 96 | function chatHeaders(apiKey: string): HeadersInit { |
| 97 | const headers: Record<string, string> = { |
| 98 | "content-type": "application/json", |
| 99 | accept: "text/event-stream", |
| 100 | }; |
| 101 | if (apiKey) headers.authorization = `Bearer ${apiKey}`; |
| 102 | return headers; |
| 103 | } |
| 104 | |
| 105 | // Upstream error bodies can echo prompts, request headers, or provider-internal |
| 106 | // detail. Drain the stream so the connection releases, and discard the content |
| 107 | // without keeping it on any field that could later flow into logs or clients. |
| 108 | async function drainAndDiscard(response: Response): Promise<void> { |
| 109 | try { |
| 110 | await response.text(); |
| 111 | } catch { |
| 112 | // ignore |
| 113 | } |
| 114 | } |
| 115 | |
| 116 | function extractOpenAiChunk(payload: unknown): string | null { |
| 117 | if (typeof payload !== "object" || payload === null) return null; |
| 118 | const choices = (payload as { choices?: unknown }).choices; |
| 119 | if (!Array.isArray(choices) || choices.length === 0) return null; |
| 120 | const first = choices[0]; |
| 121 | if (typeof first !== "object" || first === null) return null; |
| 122 | const delta = (first as { delta?: unknown }).delta; |
| 123 | if (typeof delta !== "object" || delta === null) return null; |
| 124 | const content = (delta as { content?: unknown }).content; |
| 125 | return typeof content === "string" && content.length > 0 ? content : null; |
| 126 | } |
| 127 | |
| 128 | function extractOpenAiUsage(payload: unknown): AiUsage | null { |
| 129 | if (typeof payload !== "object" || payload === null) return null; |
| 130 | const usage = (payload as { usage?: unknown }).usage; |
| 131 | return parseProviderTokenUsage(usage); |
| 132 | } |