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/apps/api/src/lib.rs

504 lines19,210 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.

API and MCP server in Rust; a public index at the API root1//! The public API: REST at api.g1t.sh and the MCP server at mcp.g1t.sh.
2//!
3//! One Worker, two hostnames. Both are thin adapters over the same
4//! operations (see [`operations::Op`]), which call the services that own
5//! the data. This Worker holds none.
6
7mod mcp;
8mod oauth;
9mod openapi;
10mod operations;
11mod rest;
12
Agents as a team: lifecycle, merge queue, billing and a new shell13use g1t_contracts::billing::FinishRunArgs;
API and MCP server in Rust; a public index at the API root14use g1t_contracts::identity::{
15 DeviceClaim, DeviceClaimArgs, DeviceStart, DeviceStartArgs, TokenArgs,
16};
Agents as a team: lifecycle, merge queue, billing and a new shell17use g1t_contracts::work::{
18 CheckRun, QueueState, ReportChecksArgs, ReportPlanArgs, ReportQueueArgs, ReportReviewArgs,
19};
20use g1t_contracts::identity::AgentScope;
21use g1t_contracts::{Failure, FailureCode, Outcome, PrincipalKind, Viewer};
API and MCP server in Rust; a public index at the API root22use serde_json::{Value, json};
23use worker::{Context, Env, Method, Request, Response, Result, event};
24
25use operations::Services;
26
27const API: &str = "https://api.g1t.sh";
28
29fn method_name(method: Method) -> &'static str {
30 match method {
31 Method::Get => "GET",
32 Method::Post => "POST",
33 Method::Patch => "PATCH",
34 Method::Put => "PUT",
35 Method::Delete => "DELETE",
36 Method::Options => "OPTIONS",
37 Method::Head => "HEAD",
38 _ => "OTHER",
39 }
40}
41
42/// An error in the shape every endpoint uses.
43fn failure(failure: &Failure) -> Result<Response> {
44 Ok(Response::from_json(&json!({ "error": failure }))?.with_status(failure.code.http_status()))
45}
46
47fn fail(code: FailureCode, message: &str) -> Result<Response> {
48 failure(&Failure {
49 code,
50 message: message.to_owned(),
51 })
52}
53
54/// A request body as JSON. An empty or malformed body is no input.
55async fn json_body(request: &mut Request) -> Value {
56 request.json().await.unwrap_or(Value::Null)
57}
58
59/// Who a request's `Authorization: Bearer g1t_…` names. A missing token is
60/// an anonymous viewer; a wrong one is refused, so that a typo does not
61/// silently look signed out.
62async fn authenticate(
63 request: &Request,
64 services: &Services,
65) -> Result<std::result::Result<Viewer, Response>> {
66 let header = request.headers().get("authorization")?.unwrap_or_default();
67 let token = match header.split_once(' ') {
68 Some((scheme, token)) if scheme.eq_ignore_ascii_case("bearer") && !token.is_empty() => {
69 token.trim()
70 }
71 _ => return Ok(Ok(None)),
72 };
73 let viewer: Viewer = g1t_kit::call(
74 &services.identity,
75 "user_for_access_token",
76 &TokenArgs {
77 token: token.to_owned(),
78 },
79 )
80 .await?;
81 if viewer.is_some() {
82 return Ok(Ok(viewer));
83 }
84 let mut response = fail(FailureCode::Unauthenticated, "Invalid access token.")?;
85 // Tells an MCP client where to sign in again.
86 response.headers_mut().set(
87 "www-authenticate",
88 &format!("{}, error=\"invalid_token\"", oauth::MCP_CHALLENGE),
89 )?;
90 Ok(Err(response))
91}
92
93/// Where everything is, for someone or something exploring the API.
94fn index() -> Value {
Agents as a team: lifecycle, merge queue, billing and a new shell95 let repo = format!("{API}/repos/{{owner}}/{{name}}");
API and MCP server in Rust; a public index at the API root96 json!({
97 "documentation_url": "https://docs.g1t.sh/api/reference/",
98 "openapi_url": format!("{API}/openapi.json"),
99 "mcp_url": "https://mcp.g1t.sh",
Agents as a team: lifecycle, merge queue, billing and a new shell100 "current_user_url": format!("{API}/user"),
101 "workspaces_url": format!("{API}/workspaces"),
102 "repositories_url": format!("{API}/repos{{?q}}"),
API and MCP server in Rust; a public index at the API root103 "repository_url": repo,
104 "repository_events_url": format!("{repo}/events{{?before}}"),
105 "labels_url": format!("{repo}/labels"),
106 "issues_url": format!("{repo}/issues{{?state,label}}"),
107 "issue_url": format!("{repo}/issues/{{number}}"),
108 "issue_comments_url": format!("{repo}/issues/{{number}}/comments"),
109 "pulls_url": format!("{repo}/pulls{{?state}}"),
110 "pull_url": format!("{repo}/pulls/{{number}}"),
111 "pull_changes_url": format!("{repo}/pulls/{{number}}/changes"),
Acceptance checks in sandboxes, line comments and review verdicts112 "pull_reviews_url": format!("{repo}/pulls/{{number}}/reviews"),
API and MCP server in Rust; a public index at the API root113 "pull_session_url": format!("{repo}/pulls/{{number}}/session{{?after}}"),
Agents as a team: lifecycle, merge queue, billing and a new shell114 "device_code_url": format!("{API}/device/code"),
115 "device_token_url": format!("{API}/device/token"),
API and MCP server in Rust; a public index at the API root116 "oauth_metadata_url": format!("{API}/.well-known/oauth-authorization-server"),
117 "git_url": "https://g1t.sh/{owner}/{name}.git",
Integrations: your own model provider, alerts that open issues, tickets agents read118 "integrations_url": format!("{API}/workspaces/{{workspace}}/integrations"),
119 "context_url": format!("{repo}/context{{?reference}}"),
120 "import_issue_url": format!("{repo}/issues/import"),
121 "hooks_url": format!("{API}/hooks/{{integration}}"),
API and MCP server in Rust; a public index at the API root122 })
123}
124
125// Signing in from a tool. Accounts are created, and passwords typed, only
126// in a browser; a tool gets its token by having a person approve a code.
127
Integrations: your own model provider, alerts that open issues, tickets agents read128/// Passes a request from an outside system to its connection, as it came:
129/// its signature covers the exact bytes of the body.
130async fn receive_hook(request: &mut Request, services: &Services, id: &str) -> Result<Response> {
131 let headers: std::collections::HashMap<String, String> = request
132 .headers()
133 .entries()
134 .map(|(name, value)| (name.to_lowercase(), value))
135 .collect();
136 let body = request.text().await.unwrap_or_default();
137 if body.len() > 1_000_000 {
138 return Ok(Response::from_json(&json!({ "message": "The body is too large." }))?.with_status(413));
139 }
140 let received: g1t_contracts::integrations::Received = g1t_kit::call(
141 &services.integrations,
142 "receive",
143 &json!({ "id": id, "headers": headers, "body": body }),
144 )
145 .await?;
146 Ok(Response::from_json(&json!({ "message": received.message }))?.with_status(received.status))
147}
148
API and MCP server in Rust; a public index at the API root149async fn device_code(request: &mut Request, services: &Services) -> Result<Response> {
150 let body = json_body(request).await;
151 let started: DeviceStart = g1t_kit::call(
152 &services.identity,
153 "device_start",
154 &DeviceStartArgs {
155 client_name: body["client_name"].as_str().unwrap_or_default().to_owned(),
156 },
157 )
158 .await?;
159 Response::from_json(&json!({
160 "device_code": started.device_code,
161 "user_code": started.user_code,
162 "verification_uri": "https://g1t.sh/device",
163 "verification_uri_complete": format!("https://g1t.sh/device?code={}", started.user_code),
164 "expires_in": started.expires_in,
165 "interval": started.interval,
166 }))
167}
168
169async fn device_token(request: &mut Request, services: &Services) -> Result<Response> {
170 let body = json_body(request).await;
171 let claim: DeviceClaim = g1t_kit::call(
172 &services.identity,
173 "device_claim",
174 &DeviceClaimArgs {
175 device_code: body["device_code"].as_str().unwrap_or_default().to_owned(),
176 },
177 )
178 .await?;
179 Response::from_json(&match claim {
180 DeviceClaim::Approved { token, user } => json!({
181 "status": "approved",
182 "token": token,
183 "username": user.username,
184 "verified": user.verified,
185 }),
186 DeviceClaim::Pending => json!({ "status": "pending" }),
187 DeviceClaim::Denied => json!({ "status": "denied" }),
188 DeviceClaim::Expired => json!({ "status": "expired" }),
189 })
190}
191
Acceptance checks in sandboxes, line comments and review verdicts192/// A sandbox reporting on its run of a pull request's acceptance checks.
193/// The run's own token, in the body, is the credential: it was given to
194/// that sandbox and to nothing else.
195async fn report_checks(
196 request: &mut Request,
197 services: &Services,
198 run_id: &str,
199) -> Result<Response> {
200 let body = json_body(request).await;
201 let reported: Outcome<CheckRun> = g1t_kit::call(
202 &services.work,
203 "report_checks",
204 &ReportChecksArgs {
205 run_id: run_id.to_owned(),
206 token: body["token"].as_str().unwrap_or_default().to_owned(),
207 results: serde_json::from_value(body["results"].clone()).unwrap_or_default(),
208 error: body["error"].as_str().map(str::to_owned),
209 skip: false,
210 },
211 )
212 .await?;
213 match reported {
214 Outcome::Ok(run) => Response::from_json(&json!({ "status": run.status })),
215 Outcome::Fail(refused) => failure(&refused),
216 }
217}
218
Agents as a team: lifecycle, merge queue, billing and a new shell219/// A sandbox reporting one tested state of a merge queue. As with checks,
220/// the entry's own token is the credential.
221async fn report_queue(
222 request: &mut Request,
223 services: &Services,
224 entry_id: &str,
225) -> Result<Response> {
226 let body = json_body(request).await;
227 let reported: Outcome<QueueState> = g1t_kit::call(
228 &services.work,
229 "report_queue",
230 &ReportQueueArgs {
231 entry_id: entry_id.to_owned(),
232 token: body["token"].as_str().unwrap_or_default().to_owned(),
233 combined_commit: body["combinedCommit"].as_str().map(str::to_owned),
234 results: serde_json::from_value(body["results"].clone()).unwrap_or_default(),
235 error: body["error"].as_str().map(str::to_owned),
236 conflict_with: body["conflictWith"].as_u64().map(|n| n as u32),
237 },
238 )
239 .await?;
240 match reported {
241 Outcome::Ok(state) => Response::from_json(&json!({ "state": state })),
242 Outcome::Fail(refused) => failure(&refused),
243 }
244}
245
246/// A sandbox reporting the review its agent wrote. As with checks, the
247/// run's own token is the credential.
248async fn report_review(
249 request: &mut Request,
250 services: &Services,
251 run_id: &str,
252) -> Result<Response> {
253 let body = json_body(request).await;
254 let reported: Outcome<bool> = g1t_kit::call(
255 &services.work,
256 "report_review",
257 &ReportReviewArgs {
258 run_id: run_id.to_owned(),
259 token: body["token"].as_str().unwrap_or_default().to_owned(),
260 verdict: serde_json::from_value(body["verdict"].clone()).unwrap_or(None),
261 body: body["body"].as_str().unwrap_or_default().to_owned(),
262 comments: serde_json::from_value(body["comments"].clone()).unwrap_or_default(),
263 model: body["model"].as_str().map(str::to_owned),
264 error: body["error"].as_str().map(str::to_owned),
265 },
266 )
267 .await?;
268 match reported {
269 Outcome::Ok(_) => Response::from_json(&json!({ "recorded": true })),
270 Outcome::Fail(refused) => failure(&refused),
271 }
272}
273
274/// A sandbox reporting the plan its agent wrote. As with checks, the
275/// plan's own token is the credential.
276async fn report_plan(
277 request: &mut Request,
278 services: &Services,
279 plan_id: &str,
280) -> Result<Response> {
281 let body = json_body(request).await;
282 let reported: Outcome<bool> = g1t_kit::call(
283 &services.work,
284 "report_plan",
285 &ReportPlanArgs {
286 plan_id: plan_id.to_owned(),
287 token: body["token"].as_str().unwrap_or_default().to_owned(),
288 summary: body["summary"].as_str().unwrap_or_default().to_owned(),
289 issues: serde_json::from_value(body["issues"].clone()).unwrap_or_default(),
290 error: body["error"].as_str().map(str::to_owned),
291 },
292 )
293 .await?;
294 match reported {
295 Outcome::Ok(_) => Response::from_json(&json!({ "recorded": true })),
296 Outcome::Fail(refused) => failure(&refused),
297 }
298}
299
300/// A sandbox reporting what its agent's run cost, so that the workspace
301/// it worked for is charged. As with checks, the run's own token is the
302/// credential.
303async fn report_usage(
304 request: &mut Request,
305 services: &Services,
306 run_id: &str,
307) -> Result<Response> {
308 let body = json_body(request).await;
309 let charged: Outcome<bool> = g1t_kit::call(
310 &services.billing,
311 "finish_run",
312 &FinishRunArgs {
313 run_id: run_id.to_owned(),
314 token: body["token"].as_str().unwrap_or_default().to_owned(),
315 cost_usd: body["cost_usd"].as_f64().unwrap_or_default(),
316 turns: body["turns"].as_u64().unwrap_or_default() as u32,
317 },
318 )
319 .await?;
320 match charged {
321 Outcome::Ok(_) => Response::from_json(&json!({ "recorded": true })),
322 Outcome::Fail(refused) => failure(&refused),
323 }
324}
325
API and MCP server in Rust; a public index at the API root326async fn respond(mut request: Request, env: &Env) -> Result<Response> {
327 let method = method_name(request.method());
328 if method == "OPTIONS" {
329 return Ok(Response::empty()?.with_status(204));
330 }
331 let url = request.url()?;
Agents as a team: lifecycle, merge queue, billing and a new shell332 // Paths carry no version. An earlier form began with `/v1`, which is
333 // still accepted so that nothing already written against it breaks.
334 let path = match url.path().strip_prefix("/v1") {
335 Some(rest) if rest.is_empty() || rest.starts_with('/') => rest.to_owned(),
336 _ => url.path().to_owned(),
337 };
API and MCP server in Rust; a public index at the API root338 let on_mcp = url.host_str().is_some_and(|host| host.starts_with("mcp."));
Agents as a team: lifecycle, merge queue, billing and a new shell339 let mut services = Services::new(env)?;
API and MCP server in Rust; a public index at the API root340
Integrations: your own model provider, alerts that open issues, tickets agents read341 // Outside systems reporting to a connection. They sign what they send
342 // with the connection's own secret, which is not a g1t token, so this
343 // comes before anything that would read one.
344 if method == "POST" && !on_mcp {
345 if let Some(id) = path.strip_prefix("/hooks/").filter(|id| !id.is_empty() && !id.contains('/')) {
346 return receive_hook(&mut request, &services, id).await;
347 }
348 }
349
API and MCP server in Rust; a public index at the API root350 let viewer = match authenticate(&request, &services).await? {
351 Ok(viewer) => viewer,
352 Err(refused) => return Ok(refused),
353 };
Agents as a team: lifecycle, merge queue, billing and a new shell354 // An agent's token: what it may do comes with it.
355 if viewer.as_ref().is_some_and(|viewer| viewer.kind == PrincipalKind::Agent) {
356 let header = request.headers().get("authorization")?.unwrap_or_default();
357 let token = header.split_once(' ').map(|(_, token)| token.trim()).unwrap_or_default();
358 let scope: Option<AgentScope> = g1t_kit::call(
359 &services.identity,
360 "agent_scope",
361 &TokenArgs {
362 token: token.to_owned(),
363 },
364 )
365 .await?;
366 // A scope is what lets an agent's token do anything at all.
367 let Some(scope) = scope else {
368 return fail(FailureCode::Unauthenticated, "Invalid access token.");
369 };
370 services.scope = Some(scope);
371 }
API and MCP server in Rust; a public index at the API root372 if let Some(response) = oauth::handle(&mut request, &services, method, &path).await? {
373 return Ok(response);
374 }
375 if on_mcp {
376 return mcp::handle(request, &services, &viewer).await;
377 }
378
379 match (method, path.trim_end_matches('/')) {
Agents as a team: lifecycle, merge queue, billing and a new shell380 ("GET", "") => return Response::from_json(&index()),
API and MCP server in Rust; a public index at the API root381 ("GET", "/openapi.json") => return Response::from_json(&openapi::document()),
Agents as a team: lifecycle, merge queue, billing and a new shell382 ("POST", "/device/code") => return device_code(&mut request, &services).await,
383 ("POST", "/device/token") => return device_token(&mut request, &services).await,
Record your own agent's sessions automatically384 // Where a pull request lives, for a tool that knows only its fork.
385 ("GET", path) if path.starts_with("/pulls/") && !path[7..].contains('/') => {
386 let located: Outcome<Value> = g1t_kit::call(
387 &services.work,
388 "locate_pull",
389 &json!({ "id": &path[7..], "viewer": viewer }),
390 )
391 .await?;
392 return match located {
393 Outcome::Ok(value) => Response::from_json(&value),
394 Outcome::Fail(refused) => failure(&refused),
395 };
396 }
Agents as a team: lifecycle, merge queue, billing and a new shell397 ("POST", path) if path.starts_with("/queue/") => {
398 let entry_id = path.trim_start_matches("/queue/").to_owned();
399 return report_queue(&mut request, &services, &entry_id).await;
400 }
401 ("POST", path) if path.starts_with("/checks/") => {
402 let run_id = path.trim_start_matches("/checks/").to_owned();
Acceptance checks in sandboxes, line comments and review verdicts403 return report_checks(&mut request, &services, &run_id).await;
404 }
Agents as a team: lifecycle, merge queue, billing and a new shell405 ("POST", path) if path.starts_with("/runs/") && path.ends_with("/usage") => {
406 let run_id = path
407 .trim_start_matches("/runs/")
408 .trim_end_matches("/usage")
409 .to_owned();
410 return report_usage(&mut request, &services, &run_id).await;
411 }
412 ("POST", path) if path.starts_with("/plans/") => {
413 let plan_id = path.trim_start_matches("/plans/").to_owned();
414 return report_plan(&mut request, &services, &plan_id).await;
415 }
416 ("POST", path) if path.starts_with("/reviews/") => {
417 let run_id = path.trim_start_matches("/reviews/").to_owned();
418 return report_review(&mut request, &services, &run_id).await;
419 }
API and MCP server in Rust; a public index at the API root420 _ => {}
421 }
422
423 let query: Vec<(String, String)> = url
424 .query_pairs()
425 .map(|(name, value)| (name.into_owned(), value.into_owned()))
426 .collect();
427 let body = if method == "GET" {
428 Value::Null
429 } else {
Agents as a team: lifecycle, merge queue, billing and a new shell430 snake_case_keys(json_body(&mut request).await)
API and MCP server in Rust; a public index at the API root431 };
432 let Some((route, input)) = rest::resolve(method, &path, &query, body) else {
433 return fail(FailureCode::NotFound, "No such endpoint.");
434 };
435 match route.op.run(&services, &viewer, &input).await? {
436 Outcome::Ok(value) => Response::from_json(&value),
437 Outcome::Fail(refused) => failure(&refused),
438 }
439}
440
Agents as a team: lifecycle, merge queue, billing and a new shell441/// Request bodies take the same keys as the MCP tools, `snake_case`; the
442/// `camelCase` that responses use is accepted too, so a client can send
443/// back what it read.
444fn snake_case_keys(body: Value) -> Value {
445 let Value::Object(fields) = body else {
446 return body;
447 };
448 let mut out = serde_json::Map::new();
449 for (key, value) in fields {
450 let mut snake = String::with_capacity(key.len() + 4);
451 for c in key.chars() {
452 if c.is_ascii_uppercase() {
453 snake.push('_');
454 snake.push(c.to_ascii_lowercase());
455 } else {
456 snake.push(c);
457 }
458 }
459 // A key given in both spellings keeps the snake_case one.
460 if snake != key && out.contains_key(&snake) {
461 continue;
462 }
463 out.insert(snake, value);
464 }
465 Value::Object(out)
466}
467
468#[cfg(test)]
469mod tests {
470 use super::snake_case_keys;
471 use serde_json::json;
472
473 #[test]
474 fn camel_case_keys_are_accepted() {
475 assert_eq!(
476 snake_case_keys(json!({ "countAgentApprovals": false, "title": "x" })),
477 json!({ "count_agent_approvals": false, "title": "x" })
478 );
479 }
480
481 #[test]
482 fn snake_case_wins_when_both_are_given() {
483 assert_eq!(
484 snake_case_keys(json!({ "keep_issue_open": true, "keepIssueOpen": false })),
485 json!({ "keep_issue_open": true })
486 );
487 }
488}
489
API and MCP server in Rust; a public index at the API root490// The API is called from browsers too: the reference's explorer, and apps
491// built on g1t. It carries no cookies, so any origin may call it.
492#[event(fetch)]
493async fn fetch(request: Request, env: Env, _ctx: Context) -> Result<Response> {
494 let mut response = respond(request, &env).await?;
495 let headers = response.headers_mut();
496 headers.set("access-control-allow-origin", "*")?;
497 headers.set(
498 "access-control-allow-headers",
499 "authorization, content-type",
500 )?;
501 headers.set("access-control-allow-methods", "GET, POST, PATCH, OPTIONS")?;
502 headers.set("access-control-expose-headers", "www-authenticate")?;
503 Ok(response)
504}