// // UpstreamCall.swift // Zyquo Router // // Author: Simon-Pierre Boucher // Mail: contact@spboucher.ai // // Executes one request against one upstream provider: endpoint + auth // construction per wire format, JSON POST for non-streaming (retries live // in StreamingService.postJSON), raw SSE event stream for streaming. // Cancellation of the calling task cancels the upstream transfer. // import Foundation struct UpstreamCall { let model: AIModel let apiKey: String private var provider: ProviderID { model.provider } enum CallError: Error { case noEndpoint(ProviderID) } // MARK: - Endpoint + auth private func urlRequest(streaming: Bool, body: [String: Any]) throws -> URLRequest { let url: URL switch provider.wireFormat { case .anthropicMessages: guard let base = model.customBaseURL ?? provider.defaultBaseURL else { throw CallError.noEndpoint(provider) } url = base.appendingPathComponent("messages") case .openAIChatCompletions where provider == .gemini && Self.usesNativeGemini: // Native generateContent (D10): Cloud's catalog stores the // OpenAI-compat base (…/v1beta/openai); derive the native root. let root = "https://generativelanguage.googleapis.com/v1beta" let verb = streaming ? "streamGenerateContent?alt=sse" : "generateContent" guard let native = URL(string: "\(root)/models/\(model.id):\(verb)") else { throw CallError.noEndpoint(provider) } url = native case .openAIChatCompletions: guard let base = model.customBaseURL ?? provider.defaultBaseURL else { throw CallError.noEndpoint(provider) } url = base.appendingPathComponent("chat/completions") } var request = URLRequest(url: url) request.httpMethod = "POST" request.setValue("application/json", forHTTPHeaderField: "Content-Type") switch provider { case .anthropic: request.setValue(apiKey, forHTTPHeaderField: "x-api-key") request.setValue("2023-06-01", forHTTPHeaderField: "anthropic-version") case .gemini: request.setValue(apiKey, forHTTPHeaderField: "x-goog-api-key") default: request.setValue("Bearer \(apiKey)", forHTTPHeaderField: "Authorization") } if streaming { request.setValue("text/event-stream", forHTTPHeaderField: "Accept") } request.httpBody = try JSONSerialization.data(withJSONObject: body) return request } /// Gemini goes through the native generateContent translation. static let usesNativeGemini = true /// Whether this call's upstream speaks the native Gemini API. var isNativeGemini: Bool { provider == .gemini && Self.usesNativeGemini } // MARK: - Execution /// Non-streaming: returns the upstream JSON object. /// StreamingService.postJSON already retries 429/5xx with backoff. func complete(body: [String: Any]) async throws -> [String: Any] { let request = try urlRequest(streaming: false, body: body) let data = try await StreamingService.postJSON(request, provider: provider) guard let json = (try? JSONSerialization.jsonObject(with: data)) as? [String: Any] else { throw ProviderError.invalidResponse(provider, detail: "response is not a JSON object") } return json } /// Streaming: raw upstream SSE events. Errors before the first event are /// retryable by the caller (never after the first forwarded byte). func stream(body: [String: Any]) throws -> AsyncThrowingStream { let request = try urlRequest(streaming: true, body: body) return StreamingService.sseEvents(for: request, provider: provider) } } // MARK: - ProviderError → OpenAI wire error extension ProviderError { /// Maps upstream failures to (HTTP status, OpenAI error type/code) per /// decision D7 — clear messages, no provider payload shapes, no keys. var openAIWire: (status: Int, type: String, code: String?, message: String) { switch self { case .invalidAPIKey(let provider): return (401, "authentication_error", "invalid_provider_key", "The stored \(provider.displayName) API key was rejected upstream. Update it in Zyquo Router → Keys.") case .missingAPIKey(let provider): return (401, "authentication_error", "missing_provider_key", "No \(provider.displayName) API key is configured. Add one in Zyquo Router → Keys.") case .rateLimited(let provider, let retryAfter): let hint = retryAfter.map { " Retry in \(Int($0.rounded()))s." } ?? "" return (429, "rate_limit_error", "upstream_rate_limited", "\(provider.displayName) rate-limited the request.\(hint)") case .badRequest(let provider, let message): return (400, "invalid_request_error", nil, "\(provider.displayName) rejected the request\(message.map { ": \($0)" } ?? ".")") case .serverError(let provider, let status, _): return (502, "api_error", "upstream_error", "\(provider.displayName) upstream error (HTTP \(status)).") case .networkError(let underlying): if (underlying as? URLError)?.code == .timedOut { return (504, "api_error", "upstream_timeout", "The upstream request timed out.") } return (502, "api_error", "upstream_unreachable", "Could not reach the upstream provider: \(underlying.localizedDescription)") case .invalidResponse(let provider, let detail): return (502, "api_error", "upstream_invalid_response", "Unexpected response from \(provider.displayName): \(detail)") case .noModelAvailable(let provider): return (404, "invalid_request_error", "model_not_found", "No model available for \(provider.displayName).") case .cancelled: return (499, "api_error", "client_disconnected", "The client disconnected.") } } }