Skip to content
319 linesCodeBlameRaw
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 * Each run is held to its cost cap here too, not only by the harness in the
19 * sandbox: every answer's cost is added to the run's count (a Durable
20 * Object per run, `run-spend.ts`), and once the run has spent its cap its
21 * requests are refused with a 402 (`spend.ts`).
22 *
23 * The same address is the AI Gateway for a workspace's own code: a request
24 * with one of the workspace's access tokens (`g1t_…`) instead of a run's,
25 * at `/anthropic` in Anthropic's format or `/openai/v1` in OpenAI's, goes
26 * to `serve.ts`, and is logged and charged to the workspace.
27 */
28import { WorkerEntrypoint } from "cloudflare:workers";
29
30import {
31 type DiscoveryResult,
32 type GatewayModel,
33 type ModelDiscoveryApi,
34 type GatewayProvider,
35 type ModelUpstream,
36 type ServiceBinding,
37 type User,
38 billingClient,
39 identityClient,
40 integrationsClient,
41} from "@g1t/contracts";
42
43import { openaiError } from "./chat";
44import { discover } from "./discover";
45import { anthropicErrorType } from "./gateway";
46import { type AnthropicRequest, StreamTranslator, errorFromChat, estimateTokens, fromChat, toChat } from "./openai";
47import { isAnswer, runMayCall, tokenReport } from "./report";
48import { type HostedRouting, presentedToken, upstreamRequest } from "./route";
49import type { RunSpend } from "./run-spend";
50import { type GatewayDeps, isOpenAiPath, serveGateway } from "./serve";
51import { capOf, capReached, ceilingMicros, chargeFor, pricesFor, tooBusy } from "./spend";
52import { measure } from "./usage";
53
54export { RunSpend } from "./run-spend";
55
56interface Env extends HostedRouting {
57 INTEGRATIONS: ServiceBinding;
58 BILLING: ServiceBinding;
59 IDENTITY: ServiceBinding;
60 /** Each run's model spend, one object per model session. */
61 RUN_SPEND: DurableObjectNamespace<RunSpend>;
62}
63
64/**
65 * How long a looked-up token is trusted before it is looked up again. Short,
66 * because a run's token is closed the moment the run ends (and a connection
67 * may be removed mid-run): the proxy refuses it again within this long.
68 */
69const REMEMBER_MS = 10_000;
70const remembered = new Map<string, { upstream: ModelUpstream | null; until: number }>();
71
72async function lookUp(env: Env, token: string): Promise<ModelUpstream | null> {
73 const now = Date.now();
74 const hit = remembered.get(token);
75 if (hit && hit.until > now) return hit.upstream;
76 const upstream = await integrationsClient(env.INTEGRATIONS).modelUpstream(token);
77 if (remembered.size > 5_000) remembered.clear();
78 remembered.set(token, { upstream, until: now + REMEMBER_MS });
79 return upstream;
80}
81
82/** An error in the shape Anthropic's API uses, which the harness understands. */
83function refuse(status: number, message: string): Response {
84 return Response.json({ type: "error", error: { type: anthropicErrorType(status), message } }, { status });
85}
86
87/** The run's spend count, by its session's id. */
88function runSpend(env: Env, upstream: ModelUpstream): DurableObjectStub<RunSpend> {
89 return env.RUN_SPEND.get(env.RUN_SPEND.idFromName(upstream.session || `${upstream.workspace}/${upstream.repo}#${upstream.number}`));
90}
91
92/** One answer's place in its run's count, settled once the answer has gone by. */
93type Held = { spend: DurableObjectStub<RunSpend>; ticket: string; requested: Partial<AnthropicRequest> | null; bodyLength: number };
94
95/**
96 * Passes an answer through and, once it has all gone by, adds its cost to
97 * the run's count and tells billing what it used. Both happen after the
98 * answer, and a report that fails is dropped: the answer never waits on
99 * it or breaks for it.
100 */
101function counted(answer: Response, upstream: ModelUpstream, env: Env, ctx: ExecutionContext, held: Held): Response {
102 const { response, tokens, model } = measure(answer);
103 ctx.waitUntil(
104 (async () => {
105 const used = await tokens;
106 const answeredBy = await model;
107 const settle = (async () => {
108 const prices = pricesFor(upstream.model ?? answeredBy ?? held.requested?.model, upstream.route, await offered(env).catch(() => []));
109 const charge = chargeFor(prices, used, answer.ok, ceilingMicros(prices, held.bodyLength, held.requested?.max_tokens));
110 await held.spend.settle(held.ticket, charge);
111 })().catch((error: unknown) => console.error("models: a run's spend was not counted", upstream.session, String(error)));
112 const report = tokenReport(upstream, answeredBy, used);
113 const reported = report ? billingClient(env.BILLING).recordTokens(report).catch(() => undefined) : Promise.resolve();
114 await Promise.all([settle, reported]);
115 })().catch(() => undefined),
116 );
117 return response;
118}
119
120/** The fields of a request body the proxy reads, or null when it is not a JSON object. */
121function parsedBody(text: string): Record<string, unknown> | null {
122 try {
123 const parsed = JSON.parse(text) as unknown;
124 return parsed && typeof parsed === "object" && !Array.isArray(parsed) ? (parsed as Record<string, unknown>) : null;
125 } catch {
126 return null;
127 }
128}
129
130/** Sends an Anthropic request to a provider that speaks OpenAI's API. */
131async function viaChat(upstream: ModelUpstream, path: string, body: AnthropicRequest | null): Promise<Response> {
132 if (!path.startsWith("/v1/messages")) return refuse(404, `${path} has no counterpart at this provider.`);
133 if (!body) return refuse(400, "The request body is not a JSON object.");
134 const model = upstream.model ?? body.model ?? "";
135 if (path.startsWith("/v1/messages/count_tokens")) {
136 return Response.json({ input_tokens: estimateTokens(body) });
137 }
138
139 const headers = new Headers({ "content-type": "application/json" });
140 if (upstream.gatewayToken) headers.set("cf-aig-authorization", `Bearer ${upstream.gatewayToken}`);
141 if (upstream.apiKey) {
142 // `authorization` means a bearer token; any other header takes the key as it is.
143 const header = upstream.authHeader ?? "authorization";
144 headers.set(header, header === "authorization" ? `Bearer ${upstream.apiKey}` : upstream.apiKey);
145 }
146 const answer = await fetch(`${(upstream.baseUrl ?? "").replace(/\/+$/, "")}/chat/completions`, {
147 method: "POST",
148 headers,
149 body: JSON.stringify(toChat(body, model, { official: upstream.official, provider: upstream.provider })),
150 });
151 if (!answer.ok) {
152 return Response.json(errorFromChat(answer.status, await answer.text()), { status: answer.status });
153 }
154 if (!body.stream) return Response.json(fromChat((await answer.json()) as Record<string, unknown>, model));
155
156 const translator = new StreamTranslator(model);
157 const decoder = new TextDecoder();
158 const encoder = new TextEncoder();
159 const translated = answer.body!.pipeThrough(
160 new TransformStream<Uint8Array, Uint8Array>({
161 transform(chunk, controller) {
162 const out = translator.push(decoder.decode(chunk, { stream: true }));
163 if (out) controller.enqueue(encoder.encode(out));
164 },
165 flush(controller) {
166 const out = translator.push(decoder.decode()) + translator.finish();
167 if (out) controller.enqueue(encoder.encode(out));
168 },
169 }),
170 );
171 return new Response(translated, {
172 headers: { "content-type": "text/event-stream", "cache-control": "no-cache" },
173 });
174}
175
176// --- The AI Gateway ------------------------------------------------------------
177
178/**
179 * What gateway requests look up, remembered as briefly as a run's token
180 * is: a deleted token, a key added under Integrations or credit just bought
181 * takes effect within `REMEMBER_MS`. The catalogue changes rarely.
182 */
183const CATALOGUE_MS = 5 * 60_000;
184const callers = new Map<string, { value: User | null; until: number }>();
185const providers = new Map<string, { value: GatewayProvider[]; until: number }>();
186const admitted = new Map<string, { value: string | null; until: number }>();
187let catalogue: { models: GatewayModel[]; until: number } | null = null;
188
189async function cached<T>(cache: Map<string, { value: T; until: number }>, key: string, read: () => Promise<T>): Promise<T> {
190 const now = Date.now();
191 const hit = cache.get(key);
192 if (hit && hit.until > now) return hit.value;
193 const value = await read();
194 if (cache.size > 5_000) cache.clear();
195 cache.set(key, { value, until: now + REMEMBER_MS });
196 return value;
197}
198
199async function offered(env: Env): Promise<GatewayModel[]> {
200 if (catalogue && catalogue.until > Date.now()) return catalogue.models;
201 const models = await billingClient(env.BILLING).gatewayModels();
202 catalogue = { models, until: Date.now() + CATALOGUE_MS };
203 return models;
204}
205
206/** What serving a gateway request reaches outside the proxy. */
207function gatewayDeps(env: Env, ctx: ExecutionContext): GatewayDeps {
208 return {
209 hosted: env,
210 caller: (token) => cached(callers, token, () => identityClient(env.IDENTITY).userForAccessToken(token)),
211 providers: (workspace) => cached(providers, workspace, () => integrationsClient(env.INTEGRATIONS).gatewayProviders(workspace)),
212 offered: () => offered(env),
213 admit: (workspace) =>
214 cached(admitted, workspace, async () => {
215 const answer = await billingClient(env.BILLING).gatewayAdmit(workspace);
216 return answer.ok ? null : answer.error.message;
217 }),
218 record: (record) => billingClient(env.BILLING).recordGateway(record),
219 fetch: (url, init) => fetch(url, init),
220 waitUntil: (promise) => ctx.waitUntil(promise),
221 };
222}
223
224// --- Discovery -----------------------------------------------------------------
225
226/** Lists every provider's models and records what changed with billing (discover.ts). */
227function checkModels(env: Env, by: string): Promise<DiscoveryResult[]> {
228 const billing = billingClient(env.BILLING);
229 return discover(env, (url, init) => fetch(url, init), (provider, models, who, error) => billing.recordDiscovery(provider, models, who, error), by);
230}
231
232/**
233 * "Check for new models" in sudo, which binds this entrypoint. Only a
234 * service binding reaches it: nothing at models.g1t.sh does.
235 */
236export class Discovery extends WorkerEntrypoint<Env> implements ModelDiscoveryApi {
237 async check(by: string): Promise<DiscoveryResult[]> {
238 return checkModels(this.env, typeof by === "string" ? by : "");
239 }
240}
241
242export default {
243 /** Once a day: the providers' model lists, against the catalogue. */
244 async scheduled(_controller: ScheduledController, env: Env, ctx: ExecutionContext): Promise<void> {
245 ctx.waitUntil(
246 checkModels(env, "schedule")
247 .then((results) => {
248 for (const result of results) {
249 console.log(`models: ${result.provider} listed ${result.listed}, new ${result.added.length}, gone ${result.deprecated.length}${result.error ? `, failed: ${result.error}` : ""}`);
250 }
251 })
252 .catch((error) => console.error("models: checking the providers' lists failed", error)),
253 );
254 },
255
256 async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
257 const url = new URL(request.url);
258 if (url.pathname === "/" || url.pathname === "") {
259 return new Response("g1t's model proxy and AI Gateway. See https://docs.g1t.sh/guides/ai-gateway/\n");
260 }
261 const token = presentedToken(request.headers);
262 // OpenAI's format is the AI Gateway's alone: runs speak Anthropic's.
263 if (isOpenAiPath(url.pathname)) {
264 if (!token?.startsWith("g1t_")) {
265 return openaiError(401, "The AI Gateway takes a workspace's access token with the models:write scope, as the API key.");
266 }
267 return serveGateway(request, token, gatewayDeps(env, ctx));
268 }
269 if (!url.pathname.startsWith("/anthropic/")) return refuse(404, "Requests go to /anthropic/v1/… or /openai/v1/….");
270 // A workspace's own access token: the AI Gateway.
271 if (token?.startsWith("g1t_")) return serveGateway(request, token, gatewayDeps(env, ctx));
272 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.");
273 const upstream = await lookUp(env, token);
274 if (!upstream) return refuse(401, "This run's model token has expired, or its model connection was removed.");
275
276 const route = url.pathname.slice("/anthropic".length);
277 const path = route + url.search;
278 if (!runMayCall(route, request.method)) return refuse(404, `A run's model token reaches only /anthropic/v1/messages and /anthropic/v1/models, not ${route}.`);
279 const hasBody = request.method !== "GET" && request.method !== "HEAD";
280 const text = hasBody ? await request.text() : null;
281 const parsed = text === null ? null : parsedBody(text);
282
283 // A model's answer costs the run: it must be under its cap to start
284 // one, and the answer's cost is added to its count once it has gone by.
285 let held: Held | null = null;
286 if (isAnswer(route) && hasBody) {
287 const spend = runSpend(env, upstream);
288 const cap = capOf(upstream);
289 let admitted;
290 try {
291 admitted = await spend.admit(cap);
292 } catch (error) {
293 console.error("models: a run's spend could not be checked", upstream.session, String(error));
294 return refuse(503, "g1t could not check this run's spending just now. Try again.");
295 }
296 if (!admitted.ok) return admitted.reason === "cap" ? capReached(cap, admitted.spent) : tooBusy();
297 held = { spend, ticket: admitted.ticket, requested: parsed as Partial<AnthropicRequest> | null, bodyLength: text?.length ?? 0 };
298 }
299 // Both routes answer in Anthropic's shape, so one reading counts either.
300 const answer = (response: Response) => (held ? counted(response, upstream, env, ctx, held) : response);
301 try {
302 if (upstream.api === "openai") return answer(await viaChat(upstream, path, parsed as AnthropicRequest | null));
303
304 const { url: target, headers } = upstreamRequest(upstream, env, path, request.headers);
305 // A route that names a model gets it for every request of the run,
306 // including the harness's small background ones.
307 let body: string | null = text;
308 if (upstream.model && parsed && path.startsWith("/v1/messages")) {
309 body = JSON.stringify({ ...parsed, model: upstream.model });
310 headers.delete("content-length");
311 }
312 return answer(await fetch(target, { method: request.method, headers, body }));
313 } catch (error) {
314 // No answer: it cost nothing, and gives its place back.
315 if (held) ctx.waitUntil(held.spend.settle(held.ticket, 0).catch(() => undefined));
316 throw error;
317 }
318 },
319} satisfies ExportedHandler<Env>;