flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/models/src/index.ts

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