Skip to content

g1t/services/models/src/index.ts

220 lines10,262 bytesCodeBlame
1/**
2 * The model proxy: every model request a g1t sandbox makes comes through
3 * here, at `https://models.g1t.sh/anthropic`.
4 *
5 * A sandbox holds a token for its one run, never a key. The proxy looks the
6 * token up and forwards the request with the credentials for that run:
7 * g1t's AI Gateway when g1t pays, or one of the workspace's own providers
8 * when it does. A provider that speaks OpenAI's API gets the request
9 * translated, and its answer translated back. So a sandbox that is tricked
10 * into printing its environment gives away a token that stops working when
11 * the run ends (the runner closes its session then, and lookups are kept
12 * only seconds), and nothing of the workspace's.
13 *
14 * Responses stream through. What each answer used is read from a copy as
15 * it passes and reported to billing afterwards, counted per run for usage
16 * views.
17 *
18 * The same address is the AI Gateway for a workspace's own code: a request
19 * with one of the workspace's access tokens (`g1t_…`) instead of a run's,
20 * at `/anthropic` in Anthropic's format or `/openai/v1` in OpenAI's, goes
21 * to `serve.ts`, and is logged and charged to the workspace.
22 */
23import {
24 type GatewayModel,
25 type GatewayProvider,
26 type ModelUpstream,
27 type ServiceBinding,
28 type User,
29 billingClient,
30 identityClient,
31 integrationsClient,
32} from "@g1t/contracts";
33
34import { openaiError } from "./chat";
35import { type AnthropicRequest, StreamTranslator, errorFromChat, estimateTokens, fromChat, toChat } from "./openai";
36import { isAnswer, tokenReport } from "./report";
37import { type HostedRouting, presentedToken, upstreamRequest } from "./route";
38import { type GatewayDeps, isOpenAiPath, serveGateway } from "./serve";
39import { measure } from "./usage";
40
41interface Env extends HostedRouting {
42 INTEGRATIONS: ServiceBinding;
43 BILLING: ServiceBinding;
44 IDENTITY: ServiceBinding;
45}
46
47/**
48 * How long a looked-up token is trusted before it is looked up again. Short,
49 * because a run's token is closed the moment the run ends (and a connection
50 * may be removed mid-run): the proxy refuses it again within this long.
51 */
52const REMEMBER_MS = 10_000;
53const remembered = new Map<string, { upstream: ModelUpstream | null; until: number }>();
54
55async function lookUp(env: Env, token: string): Promise<ModelUpstream | null> {
56 const now = Date.now();
57 const hit = remembered.get(token);
58 if (hit && hit.until > now) return hit.upstream;
59 const upstream = await integrationsClient(env.INTEGRATIONS).modelUpstream(token);
60 if (remembered.size > 5_000) remembered.clear();
61 remembered.set(token, { upstream, until: now + REMEMBER_MS });
62 return upstream;
63}
64
65/** An error in the shape Anthropic's API uses, which the harness understands. */
66function refuse(status: number, message: string): Response {
67 return Response.json(
68 { type: "error", error: { type: status === 401 ? "authentication_error" : "not_found_error", message } },
69 { status },
70 );
71}
72
73/**
74 * Passes an answer through and, once it has all gone by, tells billing what
75 * it used. Reporting happens after the answer, and a report that fails is
76 * dropped: the answer never waits on it or breaks for it.
77 */
78function counted(answer: Response, upstream: ModelUpstream, env: Env, ctx: ExecutionContext): Response {
79 const { response, tokens, model } = measure(answer);
80 ctx.waitUntil(
81 (async () => {
82 const report = tokenReport(upstream, await model, await tokens);
83 if (report) await billingClient(env.BILLING).recordTokens(report);
84 })().catch(() => undefined),
85 );
86 return response;
87}
88
89/** Sends an Anthropic request to a provider that speaks OpenAI's API. */
90async function viaChat(upstream: ModelUpstream, path: string, request: Request): Promise<Response> {
91 const body = (await request.json()) as AnthropicRequest;
92 const model = upstream.model ?? body.model ?? "";
93 if (path.startsWith("/v1/messages/count_tokens")) {
94 return Response.json({ input_tokens: estimateTokens(body) });
95 }
96 if (!path.startsWith("/v1/messages")) return refuse(404, `${path} has no counterpart at this provider.`);
97
98 const headers = new Headers({ "content-type": "application/json" });
99 if (upstream.gatewayToken) headers.set("cf-aig-authorization", `Bearer ${upstream.gatewayToken}`);
100 if (upstream.apiKey) {
101 // `authorization` means a bearer token; any other header takes the key as it is.
102 const header = upstream.authHeader ?? "authorization";
103 headers.set(header, header === "authorization" ? `Bearer ${upstream.apiKey}` : upstream.apiKey);
104 }
105 const answer = await fetch(`${(upstream.baseUrl ?? "").replace(/\/+$/, "")}/chat/completions`, {
106 method: "POST",
107 headers,
108 body: JSON.stringify(toChat(body, model, { official: upstream.official, provider: upstream.provider })),
109 });
110 if (!answer.ok) {
111 return Response.json(errorFromChat(answer.status, await answer.text()), { status: answer.status });
112 }
113 if (!body.stream) return Response.json(fromChat((await answer.json()) as Record<string, unknown>, model));
114
115 const translator = new StreamTranslator(model);
116 const decoder = new TextDecoder();
117 const encoder = new TextEncoder();
118 const translated = answer.body!.pipeThrough(
119 new TransformStream<Uint8Array, Uint8Array>({
120 transform(chunk, controller) {
121 const out = translator.push(decoder.decode(chunk, { stream: true }));
122 if (out) controller.enqueue(encoder.encode(out));
123 },
124 flush(controller) {
125 const out = translator.push(decoder.decode()) + translator.finish();
126 if (out) controller.enqueue(encoder.encode(out));
127 },
128 }),
129 );
130 return new Response(translated, {
131 headers: { "content-type": "text/event-stream", "cache-control": "no-cache" },
132 });
133}
134
135// --- The AI Gateway ------------------------------------------------------------
136
137/**
138 * What gateway requests look up, remembered as briefly as a run's token
139 * is: a deleted token, a key added under Integrations or credit just bought
140 * takes effect within `REMEMBER_MS`. The catalogue changes rarely.
141 */
142const CATALOGUE_MS = 5 * 60_000;
143const callers = new Map<string, { value: User | null; until: number }>();
144const providers = new Map<string, { value: GatewayProvider[]; until: number }>();
145const admitted = new Map<string, { value: string | null; until: number }>();
146let catalogue: { models: GatewayModel[]; until: number } | null = null;
147
148async function cached<T>(cache: Map<string, { value: T; until: number }>, key: string, read: () => Promise<T>): Promise<T> {
149 const now = Date.now();
150 const hit = cache.get(key);
151 if (hit && hit.until > now) return hit.value;
152 const value = await read();
153 if (cache.size > 5_000) cache.clear();
154 cache.set(key, { value, until: now + REMEMBER_MS });
155 return value;
156}
157
158async function offered(env: Env): Promise<GatewayModel[]> {
159 if (catalogue && catalogue.until > Date.now()) return catalogue.models;
160 const models = await billingClient(env.BILLING).gatewayModels();
161 catalogue = { models, until: Date.now() + CATALOGUE_MS };
162 return models;
163}
164
165/** What serving a gateway request reaches outside the proxy. */
166function gatewayDeps(env: Env, ctx: ExecutionContext): GatewayDeps {
167 return {
168 hosted: env,
169 caller: (token) => cached(callers, token, () => identityClient(env.IDENTITY).userForAccessToken(token)),
170 providers: (workspace) => cached(providers, workspace, () => integrationsClient(env.INTEGRATIONS).gatewayProviders(workspace)),
171 offered: () => offered(env),
172 admit: (workspace) =>
173 cached(admitted, workspace, async () => {
174 const answer = await billingClient(env.BILLING).gatewayAdmit(workspace);
175 return answer.ok ? null : answer.error.message;
176 }),
177 record: (record) => billingClient(env.BILLING).recordGateway(record),
178 fetch: (url, init) => fetch(url, init),
179 waitUntil: (promise) => ctx.waitUntil(promise),
180 };
181}
182
183export default {
184 async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
185 const url = new URL(request.url);
186 if (url.pathname === "/" || url.pathname === "") {
187 return new Response("g1t's model proxy and AI Gateway. See https://docs.g1t.sh/guides/ai-gateway/\n");
188 }
189 const token = presentedToken(request.headers);
190 // OpenAI's format is the AI Gateway's alone: runs speak Anthropic's.
191 if (isOpenAiPath(url.pathname)) {
192 if (!token?.startsWith("g1t_")) {
193 return openaiError(401, "The AI Gateway takes a workspace's access token with the models:write scope, as the API key.");
194 }
195 return serveGateway(request, token, gatewayDeps(env, ctx));
196 }
197 if (!url.pathname.startsWith("/anthropic/")) return refuse(404, "Requests go to /anthropic/v1/… or /openai/v1/….");
198 // A workspace's own access token: the AI Gateway.
199 if (token?.startsWith("g1t_")) return serveGateway(request, token, gatewayDeps(env, ctx));
200 if (!token?.startsWith("g1tm_")) return refuse(401, "This needs a g1t run's model token, or a workspace's access token for the AI Gateway.");
201 const upstream = await lookUp(env, token);
202 if (!upstream) return refuse(401, "This run's model token has expired, or its model connection was removed.");
203
204 const path = url.pathname.slice("/anthropic".length) + url.search;
205 // Both routes answer in Anthropic's shape, so one reading counts either.
206 const answer = (response: Response) => (isAnswer(url.pathname.slice("/anthropic".length)) ? counted(response, upstream, env, ctx) : response);
207 if (upstream.api === "openai") return answer(await viaChat(upstream, path, request));
208
209 const { url: target, headers } = upstreamRequest(upstream, env, path, request.headers);
210 // A route that names a model gets it for every request of the run,
211 // including the harness's small background ones.
212 let body: BodyInit | null = request.method === "GET" || request.method === "HEAD" ? null : request.body;
213 if (upstream.model && body && path.startsWith("/v1/messages")) {
214 const parsed = (await request.json()) as Record<string, unknown>;
215 body = JSON.stringify({ ...parsed, model: upstream.model });
216 headers.delete("content-length");
217 }
218 return answer(await fetch(target, { method: request.method, headers, body }));
219 },
220} satisfies ExportedHandler<Env>;