pr_01m47d24b0e6n91zwymwxg0vpx/services/models/src/index.ts

118 lines5,642 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, and nothing of the workspace's.
12 *
13 * Responses stream through.
14 */
15import { type ModelUpstream, type ServiceBinding, integrationsClient } from "@g1t/contracts";
16
17import { type AnthropicRequest, StreamTranslator, errorFromChat, estimateTokens, fromChat, toChat } from "./openai";
18import { type HostedRouting, presentedToken, upstreamRequest } from "./route";
19
20interface Env extends HostedRouting {
21 INTEGRATIONS: ServiceBinding;
22}
23
24/** How long a looked-up token is trusted before it is looked up again. */
25const REMEMBER_MS = 60_000;
26const remembered = new Map<string, { upstream: ModelUpstream | null; until: number }>();
27
28async function lookUp(env: Env, token: string): Promise<ModelUpstream | null> {
29 const now = Date.now();
30 const hit = remembered.get(token);
31 if (hit && hit.until > now) return hit.upstream;
32 const upstream = await integrationsClient(env.INTEGRATIONS).modelUpstream(token);
33 if (remembered.size > 5_000) remembered.clear();
34 remembered.set(token, { upstream, until: now + REMEMBER_MS });
35 return upstream;
36}
37
38/** An error in the shape Anthropic's API uses, which the harness understands. */
39function refuse(status: number, message: string): Response {
40 return Response.json(
41 { type: "error", error: { type: status === 401 ? "authentication_error" : "not_found_error", message } },
42 { status },
43 );
44}
45
46/** Sends an Anthropic request to a provider that speaks OpenAI's API. */
47async function viaChat(upstream: ModelUpstream, path: string, request: Request): Promise<Response> {
48 const body = (await request.json()) as AnthropicRequest;
49 const model = upstream.model ?? body.model ?? "";
50 if (path.startsWith("/v1/messages/count_tokens")) {
51 return Response.json({ input_tokens: estimateTokens(body) });
52 }
53 if (!path.startsWith("/v1/messages")) return refuse(404, `${path} has no counterpart at this provider.`);
54
55 const headers = new Headers({ "content-type": "application/json" });
56 if (upstream.gatewayToken) headers.set("cf-aig-authorization", `Bearer ${upstream.gatewayToken}`);
57 if (upstream.apiKey) {
58 // `authorization` means a bearer token; any other header takes the key as it is.
59 const header = upstream.authHeader ?? "authorization";
60 headers.set(header, header === "authorization" ? `Bearer ${upstream.apiKey}` : upstream.apiKey);
61 }
62 const answer = await fetch(`${(upstream.baseUrl ?? "").replace(/\/+$/, "")}/chat/completions`, {
63 method: "POST",
64 headers,
65 body: JSON.stringify(toChat(body, model, { official: upstream.official, provider: upstream.provider })),
66 });
67 if (!answer.ok) {
68 return Response.json(errorFromChat(answer.status, await answer.text()), { status: answer.status });
69 }
70 if (!body.stream) return Response.json(fromChat((await answer.json()) as Record<string, unknown>, model));
71
72 const translator = new StreamTranslator(model);
73 const decoder = new TextDecoder();
74 const encoder = new TextEncoder();
75 const translated = answer.body!.pipeThrough(
76 new TransformStream<Uint8Array, Uint8Array>({
77 transform(chunk, controller) {
78 const out = translator.push(decoder.decode(chunk, { stream: true }));
79 if (out) controller.enqueue(encoder.encode(out));
80 },
81 flush(controller) {
82 const out = translator.push(decoder.decode()) + translator.finish();
83 if (out) controller.enqueue(encoder.encode(out));
84 },
85 }),
86 );
87 return new Response(translated, {
88 headers: { "content-type": "text/event-stream", "cache-control": "no-cache" },
89 });
90}
91
92export default {
93 async fetch(request: Request, env: Env): Promise<Response> {
94 const url = new URL(request.url);
95 if (url.pathname === "/" || url.pathname === "") {
96 return new Response("g1t's model proxy, for g1t's sandboxes. See https://docs.g1t.sh/guides/models/\n");
97 }
98 if (!url.pathname.startsWith("/anthropic/")) return refuse(404, "Requests go to /anthropic/v1/….");
99 const token = presentedToken(request.headers);
100 if (!token?.startsWith("g1tm_")) return refuse(401, "This needs a g1t run's model token.");
101 const upstream = await lookUp(env, token);
102 if (!upstream) return refuse(401, "This run's model token has expired, or its model connection was removed.");
103
104 const path = url.pathname.slice("/anthropic".length) + url.search;
105 if (upstream.api === "openai") return viaChat(upstream, path, request);
106
107 const { url: target, headers } = upstreamRequest(upstream, env, path, request.headers);
108 // A route that names a model gets it for every request of the run,
109 // including the harness's small background ones.
110 let body: BodyInit | null = request.method === "GET" || request.method === "HEAD" ? null : request.body;
111 if (upstream.model && body && path.startsWith("/v1/messages")) {
112 const parsed = (await request.json()) as Record<string, unknown>;
113 body = JSON.stringify({ ...parsed, model: upstream.model });
114 headers.delete("content-length");
115 }
116 return fetch(target, { method: request.method, headers, body });
117 },
118} satisfies ExportedHandler<Env>;