g1t/services/models/src/index.ts

146 lines7,030 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 */
18import { type ModelUpstream, type ServiceBinding, billingClient, integrationsClient } from "@g1t/contracts";
19
20import { type AnthropicRequest, StreamTranslator, errorFromChat, estimateTokens, fromChat, toChat } from "./openai";
21import { isAnswer, tokenReport } from "./report";
22import { type HostedRouting, presentedToken, upstreamRequest } from "./route";
23import { measure } from "./usage";
24
25interface Env extends HostedRouting {
26 INTEGRATIONS: ServiceBinding;
27 BILLING: ServiceBinding;
28}
29
30/**
31 * How long a looked-up token is trusted before it is looked up again. Short,
32 * because a run's token is closed the moment the run ends (and a connection
33 * may be removed mid-run): the proxy refuses it again within this long.
34 */
35const REMEMBER_MS = 10_000;
36const remembered = new Map<string, { upstream: ModelUpstream | null; until: number }>();
37
38async function lookUp(env: Env, token: string): Promise<ModelUpstream | null> {
39 const now = Date.now();
40 const hit = remembered.get(token);
41 if (hit && hit.until > now) return hit.upstream;
42 const upstream = await integrationsClient(env.INTEGRATIONS).modelUpstream(token);
43 if (remembered.size > 5_000) remembered.clear();
44 remembered.set(token, { upstream, until: now + REMEMBER_MS });
45 return upstream;
46}
47
48/** An error in the shape Anthropic's API uses, which the harness understands. */
49function refuse(status: number, message: string): Response {
50 return Response.json(
51 { type: "error", error: { type: status === 401 ? "authentication_error" : "not_found_error", message } },
52 { status },
53 );
54}
55
56/**
57 * Passes an answer through and, once it has all gone by, tells billing what
58 * it used. Reporting happens after the answer, and a report that fails is
59 * dropped: the answer never waits on it or breaks for it.
60 */
61function counted(answer: Response, upstream: ModelUpstream, env: Env, ctx: ExecutionContext): Response {
62 const { response, tokens, model } = measure(answer);
63 ctx.waitUntil(
64 (async () => {
65 const report = tokenReport(upstream, await model, await tokens);
66 if (report) await billingClient(env.BILLING).recordTokens(report);
67 })().catch(() => undefined),
68 );
69 return response;
70}
71
72/** Sends an Anthropic request to a provider that speaks OpenAI's API. */
73async function viaChat(upstream: ModelUpstream, path: string, request: Request): Promise<Response> {
74 const body = (await request.json()) as AnthropicRequest;
75 const model = upstream.model ?? body.model ?? "";
76 if (path.startsWith("/v1/messages/count_tokens")) {
77 return Response.json({ input_tokens: estimateTokens(body) });
78 }
79 if (!path.startsWith("/v1/messages")) return refuse(404, `${path} has no counterpart at this provider.`);
80
81 const headers = new Headers({ "content-type": "application/json" });
82 if (upstream.gatewayToken) headers.set("cf-aig-authorization", `Bearer ${upstream.gatewayToken}`);
83 if (upstream.apiKey) {
84 // `authorization` means a bearer token; any other header takes the key as it is.
85 const header = upstream.authHeader ?? "authorization";
86 headers.set(header, header === "authorization" ? `Bearer ${upstream.apiKey}` : upstream.apiKey);
87 }
88 const answer = await fetch(`${(upstream.baseUrl ?? "").replace(/\/+$/, "")}/chat/completions`, {
89 method: "POST",
90 headers,
91 body: JSON.stringify(toChat(body, model, { official: upstream.official, provider: upstream.provider })),
92 });
93 if (!answer.ok) {
94 return Response.json(errorFromChat(answer.status, await answer.text()), { status: answer.status });
95 }
96 if (!body.stream) return Response.json(fromChat((await answer.json()) as Record<string, unknown>, model));
97
98 const translator = new StreamTranslator(model);
99 const decoder = new TextDecoder();
100 const encoder = new TextEncoder();
101 const translated = answer.body!.pipeThrough(
102 new TransformStream<Uint8Array, Uint8Array>({
103 transform(chunk, controller) {
104 const out = translator.push(decoder.decode(chunk, { stream: true }));
105 if (out) controller.enqueue(encoder.encode(out));
106 },
107 flush(controller) {
108 const out = translator.push(decoder.decode()) + translator.finish();
109 if (out) controller.enqueue(encoder.encode(out));
110 },
111 }),
112 );
113 return new Response(translated, {
114 headers: { "content-type": "text/event-stream", "cache-control": "no-cache" },
115 });
116}
117
118export default {
119 async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
120 const url = new URL(request.url);
121 if (url.pathname === "/" || url.pathname === "") {
122 return new Response("g1t's model proxy, for g1t's sandboxes. See https://docs.g1t.sh/guides/models/\n");
123 }
124 if (!url.pathname.startsWith("/anthropic/")) return refuse(404, "Requests go to /anthropic/v1/….");
125 const token = presentedToken(request.headers);
126 if (!token?.startsWith("g1tm_")) return refuse(401, "This needs a g1t run's model token.");
127 const upstream = await lookUp(env, token);
128 if (!upstream) return refuse(401, "This run's model token has expired, or its model connection was removed.");
129
130 const path = url.pathname.slice("/anthropic".length) + url.search;
131 // Both routes answer in Anthropic's shape, so one reading counts either.
132 const answer = (response: Response) => (isAnswer(url.pathname.slice("/anthropic".length)) ? counted(response, upstream, env, ctx) : response);
133 if (upstream.api === "openai") return answer(await viaChat(upstream, path, request));
134
135 const { url: target, headers } = upstreamRequest(upstream, env, path, request.headers);
136 // A route that names a model gets it for every request of the run,
137 // including the harness's small background ones.
138 let body: BodyInit | null = request.method === "GET" || request.method === "HEAD" ? null : request.body;
139 if (upstream.model && body && path.startsWith("/v1/messages")) {
140 const parsed = (await request.json()) as Record<string, unknown>;
141 body = JSON.stringify({ ...parsed, model: upstream.model });
142 headers.delete("content-length");
143 }
144 return answer(await fetch(target, { method: request.method, headers, body }));
145 },
146} satisfies ExportedHandler<Env>;