pr_01m47d15m3e54sn21z27rpy5n9/apps/api/src/lib.rs

615 lines24,480 bytesCodeBlame
1//! 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 blobs;
8mod mcp;
9mod oauth;
10mod openapi;
11mod operations;
12mod rest;
13
14use g1t_contracts::billing::FinishRunArgs;
15use g1t_contracts::identity::{
16 DeviceClaim, DeviceClaimArgs, DeviceStart, DeviceStartArgs, TokenArgs,
17};
18use g1t_contracts::work::{
19 CheckRun, QueueState, ReportChecksArgs, ReportPlanArgs, ReportQueueArgs, ReportReviewArgs,
20};
21use g1t_contracts::identity::AgentScope;
22use g1t_contracts::{Failure, FailureCode, Outcome, PrincipalKind, Viewer};
23use serde_json::{Value, json};
24use worker::{Context, Env, Method, Request, Response, Result, event};
25
26use operations::Services;
27
28const API: &str = "https://api.g1t.sh";
29
30fn method_name(method: Method) -> &'static str {
31 match method {
32 Method::Get => "GET",
33 Method::Post => "POST",
34 Method::Patch => "PATCH",
35 Method::Put => "PUT",
36 Method::Delete => "DELETE",
37 Method::Options => "OPTIONS",
38 Method::Head => "HEAD",
39 _ => "OTHER",
40 }
41}
42
43/// An error in the shape every endpoint uses.
44fn failure(failure: &Failure) -> Result<Response> {
45 Ok(Response::from_json(&json!({ "error": failure }))?.with_status(failure.code.http_status()))
46}
47
48fn fail(code: FailureCode, message: &str) -> Result<Response> {
49 failure(&Failure {
50 code,
51 message: message.to_owned(),
52 })
53}
54
55/// A request body as JSON. An empty or malformed body is no input.
56async fn json_body(request: &mut Request) -> Value {
57 request.json().await.unwrap_or(Value::Null)
58}
59
60/// Who a request's `Authorization: Bearer g1t_…` names. A missing token is
61/// an anonymous viewer; a wrong one is refused, so that a typo does not
62/// silently look signed out.
63async fn authenticate(
64 request: &Request,
65 services: &Services,
66) -> Result<std::result::Result<Viewer, Response>> {
67 let header = request.headers().get("authorization")?.unwrap_or_default();
68 let token = match header.split_once(' ') {
69 Some((scheme, token)) if scheme.eq_ignore_ascii_case("bearer") && !token.is_empty() => {
70 token.trim()
71 }
72 _ => return Ok(Ok(None)),
73 };
74 let viewer: Viewer = g1t_kit::call(
75 &services.identity,
76 "user_for_access_token",
77 &TokenArgs {
78 token: token.to_owned(),
79 },
80 )
81 .await?;
82 if viewer.is_some() {
83 return Ok(Ok(viewer));
84 }
85 let mut response = fail(FailureCode::Unauthenticated, "Invalid access token.")?;
86 // Tells an MCP client where to sign in again.
87 response.headers_mut().set(
88 "www-authenticate",
89 &format!("{}, error=\"invalid_token\"", oauth::MCP_CHALLENGE),
90 )?;
91 Ok(Err(response))
92}
93
94/// Where everything is, for someone or something exploring the API.
95fn index() -> Value {
96 let repo = format!("{API}/repos/{{owner}}/{{name}}");
97 json!({
98 "documentation_url": "https://docs.g1t.sh/reference/api/",
99 "openapi_url": format!("{API}/openapi.json"),
100 "mcp_url": "https://mcp.g1t.sh",
101 "current_user_url": format!("{API}/user"),
102 "workspaces_url": format!("{API}/workspaces"),
103 "repositories_url": format!("{API}/repos{{?q}}"),
104 "repository_url": repo,
105 "repository_events_url": format!("{repo}/events{{?before}}"),
106 "labels_url": format!("{repo}/labels"),
107 "issues_url": format!("{repo}/issues{{?state,label}}"),
108 "issue_url": format!("{repo}/issues/{{number}}"),
109 "issue_comments_url": format!("{repo}/issues/{{number}}/comments"),
110 "pulls_url": format!("{repo}/pulls{{?state}}"),
111 "pull_url": format!("{repo}/pulls/{{number}}"),
112 "pull_changes_url": format!("{repo}/pulls/{{number}}/changes"),
113 "pull_reviews_url": format!("{repo}/pulls/{{number}}/reviews"),
114 "pull_session_url": format!("{repo}/pulls/{{number}}/session{{?after}}"),
115 "device_code_url": format!("{API}/device/code"),
116 "device_token_url": format!("{API}/device/token"),
117 "oauth_metadata_url": format!("{API}/.well-known/oauth-authorization-server"),
118 "git_url": "https://g1t.sh/{owner}/{name}.git",
119 "integrations_url": format!("{API}/workspaces/{{workspace}}/integrations"),
120 "context_url": format!("{repo}/context{{?reference}}"),
121 "import_issue_url": format!("{repo}/issues/import"),
122 "hooks_url": format!("{API}/hooks/{{integration}}"),
123 })
124}
125
126// Signing in from a tool. Accounts are created, and passwords typed, only
127// in a browser; a tool gets its token by having a person approve a code.
128
129/// Passes a request from an outside system to its connection, as it came:
130/// its signature covers the exact bytes of the body.
131async fn receive_hook(request: &mut Request, services: &Services, id: &str) -> Result<Response> {
132 let headers: std::collections::HashMap<String, String> = request
133 .headers()
134 .entries()
135 .map(|(name, value)| (name.to_lowercase(), value))
136 .collect();
137 let body = request.text().await.unwrap_or_default();
138 if body.len() > 1_000_000 {
139 return Ok(Response::from_json(&json!({ "message": "The body is too large." }))?.with_status(413));
140 }
141 let received: g1t_contracts::integrations::Received = g1t_kit::call(
142 &services.integrations,
143 "receive",
144 &json!({ "id": id, "headers": headers, "body": body }),
145 )
146 .await?;
147 Ok(Response::from_json(&json!({ "message": received.message }))?.with_status(received.status))
148}
149
150async fn receive_stripe(request: &mut Request, env: &Env) -> Result<Response> {
151 let signature = request.headers().get("stripe-signature")?.unwrap_or_default();
152 let payload = request.text().await.unwrap_or_default();
153 if payload.len() > 1_000_000 || signature.is_empty() {
154 return Ok(Response::from_json(&json!({ "message": "Not a Stripe event." }))?.with_status(400));
155 }
156 let handled: g1t_contracts::Outcome<bool> = g1t_kit::call(
157 &env.service("BILLING")?,
158 "stripe_webhook",
159 &g1t_contracts::billing::StripeWebhookArgs { payload, signature },
160 )
161 .await?;
162 // A refusal is a 400, so Stripe shows it as failed; anything handled,
163 // or already handled, is a 200, so Stripe stops sending it.
164 Ok(match handled {
165 g1t_contracts::Outcome::Ok(_) => Response::from_json(&json!({ "received": true }))?,
166 g1t_contracts::Outcome::Fail(failure) => {
167 Response::from_json(&json!({ "message": failure.message }))?.with_status(400)
168 }
169 })
170}
171
172async fn device_code(request: &mut Request, services: &Services) -> Result<Response> {
173 let body = json_body(request).await;
174 let started: DeviceStart = g1t_kit::call(
175 &services.identity,
176 "device_start",
177 &DeviceStartArgs {
178 client_name: body["client_name"].as_str().unwrap_or_default().to_owned(),
179 },
180 )
181 .await?;
182 Response::from_json(&json!({
183 "device_code": started.device_code,
184 "user_code": started.user_code,
185 "verification_uri": "https://g1t.sh/device",
186 "verification_uri_complete": format!("https://g1t.sh/device?code={}", started.user_code),
187 "expires_in": started.expires_in,
188 "interval": started.interval,
189 }))
190}
191
192async fn device_token(request: &mut Request, services: &Services) -> Result<Response> {
193 let body = json_body(request).await;
194 let claim: DeviceClaim = g1t_kit::call(
195 &services.identity,
196 "device_claim",
197 &DeviceClaimArgs {
198 device_code: body["device_code"].as_str().unwrap_or_default().to_owned(),
199 },
200 )
201 .await?;
202 Response::from_json(&match claim {
203 DeviceClaim::Approved { token, user } => json!({
204 "status": "approved",
205 "token": token,
206 "username": user.username,
207 "verified": user.verified,
208 }),
209 DeviceClaim::Pending => json!({ "status": "pending" }),
210 DeviceClaim::Denied => json!({ "status": "denied" }),
211 DeviceClaim::Expired => json!({ "status": "expired" }),
212 })
213}
214
215/// A sandbox reporting on its run of a pull request's acceptance checks.
216/// The run's own token, in the body, is the credential: it was given to
217/// that sandbox and to nothing else.
218async fn report_checks(
219 request: &mut Request,
220 services: &Services,
221 run_id: &str,
222) -> Result<Response> {
223 let body = json_body(request).await;
224 let reported: Outcome<CheckRun> = g1t_kit::call(
225 &services.work,
226 "report_checks",
227 &ReportChecksArgs {
228 run_id: run_id.to_owned(),
229 token: body["token"].as_str().unwrap_or_default().to_owned(),
230 results: serde_json::from_value(body["results"].clone()).unwrap_or_default(),
231 error: body["error"].as_str().map(str::to_owned),
232 skip: false,
233 },
234 )
235 .await?;
236 match reported {
237 Outcome::Ok(run) => Response::from_json(&json!({ "status": run.status })),
238 Outcome::Fail(refused) => failure(&refused),
239 }
240}
241
242/// A sandbox reporting one tested state of a merge queue. As with checks,
243/// the entry's own token is the credential.
244async fn report_queue(
245 request: &mut Request,
246 services: &Services,
247 entry_id: &str,
248) -> Result<Response> {
249 let body = json_body(request).await;
250 let reported: Outcome<QueueState> = g1t_kit::call(
251 &services.work,
252 "report_queue",
253 &ReportQueueArgs {
254 entry_id: entry_id.to_owned(),
255 token: body["token"].as_str().unwrap_or_default().to_owned(),
256 combined_commit: body["combinedCommit"].as_str().map(str::to_owned),
257 results: serde_json::from_value(body["results"].clone()).unwrap_or_default(),
258 error: body["error"].as_str().map(str::to_owned),
259 conflict_with: body["conflictWith"].as_u64().map(|n| n as u32),
260 },
261 )
262 .await?;
263 match reported {
264 Outcome::Ok(state) => Response::from_json(&json!({ "state": state })),
265 Outcome::Fail(refused) => failure(&refused),
266 }
267}
268
269/// A sandbox reporting the review its agent wrote. As with checks, the
270/// run's own token is the credential.
271async fn report_review(
272 request: &mut Request,
273 services: &Services,
274 run_id: &str,
275) -> Result<Response> {
276 let body = json_body(request).await;
277 let reported: Outcome<bool> = g1t_kit::call(
278 &services.work,
279 "report_review",
280 &ReportReviewArgs {
281 run_id: run_id.to_owned(),
282 token: body["token"].as_str().unwrap_or_default().to_owned(),
283 verdict: serde_json::from_value(body["verdict"].clone()).unwrap_or(None),
284 body: body["body"].as_str().unwrap_or_default().to_owned(),
285 comments: serde_json::from_value(body["comments"].clone()).unwrap_or_default(),
286 model: body["model"].as_str().map(str::to_owned),
287 error: body["error"].as_str().map(str::to_owned),
288 },
289 )
290 .await?;
291 match reported {
292 Outcome::Ok(_) => Response::from_json(&json!({ "recorded": true })),
293 Outcome::Fail(refused) => failure(&refused),
294 }
295}
296
297/// A sandbox reporting the plan its agent wrote. As with checks, the
298/// plan's own token is the credential.
299async fn report_plan(
300 request: &mut Request,
301 services: &Services,
302 plan_id: &str,
303) -> Result<Response> {
304 let body = json_body(request).await;
305 let reported: Outcome<bool> = g1t_kit::call(
306 &services.work,
307 "report_plan",
308 &ReportPlanArgs {
309 plan_id: plan_id.to_owned(),
310 token: body["token"].as_str().unwrap_or_default().to_owned(),
311 summary: body["summary"].as_str().unwrap_or_default().to_owned(),
312 issues: serde_json::from_value(body["issues"].clone()).unwrap_or_default(),
313 error: body["error"].as_str().map(str::to_owned),
314 },
315 )
316 .await?;
317 match reported {
318 Outcome::Ok(_) => Response::from_json(&json!({ "recorded": true })),
319 Outcome::Fail(refused) => failure(&refused),
320 }
321}
322
323/// A sandbox reporting what its agent's run cost, so that the workspace
324/// it worked for is charged. As with checks, the run's own token is the
325/// credential.
326async fn report_usage(
327 request: &mut Request,
328 services: &Services,
329 run_id: &str,
330) -> Result<Response> {
331 let body = json_body(request).await;
332 let charged: Outcome<bool> = g1t_kit::call(
333 &services.billing,
334 "finish_run",
335 &FinishRunArgs {
336 run_id: run_id.to_owned(),
337 token: body["token"].as_str().unwrap_or_default().to_owned(),
338 cost_usd: body["cost_usd"].as_f64().unwrap_or_default(),
339 turns: body["turns"].as_u64().unwrap_or_default() as u32,
340 },
341 )
342 .await?;
343 match charged {
344 Outcome::Ok(_) => Response::from_json(&json!({ "recorded": true })),
345 Outcome::Fail(refused) => failure(&refused),
346 }
347}
348
349async fn respond(mut request: Request, env: &Env) -> Result<Response> {
350 let method = method_name(request.method());
351 if method == "OPTIONS" {
352 return Ok(Response::empty()?.with_status(204));
353 }
354 let url = request.url()?;
355 // Paths carry no version. An earlier form began with `/v1`, which is
356 // still accepted so that nothing already written against it breaks.
357 let path = match url.path().strip_prefix("/v1") {
358 Some(rest) if rest.is_empty() || rest.starts_with('/') => rest.to_owned(),
359 _ => url.path().to_owned(),
360 };
361 let on_mcp = url.host_str().is_some_and(|host| host.starts_with("mcp."));
362 let mut services = Services::new(env)?;
363
364 // Stripe reporting to billing. Signed with the secret of the endpoint
365 // billing registered; the body goes through exactly as received, since
366 // the signature covers its bytes.
367 if method == "POST" && !on_mcp && path == "/stripe/webhook" {
368 return receive_stripe(&mut request, env).await;
369 }
370
371 // Outside systems reporting to a connection. They sign what they send
372 // with the connection's own secret, which is not a g1t token, so this
373 // comes before anything that would read one.
374 if method == "POST" && !on_mcp
375 && let Some(id) = path.strip_prefix("/hooks/").filter(|id| !id.is_empty() && !id.contains('/')) {
376 return receive_hook(&mut request, &services, id).await;
377 }
378
379 // A sandbox building a deployment, reporting with its build's token,
380 // which is not a g1t token. The body goes through as it is: it can
381 // carry a Worker's bundled code.
382 if method == "POST" && !on_mcp
383 && let Some(rest) = path.strip_prefix("/deployments/jobs/")
384 {
385 let target = format!("https://deployments/jobs/{rest}");
386 let body = request.bytes().await?;
387 let headers = worker::Headers::new();
388 headers.set("content-type", "application/json")?;
389 let mut init = worker::RequestInit::new();
390 init.with_method(Method::Post)
391 .with_headers(headers)
392 .with_body(Some(worker::js_sys::Uint8Array::from(body.as_slice()).into()));
393 let mut answer = env
394 .service("DEPLOYMENTS")?
395 .fetch_request(Request::new_with_init(&target, &init)?)
396 .await?;
397 // A fresh response: a fetched one's headers cannot be changed, and
398 // every response gets the API's own on the way out.
399 let status = answer.status_code();
400 return Ok(Response::from_bytes(answer.bytes().await?)?
401 .with_status(status)
402 .with_headers({
403 let headers = worker::Headers::new();
404 headers.set("content-type", "application/json")?;
405 headers
406 }));
407 }
408
409 // A sandbox's artifacts and cache, with its job's token, which is not a
410 // g1t token either.
411 if !on_mcp
412 && let Some(rest) = path.strip_prefix("/actions/jobs/")
413 && (rest.contains("/artifacts") || rest.ends_with("/cache"))
414 {
415 let rest = rest.to_owned();
416 return blobs::for_job(request, env, &services, method, &rest).await;
417 }
418
419 let viewer = match authenticate(&request, &services).await? {
420 Ok(viewer) => viewer,
421 Err(refused) => return Ok(refused),
422 };
423 // An agent's token: what it may do comes with it.
424 if viewer.as_ref().is_some_and(|viewer| viewer.kind == PrincipalKind::Agent) {
425 let header = request.headers().get("authorization")?.unwrap_or_default();
426 let token = header.split_once(' ').map(|(_, token)| token.trim()).unwrap_or_default();
427 let scope: Option<AgentScope> = g1t_kit::call(
428 &services.identity,
429 "agent_scope",
430 &TokenArgs {
431 token: token.to_owned(),
432 },
433 )
434 .await?;
435 // A scope is what lets an agent's token do anything at all.
436 let Some(scope) = scope else {
437 return fail(FailureCode::Unauthenticated, "Invalid access token.");
438 };
439 services.scope = Some(scope);
440 }
441 if let Some(response) = oauth::handle(&mut request, &services, method, &path).await? {
442 return Ok(response);
443 }
444 if on_mcp {
445 return mcp::handle(request, &services, &viewer).await;
446 }
447
448 match (method, path.trim_end_matches('/')) {
449 ("GET", "") => return Response::from_json(&index()),
450 ("GET", "/openapi.json") => return Response::from_json(&openapi::document()),
451 // A run's artifacts: listed, or one downloaded.
452 ("GET", path) if path.starts_with("/repos/") && path.contains("/actions/runs/") && path.contains("/artifacts") => {
453 let parts: Vec<&str> = path.trim_start_matches("/repos/").split('/').collect();
454 if let [owner, repo, "actions", "runs", run, "artifacts", rest @ ..] = parts.as_slice() {
455 return match rest {
456 [] => {
457 let seen: Outcome<Value> = g1t_kit::call(
458 &services.actions,
459 "run",
460 &json!({ "repo": { "namespace": owner, "name": repo }, "viewer": viewer, "id": run }),
461 )
462 .await?;
463 match seen {
464 Outcome::Ok(_) => Response::from_json(&blobs::of_run(env, run).await?),
465 Outcome::Fail(refused) => failure(&refused),
466 }
467 }
468 [name] => blobs::download(env, &services, &viewer, owner, repo, run, name).await,
469 _ => fail(FailureCode::NotFound, "No such endpoint."),
470 };
471 }
472 }
473 ("POST", "/device/code") => return device_code(&mut request, &services).await,
474 ("POST", "/device/token") => return device_token(&mut request, &services).await,
475 // Where a pull request lives, for a tool that knows only its fork.
476 ("GET", path) if path.starts_with("/pulls/") && !path[7..].contains('/') => {
477 let located: Outcome<Value> = g1t_kit::call(
478 &services.work,
479 "locate_pull",
480 &json!({ "id": &path[7..], "viewer": viewer }),
481 )
482 .await?;
483 return match located {
484 Outcome::Ok(value) => Response::from_json(&value),
485 Outcome::Fail(refused) => failure(&refused),
486 };
487 }
488 ("POST", path) if path.starts_with("/queue/") => {
489 let entry_id = path.trim_start_matches("/queue/").to_owned();
490 return report_queue(&mut request, &services, &entry_id).await;
491 }
492 // A sandbox running a GitHub Actions job: fetching the job, and
493 // reporting how it goes. The job's own token is the credential.
494 ("POST", path) if path.starts_with("/actions/jobs/") => {
495 let rest = path.trim_start_matches("/actions/jobs/");
496 let (job, method) = match rest.strip_suffix("/spec") {
497 Some(job) => (job.to_owned(), "job_spec"),
498 None => (rest.to_owned(), "job_report"),
499 };
500 let body = json_body(&mut request).await;
501 let answered: Outcome<Value> = g1t_kit::call(
502 &services.actions,
503 method,
504 &json!({ "job": job, "token": body["token"], "report": body["report"] }),
505 )
506 .await?;
507 return match answered {
508 Outcome::Ok(value) => Response::from_json(&value),
509 Outcome::Fail(refused) => failure(&refused),
510 };
511 }
512 ("POST", path) if path.starts_with("/checks/") => {
513 let run_id = path.trim_start_matches("/checks/").to_owned();
514 return report_checks(&mut request, &services, &run_id).await;
515 }
516 ("POST", path) if path.starts_with("/runs/") && path.ends_with("/usage") => {
517 let run_id = path
518 .trim_start_matches("/runs/")
519 .trim_end_matches("/usage")
520 .to_owned();
521 return report_usage(&mut request, &services, &run_id).await;
522 }
523 ("POST", path) if path.starts_with("/plans/") => {
524 let plan_id = path.trim_start_matches("/plans/").to_owned();
525 return report_plan(&mut request, &services, &plan_id).await;
526 }
527 ("POST", path) if path.starts_with("/reviews/") => {
528 let run_id = path.trim_start_matches("/reviews/").to_owned();
529 return report_review(&mut request, &services, &run_id).await;
530 }
531 _ => {}
532 }
533
534 let query: Vec<(String, String)> = url
535 .query_pairs()
536 .map(|(name, value)| (name.into_owned(), value.into_owned()))
537 .collect();
538 let body = if method == "GET" {
539 Value::Null
540 } else {
541 snake_case_keys(json_body(&mut request).await)
542 };
543 let Some((route, input)) = rest::resolve(method, &path, &query, body) else {
544 return fail(FailureCode::NotFound, "No such endpoint.");
545 };
546 match route.op.run(&services, &viewer, &input).await? {
547 Outcome::Ok(value) => Response::from_json(&value),
548 Outcome::Fail(refused) => failure(&refused),
549 }
550}
551
552/// Request bodies take the same keys as the MCP tools, `snake_case`; the
553/// `camelCase` that responses use is accepted too, so a client can send
554/// back what it read.
555fn snake_case_keys(body: Value) -> Value {
556 let Value::Object(fields) = body else {
557 return body;
558 };
559 let mut out = serde_json::Map::new();
560 for (key, value) in fields {
561 let mut snake = String::with_capacity(key.len() + 4);
562 for c in key.chars() {
563 if c.is_ascii_uppercase() {
564 snake.push('_');
565 snake.push(c.to_ascii_lowercase());
566 } else {
567 snake.push(c);
568 }
569 }
570 // A key given in both spellings keeps the snake_case one.
571 if snake != key && out.contains_key(&snake) {
572 continue;
573 }
574 out.insert(snake, value);
575 }
576 Value::Object(out)
577}
578
579#[cfg(test)]
580mod tests {
581 use super::snake_case_keys;
582 use serde_json::json;
583
584 #[test]
585 fn camel_case_keys_are_accepted() {
586 assert_eq!(
587 snake_case_keys(json!({ "countAgentApprovals": false, "title": "x" })),
588 json!({ "count_agent_approvals": false, "title": "x" })
589 );
590 }
591
592 #[test]
593 fn snake_case_wins_when_both_are_given() {
594 assert_eq!(
595 snake_case_keys(json!({ "keep_issue_open": true, "keepIssueOpen": false })),
596 json!({ "keep_issue_open": true })
597 );
598 }
599}
600
601// The API is called from browsers too: the reference's explorer, and apps
602// built on g1t. It carries no cookies, so any origin may call it.
603#[event(fetch)]
604async fn fetch(request: Request, env: Env, _ctx: Context) -> Result<Response> {
605 let mut response = respond(request, &env).await?;
606 let headers = response.headers_mut();
607 headers.set("access-control-allow-origin", "*")?;
608 headers.set(
609 "access-control-allow-headers",
610 "authorization, content-type",
611 )?;
612 headers.set("access-control-allow-methods", "GET, POST, PATCH, OPTIONS")?;
613 headers.set("access-control-expose-headers", "www-authenticate")?;
614 Ok(response)
615}