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

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