Skip to content

g1t/services/models/src/serve.ts

278 lines12,114 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

AI Gateway: OpenAI's format, open models, and your own providers1/**
2 * Serving one AI Gateway request, in either format, to whichever provider
3 * its model goes to: authenticate the token, route the model, admit it on
4 * g1t's key, translate when the caller's format and the provider's API
5 * differ, stream the answer back, and log (and on g1t's key charge) what
6 * it used.
7 *
8 * Everything outside the worker comes in through `GatewayDeps`, so the
9 * whole path runs in tests with no network.
10 */
11
12import type { GatewayModel, GatewayProvider, GatewayRecord, User } from "@g1t/contracts";
13
14import { type Kind, listModels, routeModel } from "./catalogue.ts";
15import { ChatStreamTranslator, Untranslatable, anthropicToChat, chatToAnthropic, openaiError } from "./chat.ts";
16import {
17 type Caller,
18 type Format,
19 ROUTES,
20 type Target,
21 anthropicError,
22 anthropicErrorType,
23 callerOf,
24 errorMessage,
25 gatewayOperation,
26 gatewayRecord,
27 hostedTarget,
28 ownTarget,
29 requestId,
30 scrub,
31 sessionOf,
32 targetUrl,
33 unpriced,
34 unpricedChat,
35} from "./gateway.ts";
36import { type AnthropicRequest, StreamTranslator, estimateTokens, fromChat, toChat } from "./openai.ts";
37import type { HostedRouting } from "./route.ts";
38import { NO_TOKENS, type Tokens, measure } from "./usage.ts";
39
40type Json = Record<string, unknown>;
41
42/** What serving needs from outside: who a token is, the workspace's providers, billing, the network. */
43export type GatewayDeps = {
44 hosted: HostedRouting;
45 caller(token: string): Promise<User | null>;
46 providers(workspace: string): Promise<GatewayProvider[]>;
47 offered(): Promise<GatewayModel[]>;
48 /** Why a workspace may not use g1t's models now, or null. */
49 admit(workspace: string): Promise<string | null>;
50 record(record: GatewayRecord): Promise<unknown>;
51 fetch(url: string, init: RequestInit): Promise<Response>;
52 waitUntil(promise: Promise<unknown>): void;
53};
54
55/** Headers of a provider's answer that reach the caller. */
56const KEPT = ["content-type", "cache-control", "retry-after", "request-id", "x-request-id"];
57
58function answerHeaders(upstream: Headers | null, id: string, contentType?: string): Headers {
59 const headers = new Headers();
60 for (const name of KEPT) {
61 const value = upstream?.get(name);
62 if (value) headers.set(name, value);
63 }
64 if (contentType) headers.set("content-type", contentType);
65 headers.set("x-g1t-request-id", id);
66 return headers;
67}
68
69/** An error in the caller's format. */
70function failure(format: Format, status: number, message: string, id?: string): Response {
71 const response = format === "openai" ? openaiError(status, message) : anthropicError(status, anthropicErrorType(status), message);
72 if (id) response.headers.set("x-g1t-request-id", id);
73 return response;
74}
75
76/** The error type a provider gave, when it is one of the caller's format's. */
77function upstreamType(body: string): string | null {
78 try {
79 const parsed = JSON.parse(body) as { error?: { type?: unknown } };
80 return typeof parsed.error?.type === "string" ? parsed.error.type : null;
81 } catch {
82 return null;
83 }
84}
85
86/** Streams a body through a translator, a chunk at a time. */
87function translated(body: ReadableStream<Uint8Array>, push: (text: string) => string, finish: () => string): ReadableStream<Uint8Array> {
88 const decoder = new TextDecoder();
89 const encoder = new TextEncoder();
90 return body.pipeThrough(
91 new TransformStream<Uint8Array, Uint8Array>({
92 transform(chunk, controller) {
93 const out = push(decoder.decode(chunk, { stream: true }));
94 if (out) controller.enqueue(encoder.encode(out));
95 },
96 flush(controller) {
97 const out = push(decoder.decode()) + finish();
98 if (out) controller.enqueue(encoder.encode(out));
99 },
100 }),
101 );
102}
103
104/** Serves an AI Gateway request; `token` is the workspace's access token it carried. */
105export async function serveGateway(request: Request, token: string, deps: GatewayDeps): Promise<Response> {
106 const started = Date.now();
107 const url = new URL(request.url);
108 const format: Format = url.pathname.startsWith("/openai/") ? "openai" : "anthropic";
109 const who = callerOf(await deps.caller(token));
110 if (!("caller" in who)) return failure(format, who.status, who.message);
111 const caller: Caller = who.caller;
112 const route = gatewayOperation(url.pathname, request.method);
113 if (!route) return failure(format, 404, ROUTES[format]);
114 const { op } = route;
115
116 if (op === "models") {
117 const [providers, offered] = await Promise.all([deps.providers(caller.workspace), deps.offered()]);
118 return Response.json({ object: "list", data: listModels(providers, offered) });
119 }
120
121 let body: Json;
122 try {
123 body = (await request.json()) as Json;
124 if (!body || typeof body !== "object" || Array.isArray(body)) throw new Error("not an object");
125 } catch {
126 return failure(format, 400, "The request body is not a JSON object.");
127 }
128 const requested = typeof body.model === "string" ? body.model.trim() : "";
129 const streamed = body.stream === true;
130 const id = requestId();
131 const kind: Kind = op === "embeddings" ? "embeddings" : "chat";
132
133 // Every request is logged once it is known whose it is. Counting tokens
134 // is a question about a request, not one, and is not.
135 const log = (input: { status: number; target?: Target | null; model?: string; tokens?: Tokens; error?: string | null }) => {
136 if (op === "count_tokens") return Promise.resolve();
137 const target = input.target ?? null;
138 const record = gatewayRecord({
139 id,
140 caller,
141 model: input.model ?? requested,
142 tokens: input.tokens ?? NO_TOKENS,
143 status: input.status,
144 ownKey: target?.ownKey ?? false,
145 streamed,
146 durationMs: Date.now() - started,
147 error: input.error ? scrub(input.error, target?.secrets ?? []) : null,
148 format,
149 provider: target?.provider ?? "",
150 connection: target?.connection ?? null,
151 });
152 return deps
153 .record(record)
154 .then(() => undefined)
155 .catch(() => undefined);
156 };
157 const refuse = (status: number, message: string, target?: Target | null) => {
158 deps.waitUntil(log({ status, target, error: message }));
159 return failure(format, status, message, id);
160 };
161
162 const [providers, offered] = await Promise.all([deps.providers(caller.workspace), deps.offered()]);
163 const routed = routeModel(requested, kind, providers, offered);
164 if (routed.to === "none") return refuse(routed.status, routed.message);
165
166 let target: Target | null;
167 if (routed.to === "g1t") {
168 const why = format === "anthropic" ? unpriced(body) : unpricedChat(body);
169 if (why) return refuse(400, why);
170 const refusal = await deps.admit(caller.workspace);
171 if (refusal) return refuse(402, refusal);
172 target = hostedTarget(deps.hosted, routed.entry.provider, request.headers, caller, sessionOf(caller.tokenId, new Date()));
173 if (!target) return refuse(503, `${routed.entry.name} is not available on this g1t: it has no way to ${routed.entry.provider}.`);
174 } else {
175 target = ownTarget(routed.provider, request.headers);
176 if (!target.base) return refuse(400, `${routed.provider.name} has no address. Give it one under Integrations.`, target);
177 }
178 // On g1t's key the catalogue's id is what billing prices; on the
179 // workspace's own, the model that answered.
180 const pricedAs = routed.to === "g1t" ? routed.entry.model : routed.model;
181 const shown = requested;
182
183 // What goes upstream, in the provider's API.
184 let upstreamBody: Json;
185 if (format === "anthropic" && target.api === "anthropic") {
186 upstreamBody = { ...body, model: routed.model };
187 } else if (format === "anthropic") {
188 if (op === "count_tokens") return Response.json({ input_tokens: estimateTokens(body as AnthropicRequest) });
189 upstreamBody = toChat(body as AnthropicRequest, routed.model, target.dialect);
190 } else if (target.api === "openai") {
191 upstreamBody = { ...body, model: routed.model };
192 // Usage at the end of a stream, to count it by; Mistral refuses the option.
193 if (op === "chat" && streamed && target.provider !== "mistral") {
194 upstreamBody.stream_options = { ...((body.stream_options as Json | undefined) ?? {}), include_usage: true };
195 }
196 } else {
197 try {
198 upstreamBody = chatToAnthropic(body, routed.model);
199 } catch (error) {
200 if (error instanceof Untranslatable) return refuse(400, error.message, target);
201 throw error;
202 }
203 }
204 const upstreamOp = target.api === "anthropic" ? (op === "count_tokens" ? "count_tokens" : "messages") : op === "embeddings" ? "embeddings" : "chat";
205
206 let answer: Response;
207 try {
208 answer = await deps.fetch(targetUrl(target, upstreamOp), { method: "POST", headers: target.headers, body: JSON.stringify(upstreamBody) });
209 } catch {
210 return refuse(502, `${target.connection ?? "The model provider"} could not be reached.`, target);
211 }
212
213 if (!answer.ok) {
214 const text = scrub(await answer.text(), target.secrets);
215 let message = errorMessage(answer.status, text);
216 if (target.ownKey && (answer.status === 401 || answer.status === 403)) {
217 message = `${target.connection} refused the workspace's key (${answer.status}): ${message} Check it under Integrations.`;
218 }
219 deps.waitUntil(log({ status: answer.status, target, error: message }));
220 // An error already in the caller's format keeps its type.
221 const native = (format === "anthropic") === (target.api === "anthropic") ? upstreamType(text) : null;
222 if (native && format === "anthropic") {
223 return new Response(JSON.stringify({ type: "error", error: { type: native, message } }), {
224 status: answer.status,
225 headers: answerHeaders(answer.headers, id, "application/json"),
226 });
227 }
228 const response = failure(format, answer.status, message, id);
229 const retry = answer.headers.get("retry-after");
230 if (retry) response.headers.set("retry-after", retry);
231 return response;
232 }
233
234 if (op === "count_tokens") {
235 return new Response(answer.body, { status: answer.status, headers: answerHeaders(answer.headers, id) });
236 }
237
238 // What it used, read from the provider's own answer as it passes.
239 const measured = measure(answer, target.api === "anthropic" ? "anthropic" : "openai");
240 const settle = target;
241 deps.waitUntil(
242 (async () => {
243 const tokens = await measured.tokens;
244 const answeredBy = await measured.model;
245 await log({ status: answer.status, target: settle, model: routed.to === "g1t" ? pricedAs : (answeredBy ?? pricedAs), tokens });
246 })().catch(() => undefined),
247 );
248 const passed = measured.response;
249 const eventStream = (passed.headers.get("content-type") ?? "").includes("text/event-stream");
250
251 // Same API both sides: the answer as it is.
252 if ((format === "anthropic") === (target.api === "anthropic")) {
253 return new Response(passed.body, { status: passed.status, headers: answerHeaders(passed.headers, id) });
254 }
255 const sse = "text/event-stream";
256 if (format === "anthropic") {
257 if (!eventStream) {
258 const whole = (await passed.json()) as Json;
259 return new Response(JSON.stringify(fromChat(whole, shown)), { headers: answerHeaders(passed.headers, id, "application/json") });
260 }
261 const translator = new StreamTranslator(shown);
262 const stream = translated(passed.body!, (text) => translator.push(text), () => translator.finish());
263 return new Response(stream, { headers: answerHeaders(null, id, sse) });
264 }
265 if (!eventStream) {
266 const whole = (await passed.json()) as Json;
267 return new Response(JSON.stringify(anthropicToChat(whole, shown)), { headers: answerHeaders(passed.headers, id, "application/json") });
268 }
269 const includeUsage = (body.stream_options as Json | undefined)?.include_usage === true;
270 const translator = new ChatStreamTranslator(shown, includeUsage);
271 const stream = translated(passed.body!, (text) => translator.push(text), () => translator.finish());
272 return new Response(stream, { headers: answerHeaders(null, id, sse) });
273}
274
275/** Whether a path is one of the AI Gateway's OpenAI-format routes. */
276export function isOpenAiPath(path: string): boolean {
277 return path === "/openai" || path.startsWith("/openai/");
278}