From 1936a22c95d87b49178ca0883c2ac9e70448bf43 Mon Sep 17 00:00:00 2001 From: RainbowBird Date: Tue, 2 Jun 2026 15:28:25 +0800 Subject: [PATCH] refactor(server): split openai v1 route pipeline Separate OpenAI-compatible route wiring from chat, speech, catalog, billing, and telemetry pipeline code. Signed-off-by: RainbowBird Commit-Message-Assisted-by: Claude (via Claude Code) --- apps/server/package.json | 1 + apps/server/src/routes/openai/v1/billing.ts | 173 ++++ apps/server/src/routes/openai/v1/catalog.ts | 160 ++++ apps/server/src/routes/openai/v1/chat.ts | 456 ++++++++++ apps/server/src/routes/openai/v1/guards.ts | 19 + apps/server/src/routes/openai/v1/index.ts | 851 +----------------- apps/server/src/routes/openai/v1/response.ts | 15 + apps/server/src/routes/openai/v1/speech.ts | 152 ++++ apps/server/src/routes/openai/v1/telemetry.ts | 188 ++++ apps/server/src/routes/openai/v1/types.ts | 36 + pnpm-lock.yaml | 3 + 11 files changed, 1237 insertions(+), 817 deletions(-) create mode 100644 apps/server/src/routes/openai/v1/billing.ts create mode 100644 apps/server/src/routes/openai/v1/catalog.ts create mode 100644 apps/server/src/routes/openai/v1/chat.ts create mode 100644 apps/server/src/routes/openai/v1/guards.ts create mode 100644 apps/server/src/routes/openai/v1/response.ts create mode 100644 apps/server/src/routes/openai/v1/speech.ts create mode 100644 apps/server/src/routes/openai/v1/telemetry.ts create mode 100644 apps/server/src/routes/openai/v1/types.ts diff --git a/apps/server/package.json b/apps/server/package.json index 8daa135f8..1797d653f 100644 --- a/apps/server/package.json +++ b/apps/server/package.json @@ -58,6 +58,7 @@ "injeca": "catalog:", "ioredis": "^5.10.1", "jose": "catalog:", + "ofetch": "catalog:", "pg": "^8.20.0", "posthog-node": "catalog:", "resend": "^6.12.2", diff --git a/apps/server/src/routes/openai/v1/billing.ts b/apps/server/src/routes/openai/v1/billing.ts new file mode 100644 index 000000000..e6a4e6d8f --- /dev/null +++ b/apps/server/src/routes/openai/v1/billing.ts @@ -0,0 +1,173 @@ +import type { RevenueMetrics } from '../../../otel' +import type { ConfigKVService } from '../../../services/adapters/config-kv' +import type { UsageInfo } from '../../../services/domain/billing/billing' +import type { BillingService } from '../../../services/domain/billing/billing-service' +import type { FluxMeter } from '../../../services/domain/billing/flux-meter' +import type { FluxService } from '../../../services/domain/flux' + +import { calculateFluxFromUsage } from '../../../services/domain/billing/billing' +import { createPaymentRequiredError } from '../../../utils/error' +import { GEN_AI_ATTR_REQUEST_MODEL } from '../../../utils/observability' + +export interface ChatFluxDebitInput extends UsageInfo { + billingService: BillingService + revenue?: RevenueMetrics | null + userId: string + requestId: string + model: string + amount: number + stage: 'streaming' | 'non_streaming' + logger: { + withFields: (fields: Record) => { + warn: (message: string) => void + } + } +} + +export interface ChatBillingPolicy { + fallbackRate: number + fluxPer1kTokens?: number +} + +export interface TtsBillingAuthorization { + balance: number + inputChars: number +} + +export interface OpenAiRouteBilling { + authorizeChat: (userId: string) => Promise + authorizeTts: (userId: string, inputText: string) => Promise + priceChatUsage: (usage: UsageInfo, policy: ChatBillingPolicy) => number + recordChatDebitFailure: (input: { + amount: number + model: string + stage: 'streaming' | 'non_streaming' + }) => void + settleChat: (input: Omit) => Promise + settleTts: (input: { + userId: string + inputText: string + currentBalance: number + requestId: string + model: string + }) => Promise<{ fluxDebited: number }> +} + +export function createOpenAiRouteBilling(deps: { + billingService: BillingService + configKV: ConfigKVService + fluxService: FluxService + revenue?: RevenueMetrics | null + ttsMeter: FluxMeter +}): OpenAiRouteBilling { + // NOTICE: Billing is best-effort — chat flux is debited AFTER the LLM + // response is sent. This is a deliberate tradeoff: users get lower latency + // and uninterrupted streaming, at the cost of a small revenue leak when + // debit fails (e.g. DB timeout). Failed debits are logged at error level by + // the settlement call site. + // + // Pre-flight gates on `balance >= fallbackRate` (not just `> 0`) because + // streaming providers that don't echo `usage` cause every billable request + // to fall back to `FLUX_PER_REQUEST`. Without this gate, a user sitting on + // `0 < balance < fallbackRate` could spawn N parallel requests that each + // pass the loose `>0` check, complete the stream, and race on the debit. + async function authorizeChat(userId: string): Promise { + const fallbackRate = await deps.configKV.getOrThrow('FLUX_PER_REQUEST') + const fluxPer1kTokens = await deps.configKV.get('FLUX_PER_1K_TOKENS') + + const flux = await deps.fluxService.getFlux(userId) + if (flux.flux < fallbackRate) { + throw createPaymentRequiredError('Insufficient flux') + } + + return { fallbackRate, fluxPer1kTokens } + } + + function priceChatUsage(usage: UsageInfo, policy: ChatBillingPolicy): number { + if (policy.fluxPer1kTokens == null) + return policy.fallbackRate + return calculateFluxFromUsage(usage, policy.fluxPer1kTokens, policy.fallbackRate) + } + + async function settleChat(input: Omit): Promise { + return debitChatFlux({ + ...input, + billingService: deps.billingService, + revenue: deps.revenue, + }) + } + + function recordChatDebitFailure(input: { + amount: number + model: string + stage: 'streaming' | 'non_streaming' + }): void { + deps.revenue?.fluxUnbilled.add(input.amount, { + [GEN_AI_ATTR_REQUEST_MODEL]: input.model, + reason: 'debit_failed', + stage: input.stage, + }) + } + + async function authorizeTts(userId: string, inputText: string): Promise { + const flux = await deps.fluxService.getFlux(userId) + if (flux.flux <= 0) { + throw createPaymentRequiredError('Insufficient flux') + } + + // Pre-flight: refuse before hitting upstream if this segment would push the + // user past their balance. Cheap-path requests below the Flux threshold + // still pass when the user has at least 1 Flux. + await deps.ttsMeter.assertCanAfford(userId, inputText.length, flux.flux) + return { balance: flux.flux, inputChars: inputText.length } + } + + async function settleTts(input: { + userId: string + inputText: string + currentBalance: number + requestId: string + model: string + }) { + return deps.ttsMeter.accumulate({ + userId: input.userId, + units: input.inputText.length, + currentBalance: input.currentBalance, + requestId: input.requestId, + metadata: { model: input.model }, + }) + } + + return { authorizeChat, authorizeTts, priceChatUsage, recordChatDebitFailure, settleChat, settleTts } +} + +export async function debitChatFlux(input: ChatFluxDebitInput): Promise { + const result = await input.billingService.consumeFluxForLLM({ + userId: input.userId, + amount: input.amount, + requestId: input.requestId, + description: 'llm_request', + model: input.model, + promptTokens: input.promptTokens, + completionTokens: input.completionTokens, + }) + + if (result.charged < result.requested) { + input.revenue?.fluxUnbilled.add(result.requested - result.charged, { + [GEN_AI_ATTR_REQUEST_MODEL]: input.model, + reason: 'partial_debit_drained', + stage: input.stage, + }) + input.logger.withFields({ + userId: input.userId, + requestId: input.requestId, + requested: result.requested, + charged: result.charged, + unbilled: result.requested - result.charged, + }).warn(input.stage === 'streaming' + ? 'Partial debit after streaming — flux drained to zero' + : 'Partial debit on non-streaming completion — flux drained to zero') + } + + return result.charged +} diff --git a/apps/server/src/routes/openai/v1/catalog.ts b/apps/server/src/routes/openai/v1/catalog.ts new file mode 100644 index 000000000..b61ba080a --- /dev/null +++ b/apps/server/src/routes/openai/v1/catalog.ts @@ -0,0 +1,160 @@ +import type { Context, Handler } from 'hono' + +import type { HonoEnv } from '../../../types/hono' +import type { V1RouteDeps } from './types' + +import { useLogger } from '@guiiai/logg' +import { ofetch } from 'ofetch' + +import { createBadGatewayError, createBadRequestError, createServiceUnavailableError } from '../../../utils/error' + +export interface AudioCatalogHandlers { + handleListStreamingTTSModels: Handler + handleListStreamingVoices: Handler + handleListTTSModels: Handler + handleListVoices: Handler +} + +export function createAudioCatalogHandlers(deps: V1RouteDeps): AudioCatalogHandlers { + const logger = useLogger('v1-completions').useGlobalConfig() + + async function handleListVoices(c: Context) { + // Voice catalogs are per-model. Live providers (Azure) call upstream + // via unspeech; static providers (cosyvoice, volcengine) return their + // bundled JSON. The Redis cache + invalidation lives one layer down + // in the router so route-level changes don't leak into the cache + // contract. Recommended map stays in configKV so operators can edit it + // without a deploy. + // + // No implicit fallback: an empty `?model=` is a client bug (the UI is + // expected to pass either an explicit model id or the `auto` alias) and + // returns 400 instead of silently resolving to DEFAULT_TTS_MODEL. + const requested = c.req.query('model') + if (requested === undefined || requested === '') + throw createBadRequestError('audio voices: ?model= is required (use `auto` to defer to DEFAULT_TTS_MODEL)', 'MISSING_MODEL') + + const model = requested === 'auto' + ? await deps.configKV.getOrThrow('DEFAULT_TTS_MODEL') + : requested + + const voices = await deps.llmRouter.listTtsVoices(model) + const recommended = (await deps.configKV.getOptional('DEFAULT_TTS_VOICES'))?.[model] ?? {} + // Debug level: high-frequency catalog poll from UI selectors, no + // billing / user-facing side effect — useful only when debugging + // voice-picker drift, never as a permanent audit trail line. + logger.withFields({ model, voiceCount: voices.length }).debug('list tts voices') + return Response.json({ voices, recommended }) + } + + /** + * Voice catalog for the streaming TTS provider (`/audio/speech/ws`). + * + * Errors propagate verbatim: missing config → 503, malformed upstream + * URL → 502, unspeech network failure → 502, unspeech non-2xx → 502. + * No empty-array fallback — the UI surfaces a real failure state. + */ + async function handleListStreamingVoices(c: Context) { + const unspeech = await deps.configKV.getOptional('UNSPEECH_UPSTREAM') + if (!unspeech?.streaming?.baseURL) + throw createServiceUnavailableError('streaming tts upstream not configured', 'STREAMING_TTS_NOT_CONFIGURED') + + // Pass through the api_resource_id (e.g. `seed-tts-2.0`). unspeech + // filters the embedded Volcengine catalogue server-side; absent model + // means "return everything streaming-safe". + const model = c.req.query('model') + + let voicesURL: string + try { + const u = new URL(unspeech.restBaseURL) + u.pathname = '/api/voices' + const params = new URLSearchParams({ provider: 'volcengine' }) + if (model) + params.set('model', model) + u.search = `?${params.toString()}` + voicesURL = u.toString() + } + catch (err) { + logger.withError(err).withFields({ restBaseURL: unspeech.restBaseURL }).warn('streaming-voices: bad UNSPEECH_UPSTREAM.restBaseURL') + throw createBadGatewayError('UNSPEECH_UPSTREAM.restBaseURL is malformed') + } + + let res: Awaited> + try { + res = await ofetch.raw(voicesURL, { + ignoreResponseError: true, + timeout: 5000, + }) + } + catch (err) { + logger.withError(err).withFields({ voicesURL }).warn('streaming-voices: unspeech fetch failed') + throw createBadGatewayError('streaming voices upstream fetch failed') + } + + if (!res.ok) { + let snippet = '' + if (typeof res._data === 'string') { + snippet = res._data + } + else if (res._data != null) { + try { + snippet = JSON.stringify(res._data) + } + catch { + snippet = String(res._data) + } + } + logger.withFields({ voicesURL, status: res.status, snippet: snippet.slice(0, 256) }).warn('streaming-voices: unspeech non-2xx') + throw createBadGatewayError(`streaming voices upstream ${res.status}`, { lastStatusCode: res.status }) + } + + const data = res._data as { voices: unknown[] } + if (!Array.isArray(data.voices)) + throw createBadGatewayError('streaming voices upstream missing voices[]') + + const recommended = model + ? ((await deps.configKV.getOptional('DEFAULT_TTS_VOICES'))?.[model] ?? {}) + : {} + return Response.json({ voices: data.voices, recommended }) + } + + async function handleListTTSModels(_c: Context) { + // Surface the concrete TTS models the operator has configured. The UI + // should select an explicit model id so voice catalog requests stay + // model-scoped instead of hiding behind DEFAULT_TTS_MODEL. + const config = await deps.configKV.getOrThrow('LLM_ROUTER_CONFIG') + // `LLM_ROUTER_CONFIG` is `optional()` at the schema, so its inferred type + // tolerates `undefined`. `getOrThrow` already throws on missing entries, + // so by this line we know `config` is present — the `?.` here is purely + // a TS narrowing aid. + const modelIds = Object.keys(config?.tts?.models ?? {}).sort() + return Response.json({ + models: modelIds.map(id => ({ id, name: id })), + }) + } + + async function handleListStreamingTTSModels(_c: Context) { + const unspeech = await deps.configKV.getOptional('UNSPEECH_UPSTREAM') + const models = unspeech?.streaming?.models ?? [] + // `available` is the operator-controlled visibility switch the client gates + // the streaming provider on. It tracks whether `UNSPEECH_UPSTREAM.streaming` + // is configured at all — not whether `models[]` happens to be empty — so an + // operator who has wired the upstream but not yet curated models still + // surfaces the provider rather than silently hiding it. + return Response.json({ + available: !!unspeech?.streaming?.baseURL, + models: models.map(m => ({ + id: m.id, + name: m.name ?? m.id, + description: m.description, + })), + default: unspeech?.streaming?.defaultModel ?? null, + }) + } + + return { + handleListStreamingTTSModels, + handleListStreamingVoices, + handleListTTSModels, + handleListVoices, + } +} diff --git a/apps/server/src/routes/openai/v1/chat.ts b/apps/server/src/routes/openai/v1/chat.ts new file mode 100644 index 000000000..ac350b671 --- /dev/null +++ b/apps/server/src/routes/openai/v1/chat.ts @@ -0,0 +1,456 @@ +import type { Context, Handler } from 'hono' + +import type { UsageInfo } from '../../../services/domain/billing/billing' +import type { HonoEnv } from '../../../types/hono' +import type { V1RouteDeps } from './types' + +import { useLogger } from '@guiiai/logg' + +import { captureSafe } from '../../../services/adapters/posthog' +import { extractUsageFromBody } from '../../../services/domain/billing/billing' +import { nanoid } from '../../../utils/id' +import { createOpenAiRouteBilling } from './billing' +import { buildSafeResponseHeaders } from './response' +import { createRouteTelemetry, newRouteContext } from './telemetry' + +type ChatBilling = ReturnType +type ChatBillingPolicy = Awaited> +type RouteTelemetry = ReturnType + +export function createChatCompletionHandler(deps: V1RouteDeps): Handler { + const logger = useLogger('v1-completions').useGlobalConfig() + const telemetry = createRouteTelemetry({ + genAi: deps.genAi, + requestLogService: deps.requestLogService, + }) + const billing = createOpenAiRouteBilling(deps) + + return async function handleCompletion(c: Context) { + const user = c.get('user')! + // Generated up-front so incoming, completion, partial-debit, debit-failure, + // and request-log entries all carry the same correlation id. Re-used as + // the billing requestId (both streaming and non-streaming branches) for + // DB-level idempotency. + const requestId = nanoid() + + const billingPolicy = await billing.authorizeChat(user.id) + + const body = await c.req.json() + let requestModel = body.model || 'auto' + + if (requestModel === 'auto') { + requestModel = await deps.configKV.getOrThrow('DEFAULT_CHAT_MODEL') + } + + const stream = !!body.stream + logger.withFields({ + requestId, + userId: user.id, + model: requestModel, + stream, + messageCount: Array.isArray(body.messages) ? body.messages.length : undefined, + }).log('chat completion request') + + // Server-connection attrs come from the router (which knows the actual + // upstream baseURL it dispatched to) — it enriches the active span with + // its own `airi.gen_ai.gateway.*` attrs on success. + const span = telemetry.startChatSpan({ model: requestModel, stream }) + + const startedAt = Date.now() + + // Router throws ApiError (502/503/504/400) on full exhaustion or unknown + // model. We do NOT catch here — global app.onError renders the ApiError + // shape. Span is closed inside the catch so failures show up in traces. + // NOTICE: + // Propagate the client disconnect signal so an upstream LLM call doesn't + // keep generating tokens (and burning paid upstream quota) after the + // caller hangs up. Without this the streaming-cancel path records + // fluxConsumed: 0 while real cost was incurred — a silent revenue leak. + // Source: codex review 2026-05-15 HIGH #1. + const clientAbort = c.req.raw.signal + const routeCtx = newRouteContext() + let response: Response + try { + response = await telemetry.runWithSpan(span, () => + deps.llmRouter.route({ modelName: requestModel, body, headers: {}, abortSignal: clientAbort }, routeCtx)) + } + catch (err) { + telemetry.failSpan(span, 'Router exhausted or unknown model') + deps.llmTracing.startChatGeneration({ + input: body.messages, + model: routeCtx.upstreamModel ?? requestModel, + requestId, + stream, + userId: user.id, + sessionId: c.req.header('x-airi-session-id'), + }).fail('Router exhausted or unknown model') + telemetry.recordMetrics({ model: requestModel, status: 502, type: 'chat', provider: routeCtx.provider, durationMs: Date.now() - startedAt, fluxConsumed: 0 }) + throw err + } + + const durationMs = Date.now() - startedAt + telemetry.setHttpStatus(span, response.status) + const langfuseModel = routeCtx.upstreamModel ?? requestModel + + // Langfuse LLM-native generation: per-request prompt/completion record + // (input/output/model/usage) powering prompt trace, eval, and per-user/ + // session cost. Use the router-resolved upstream model, not the client + // alias (`auto` / `chat-auto`), so Langfuse model-cost grouping matches the + // provider model that actually generated the tokens. + const generationTrace = deps.llmTracing.startChatGeneration({ + input: body.messages, + model: langfuseModel, + requestId, + stream, + userId: user.id, + sessionId: c.req.header('x-airi-session-id'), + }) + + if (!response.ok) { + telemetry.failSpan(span, `Gateway ${response.status}`) + generationTrace.fail(`Gateway ${response.status}`) + telemetry.recordMetrics({ model: requestModel, status: response.status, type: 'chat', provider: routeCtx.provider, durationMs, fluxConsumed: 0 }) + // Emit server-side so funnels see real HTTP status — the client only + // ever observes "stream closed" and cannot tell 401 / 429 / 5xx apart. + void captureSafe(deps.posthog ?? null, { + distinctId: user.id, + event: 'llm_request_failed', + properties: { + model: requestModel, + http_status: response.status, + duration_ms: durationMs, + stream: !!body.stream, + }, + }) + + logger.withFields({ requestId, userId: user.id, model: requestModel, status: response.status, durationMs }) + .warn('chat completion delivered with upstream error status') + + return new Response(response.body, { + status: response.status, + headers: buildSafeResponseHeaders(response), + }) + } + + if (stream) { + return streamChatCompletion({ + deps, + response, + generationTrace, + span, + startedAt, + durationMs, + requestId, + userId: user.id, + requestModel, + routeCtxProvider: routeCtx.provider, + billing, + billingPolicy, + telemetry, + logger, + }) + } + + return completeNonStreamingChat({ + c, + deps, + response, + generationTrace, + span, + durationMs, + requestId, + userId: user.id, + requestModel, + routeCtxProvider: routeCtx.provider, + billing, + billingPolicy, + telemetry, + logger, + }) + } +} + +function streamChatCompletion(input: { + deps: V1RouteDeps + response: Response + generationTrace: ReturnType + span: Parameters[0] + startedAt: number + durationMs: number + requestId: string + userId: string + requestModel: string + routeCtxProvider: string + billing: ChatBilling + billingPolicy: ChatBillingPolicy + telemetry: RouteTelemetry + logger: ReturnType +}) { + // Streaming: return response immediately, bill after stream ends + const { readable, writable } = new TransformStream() + const reader = input.response.body!.getReader() + const writer = writable.getWriter() + const decoder = new TextDecoder() + // Buffer last 2KB to handle chunk boundary splits for usage extraction + let tailBuffer = '' + let streamCompleted = false + let streamInterrupted = false + // First-chunk timestamp for gen_ai.client.first_token.duration. Latched + // on the first byte from upstream — captures perceived "time to first + // token" for streaming clients. NaN until the first chunk lands so + // `Number.isFinite` gates the histogram record. + let firstChunkAt = Number.NaN + + // Process stream in background + ;(async () => { + try { + while (true) { + const { done, value } = await reader.read() + if (done) { + streamCompleted = true + break + } + if (!Number.isFinite(firstChunkAt)) { + firstChunkAt = Date.now() + input.telemetry.recordFirstToken({ + firstChunkAt, + model: input.requestModel, + provider: input.routeCtxProvider, + startedAt: input.startedAt, + }) + } + await writer.write(value) + const text = decoder.decode(value, { stream: true }) + tailBuffer = (tailBuffer + text).slice(-2048) + // Accumulate the assistant completion for the Langfuse trace output + // (no-op when tracing is off). Module owns SSE parsing + the cap. + input.generationTrace.appendStreamChunk(text) + } + } + catch (err) { + streamInterrupted = true + input.telemetry.recordStreamInterrupted({ + model: input.requestModel, + span: input.span, + stage: Number.isFinite(firstChunkAt) ? 'mid_stream' : 'before_first_chunk', + }) + + try { + await writer.abort(err) + } + catch (abortErr) { + input.logger.withError(abortErr).warn('Failed to abort stream writer after upstream interruption') + } + + input.logger.withError(err).warn('Upstream stream interrupted before completion') + return + } + finally { + if (streamInterrupted) { + input.telemetry.endSpan(input.span) + input.generationTrace.fail('Gateway stream interrupted') + input.telemetry.recordMetrics({ model: input.requestModel, status: input.response.status, type: 'chat', provider: input.routeCtxProvider, durationMs: input.durationMs, fluxConsumed: 0 }) + } + else if (streamCompleted) { + try { + await writer.close() + } + catch (err) { + input.logger.withError(err).warn('Failed to close stream writer') + } + + let usage: UsageInfo = {} + try { + const lines = tailBuffer.split('\n').filter(l => l.startsWith('data: ') && !l.includes('[DONE]')) + const lastDataLine = lines.at(-1) + if (lastDataLine) { + const json = JSON.parse(lastDataLine.slice(6)) + usage = extractUsageFromBody(json) + } + } + catch (err) { input.logger.withError(err).warn('Failed to extract usage from stream, falling back to flat rate') } + + const fluxConsumed = input.billing.priceChatUsage(usage, input.billingPolicy) + + input.telemetry.recordUsageOnSpan(input.span, { ...usage, fluxConsumed }) + input.telemetry.endSpan(input.span) + // Streaming output comes from appendStreamChunk above, so succeed + // omits it and the module uses the assembled assistant text. + input.generationTrace.succeed({ + promptTokens: usage.promptTokens, + completionTokens: usage.completionTokens, + fluxConsumed, + }) + input.telemetry.recordMetrics({ model: input.requestModel, status: input.response.status, type: 'chat', provider: input.routeCtxProvider, durationMs: input.durationMs, fluxConsumed, ...usage }) + + // Debit flux via DB transaction (source of truth) + // NOTICE: streaming response is already sent, so we cannot reject on failure. + // Log at error level so unpaid usage is visible in monitoring/alerts. + // + // `consumeFluxForLLM` now drains to zero on partial balance instead + // of throwing — the catch path only fires on `balance <= 0` (post- + // race) or real DB errors. Partial debits are signalled via the + // returned `charged < requested` and accounted to the same + // `fluxUnbilled` counter (different `reason` label). + let actualCharged = 0 + try { + actualCharged = await input.billing.settleChat({ + userId: input.userId, + amount: fluxConsumed, + requestId: input.requestId, + model: input.requestModel, + stage: 'streaming', + logger: input.logger, + ...usage, + }) + } + catch (err) { + // Real revenue leak: streaming response already sent (HTTP 200, + // tokens delivered), so this catch produces no 5xx and no DB + // latency spike on the request path. Without a dedicated counter, + // the failure is silent. Page on any sustained `increase()`. + input.billing.recordChatDebitFailure({ amount: fluxConsumed, model: input.requestModel, stage: 'streaming' }) + input.logger.withError(err).withFields({ userId: input.userId, fluxConsumed, requestId: input.requestId }).error('Failed to debit flux after streaming — unpaid usage') + } + + input.telemetry.recordRequestLog({ + userId: input.userId, + model: input.requestModel, + status: input.response.status, + durationMs: input.durationMs, + fluxConsumed: actualCharged, + promptTokens: usage.promptTokens, + completionTokens: usage.completionTokens, + }) + + void captureSafe(input.deps.posthog ?? null, { + distinctId: input.userId, + event: 'llm_request_succeeded', + properties: { + model: input.requestModel, + http_status: input.response.status, + duration_ms: input.durationMs, + prompt_tokens: usage.promptTokens ?? 0, + completion_tokens: usage.completionTokens ?? 0, + flux_consumed: actualCharged, + stream: true, + stream_interrupted: streamInterrupted, + }, + }) + + input.logger.withFields({ + requestId: input.requestId, + userId: input.userId, + model: input.requestModel, + status: input.response.status, + durationMs: input.durationMs, + promptTokens: usage.promptTokens, + completionTokens: usage.completionTokens, + fluxConsumed: actualCharged, + stream: true, + }).log('chat completion delivered') + } + } + })() + + return new Response(readable, { + status: input.response.status, + headers: buildSafeResponseHeaders(input.response), + }) +} + +async function completeNonStreamingChat(input: { + c: Context + deps: V1RouteDeps + response: Response + generationTrace: ReturnType + span: Parameters[0] + durationMs: number + requestId: string + userId: string + requestModel: string + routeCtxProvider: string + billing: ChatBilling + billingPolicy: ChatBillingPolicy + telemetry: RouteTelemetry + logger: ReturnType +}) { + // Non-streaming: parse response, bill, then return. + // Parse failure (malformed upstream JSON) must close both span and the + // Langfuse generation before bubbling up — otherwise the trace leaks. + // Mirrors the error-branch shape used above (router throw / !response.ok). + let responseBody + try { + responseBody = await input.response.json() + } + catch (err) { + input.telemetry.failSpan(input.span, 'Failed to parse upstream response body') + input.generationTrace.fail('Failed to parse upstream response body') + input.telemetry.recordMetrics({ model: input.requestModel, status: input.response.status, type: 'chat', provider: input.routeCtxProvider, durationMs: input.durationMs, fluxConsumed: 0 }) + throw err + } + const usage = extractUsageFromBody(responseBody) + const fluxConsumed = input.billing.priceChatUsage(usage, input.billingPolicy) + + input.telemetry.recordUsageOnSpan(input.span, { ...usage, fluxConsumed }) + input.telemetry.endSpan(input.span) + input.generationTrace.succeed({ + output: responseBody, + promptTokens: usage.promptTokens, + completionTokens: usage.completionTokens, + fluxConsumed, + }) + input.telemetry.recordMetrics({ model: input.requestModel, status: input.response.status, type: 'chat', provider: input.routeCtxProvider, durationMs: input.durationMs, fluxConsumed, ...usage }) + + // Debit flux via DB transaction (source of truth). + // The upstream call has already happened (cost incurred), so partial + // debit + `fluxUnbilled` is the only sane recovery — same shape as the + // streaming path. `balance <= 0` still throws and bubbles up as 402. + const actualCharged = await input.billing.settleChat({ + userId: input.userId, + amount: fluxConsumed, + requestId: input.requestId, + model: input.requestModel, + stage: 'non_streaming', + logger: input.logger, + ...usage, + }) + + input.telemetry.recordRequestLog({ + userId: input.userId, + model: input.requestModel, + status: input.response.status, + durationMs: input.durationMs, + fluxConsumed: actualCharged, + promptTokens: usage.promptTokens, + completionTokens: usage.completionTokens, + }) + + void captureSafe(input.deps.posthog ?? null, { + distinctId: input.userId, + event: 'llm_request_succeeded', + properties: { + model: input.requestModel, + http_status: input.response.status, + duration_ms: input.durationMs, + prompt_tokens: usage.promptTokens ?? 0, + completion_tokens: usage.completionTokens ?? 0, + flux_consumed: actualCharged, + stream: false, + }, + }) + + input.logger.withFields({ + requestId: input.requestId, + userId: input.userId, + model: input.requestModel, + status: input.response.status, + durationMs: input.durationMs, + promptTokens: usage.promptTokens, + completionTokens: usage.completionTokens, + fluxConsumed: actualCharged, + stream: false, + }).log('chat completion delivered') + + return input.c.json(responseBody) +} diff --git a/apps/server/src/routes/openai/v1/guards.ts b/apps/server/src/routes/openai/v1/guards.ts new file mode 100644 index 000000000..72c84794e --- /dev/null +++ b/apps/server/src/routes/openai/v1/guards.ts @@ -0,0 +1,19 @@ +import type { V1RouteDeps } from './types' + +import { authGuard } from '../../../middlewares/auth' +import { configGuard } from '../../../middlewares/config-guard' +import { rateLimiter } from '../../../middlewares/rate-limit' + +export function createV1RouteGuards(deps: V1RouteDeps) { + return { + authGuard, + chatGuard: configGuard(deps.configKV, ['FLUX_PER_REQUEST'], 'Service is not available yet'), + completionsRateLimit: rateLimiter({ + max: 60, + windowSec: 60, + metrics: deps.rateLimitMetrics, + routeLabel: 'openai.completions', + }), + ttsGuard: configGuard(deps.configKV, ['FLUX_PER_1K_CHARS_TTS'], 'TTS service is not available yet'), + } +} diff --git a/apps/server/src/routes/openai/v1/index.ts b/apps/server/src/routes/openai/v1/index.ts index 9cdd07c48..b74228551 100644 --- a/apps/server/src/routes/openai/v1/index.ts +++ b/apps/server/src/routes/openai/v1/index.ts @@ -1,93 +1,22 @@ -import type { Context } from 'hono' import type { PostHog } from 'posthog-node' import type { GenAiMetrics, RateLimitMetrics, RevenueMetrics } from '../../../otel' import type { ConfigKVService } from '../../../services/adapters/config-kv' -import type { UsageInfo } from '../../../services/domain/billing/billing' import type { BillingService } from '../../../services/domain/billing/billing-service' import type { FluxMeter } from '../../../services/domain/billing/flux-meter' import type { FluxService } from '../../../services/domain/flux' -import type { LlmRouteContext, LlmRouterService } from '../../../services/domain/llm-router' -import type { ChatGenerationTrace, TtsGenerationTrace } from '../../../services/domain/llm-tracing' +import type { LlmRouterService } from '../../../services/domain/llm-router' import type { RequestLogService } from '../../../services/domain/request-log' import type { HonoEnv } from '../../../types/hono' +import type { LlmTracingDeps } from './types' -import { useLogger } from '@guiiai/logg' -import { context, SpanStatusCode, trace } from '@opentelemetry/api' import { Hono } from 'hono' -import { authGuard } from '../../../middlewares/auth' -import { configGuard } from '../../../middlewares/config-guard' -import { rateLimiter } from '../../../middlewares/rate-limit' -import { captureSafe } from '../../../services/adapters/posthog' -import { calculateFluxFromUsage, extractUsageFromBody } from '../../../services/domain/billing/billing' -import { startChatGeneration, startTtsGeneration } from '../../../services/domain/llm-tracing' -import { createBadGatewayError, createBadRequestError, createPaymentRequiredError, createServiceUnavailableError } from '../../../utils/error' -import { nanoid } from '../../../utils/id' -import { - AIRI_ATTR_BILLING_FLUX_CONSUMED, - AIRI_ATTR_GEN_AI_OPERATION_KIND, - AIRI_ATTR_GEN_AI_STREAM, - AIRI_ATTR_GEN_AI_STREAM_INTERRUPTED, - GEN_AI_ATTR_OPERATION_NAME, - GEN_AI_ATTR_REQUEST_MODEL, - GEN_AI_ATTR_USAGE_INPUT_TOKENS, - GEN_AI_ATTR_USAGE_OUTPUT_TOKENS, -} from '../../../utils/observability' - -const tracer = trace.getTracer('v1-completions') - -interface LlmTracingDeps { - startChatGeneration: (input: Parameters[0]) => ChatGenerationTrace - startTtsGeneration: (input: Parameters[0]) => TtsGenerationTrace -} - -const SAFE_RESPONSE_HEADERS = new Set([ - 'content-type', - 'content-length', - 'transfer-encoding', - 'cache-control', -]) - -function buildSafeResponseHeaders(response: Response): Headers { - const headers = new Headers() - response.headers.forEach((value, key) => { - if (SAFE_RESPONSE_HEADERS.has(key.toLowerCase())) - headers.set(key, value) - }) - return headers -} - -function getLlmMetricAttributes(opts: { model: string, type: string, status: number, provider: string }): Record { - // `provider` is the upstream the router actually used (winning upstream on - // success, last-tried on exhaustion), so per-provider rollups in Grafana - // line up with each vendor's own console. Same label name as the gateway - // error counters (`airi_gen_ai_gateway_upstream_errors{provider}`) so the - // two can be compared/joined. - if (opts.type === 'chat') { - return { - [GEN_AI_ATTR_REQUEST_MODEL]: opts.model, - [GEN_AI_ATTR_OPERATION_NAME]: 'chat', - 'http.response.status_code': opts.status, - 'provider': opts.provider, - } - } - - return { - [GEN_AI_ATTR_REQUEST_MODEL]: opts.model, - [AIRI_ATTR_GEN_AI_OPERATION_KIND]: opts.type, - 'http.response.status_code': opts.status, - 'provider': opts.provider, - } -} - -// Fresh per-request context handed to `llmRouter.route` / `routeTts` so the -// router can report back which upstream it used (for the `provider` metric -// label). Must be created per request — never shared — because the route -// closures live at factory scope across concurrent requests. -function newRouteContext(): LlmRouteContext { - return { provider: 'unknown', triedUpstreams: 0, triedKeys: 0, lastStatus: null } -} +import { createAudioCatalogHandlers } from './catalog' +import { createChatCompletionHandler } from './chat' +import { createV1RouteGuards } from './guards' +import { createSpeechHandler } from './speech' +import { defaultLlmTracing } from './types' export function createV1Routes( fluxService: FluxService, @@ -100,742 +29,30 @@ export function createV1Routes( revenue?: RevenueMetrics | null, rateLimitMetrics?: RateLimitMetrics | null, posthog?: PostHog | null, - llmTracing: LlmTracingDeps = { startChatGeneration, startTtsGeneration }, + llmTracing: LlmTracingDeps = defaultLlmTracing, ) { - const logger = useLogger('v1-completions').useGlobalConfig() - // TODO: Extract this compat route into smaller facades/modules. - // It currently mixes auth, rate limiting, proxying, billing, telemetry, and event publishing in one transport layer entrypoint. - - function recordMetrics(opts: { model: string, status: number, type: string, provider: string, durationMs: number, fluxConsumed: number, promptTokens?: number, completionTokens?: number }) { - if (!genAi) - return - const attrs = getLlmMetricAttributes(opts) - genAi.operationCount.add(1, attrs) - genAi.operationDuration.record(opts.durationMs / 1000, attrs) - genAi.fluxConsumed.add(opts.fluxConsumed, attrs) - if (opts.promptTokens != null) - genAi.tokenUsageInput.add(opts.promptTokens, attrs) - if (opts.completionTokens != null) - genAi.tokenUsageOutput.add(opts.completionTokens, attrs) + const deps = { + fluxService, + billingService, + configKV, + requestLogService, + ttsMeter, + llmRouter, + genAi, + revenue, + rateLimitMetrics, + posthog, + llmTracing, } - - function recordRequestLog(entry: { userId: string, model: string, status: number, durationMs: number, fluxConsumed: number, promptTokens?: number, completionTokens?: number }) { - // Best-effort: a failed request log must not surface to the user — the - // upstream LLM response has already been delivered (or is mid-stream) by - // the time we get here. Log loss is observability-only. - requestLogService.logRequest(entry).catch(err => logger.withError(err).warn('Failed to write llm_request_log row')) - } - - // NOTICE: Billing is best-effort — flux is debited AFTER the LLM response is sent. - // This is a deliberate tradeoff: users get lower latency and uninterrupted streaming, - // at the cost of a small revenue leak when debit fails (e.g. DB timeout). - // Failed debits are logged at error level for monitoring/alerting. - // A pre-debit model would require holding the response until billing confirms, - // which adds latency and complicates streaming. We accept the leak for now. - // - // Pre-flight gates on `balance >= fallbackRate` (not just `> 0`) because - // streaming providers that don't echo `usage` cause every billable request - // to fall back to `FLUX_PER_REQUEST`. Without this gate, a user sitting on - // `0 < balance < fallbackRate` could spawn N parallel requests that each - // pass the loose `>0` check, complete the stream, and race on the debit — - // first wins, rest land in the partial-debit / catch path unbilled. With - // the gate, concurrent requests are rejected before the upstream call. - async function handleCompletion(c: Context) { - const user = c.get('user')! - // Generated up-front so incoming, completion, partial-debit, debit-failure, - // and request-log entries all carry the same correlation id. Re-used as - // the billing requestId (both streaming and non-streaming branches) for - // DB-level idempotency. - const requestId = nanoid() - - // Read billing rates before pre-flight so the gate can compare against - // the realistic per-request cost (fallback rate), not just `> 0`. - const fallbackRate = await configKV.getOrThrow('FLUX_PER_REQUEST') - const fluxPer1kTokens = await configKV.get('FLUX_PER_1K_TOKENS') - - const flux = await fluxService.getFlux(user.id) - if (flux.flux < fallbackRate) { - throw createPaymentRequiredError('Insufficient flux') - } - - const body = await c.req.json() - let requestModel = body.model || 'auto' - - if (requestModel === 'auto') { - requestModel = await configKV.getOrThrow('DEFAULT_CHAT_MODEL') - } - - const stream = !!body.stream - logger.withFields({ - requestId, - userId: user.id, - model: requestModel, - stream, - messageCount: Array.isArray(body.messages) ? body.messages.length : undefined, - }).log('chat completion request') - - // Server-connection attrs come from the router (which knows the actual - // upstream baseURL it dispatched to) — it enriches the active span with - // its own `airi.gen_ai.gateway.*` attrs on success. - const span = tracer.startSpan('llm.gateway.chat', { - attributes: { - [GEN_AI_ATTR_OPERATION_NAME]: 'chat', - [GEN_AI_ATTR_REQUEST_MODEL]: requestModel, - [AIRI_ATTR_GEN_AI_STREAM]: stream, - }, - }) - - const startedAt = Date.now() - - // Router throws ApiError (502/503/504/400) on full exhaustion or unknown - // model. We do NOT catch here — global app.onError renders the ApiError - // shape. Span is closed inside the catch so failures show up in traces. - // NOTICE: - // Propagate the client disconnect signal so an upstream LLM call doesn't - // keep generating tokens (and burning paid upstream quota) after the - // caller hangs up. Without this the streaming-cancel path records - // fluxConsumed: 0 while real cost was incurred — a silent revenue leak. - // Source: codex review 2026-05-15 HIGH #1. - const clientAbort = c.req.raw.signal - const routeCtx = newRouteContext() - let response: Response - try { - response = await context.with(trace.setSpan(context.active(), span), () => - llmRouter.route({ modelName: requestModel, body, headers: {}, abortSignal: clientAbort }, routeCtx)) - } - catch (err) { - span.setStatus({ code: SpanStatusCode.ERROR, message: 'Router exhausted or unknown model' }) - span.end() - llmTracing.startChatGeneration({ - input: body.messages, - model: routeCtx.upstreamModel ?? requestModel, - requestId, - stream, - userId: user.id, - sessionId: c.req.header('x-airi-session-id'), - }).fail('Router exhausted or unknown model') - recordMetrics({ model: requestModel, status: 502, type: 'chat', provider: routeCtx.provider, durationMs: Date.now() - startedAt, fluxConsumed: 0 }) - throw err - } - - const durationMs = Date.now() - startedAt - span.setAttribute('http.response.status_code', response.status) - const langfuseModel = routeCtx.upstreamModel ?? requestModel - - // Langfuse LLM-native generation: per-request prompt/completion record - // (input/output/model/usage) powering prompt trace, eval, and per-user/ - // session cost. Use the router-resolved upstream model, not the client - // alias (`auto` / `chat-auto`), so Langfuse model-cost grouping matches the - // provider model that actually generated the tokens. - const generationTrace = llmTracing.startChatGeneration({ - input: body.messages, - model: langfuseModel, - requestId, - stream, - userId: user.id, - sessionId: c.req.header('x-airi-session-id'), - }) - - if (!response.ok) { - span.setStatus({ code: SpanStatusCode.ERROR, message: `Gateway ${response.status}` }) - span.end() - generationTrace.fail(`Gateway ${response.status}`) - recordMetrics({ model: requestModel, status: response.status, type: 'chat', provider: routeCtx.provider, durationMs, fluxConsumed: 0 }) - // Emit server-side so funnels see real HTTP status — the client only - // ever observes "stream closed" and cannot tell 401 / 429 / 5xx apart. - void captureSafe(posthog ?? null, { - distinctId: user.id, - event: 'llm_request_failed', - properties: { - model: requestModel, - http_status: response.status, - duration_ms: durationMs, - stream: !!body.stream, - }, - }) - - logger.withFields({ requestId, userId: user.id, model: requestModel, status: response.status, durationMs }) - .warn('chat completion delivered with upstream error status') - - return new Response(response.body, { - status: response.status, - headers: buildSafeResponseHeaders(response), - }) - } - - // Post-billing: parse usage and charge after successful response. - // `fallbackRate` / `fluxPer1kTokens` were hoisted to the top of this - // function so the pre-flight gate can use them too. - - if (stream) { - // Streaming: return response immediately, bill after stream ends - const { readable, writable } = new TransformStream() - const reader = response.body!.getReader() - const writer = writable.getWriter() - const decoder = new TextDecoder() - // Buffer last 2KB to handle chunk boundary splits for usage extraction - let tailBuffer = '' - let streamCompleted = false - let streamInterrupted = false - // First-chunk timestamp for gen_ai.client.first_token.duration. Latched - // on the first byte from upstream — captures perceived "time to first - // token" for streaming clients. NaN until the first chunk lands so - // `Number.isFinite` gates the histogram record. - let firstChunkAt = Number.NaN - - // Process stream in background - ;(async () => { - try { - while (true) { - const { done, value } = await reader.read() - if (done) { - streamCompleted = true - break - } - if (!Number.isFinite(firstChunkAt)) { - firstChunkAt = Date.now() - genAi?.firstTokenDuration.record((firstChunkAt - startedAt) / 1000, { - [GEN_AI_ATTR_REQUEST_MODEL]: requestModel, - [GEN_AI_ATTR_OPERATION_NAME]: 'chat', - provider: routeCtx.provider, - }) - } - await writer.write(value) - const text = decoder.decode(value, { stream: true }) - tailBuffer = (tailBuffer + text).slice(-2048) - // Accumulate the assistant completion for the Langfuse trace output - // (no-op when tracing is off). Module owns SSE parsing + the cap. - generationTrace.appendStreamChunk(text) - } - } - catch (err) { - streamInterrupted = true - span.setStatus({ code: SpanStatusCode.ERROR, message: 'Gateway stream interrupted' }) - span.setAttribute(AIRI_ATTR_GEN_AI_STREAM_INTERRUPTED, true) - // Counter so alerts/dashboards can fire on interrupted streams; the - // span attribute alone only shows up in trace search, not metrics. - genAi?.streamInterrupted.add(1, { - [GEN_AI_ATTR_REQUEST_MODEL]: requestModel, - stage: Number.isFinite(firstChunkAt) ? 'mid_stream' : 'before_first_chunk', - }) - - try { - await writer.abort(err) - } - catch (abortErr) { - logger.withError(abortErr).warn('Failed to abort stream writer after upstream interruption') - } - - logger.withError(err).warn('Upstream stream interrupted before completion') - return - } - finally { - if (streamInterrupted) { - span.end() - generationTrace.fail('Gateway stream interrupted') - recordMetrics({ model: requestModel, status: response.status, type: 'chat', provider: routeCtx.provider, durationMs, fluxConsumed: 0 }) - } - else if (streamCompleted) { - try { - await writer.close() - } - catch (err) { - logger.withError(err).warn('Failed to close stream writer') - } - - let usage: UsageInfo = {} - try { - const lines = tailBuffer.split('\n').filter(l => l.startsWith('data: ') && !l.includes('[DONE]')) - const lastDataLine = lines.at(-1) - if (lastDataLine) { - const json = JSON.parse(lastDataLine.slice(6)) - usage = extractUsageFromBody(json) - } - } - catch (err) { logger.withError(err).warn('Failed to extract usage from stream, falling back to flat rate') } - - const fluxConsumed = calculateFluxFromUsage(usage, fluxPer1kTokens, fallbackRate) - - span.setAttributes({ - [GEN_AI_ATTR_USAGE_INPUT_TOKENS]: usage.promptTokens ?? 0, - [GEN_AI_ATTR_USAGE_OUTPUT_TOKENS]: usage.completionTokens ?? 0, - [AIRI_ATTR_BILLING_FLUX_CONSUMED]: fluxConsumed, - }) - span.end() - // Streaming output comes from appendStreamChunk above, so succeed - // omits it and the module uses the assembled assistant text. - generationTrace.succeed({ - promptTokens: usage.promptTokens, - completionTokens: usage.completionTokens, - fluxConsumed, - }) - recordMetrics({ model: requestModel, status: response.status, type: 'chat', provider: routeCtx.provider, durationMs, fluxConsumed, ...usage }) - - // Debit flux via DB transaction (source of truth) - // NOTICE: streaming response is already sent, so we cannot reject on failure. - // Log at error level so unpaid usage is visible in monitoring/alerts. - // - // `consumeFluxForLLM` now drains to zero on partial balance instead - // of throwing — the catch path only fires on `balance <= 0` (post- - // race) or real DB errors. Partial debits are signalled via the - // returned `charged < requested` and accounted to the same - // `fluxUnbilled` counter (different `reason` label). - let actualCharged = 0 - try { - const result = await billingService.consumeFluxForLLM({ - userId: user.id, - amount: fluxConsumed, - requestId, - description: 'llm_request', - model: requestModel, - promptTokens: usage.promptTokens, - completionTokens: usage.completionTokens, - }) - actualCharged = result.charged - if (result.charged < result.requested) { - revenue?.fluxUnbilled.add(result.requested - result.charged, { - [GEN_AI_ATTR_REQUEST_MODEL]: requestModel, - reason: 'partial_debit_drained', - stage: 'streaming', - }) - logger.withFields({ - userId: user.id, - requestId, - requested: result.requested, - charged: result.charged, - unbilled: result.requested - result.charged, - }).warn('Partial debit after streaming — flux drained to zero') - } - } - catch (err) { - // Real revenue leak: streaming response already sent (HTTP 200, - // tokens delivered), so this catch produces no 5xx and no DB - // latency spike on the request path. Without a dedicated counter, - // the failure is silent. Page on any sustained `increase()`. - revenue?.fluxUnbilled.add(fluxConsumed, { - [GEN_AI_ATTR_REQUEST_MODEL]: requestModel, - reason: 'debit_failed', - stage: 'streaming', - }) - logger.withError(err).withFields({ userId: user.id, fluxConsumed, requestId }).error('Failed to debit flux after streaming — unpaid usage') - } - - recordRequestLog({ - userId: user.id, - model: requestModel, - status: response.status, - durationMs, - fluxConsumed: actualCharged, - promptTokens: usage.promptTokens, - completionTokens: usage.completionTokens, - }) - - void captureSafe(posthog ?? null, { - distinctId: user.id, - event: 'llm_request_succeeded', - properties: { - model: requestModel, - http_status: response.status, - duration_ms: durationMs, - prompt_tokens: usage.promptTokens ?? 0, - completion_tokens: usage.completionTokens ?? 0, - flux_consumed: actualCharged, - stream: true, - stream_interrupted: streamInterrupted, - }, - }) - - logger.withFields({ - requestId, - userId: user.id, - model: requestModel, - status: response.status, - durationMs, - promptTokens: usage.promptTokens, - completionTokens: usage.completionTokens, - fluxConsumed: actualCharged, - stream: true, - }).log('chat completion delivered') - } - } - })() - - return new Response(readable, { - status: response.status, - headers: buildSafeResponseHeaders(response), - }) - } - - // Non-streaming: parse response, bill, then return. - // Parse failure (malformed upstream JSON) must close both span and the - // Langfuse generation before bubbling up — otherwise the trace leaks. - // Mirrors the error-branch shape used above (router throw / !response.ok). - let responseBody - try { - responseBody = await response.json() - } - catch (err) { - span.setStatus({ code: SpanStatusCode.ERROR, message: 'Failed to parse upstream response body' }) - span.end() - generationTrace.fail('Failed to parse upstream response body') - recordMetrics({ model: requestModel, status: response.status, type: 'chat', provider: routeCtx.provider, durationMs, fluxConsumed: 0 }) - throw err - } - const usage = extractUsageFromBody(responseBody) - const fluxConsumed = calculateFluxFromUsage(usage, fluxPer1kTokens, fallbackRate) - - span.setAttributes({ - [GEN_AI_ATTR_USAGE_INPUT_TOKENS]: usage.promptTokens ?? 0, - [GEN_AI_ATTR_USAGE_OUTPUT_TOKENS]: usage.completionTokens ?? 0, - [AIRI_ATTR_BILLING_FLUX_CONSUMED]: fluxConsumed, - }) - span.end() - generationTrace.succeed({ - output: responseBody, - promptTokens: usage.promptTokens, - completionTokens: usage.completionTokens, - fluxConsumed, - }) - recordMetrics({ model: requestModel, status: response.status, type: 'chat', provider: routeCtx.provider, durationMs, fluxConsumed, ...usage }) - - // Debit flux via DB transaction (source of truth). - // The upstream call has already happened (cost incurred), so partial - // debit + `fluxUnbilled` is the only sane recovery — same shape as the - // streaming path. `balance <= 0` still throws and bubbles up as 402. - const result = await billingService.consumeFluxForLLM({ - userId: user.id, - amount: fluxConsumed, - requestId, - description: 'llm_request', - model: requestModel, - promptTokens: usage.promptTokens, - completionTokens: usage.completionTokens, - }) - if (result.charged < result.requested) { - revenue?.fluxUnbilled.add(result.requested - result.charged, { - [GEN_AI_ATTR_REQUEST_MODEL]: requestModel, - reason: 'partial_debit_drained', - stage: 'non_streaming', - }) - logger.withFields({ - userId: user.id, - requestId, - requested: result.requested, - charged: result.charged, - unbilled: result.requested - result.charged, - }).warn('Partial debit on non-streaming completion — flux drained to zero') - } - - recordRequestLog({ - userId: user.id, - model: requestModel, - status: response.status, - durationMs, - fluxConsumed: result.charged, - promptTokens: usage.promptTokens, - completionTokens: usage.completionTokens, - }) - - void captureSafe(posthog ?? null, { - distinctId: user.id, - event: 'llm_request_succeeded', - properties: { - model: requestModel, - http_status: response.status, - duration_ms: durationMs, - prompt_tokens: usage.promptTokens ?? 0, - completion_tokens: usage.completionTokens ?? 0, - flux_consumed: result.charged, - stream: false, - }, - }) - - logger.withFields({ - requestId, - userId: user.id, - model: requestModel, - status: response.status, - durationMs, - promptTokens: usage.promptTokens, - completionTokens: usage.completionTokens, - fluxConsumed: result.charged, - stream: false, - }).log('chat completion delivered') - - return c.json(responseBody) - } - - async function handleTTS(c: Context) { - const user = c.get('user')! - const requestId = nanoid() - const flux = await fluxService.getFlux(user.id) - if (flux.flux <= 0) { - throw createPaymentRequiredError('Insufficient flux') - } - - const body = await c.req.json() - let requestModel = body.model || 'auto' - // NOTICE: Guard against non-string body.input — upstream would reject it - // anyway, but billing math (.length → INCRBY) turns NaN into a Redis error. - const inputText: string = typeof body.input === 'string' ? body.input : '' - - if (requestModel === 'auto') { - requestModel = await configKV.getOrThrow('DEFAULT_TTS_MODEL') - } - - logger.withFields({ - requestId, - userId: user.id, - model: requestModel, - inputChars: inputText.length, - voice: typeof body.voice === 'string' ? body.voice : undefined, - }).log('tts speech request') - - // Pre-flight: refuse before hitting upstream if this segment would push the - // user past their balance. Cheap-path requests below the Flux threshold - // still pass when the user has at least 1 Flux. - await ttsMeter.assertCanAfford(user.id, inputText.length, flux.flux) - - // Map OpenAI-shaped /audio/speech body → adapter-neutral TtsInput. Speed - // / response_format / extra fields stay in adapterParams for adapters that - // care (Azure SSML rate, Volcengine audio_params, etc.). - const ttsInput = { - text: inputText, - voice: typeof body.voice === 'string' ? body.voice : undefined, - speed: typeof body.speed === 'number' ? body.speed : undefined, - responseFormat: typeof body.response_format === 'string' ? body.response_format : undefined, - } - const generationTrace = llmTracing.startTtsGeneration({ - input: ttsInput, - model: requestModel, - requestId, - userId: user.id, - sessionId: c.req.header('x-airi-session-id'), - }) - - const span = tracer.startSpan('llm.gateway.tts', { - attributes: { - [GEN_AI_ATTR_REQUEST_MODEL]: requestModel, - [AIRI_ATTR_GEN_AI_OPERATION_KIND]: 'text_to_speech', - }, - }) - - const startedAt = Date.now() - - const routeCtx = newRouteContext() - let response: Response - try { - response = await context.with(trace.setSpan(context.active(), span), () => - llmRouter.routeTts({ modelName: requestModel, input: ttsInput, abortSignal: c.req.raw.signal }, routeCtx)) - } - catch (err) { - span.setStatus({ code: SpanStatusCode.ERROR, message: 'TTS router exhausted or unknown model' }) - span.end() - generationTrace.fail('TTS router exhausted or unknown model') - recordMetrics({ model: requestModel, status: 502, type: 'tts', provider: routeCtx.provider, durationMs: Date.now() - startedAt, fluxConsumed: 0 }) - throw err - } - - const durationMs = Date.now() - startedAt - span.setAttribute('http.response.status_code', response.status) - - if (!response.ok) { - span.setStatus({ code: SpanStatusCode.ERROR, message: `Gateway ${response.status}` }) - span.end() - generationTrace.fail(`Gateway ${response.status}`) - recordMetrics({ model: requestModel, status: response.status, type: 'tts', provider: routeCtx.provider, durationMs, fluxConsumed: 0 }) - logger.withFields({ requestId, userId: user.id, model: requestModel, status: response.status, durationMs }) - .warn('tts speech delivered with upstream error status') - return new Response(response.body, { - status: response.status, - headers: buildSafeResponseHeaders(response), - }) - } - - // Debt-ledger billing: accumulate chars in Redis; only debit when we - // cross a whole-Flux boundary. Sub-threshold requests cost 0 Flux at this - // call site — the cost is realised on a later request that crosses. - // - // Wrapped in try/finally so a Redis blip inside `accumulate()` (or any - // throw before `span.end()`) doesn't leak the active span. Falling-through - // to `throw` reaches the global ApiError handler — billing failure on a - // 200 upstream is rare but observable, and a dropped span would have - // hidden it. - let fluxConsumed = 0 - try { - const result = await ttsMeter.accumulate({ - userId: user.id, - units: inputText.length, - currentBalance: flux.flux, - requestId, - metadata: { model: requestModel }, - }) - fluxConsumed = result.fluxDebited - span.setAttribute(AIRI_ATTR_BILLING_FLUX_CONSUMED, fluxConsumed) - generationTrace.succeed({ - inputChars: inputText.length, - fluxConsumed, - output: { contentType: response.headers.get('content-type') }, - }) - } - catch (err) { - generationTrace.fail('TTS billing failed') - throw err - } - finally { - span.end() - } - recordMetrics({ model: requestModel, status: response.status, type: 'tts', provider: routeCtx.provider, durationMs, fluxConsumed }) - - recordRequestLog({ - userId: user.id, - model: requestModel, - status: response.status, - durationMs, - fluxConsumed, - }) - - logger.withFields({ - requestId, - userId: user.id, - model: requestModel, - status: response.status, - durationMs, - inputChars: inputText.length, - fluxConsumed, - }).log('tts speech delivered') - - return new Response(response.body, { - status: response.status, - headers: buildSafeResponseHeaders(response), - }) - } - - async function handleListVoices(c: Context) { - // Voice catalogs are per-model. Live providers (Azure) call upstream - // via unspeech; static providers (cosyvoice, volcengine) return their - // bundled JSON. The Redis cache + invalidation lives one layer down - // in the router so route-level changes don't leak into the cache - // contract. Recommended map stays in configKV so operators can edit it - // without a deploy. - // - // No implicit fallback: an empty `?model=` is a client bug (the UI is - // expected to pass either an explicit model id or the `auto` alias) and - // returns 400 instead of silently resolving to DEFAULT_TTS_MODEL. - const requested = c.req.query('model') - if (requested === undefined || requested === '') - throw createBadRequestError('audio voices: ?model= is required (use `auto` to defer to DEFAULT_TTS_MODEL)', 'MISSING_MODEL') - - const model = requested === 'auto' - ? await configKV.getOrThrow('DEFAULT_TTS_MODEL') - : requested - - const voices = await llmRouter.listTtsVoices(model) - const recommended = (await configKV.getOptional('DEFAULT_TTS_VOICES'))?.[model] ?? {} - // Debug level: high-frequency catalog poll from UI selectors, no - // billing / user-facing side effect — useful only when debugging - // voice-picker drift, never as a permanent audit trail line. - logger.withFields({ model, voiceCount: voices.length }).debug('list tts voices') - return Response.json({ voices, recommended }) - } - - /** - * Voice catalog for the streaming TTS provider (`/audio/speech/ws`). - * - * Errors propagate verbatim: missing config → 503, malformed upstream - * URL → 502, unspeech network failure → 502, unspeech non-2xx → 502. - * No empty-array fallback — the UI surfaces a real failure state. - */ - async function handleListStreamingVoices(c: Context) { - const unspeech = await configKV.getOptional('UNSPEECH_UPSTREAM') - if (!unspeech?.streaming?.baseURL) - throw createServiceUnavailableError('streaming tts upstream not configured', 'STREAMING_TTS_NOT_CONFIGURED') - - // Pass through the api_resource_id (e.g. `seed-tts-2.0`). unspeech - // filters the embedded Volcengine catalogue server-side; absent model - // means "return everything streaming-safe". - const model = c.req.query('model') - - let voicesURL: string - try { - const u = new URL(unspeech.restBaseURL) - u.pathname = '/api/voices' - const params = new URLSearchParams({ provider: 'volcengine' }) - if (model) - params.set('model', model) - u.search = `?${params.toString()}` - voicesURL = u.toString() - } - catch (err) { - logger.withError(err).withFields({ restBaseURL: unspeech.restBaseURL }).warn('streaming-voices: bad UNSPEECH_UPSTREAM.restBaseURL') - throw createBadGatewayError('UNSPEECH_UPSTREAM.restBaseURL is malformed') - } - - let res: Response - try { - res = await globalThis.fetch(voicesURL, { - signal: AbortSignal.timeout(5000), - }) - } - catch (err) { - logger.withError(err).withFields({ voicesURL }).warn('streaming-voices: unspeech fetch failed') - throw createBadGatewayError('streaming voices upstream fetch failed') - } - - if (!res.ok) { - const snippet = await res.text().catch(() => '') - logger.withFields({ voicesURL, status: res.status, snippet: snippet.slice(0, 256) }).warn('streaming-voices: unspeech non-2xx') - throw createBadGatewayError(`streaming voices upstream ${res.status}`, { lastStatusCode: res.status }) - } - - const data = await res.json() as { voices: unknown[] } - if (!Array.isArray(data.voices)) - throw createBadGatewayError('streaming voices upstream missing voices[]') - - const recommended = model - ? ((await configKV.getOptional('DEFAULT_TTS_VOICES'))?.[model] ?? {}) - : {} - return Response.json({ voices: data.voices, recommended }) - } - - async function handleListTTSModels(_c: Context) { - // Surface the concrete TTS models the operator has configured. The UI - // should select an explicit model id so voice catalog requests stay - // model-scoped instead of hiding behind DEFAULT_TTS_MODEL. - const config = await configKV.getOrThrow('LLM_ROUTER_CONFIG') - // `LLM_ROUTER_CONFIG` is `optional()` at the schema, so its inferred type - // tolerates `undefined`. `getOrThrow` already throws on missing entries, - // so by this line we know `config` is present — the `?.` here is purely - // a TS narrowing aid. - const modelIds = Object.keys(config?.tts?.models ?? {}).sort() - return Response.json({ - models: modelIds.map(id => ({ id, name: id })), - }) - } - - async function handleListStreamingTTSModels(_c: Context) { - const unspeech = await configKV.getOptional('UNSPEECH_UPSTREAM') - const models = unspeech?.streaming?.models ?? [] - // `available` is the operator-controlled visibility switch the client gates - // the streaming provider on. It tracks whether `UNSPEECH_UPSTREAM.streaming` - // is configured at all — not whether `models[]` happens to be empty — so an - // operator who has wired the upstream but not yet curated models still - // surfaces the provider rather than silently hiding it. - return Response.json({ - available: !!unspeech?.streaming?.baseURL, - models: models.map(m => ({ - id: m.id, - name: m.name ?? m.id, - description: m.description, - })), - default: unspeech?.streaming?.defaultModel ?? null, - }) - } - - const chatGuard = configGuard(configKV, ['FLUX_PER_REQUEST'], 'Service is not available yet') - const ttsGuard = configGuard(configKV, ['FLUX_PER_1K_CHARS_TTS'], 'TTS service is not available yet') - - const completionsRateLimit = rateLimiter({ max: 60, windowSec: 60, metrics: rateLimitMetrics, routeLabel: 'openai.completions' }) + const guards = createV1RouteGuards(deps) + const handleCompletion = createChatCompletionHandler(deps) + const handleTTS = createSpeechHandler(deps) + const { + handleListStreamingTTSModels, + handleListStreamingVoices, + handleListTTSModels, + handleListVoices, + } = createAudioCatalogHandlers(deps) // OpenAI-compatible surface (mounted at /api/v1/openai). Only routes that // mirror an actual OpenAI public endpoint belong here. Audio used to live @@ -844,17 +61,17 @@ export function createV1Routes( // OpenAI — keeping them here mislabelled the surface, so audio now mounts // at /api/v1/audio (see `audioRoutes` below). const openaiRoutes = new Hono() - .use('*', authGuard) - .post('/chat/completions', completionsRateLimit, chatGuard, handleCompletion) - .post('/chat/completion', completionsRateLimit, chatGuard, handleCompletion) + .use('*', guards.authGuard) + .post('/chat/completions', guards.completionsRateLimit, guards.chatGuard, handleCompletion) + .post('/chat/completion', guards.completionsRateLimit, guards.chatGuard, handleCompletion) // AIRI audio surface (mounted at /api/v1/audio). Lives outside /openai/ so // the `/voices`, `/voices/streaming`, and `/models` extensions aren't // misread as OpenAI-compatible. `/audio/speech/ws` is registered // separately in app.ts because it needs the WebSocket upgrade middleware. const audioRoutes = new Hono() - .use('*', authGuard) - .post('/speech', ttsGuard, handleTTS) + .use('*', guards.authGuard) + .post('/speech', guards.ttsGuard, handleTTS) .get('/voices', handleListVoices) .get('/voices/streaming', handleListStreamingVoices) .get('/models', handleListTTSModels) diff --git a/apps/server/src/routes/openai/v1/response.ts b/apps/server/src/routes/openai/v1/response.ts new file mode 100644 index 000000000..668a39f54 --- /dev/null +++ b/apps/server/src/routes/openai/v1/response.ts @@ -0,0 +1,15 @@ +const SAFE_RESPONSE_HEADERS = new Set([ + 'content-type', + 'content-length', + 'transfer-encoding', + 'cache-control', +]) + +export function buildSafeResponseHeaders(response: Response): Headers { + const headers = new Headers() + response.headers.forEach((value, key) => { + if (SAFE_RESPONSE_HEADERS.has(key.toLowerCase())) + headers.set(key, value) + }) + return headers +} diff --git a/apps/server/src/routes/openai/v1/speech.ts b/apps/server/src/routes/openai/v1/speech.ts new file mode 100644 index 000000000..11c4486e3 --- /dev/null +++ b/apps/server/src/routes/openai/v1/speech.ts @@ -0,0 +1,152 @@ +import type { Context, Handler } from 'hono' + +import type { HonoEnv } from '../../../types/hono' +import type { V1RouteDeps } from './types' + +import { useLogger } from '@guiiai/logg' + +import { nanoid } from '../../../utils/id' +import { createOpenAiRouteBilling } from './billing' +import { buildSafeResponseHeaders } from './response' +import { createRouteTelemetry, newRouteContext } from './telemetry' + +export function createSpeechHandler(deps: V1RouteDeps): Handler { + const logger = useLogger('v1-completions').useGlobalConfig() + const telemetry = createRouteTelemetry({ + genAi: deps.genAi, + requestLogService: deps.requestLogService, + }) + const billing = createOpenAiRouteBilling(deps) + + return async function handleTTS(c: Context) { + const user = c.get('user')! + const requestId = nanoid() + + const body = await c.req.json() + let requestModel = body.model || 'auto' + // NOTICE: Guard against non-string body.input — upstream would reject it + // anyway, but billing math (.length → INCRBY) turns NaN into a Redis error. + const inputText: string = typeof body.input === 'string' ? body.input : '' + + if (requestModel === 'auto') { + requestModel = await deps.configKV.getOrThrow('DEFAULT_TTS_MODEL') + } + + logger.withFields({ + requestId, + userId: user.id, + model: requestModel, + inputChars: inputText.length, + voice: typeof body.voice === 'string' ? body.voice : undefined, + }).log('tts speech request') + + const billingAuthorization = await billing.authorizeTts(user.id, inputText) + + // Map OpenAI-shaped /audio/speech body → adapter-neutral TtsInput. Speed + // / response_format / extra fields stay in adapterParams for adapters that + // care (Azure SSML rate, Volcengine audio_params, etc.). + const ttsInput = { + text: inputText, + voice: typeof body.voice === 'string' ? body.voice : undefined, + speed: typeof body.speed === 'number' ? body.speed : undefined, + responseFormat: typeof body.response_format === 'string' ? body.response_format : undefined, + } + const generationTrace = deps.llmTracing.startTtsGeneration({ + input: ttsInput, + model: requestModel, + requestId, + userId: user.id, + sessionId: c.req.header('x-airi-session-id'), + }) + + const span = telemetry.startTtsSpan({ model: requestModel }) + + const startedAt = Date.now() + + const routeCtx = newRouteContext() + let response: Response + try { + response = await telemetry.runWithSpan(span, () => + deps.llmRouter.routeTts({ modelName: requestModel, input: ttsInput, abortSignal: c.req.raw.signal }, routeCtx)) + } + catch (err) { + telemetry.failSpan(span, 'TTS router exhausted or unknown model') + generationTrace.fail('TTS router exhausted or unknown model') + telemetry.recordMetrics({ model: requestModel, status: 502, type: 'tts', provider: routeCtx.provider, durationMs: Date.now() - startedAt, fluxConsumed: 0 }) + throw err + } + + const durationMs = Date.now() - startedAt + telemetry.setHttpStatus(span, response.status) + + if (!response.ok) { + telemetry.failSpan(span, `Gateway ${response.status}`) + generationTrace.fail(`Gateway ${response.status}`) + telemetry.recordMetrics({ model: requestModel, status: response.status, type: 'tts', provider: routeCtx.provider, durationMs, fluxConsumed: 0 }) + logger.withFields({ requestId, userId: user.id, model: requestModel, status: response.status, durationMs }) + .warn('tts speech delivered with upstream error status') + return new Response(response.body, { + status: response.status, + headers: buildSafeResponseHeaders(response), + }) + } + + // Debt-ledger billing: accumulate chars in Redis; only debit when we + // cross a whole-Flux boundary. Sub-threshold requests cost 0 Flux at this + // call site — the cost is realised on a later request that crosses. + // + // Wrapped in try/finally so a Redis blip inside `accumulate()` (or any + // throw before `span.end()`) doesn't leak the active span. Falling-through + // to `throw` reaches the global ApiError handler — billing failure on a + // 200 upstream is rare but observable, and a dropped span would have + // hidden it. + let fluxConsumed = 0 + try { + const result = await billing.settleTts({ + userId: user.id, + inputText, + currentBalance: billingAuthorization.balance, + requestId, + model: requestModel, + }) + fluxConsumed = result.fluxDebited + telemetry.recordTtsBillingOnSpan(span, fluxConsumed) + generationTrace.succeed({ + inputChars: inputText.length, + fluxConsumed, + output: { contentType: response.headers.get('content-type') }, + }) + } + catch (err) { + generationTrace.fail('TTS billing failed') + throw err + } + finally { + telemetry.endSpan(span) + } + telemetry.recordMetrics({ model: requestModel, status: response.status, type: 'tts', provider: routeCtx.provider, durationMs, fluxConsumed }) + + telemetry.recordRequestLog({ + userId: user.id, + model: requestModel, + status: response.status, + durationMs, + fluxConsumed, + }) + + logger.withFields({ + requestId, + userId: user.id, + model: requestModel, + status: response.status, + durationMs, + inputChars: inputText.length, + fluxConsumed, + }).log('tts speech delivered') + + return new Response(response.body, { + status: response.status, + headers: buildSafeResponseHeaders(response), + }) + } +} diff --git a/apps/server/src/routes/openai/v1/telemetry.ts b/apps/server/src/routes/openai/v1/telemetry.ts new file mode 100644 index 000000000..65cc074d7 --- /dev/null +++ b/apps/server/src/routes/openai/v1/telemetry.ts @@ -0,0 +1,188 @@ +import type { GenAiMetrics } from '../../../otel' +import type { UsageInfo } from '../../../services/domain/billing/billing' +import type { LlmRouteContext } from '../../../services/domain/llm-router' +import type { RequestLogService } from '../../../services/domain/request-log' + +import { useLogger } from '@guiiai/logg' +import { context, SpanStatusCode, trace } from '@opentelemetry/api' + +import { + AIRI_ATTR_BILLING_FLUX_CONSUMED, + AIRI_ATTR_GEN_AI_OPERATION_KIND, + AIRI_ATTR_GEN_AI_STREAM, + AIRI_ATTR_GEN_AI_STREAM_INTERRUPTED, + GEN_AI_ATTR_OPERATION_NAME, + GEN_AI_ATTR_REQUEST_MODEL, + GEN_AI_ATTR_USAGE_INPUT_TOKENS, + GEN_AI_ATTR_USAGE_OUTPUT_TOKENS, +} from '../../../utils/observability' + +export const tracer = trace.getTracer('v1-completions') + +export type GatewaySpan = ReturnType + +export interface OperationMetricsInput extends UsageInfo { + model: string + status: number + type: string + provider: string + durationMs: number + fluxConsumed: number +} + +export interface RequestLogInput extends UsageInfo { + userId: string + model: string + status: number + durationMs: number + fluxConsumed: number +} + +export function getLlmMetricAttributes(opts: { model: string, type: string, status: number, provider: string }): Record { + // `provider` is the upstream the router actually used (winning upstream on + // success, last-tried on exhaustion), so per-provider rollups in Grafana + // line up with each vendor's own console. Same label name as the gateway + // error counters (`airi_gen_ai_gateway_upstream_errors{provider}`) so the + // two can be compared/joined. + if (opts.type === 'chat') { + return { + [GEN_AI_ATTR_REQUEST_MODEL]: opts.model, + [GEN_AI_ATTR_OPERATION_NAME]: 'chat', + 'http.response.status_code': opts.status, + 'provider': opts.provider, + } + } + + return { + [GEN_AI_ATTR_REQUEST_MODEL]: opts.model, + [AIRI_ATTR_GEN_AI_OPERATION_KIND]: opts.type, + 'http.response.status_code': opts.status, + 'provider': opts.provider, + } +} + +// Fresh per-request context handed to `llmRouter.route` / `routeTts` so the +// router can report back which upstream it used (for the `provider` metric +// label). Must be created per request — never shared — because the route +// closures live at factory scope across concurrent requests. +export function newRouteContext(): LlmRouteContext { + return { provider: 'unknown', triedUpstreams: 0, triedKeys: 0, lastStatus: null } +} + +export function createRouteTelemetry(deps: { + genAi?: GenAiMetrics | null + requestLogService: RequestLogService +}) { + const logger = useLogger('v1-completions').useGlobalConfig() + + function recordMetrics(opts: OperationMetricsInput) { + if (!deps.genAi) + return + const attrs = getLlmMetricAttributes(opts) + deps.genAi.operationCount.add(1, attrs) + deps.genAi.operationDuration.record(opts.durationMs / 1000, attrs) + deps.genAi.fluxConsumed.add(opts.fluxConsumed, attrs) + if (opts.promptTokens != null) + deps.genAi.tokenUsageInput.add(opts.promptTokens, attrs) + if (opts.completionTokens != null) + deps.genAi.tokenUsageOutput.add(opts.completionTokens, attrs) + } + + function recordRequestLog(entry: RequestLogInput) { + // Best-effort: a failed request log must not surface to the user — the + // upstream LLM response has already been delivered (or is mid-stream) by + // the time we get here. Log loss is observability-only. + deps.requestLogService.logRequest(entry).catch(err => logger.withError(err).warn('Failed to write llm_request_log row')) + } + + function startChatSpan(input: { model: string, stream: boolean }): GatewaySpan { + return tracer.startSpan('llm.gateway.chat', { + attributes: { + [GEN_AI_ATTR_OPERATION_NAME]: 'chat', + [GEN_AI_ATTR_REQUEST_MODEL]: input.model, + [AIRI_ATTR_GEN_AI_STREAM]: input.stream, + }, + }) + } + + function startTtsSpan(input: { model: string }): GatewaySpan { + return tracer.startSpan('llm.gateway.tts', { + attributes: { + [GEN_AI_ATTR_REQUEST_MODEL]: input.model, + [AIRI_ATTR_GEN_AI_OPERATION_KIND]: 'text_to_speech', + }, + }) + } + + async function runWithSpan(span: GatewaySpan, work: () => Promise): Promise { + return context.with(trace.setSpan(context.active(), span), work) + } + + function setHttpStatus(span: GatewaySpan, status: number): void { + span.setAttribute('http.response.status_code', status) + } + + function failSpan(span: GatewaySpan, message: string): void { + span.setStatus({ code: SpanStatusCode.ERROR, message }) + span.end() + } + + function endSpan(span: GatewaySpan): void { + span.end() + } + + function recordUsageOnSpan(span: GatewaySpan, input: UsageInfo & { fluxConsumed: number }): void { + span.setAttributes({ + [GEN_AI_ATTR_USAGE_INPUT_TOKENS]: input.promptTokens ?? 0, + [GEN_AI_ATTR_USAGE_OUTPUT_TOKENS]: input.completionTokens ?? 0, + [AIRI_ATTR_BILLING_FLUX_CONSUMED]: input.fluxConsumed, + }) + } + + function recordTtsBillingOnSpan(span: GatewaySpan, fluxConsumed: number): void { + span.setAttribute(AIRI_ATTR_BILLING_FLUX_CONSUMED, fluxConsumed) + } + + function recordFirstToken(input: { + model: string + provider: string + startedAt: number + firstChunkAt: number + }): void { + deps.genAi?.firstTokenDuration.record((input.firstChunkAt - input.startedAt) / 1000, { + [GEN_AI_ATTR_REQUEST_MODEL]: input.model, + [GEN_AI_ATTR_OPERATION_NAME]: 'chat', + provider: input.provider, + }) + } + + function recordStreamInterrupted(input: { + model: string + stage: 'mid_stream' | 'before_first_chunk' + span: GatewaySpan + }): void { + input.span.setStatus({ code: SpanStatusCode.ERROR, message: 'Gateway stream interrupted' }) + input.span.setAttribute(AIRI_ATTR_GEN_AI_STREAM_INTERRUPTED, true) + // Counter so alerts/dashboards can fire on interrupted streams; the + // span attribute alone only shows up in trace search, not metrics. + deps.genAi?.streamInterrupted.add(1, { + [GEN_AI_ATTR_REQUEST_MODEL]: input.model, + stage: input.stage, + }) + } + + return { + endSpan, + failSpan, + recordFirstToken, + recordMetrics, + recordRequestLog, + recordStreamInterrupted, + recordTtsBillingOnSpan, + recordUsageOnSpan, + runWithSpan, + setHttpStatus, + startChatSpan, + startTtsSpan, + } +} diff --git a/apps/server/src/routes/openai/v1/types.ts b/apps/server/src/routes/openai/v1/types.ts new file mode 100644 index 000000000..f453e9ba1 --- /dev/null +++ b/apps/server/src/routes/openai/v1/types.ts @@ -0,0 +1,36 @@ +import type { PostHog } from 'posthog-node' + +import type { GenAiMetrics, RateLimitMetrics, RevenueMetrics } from '../../../otel' +import type { ConfigKVService } from '../../../services/adapters/config-kv' +import type { BillingService } from '../../../services/domain/billing/billing-service' +import type { FluxMeter } from '../../../services/domain/billing/flux-meter' +import type { FluxService } from '../../../services/domain/flux' +import type { LlmRouterService } from '../../../services/domain/llm-router' +import type { ChatGenerationTrace, TtsGenerationTrace } from '../../../services/domain/llm-tracing' +import type { RequestLogService } from '../../../services/domain/request-log' + +import { startChatGeneration, startTtsGeneration } from '../../../services/domain/llm-tracing' + +export interface LlmTracingDeps { + startChatGeneration: (input: Parameters[0]) => ChatGenerationTrace + startTtsGeneration: (input: Parameters[0]) => TtsGenerationTrace +} + +export interface V1RouteDeps { + fluxService: FluxService + billingService: BillingService + configKV: ConfigKVService + requestLogService: RequestLogService + ttsMeter: FluxMeter + llmRouter: LlmRouterService + genAi?: GenAiMetrics | null + revenue?: RevenueMetrics | null + rateLimitMetrics?: RateLimitMetrics | null + posthog?: PostHog | null + llmTracing: LlmTracingDeps +} + +export const defaultLlmTracing: LlmTracingDeps = { + startChatGeneration, + startTtsGeneration, +} diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 76cfb8e60..afb344ad7 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -758,6 +758,9 @@ importers: jose: specifier: 'catalog:' version: 6.2.2 + ofetch: + specifier: 'catalog:' + version: 1.5.1 pg: specifier: ^8.20.0 version: 8.20.0