Skip to content

g1t/services/models/src/index.ts

256 lines11,832 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Integrations: your own model provider, alerts that open issues, tickets agents read1/**
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:
Models per workspace: several providers, routed by kind of work7 * 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
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily11 * the run ends (the runner closes its session then, and lookups are kept
12 * only seconds), and nothing of the workspace's.
Integrations: your own model provider, alerts that open issues, tickets agents read13 *
Mission control shows model usage, yours and the workspace's: tokens, cost, active days, cache share, each day, and the mix14 * 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.
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens17 *
18 * The same address is the AI Gateway for a workspace's own code: a request
AI Gateway: OpenAI's format, open models, and your own providers19 * 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.
Integrations: your own model provider, alerts that open issues, tickets agents read22 */
Models: g1t keeps up with new models, and staff choose each default in sudo23import { WorkerEntrypoint } from "cloudflare:workers";
24
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens25import {
Models: g1t keeps up with new models, and staff choose each default in sudo26 type DiscoveryResult,
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens27 type GatewayModel,
Models: g1t keeps up with new models, and staff choose each default in sudo28 type ModelDiscoveryApi,
AI Gateway: OpenAI's format, open models, and your own providers29 type GatewayProvider,
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens30 type ModelUpstream,
31 type ServiceBinding,
32 type User,
33 billingClient,
34 identityClient,
35 integrationsClient,
36} from "@g1t/contracts";
Integrations: your own model provider, alerts that open issues, tickets agents read37
AI Gateway: OpenAI's format, open models, and your own providers38import { openaiError } from "./chat";
Models: g1t keeps up with new models, and staff choose each default in sudo39import { discover } from "./discover";
Models per workspace: several providers, routed by kind of work40import { type AnthropicRequest, StreamTranslator, errorFromChat, estimateTokens, fromChat, toChat } from "./openai";
Mission control shows model usage, yours and the workspace's: tokens, cost, active days, cache share, each day, and the mix41import { isAnswer, tokenReport } from "./report";
Integrations: your own model provider, alerts that open issues, tickets agents read42import { type HostedRouting, presentedToken, upstreamRequest } from "./route";
AI Gateway: OpenAI's format, open models, and your own providers43import { type GatewayDeps, isOpenAiPath, serveGateway } from "./serve";
44import { measure } from "./usage";
Integrations: your own model provider, alerts that open issues, tickets agents read45
46interface Env extends HostedRouting {
47 INTEGRATIONS: ServiceBinding;
Mission control shows model usage, yours and the workspace's: tokens, cost, active days, cache share, each day, and the mix48 BILLING: ServiceBinding;
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens49 IDENTITY: ServiceBinding;
Integrations: your own model provider, alerts that open issues, tickets agents read50}
51
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily52/**
53 * How long a looked-up token is trusted before it is looked up again. Short,
54 * because a run's token is closed the moment the run ends (and a connection
55 * may be removed mid-run): the proxy refuses it again within this long.
56 */
57const REMEMBER_MS = 10_000;
Integrations: your own model provider, alerts that open issues, tickets agents read58const remembered = new Map<string, { upstream: ModelUpstream | null; until: number }>();
59
60async function lookUp(env: Env, token: string): Promise<ModelUpstream | null> {
61 const now = Date.now();
62 const hit = remembered.get(token);
63 if (hit && hit.until > now) return hit.upstream;
64 const upstream = await integrationsClient(env.INTEGRATIONS).modelUpstream(token);
65 if (remembered.size > 5_000) remembered.clear();
66 remembered.set(token, { upstream, until: now + REMEMBER_MS });
67 return upstream;
68}
69
70/** An error in the shape Anthropic's API uses, which the harness understands. */
71function refuse(status: number, message: string): Response {
72 return Response.json(
73 { type: "error", error: { type: status === 401 ? "authentication_error" : "not_found_error", message } },
74 { status },
75 );
76}
77
Mission control shows model usage, yours and the workspace's: tokens, cost, active days, cache share, each day, and the mix78/**
79 * Passes an answer through and, once it has all gone by, tells billing what
80 * it used. Reporting happens after the answer, and a report that fails is
81 * dropped: the answer never waits on it or breaks for it.
82 */
83function counted(answer: Response, upstream: ModelUpstream, env: Env, ctx: ExecutionContext): Response {
84 const { response, tokens, model } = measure(answer);
85 ctx.waitUntil(
86 (async () => {
87 const report = tokenReport(upstream, await model, await tokens);
88 if (report) await billingClient(env.BILLING).recordTokens(report);
89 })().catch(() => undefined),
90 );
91 return response;
92}
93
Models per workspace: several providers, routed by kind of work94/** Sends an Anthropic request to a provider that speaks OpenAI's API. */
95async function viaChat(upstream: ModelUpstream, path: string, request: Request): Promise<Response> {
96 const body = (await request.json()) as AnthropicRequest;
97 const model = upstream.model ?? body.model ?? "";
98 if (path.startsWith("/v1/messages/count_tokens")) {
99 return Response.json({ input_tokens: estimateTokens(body) });
100 }
101 if (!path.startsWith("/v1/messages")) return refuse(404, `${path} has no counterpart at this provider.`);
102
103 const headers = new Headers({ "content-type": "application/json" });
Model providers: gateway tokens for endpoints, tidier rows, and the docs104 if (upstream.gatewayToken) headers.set("cf-aig-authorization", `Bearer ${upstream.gatewayToken}`);
Models per workspace: several providers, routed by kind of work105 if (upstream.apiKey) {
A catalogue of model providers, and settings that feel like settings106 // `authorization` means a bearer token; any other header takes the key as it is.
107 const header = upstream.authHeader ?? "authorization";
108 headers.set(header, header === "authorization" ? `Bearer ${upstream.apiKey}` : upstream.apiKey);
Models per workspace: several providers, routed by kind of work109 }
110 const answer = await fetch(`${(upstream.baseUrl ?? "").replace(/\/+$/, "")}/chat/completions`, {
111 method: "POST",
112 headers,
A catalogue of model providers, and settings that feel like settings113 body: JSON.stringify(toChat(body, model, { official: upstream.official, provider: upstream.provider })),
Models per workspace: several providers, routed by kind of work114 });
115 if (!answer.ok) {
116 return Response.json(errorFromChat(answer.status, await answer.text()), { status: answer.status });
117 }
118 if (!body.stream) return Response.json(fromChat((await answer.json()) as Record<string, unknown>, model));
119
120 const translator = new StreamTranslator(model);
121 const decoder = new TextDecoder();
122 const encoder = new TextEncoder();
123 const translated = answer.body!.pipeThrough(
124 new TransformStream<Uint8Array, Uint8Array>({
125 transform(chunk, controller) {
126 const out = translator.push(decoder.decode(chunk, { stream: true }));
127 if (out) controller.enqueue(encoder.encode(out));
128 },
129 flush(controller) {
130 const out = translator.push(decoder.decode()) + translator.finish();
131 if (out) controller.enqueue(encoder.encode(out));
132 },
133 }),
134 );
135 return new Response(translated, {
136 headers: { "content-type": "text/event-stream", "cache-control": "no-cache" },
137 });
138}
139
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens140// --- The AI Gateway ------------------------------------------------------------
141
142/**
143 * What gateway requests look up, remembered as briefly as a run's token
144 * is: a deleted token, a key added under Integrations or credit just bought
145 * takes effect within `REMEMBER_MS`. The catalogue changes rarely.
146 */
147const CATALOGUE_MS = 5 * 60_000;
148const callers = new Map<string, { value: User | null; until: number }>();
AI Gateway: OpenAI's format, open models, and your own providers149const providers = new Map<string, { value: GatewayProvider[]; until: number }>();
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens150const admitted = new Map<string, { value: string | null; until: number }>();
151let catalogue: { models: GatewayModel[]; until: number } | null = null;
152
153async function cached<T>(cache: Map<string, { value: T; until: number }>, key: string, read: () => Promise<T>): Promise<T> {
154 const now = Date.now();
155 const hit = cache.get(key);
156 if (hit && hit.until > now) return hit.value;
157 const value = await read();
158 if (cache.size > 5_000) cache.clear();
159 cache.set(key, { value, until: now + REMEMBER_MS });
160 return value;
161}
162
163async function offered(env: Env): Promise<GatewayModel[]> {
164 if (catalogue && catalogue.until > Date.now()) return catalogue.models;
165 const models = await billingClient(env.BILLING).gatewayModels();
166 catalogue = { models, until: Date.now() + CATALOGUE_MS };
167 return models;
168}
169
AI Gateway: OpenAI's format, open models, and your own providers170/** What serving a gateway request reaches outside the proxy. */
171function gatewayDeps(env: Env, ctx: ExecutionContext): GatewayDeps {
172 return {
173 hosted: env,
174 caller: (token) => cached(callers, token, () => identityClient(env.IDENTITY).userForAccessToken(token)),
175 providers: (workspace) => cached(providers, workspace, () => integrationsClient(env.INTEGRATIONS).gatewayProviders(workspace)),
176 offered: () => offered(env),
177 admit: (workspace) =>
178 cached(admitted, workspace, async () => {
179 const answer = await billingClient(env.BILLING).gatewayAdmit(workspace);
180 return answer.ok ? null : answer.error.message;
181 }),
182 record: (record) => billingClient(env.BILLING).recordGateway(record),
183 fetch: (url, init) => fetch(url, init),
184 waitUntil: (promise) => ctx.waitUntil(promise),
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens185 };
186}
187
Models: g1t keeps up with new models, and staff choose each default in sudo188// --- Discovery -----------------------------------------------------------------
189
190/** Lists every provider's models and records what changed with billing (discover.ts). */
191function checkModels(env: Env, by: string): Promise<DiscoveryResult[]> {
192 const billing = billingClient(env.BILLING);
193 return discover(env, (url, init) => fetch(url, init), (provider, models, who, error) => billing.recordDiscovery(provider, models, who, error), by);
194}
195
196/**
197 * "Check for new models" in sudo, which binds this entrypoint. Only a
198 * service binding reaches it: nothing at models.g1t.sh does.
199 */
200export class Discovery extends WorkerEntrypoint<Env> implements ModelDiscoveryApi {
201 async check(by: string): Promise<DiscoveryResult[]> {
202 return checkModels(this.env, typeof by === "string" ? by : "");
203 }
204}
205
Integrations: your own model provider, alerts that open issues, tickets agents read206export default {
Models: g1t keeps up with new models, and staff choose each default in sudo207 /** Once a day: the providers' model lists, against the catalogue. */
208 async scheduled(_controller: ScheduledController, env: Env, ctx: ExecutionContext): Promise<void> {
209 ctx.waitUntil(
210 checkModels(env, "schedule")
211 .then((results) => {
212 for (const result of results) {
213 console.log(`models: ${result.provider} listed ${result.listed}, new ${result.added.length}, gone ${result.deprecated.length}${result.error ? `, failed: ${result.error}` : ""}`);
214 }
215 })
216 .catch((error) => console.error("models: checking the providers' lists failed", error)),
217 );
218 },
219
Mission control shows model usage, yours and the workspace's: tokens, cost, active days, cache share, each day, and the mix220 async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
Integrations: your own model provider, alerts that open issues, tickets agents read221 const url = new URL(request.url);
222 if (url.pathname === "/" || url.pathname === "") {
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens223 return new Response("g1t's model proxy and AI Gateway. See https://docs.g1t.sh/guides/ai-gateway/\n");
Integrations: your own model provider, alerts that open issues, tickets agents read224 }
225 const token = presentedToken(request.headers);
AI Gateway: OpenAI's format, open models, and your own providers226 // OpenAI's format is the AI Gateway's alone: runs speak Anthropic's.
227 if (isOpenAiPath(url.pathname)) {
228 if (!token?.startsWith("g1t_")) {
229 return openaiError(401, "The AI Gateway takes a workspace's access token with the models:write scope, as the API key.");
230 }
231 return serveGateway(request, token, gatewayDeps(env, ctx));
232 }
233 if (!url.pathname.startsWith("/anthropic/")) return refuse(404, "Requests go to /anthropic/v1/… or /openai/v1/….");
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens234 // A workspace's own access token: the AI Gateway.
AI Gateway: OpenAI's format, open models, and your own providers235 if (token?.startsWith("g1t_")) return serveGateway(request, token, gatewayDeps(env, ctx));
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens236 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.");
Integrations: your own model provider, alerts that open issues, tickets agents read237 const upstream = await lookUp(env, token);
238 if (!upstream) return refuse(401, "This run's model token has expired, or its model connection was removed.");
239
240 const path = url.pathname.slice("/anthropic".length) + url.search;
Mission control shows model usage, yours and the workspace's: tokens, cost, active days, cache share, each day, and the mix241 // Both routes answer in Anthropic's shape, so one reading counts either.
242 const answer = (response: Response) => (isAnswer(url.pathname.slice("/anthropic".length)) ? counted(response, upstream, env, ctx) : response);
243 if (upstream.api === "openai") return answer(await viaChat(upstream, path, request));
Models per workspace: several providers, routed by kind of work244
Integrations: your own model provider, alerts that open issues, tickets agents read245 const { url: target, headers } = upstreamRequest(upstream, env, path, request.headers);
Models per workspace: several providers, routed by kind of work246 // A route that names a model gets it for every request of the run,
247 // including the harness's small background ones.
248 let body: BodyInit | null = request.method === "GET" || request.method === "HEAD" ? null : request.body;
249 if (upstream.model && body && path.startsWith("/v1/messages")) {
250 const parsed = (await request.json()) as Record<string, unknown>;
251 body = JSON.stringify({ ...parsed, model: upstream.model });
252 headers.delete("content-length");
253 }
Mission control shows model usage, yours and the workspace's: tokens, cost, active days, cache share, each day, and the mix254 return answer(await fetch(target, { method: request.method, headers, body }));
Integrations: your own model provider, alerts that open issues, tickets agents read255 },
256} satisfies ExportedHandler<Env>;

This file's history is long; its oldest lines are credited to the oldest commit read.