// // ChatCompletionsRoute.swift // Zyquo Router // // Author: Simon-Pierre Boucher // Mail: contact@spboucher.ai // // POST /v1/chat/completions — the core of the router. Parses the OpenAI // request, resolves the model (aliases, fallback chains), checks // capabilities, calls the upstream through the right translation path, // and returns a spec-exact response or byte-exact chunk stream. Upstream // failures become OpenAI-shaped errors (D7); transient ones retry before // the first forwarded byte (D8); usage is metered, estimated-and-flagged // when the upstream reports none (D9). // import Foundation import NIOHTTP1 struct ChatCompletionsRoute { let router: RequestRouter let providerKey: @Sendable (ProviderID) -> String? let usageMeter: UsageMeter let requestLog: RequestLogStore let retryPolicy: RetryPolicy /// Cap on stored body previews in the request log (inspector display). private static let previewLimit = 20_000 static func prettyJSON(_ data: Data) -> String { guard let object = try? JSONSerialization.jsonObject(with: data), let pretty = try? JSONSerialization.data(withJSONObject: object, options: [.prettyPrinted, .sortedKeys]) else { return String(decoding: data.prefix(previewLimit), as: UTF8.self) } return String(decoding: pretty.prefix(previewLimit), as: UTF8.self) } func handle(_ request: RouteRequest, localKey: APIKeyRecord?) async -> RouteResult { // Parse. let chat: ChatCompletionRequest do { chat = try ChatCompletionRequest(body: request.body) } catch let error as ChatCompletionRequest.ParseError { return OpenAIError.response( status: .badRequest, message: error.localizedDescription, type: "invalid_request_error", param: error.param ) } catch { return OpenAIError.response( status: .badRequest, message: "Malformed JSON body.", type: "invalid_request_error" ) } // Resolve the primary model + any configured fallback chain. var candidates: [RequestRouter.Resolution] = [] do { let primary = try router.resolve(chat.model) candidates.append(primary) for fallbackID in router.fallbackChains[primary.namespacedID] ?? [] { if let fallback = try? router.resolve(fallbackID) { candidates.append(fallback) } } } catch let error as RequestRouter.RoutingError { return Routes.routingErrorResponse(error) } catch { return OpenAIError.response(status: .internalServerError, message: "Internal error.", type: "server_error") } // Per-key model allow-list. if let allowed = localKey?.allowedModels, let primary = candidates.first, !allowed.contains(primary.namespacedID) { return OpenAIError.response( status: .forbidden, message: "This API key is not allowed to use `\(candidates[0].namespacedID)`.", type: "permission_error", code: "model_not_allowed" ) } // Try candidates in order; report the actually-used model honestly. let requestPreview = Self.prettyJSON(request.body) var lastFailure: ProviderError = .noModelAvailable(candidates[0].model.provider) for (index, resolution) in candidates.enumerated() { let isLastCandidate = index == candidates.count - 1 switch await attempt(chat: chat, resolution: resolution, localKey: localKey, requestPreview: requestPreview) { case .success(let result): return result case .failure(let error): lastFailure = error if isLastCandidate || !shouldFallback(on: error) { return errorResponse(for: error) } } } return errorResponse(for: lastFailure) } // MARK: - One candidate attempt private enum AttemptOutcome { case success(RouteResult) case failure(ProviderError) } private func attempt( chat: ChatCompletionRequest, resolution: RequestRouter.Resolution, localKey: APIKeyRecord?, requestPreview: String ) async -> AttemptOutcome { let model = resolution.model // Capability gates — clear OpenAI errors instead of upstream 400s. if chat.hasTools, !model.capabilities.tools { return .failure(.badRequest(model.provider, message: "`\(resolution.namespacedID)` does not support tools/function calling")) } if chat.hasImageContent, !model.capabilities.vision { return .failure(.badRequest(model.provider, message: "`\(resolution.namespacedID)` does not accept image input")) } guard let apiKey = providerKey(model.provider) else { return .failure(.missingAPIKey(model.provider)) } let call = UpstreamCall(model: model, apiKey: apiKey) if call.isNativeGemini, GeminiTranslator.hasRemoteImageURL(chat) { return .failure(.badRequest(model.provider, message: "Gemini requires images as base64 data URIs — remote image URLs are not fetched by the router")) } // A model that rejects non-streaming calls is transparently streamed // and aggregated when the client asked for a buffered response. let clientWantsStream = chat.stream let mustStreamUpstream = model.parameterSupport.requiresStreaming || CompatAdjuster.requiresStreamingOverride(model) if clientWantsStream { return await streamingAttempt(chat: chat, resolution: resolution, call: call, localKey: localKey, requestPreview: requestPreview) } if mustStreamUpstream { return await aggregatedStreamingAttempt(chat: chat, resolution: resolution, call: call, localKey: localKey, requestPreview: requestPreview) } return await bufferedAttempt(chat: chat, resolution: resolution, call: call, localKey: localKey, requestPreview: requestPreview) } private func shouldFallback(on error: ProviderError) -> Bool { switch error { case .rateLimited, .serverError, .networkError, .invalidResponse, .missingAPIKey, .invalidAPIKey: return true case .badRequest, .noModelAvailable, .cancelled: return false } } private func errorResponse(for error: ProviderError) -> RouteResult { let wire = error.openAIWire var extraHeaders: [(String, String)] = [] if case .rateLimited(_, let retryAfter) = error, let retryAfter { extraHeaders.append(("Retry-After", String(Int(retryAfter.rounded())))) } return OpenAIError.response( status: HTTPResponseStatus(statusCode: wire.status), message: wire.message, type: wire.type, code: wire.code, extraHeaders: extraHeaders ) } // MARK: - Buffered (non-streaming) private func bufferedAttempt( chat: ChatCompletionRequest, resolution: RequestRouter.Resolution, call: UpstreamCall, localKey: APIKeyRecord?, requestPreview: String ) async -> AttemptOutcome { let started = Date() let emitter = ChunkEmitter(model: resolution.namespacedID) do { let response: [String: Any] switch upstreamKind(call) { case .anthropic: let body = AnthropicTranslator.buildRequest(chat, model: resolution.model) response = AnthropicTranslator.translateResponse(try await call.complete(body: body), emitter: emitter) case .gemini: let body = GeminiTranslator.buildRequest(chat, model: resolution.model) let upstream = try await call.complete(body: body) if let blocked = GeminiTranslator.blockReason(upstream) { return .failure(.badRequest(.gemini, message: "Gemini blocked the prompt (reason: \(blocked))")) } response = GeminiTranslator.translateResponse(upstream, emitter: emitter) case .compat: let body = CompatAdjuster.adjustRequest(chat.raw, model: resolution.model, stream: false) var normalized = CompatAdjuster.normalizeResponse( try await call.complete(body: body), namespacedModel: resolution.namespacedID, provider: resolution.model.provider ) if normalized["usage"] == nil { normalized["usage"] = estimatedUsage(chat: chat, outputText: Self.responseText(normalized)) } response = normalized } let body = ChunkEmitter.serialize(response) await meter( response: response, resolution: resolution, localKey: localKey, started: started, streamed: false, status: 200, requestBody: requestPreview, responseBody: Self.prettyJSON(body) ) return .success(.complete( status: .ok, headers: [("Content-Type", "application/json")], body: body )) } catch let error as ProviderError { await meterFailure(resolution: resolution, localKey: localKey, started: started, streamed: false, error: error, requestBody: requestPreview) return .failure(error) } catch { let wrapped = ProviderError.networkError(underlying: error) await meterFailure(resolution: resolution, localKey: localKey, started: started, streamed: false, error: wrapped, requestBody: requestPreview) return .failure(wrapped) } } // MARK: - Streaming private func streamingAttempt( chat: ChatCompletionRequest, resolution: RequestRouter.Resolution, call: UpstreamCall, localKey: APIKeyRecord?, requestPreview: String ) async -> AttemptOutcome { // Pre-flight retries: transient failures before ANY byte reaches the // client are retried/fallback-able. Open the upstream stream and pull // its first event before committing to the client response. let started = Date() var attempt = 1 while true { do { let (events, firstEvent) = try await openUpstreamStream(chat: chat, call: call, resolution: resolution) return .success(streamResult( chat: chat, resolution: resolution, call: call, localKey: localKey, events: events, firstEvent: firstEvent, started: started, ttfb: Date().timeIntervalSince(started), requestPreview: requestPreview )) } catch let error as ProviderError { let retryable: Bool switch error { case .rateLimited: retryable = true case .serverError(_, let status, _): retryable = retryPolicy.shouldRetry(status: status, attempt: attempt) case .networkError: retryable = attempt < retryPolicy.maxAttempts default: retryable = false } var retryAfter: TimeInterval? if case .rateLimited(_, let after) = error { retryAfter = after } guard retryable, attempt < retryPolicy.maxAttempts else { return .failure(error) } try? await Task.sleep(nanoseconds: UInt64(retryPolicy.delay(attempt: attempt, retryAfter: retryAfter) * 1_000_000_000)) attempt += 1 } catch { return .failure(.networkError(underlying: error)) } } } /// Opens the upstream stream and awaits its first event so that upstream /// HTTP errors surface here (retryable) instead of mid-client-stream. private func openUpstreamStream( chat: ChatCompletionRequest, call: UpstreamCall, resolution: RequestRouter.Resolution ) async throws -> (AsyncThrowingStream.AsyncIterator, SSEEvent?) { let body: [String: Any] switch upstreamKind(call) { case .anthropic: body = AnthropicTranslator.buildRequest(chat, model: resolution.model) case .gemini: body = GeminiTranslator.buildRequest(chat, model: resolution.model) case .compat: body = CompatAdjuster.adjustRequest(chat.raw, model: resolution.model, stream: true) } var iterator = try call.stream(body: body).makeAsyncIterator() let first = try await iterator.next() return (iterator, first) } private func streamResult( chat: ChatCompletionRequest, resolution: RequestRouter.Resolution, call: UpstreamCall, localKey: APIKeyRecord?, events: AsyncThrowingStream.AsyncIterator, firstEvent: SSEEvent?, started: Date, ttfb: TimeInterval, requestPreview: String ) -> RouteResult { return .stream(status: .ok, headers: []) { writer in await self.usageMeter.streamBegan() defer { Task { await self.usageMeter.streamEnded() } } var iterator = events var next = firstEvent var usage: [String: Any]? var finishReason: String? var preview = "" func iterate(_ handle: (SSEEvent) async throws -> Void) async throws { while let event = next { try await handle(event) next = try await iterator.next() } } do { switch self.upstreamKind(call) { case .anthropic: var machine = AnthropicTranslator.StreamMachine( emitter: ChunkEmitter(model: resolution.namespacedID), includeUsage: chat.includeUsage ) try await iterate { event in let (payloads, _) = machine.consume(event) for payload in payloads { try await writer.send(raw: payload) } } usage = UsageBuilder.build( promptTokens: machine.promptTokens, completionTokens: machine.completionTokens, cachedTokens: machine.cachedTokens > 0 ? machine.cachedTokens : nil ) finishReason = machine.finishReasonSent preview = machine.textPreview case .gemini: var machine = GeminiTranslator.StreamMachine( emitter: ChunkEmitter(model: resolution.namespacedID), includeUsage: chat.includeUsage ) try await iterate { event in for payload in machine.consume(event) { try await writer.send(raw: payload) } } for payload in machine.finalPayloads() { try await writer.send(raw: payload) } usage = machine.lastUsage.map(GeminiTranslator.normalizedUsage) finishReason = "stop" preview = machine.textPreview case .compat: // Spec-discipline guards for deviant upstreams: synthesize // the role delta if the first chunk lacks it, and turn a // missing finish_reason (Perplexity puts it only in its // non-spec `.done` summary event) into a proper finish // chunk before usage/[DONE]. let emitter = ChunkEmitter(model: resolution.namespacedID) var usageChunkForwarded = false var roleForwarded = false var doneEventFinish: String? try await iterate { event in if event.data == "[DONE]" { return } guard let json = (try? JSONSerialization.jsonObject(with: Data(event.data.utf8))) as? [String: Any] else { return } guard let chunk = CompatAdjuster.normalizeChunk( json, namespacedModel: resolution.namespacedID, provider: resolution.model.provider, clientWantsUsage: chat.includeUsage ) else { // Swallowed event (usage-only chunk, Perplexity // `.done` summary): capture its usage + finish. if let chunkUsage = json["usage"] as? [String: Any] { usage = chunkUsage } if let finish = ((json["choices"] as? [[String: Any]])?.first?["finish_reason"] as? String) { doneEventFinish = CompatAdjuster.normalizeFinishReason(finish) } return } if let chunkUsage = chunk["usage"] as? [String: Any] { usage = chunkUsage if (chunk["choices"] as? [Any])?.isEmpty == true { usageChunkForwarded = true } } if let choices = chunk["choices"] as? [[String: Any]] { for choice in choices { if let finish = choice["finish_reason"] as? String { finishReason = finish } if let delta = choice["delta"] as? [String: Any] { if !roleForwarded, delta["role"] == nil, !(chunk["choices"] as? [Any] ?? []).isEmpty { try await writer.send(raw: emitter.roleChunk()) } roleForwarded = true if let content = delta["content"] as? String, preview.count < Self.previewLimit { preview += content } } } } try await writer.send(raw: ChunkEmitter.serialize(chunk)) } if finishReason == nil { // Even an empty stream must open with the role delta. if !roleForwarded { roleForwarded = true try await writer.send(raw: emitter.roleChunk()) } let reason = doneEventFinish ?? "stop" finishReason = reason try await writer.send(raw: emitter.finishChunk(reason: reason)) } // Client asked for usage but no usage chunk was forwarded // (upstream sent none, or only in a swallowed event). if chat.includeUsage, !usageChunkForwarded { let payload = usage ?? self.estimatedUsage(chat: chat, outputText: preview) usage = payload try await writer.send(raw: emitter.usageChunk(payload)) } } try await writer.sendDone() _ = finishReason await self.meter( usageDict: usage, resolution: resolution, localKey: localKey, started: started, streamed: true, status: 200, ttfb: ttfb, requestBody: requestPreview, responseBody: preview ) } catch let error as ProviderError { // Mid-stream failure: never retry (bytes were forwarded). // Emit a LiteLLM-style error frame, then terminate. await self.meterFailure(resolution: resolution, localKey: localKey, started: started, streamed: true, error: error, requestBody: requestPreview) let wire = error.openAIWire let frame = OpenAIError(error: .init(message: wire.message, type: wire.type, param: nil, code: wire.code)) try? await writer.send(raw: (try? JSONEncoder().encode(frame)) ?? Data()) try? await writer.sendDone() } } } /// Client asked non-streaming but the model only streams: aggregate. private func aggregatedStreamingAttempt( chat: ChatCompletionRequest, resolution: RequestRouter.Resolution, call: UpstreamCall, localKey: APIKeyRecord?, requestPreview: String ) async -> AttemptOutcome { let started = Date() do { var content = "" var reasoning = "" var toolCalls: [Int: [String: Any]] = [:] var finishReason = "stop" var usage: [String: Any]? // Produce normalized OpenAI chunk dicts from whichever upstream // wire this model speaks, then fold them into one completion. let kind = upstreamKind(call) let emitterForMachines = ChunkEmitter(model: resolution.namespacedID) var anthropicMachine = AnthropicTranslator.StreamMachine(emitter: emitterForMachines, includeUsage: true) var geminiMachine = GeminiTranslator.StreamMachine(emitter: emitterForMachines, includeUsage: true) let body: [String: Any] switch kind { case .anthropic: body = AnthropicTranslator.buildRequest(chat, model: resolution.model) case .gemini: body = GeminiTranslator.buildRequest(chat, model: resolution.model) case .compat: body = CompatAdjuster.adjustRequest(chat.raw, model: resolution.model, stream: true) } var chunks: [[String: Any]] = [] for try await event in try call.stream(body: body) { if event.data == "[DONE]" { break } switch kind { case .anthropic: chunks.append(contentsOf: anthropicMachine.consume(event).payloads.compactMap { (try? JSONSerialization.jsonObject(with: $0)) as? [String: Any] }) case .gemini: chunks.append(contentsOf: geminiMachine.consume(event).compactMap { (try? JSONSerialization.jsonObject(with: $0)) as? [String: Any] }) case .compat: if let json = (try? JSONSerialization.jsonObject(with: Data(event.data.utf8))) as? [String: Any], let chunk = CompatAdjuster.normalizeChunk( json, namespacedModel: resolution.namespacedID, provider: resolution.model.provider, clientWantsUsage: true ) { chunks.append(chunk) } } } if kind == .gemini { chunks.append(contentsOf: geminiMachine.finalPayloads().compactMap { (try? JSONSerialization.jsonObject(with: $0)) as? [String: Any] }) } for chunk in chunks { if let chunkUsage = chunk["usage"] as? [String: Any] { usage = chunkUsage } for choice in chunk["choices"] as? [[String: Any]] ?? [] { if let finish = choice["finish_reason"] as? String { finishReason = finish } guard let delta = choice["delta"] as? [String: Any] else { continue } content += delta["content"] as? String ?? "" reasoning += delta["reasoning_content"] as? String ?? "" for call in delta["tool_calls"] as? [[String: Any]] ?? [] { let index = call["index"] as? Int ?? 0 var existing = toolCalls[index] ?? ["type": "function", "function": ["name": "", "arguments": ""]] if let id = call["id"] as? String { existing["id"] = id } if let function = call["function"] as? [String: Any] { var merged = existing["function"] as? [String: Any] ?? [:] if let name = function["name"] as? String, !name.isEmpty { merged["name"] = name } merged["arguments"] = (merged["arguments"] as? String ?? "") + (function["arguments"] as? String ?? "") existing["function"] = merged } toolCalls[index] = existing } } } var message: [String: Any] = ["role": "assistant", "content": content] if !reasoning.isEmpty { message["reasoning_content"] = reasoning } if !toolCalls.isEmpty { message["tool_calls"] = toolCalls.sorted { $0.key < $1.key }.map(\.value) if finishReason == "stop" { finishReason = "tool_calls" } } let emitter = ChunkEmitter(model: resolution.namespacedID) let response = emitter.completion( message: message, finishReason: finishReason, usage: usage ?? estimatedUsage(chat: chat, outputText: content + reasoning) ) let responseData = ChunkEmitter.serialize(response) await meter( response: response, resolution: resolution, localKey: localKey, started: started, streamed: false, status: 200, requestBody: requestPreview, responseBody: Self.prettyJSON(responseData) ) return .success(.complete(status: .ok, headers: [("Content-Type", "application/json")], body: responseData)) } catch let error as ProviderError { await meterFailure(resolution: resolution, localKey: localKey, started: started, streamed: false, error: error, requestBody: requestPreview) return .failure(error) } catch { return .failure(.networkError(underlying: error)) } } // MARK: - Shared helpers private enum UpstreamKind { case anthropic, gemini, compat } private func upstreamKind(_ call: UpstreamCall) -> UpstreamKind { if call.model.provider == .anthropic { return .anthropic } if call.isNativeGemini { return .gemini } return .compat } private func estimatedUsage(chat: ChatCompletionRequest, outputText: String) -> [String: Any] { let promptText = chat.messages.map(\.flattenedText).joined(separator: "\n") return UsageBuilder.build( promptTokens: UsageBuilder.estimateTokens(promptText), completionTokens: UsageBuilder.estimateTokens(outputText), estimated: true ) } private static func responseText(_ response: [String: Any]) -> String { ((response["choices"] as? [[String: Any]])?.first?["message"] as? [String: Any])?["content"] as? String ?? "" } private func meter( response: [String: Any]? = nil, usageDict: [String: Any]? = nil, resolution: RequestRouter.Resolution, localKey: APIKeyRecord?, started: Date, streamed: Bool, status: Int, ttfb: TimeInterval? = nil, requestBody: String = "", responseBody: String = "" ) async { let usage = usageDict ?? response?["usage"] as? [String: Any] ?? [:] let prompt = usage["prompt_tokens"] as? Int ?? 0 let completion = usage["completion_tokens"] as? Int ?? 0 let estimated = (usage["x_zyquo"] as? [String: Any])?["usage_estimated"] as? Bool ?? false let reasoningTokens = (usage["completion_tokens_details"] as? [String: Any])?["reasoning_tokens"] as? Int let tokens = TokenUsage(inputTokens: prompt, outputTokens: completion, reasoningTokens: reasoningTokens) let cost = resolution.model.pricing?.cost(inputTokens: prompt, outputTokens: completion) let latency = Date().timeIntervalSince(started) await usageMeter.record(UsageRecord( namespacedModelID: resolution.namespacedID, provider: resolution.model.provider, localKeyName: localKey?.name, usage: tokens, usageEstimated: estimated, estimatedCost: cost, latency: latency, streamed: streamed, status: status )) await requestLog.append(RequestLogEntry( namespacedModelID: resolution.namespacedID, provider: resolution.model.provider, status: status, streamed: streamed, latency: latency, upstreamTTFB: ttfb, usage: tokens, usageEstimated: estimated, estimatedCost: cost, localKeyName: localKey?.name, requestBody: requestBody, responseBody: responseBody )) } private func meterFailure( resolution: RequestRouter.Resolution, localKey: APIKeyRecord?, started: Date, streamed: Bool, error: ProviderError, requestBody: String = "" ) async { let latency = Date().timeIntervalSince(started) let status = error.openAIWire.status await usageMeter.record(UsageRecord( namespacedModelID: resolution.namespacedID, provider: resolution.model.provider, localKeyName: localKey?.name, latency: latency, streamed: streamed, status: status )) await requestLog.append(RequestLogEntry( namespacedModelID: resolution.namespacedID, provider: resolution.model.provider, status: status, streamed: streamed, latency: latency, localKeyName: localKey?.name, errorMessage: error.openAIWire.message, requestBody: requestBody )) } }