Commit

Rust repos service with shipping; pull requests kept in the model

- services/repos rewritten in Rust: registry, contents, forks, git over HTTPS, and landing an attempt on main by relaying a pack from the fork - landing refuses an attempt that is behind main, following every parent when checking ancestry - work: ship_attempt closes the intent; exposed in the API, MCP and site - kit: call Worker RPC stubs correctly from Rust - docs and marketing: shipping documented; intents framed as building on pull requests

syntaqxcommitted Parent91f9980Browse files
47 files+1913−7870/47 viewed
+11−0
851851 ]
852852
853853 [[package]]
854+name = "g1t-repos"
855+version = "0.1.0"
856+dependencies = [
857+ "g1t-contracts",
858+ "g1t-kit",
859+ "serde",
860+ "serde_json",
861+ "worker",
862+]
863+
864+[[package]]
854865 name = "g1t-sshd"
855866 version = "0.1.0"
856867 dependencies = [
+1−1
11 [workspace]
22 resolver = "3"
3−members = ["crates/*", "services/identity"]
3+members = ["crates/*", "services/identity", "services/repos"]
44
55 [workspace.package]
66 edition = "2024"
+9−4
2020 - Accounts, access tokens, public and private repositories.
2121 - Git over HTTPS, including creating a repository by pushing to it.
2222 - Intents, attempts (a copy-on-write fork each) and recorded agent sessions.
23+- Shipping: landing an attempt on `main`, refused when the attempt is behind
24+ so that no commit is ever lost.
25+- Registration with email verification, and password reset.
2326 - A REST API and an MCP server over the same operations.
2427 - An event bus: every state change is published, logged and delivered to
2528 subscribers.
2629
27−Not built yet: shipping an attempt through a landing queue, running
28−acceptance checks, hosted agents, git over SSH. See the build order in the
30+Not built yet: diffs and review on the site, pull requests from branches,
31+server-side merging, running acceptance checks, hosted agents, git over SSH. See the build order in the
2932 plan.
3033
3134 ## Try it
5053 | `apps/web` | The site: server-rendered React on a Worker. Holds no data. |
5154 | `apps/api` | REST API and MCP server. |
5255 | `services/identity` | Accounts, sessions, keys and tokens. Rust. |
53−| `services/repos` | Repository registry, contents, forks, git over HTTPS. |
56+| `services/repos` | Repository registry, contents, forks, landing, git over HTTPS. Rust. |
5457 | `services/work` | Intents, attempts and sessions. |
5558 | `services/events` | The event bus and its log. |
5659 | `crates/contracts` | Types and service interfaces for the Rust services. |
6063
6164 Each service is its own Worker with its own database. They call each other
6265 through service bindings and react to each other through events. Services
63−are being moved from TypeScript to Rust one at a time; identity is done.
66+are being moved from TypeScript to Rust one at a time; identity and repos are
67+done.
6468
6569 ## Run your own
6670
8690
8791 ```sh
8892 (cd services/identity && npx wrangler deploy)
93+(cd services/repos && npx wrangler deploy)
8994 npm run deploy
9095 ```
9196
+23−5
11 import { Hono } from "hono";
22
3−import { type Viewer, httpStatus, identityClient } from "@g1t/contracts";
3+import {
4+ type ServiceBinding,
5+ type Viewer,
6+ httpStatus,
7+ identityClient,
8+ reposClient,
9+} from "@g1t/contracts";
410
511 import { handleMcp } from "./mcp";
612 import { type ApiEnv, operations, operationsByName } from "./operations";
713
814 type Input = Record<string, unknown>;
9−type App = { Bindings: ApiEnv; Variables: { viewer: Viewer } };
15+/** The Worker's raw bindings; Rust services are reached through clients. */
16+type Bindings = Omit<ApiEnv, "IDENTITY" | "REPOS"> & {
17+ IDENTITY: ServiceBinding;
18+ REPOS: ServiceBinding;
19+};
20+type App = { Bindings: Bindings; Variables: { viewer: Viewer; services: ApiEnv } };
1021
1122 /**
1223 * REST routes. Each maps an HTTP request onto one operation; `input` builds
3243 { method: "POST", path: "/v1/attempts/:attempt_id/session", operation: "record_session", input: (p, _q, b) => ({ ...b, ...p }) },
3344 { method: "POST", path: "/v1/attempts/:attempt_id/submit", operation: "submit_attempt", input: (p, _q, b) => ({ ...b, ...p }) },
3445 { method: "POST", path: "/v1/attempts/:attempt_id/abandon", operation: "abandon_attempt", input: (p) => p },
46+ { method: "POST", path: "/v1/attempts/:attempt_id/ship", operation: "ship_attempt", input: (p) => p },
3547 ];
3648
3749 function repo(params: Record<string, string>): Input {
4456 // wrong one is rejected so a typo does not silently look signed out.
4557 app.use(async (c, next) => {
4658 const [scheme, token] = (c.req.header("authorization") ?? "").split(" ");
59+ const services: ApiEnv = {
60+ ...c.env,
61+ IDENTITY: identityClient(c.env.IDENTITY),
62+ REPOS: reposClient(c.env.REPOS),
63+ };
64+ c.set("services", services);
4765 let viewer: Viewer = null;
4866 if (scheme?.toLowerCase() === "bearer" && token) {
49− viewer = await identityClient(c.env.IDENTITY).userForAccessToken(token);
67+ viewer = await services.IDENTITY.userForAccessToken(token);
5068 if (!viewer) {
5169 return c.json(
5270 { error: { code: "unauthenticated", message: "Invalid access token." } },
6078
6179 app.all("*", async (c, next) => {
6280 if (new URL(c.req.url).hostname.startsWith("mcp.")) {
63− return handleMcp(c.req.raw, c.env, c.get("viewer"));
81+ return handleMcp(c.req.raw, c.get("services"), c.get("viewer"));
6482 }
6583 await next();
6684 });
86104 }
87105 }
88106 const input = route.input?.(c.req.param(), c.req.query(), body) ?? {};
89− const outcome = await operation.run(c.env, c.get("viewer"), input);
107+ const outcome = await operation.run(c.get("services"), c.get("viewer"), input);
90108 if (outcome.ok) return c.json(outcome.value);
91109 return c.json(
92110 { error: outcome.error },
+15−2
11 import {
22 type EventsApi,
3− type ServiceBinding,
3+ type IdentityApi,
44 type NewSessionEntry,
55 type RepoPath,
66 type ReposApi,
1313 } from "@g1t/contracts";
1414
1515 export interface ApiEnv {
16− IDENTITY: ServiceBinding;
16+ IDENTITY: IdentityApi;
1717 REPOS: ReposApi;
1818 WORK: WorkApi;
1919 EVENTS: EventsApi;
306306 ),
307307 },
308308 {
309+ name: "ship_attempt",
310+ description:
311+ "Land an attempt on the repository's main branch and close its intent. Only the repository's owner can ship. Fails if main has moved since the attempt started; the attempt must then pull main into its fork and push before shipping again.",
312+ input: {
313+ type: "object",
314+ properties: { attempt_id: { type: "string" } },
315+ required: ["attempt_id"],
316+ },
317+ run: authed((env, user, input) =>
318+ env.WORK.shipAttempt(user, text(input, "attempt_id")),
319+ ),
320+ },
321+ {
309322 name: "list_events",
310323 description:
311324 "The timeline of a repository: pushes, intents, attempts and session activity, newest first.",
+6−5
5454 }
5555
5656 const COMPARISON: [string, string, string][] = [
57− ["Unit of work", "A pull request: one author, one change", "An intent: one goal, any number of attempts"],
57+ ["Unit of work", "A pull request: one author, one change", "An intent: one goal, any number of attempts. One attempt is a pull request"],
5858 ["Where agents work", "Branches and local worktrees", "A server-side fork per attempt"],
5959 ["Why a change was made", "A commit message, if you are lucky", "The agent's full session, kept with the code"],
6060 ["Connecting an agent", "A vendor integration", "Any MCP client, or plain HTTP"],
110110 <section className="mx-auto max-w-6xl px-4 py-24">
111111 <p className="text-sm font-medium text-accent">How it works</p>
112112 <h2 className="mt-2 max-w-2xl text-3xl font-semibold tracking-tight text-balance sm:text-4xl">
113− Pull requests assume one author. Agents arrive by the hundred.
113+ Everything you know about git still works. It just stops assuming
114+ one author at a time.
114115 </h2>
115116 <div className="mt-12 grid gap-4 md:grid-cols-6">
116117 <FeatureCard title="Start with an intent" illustration={<IntentIllustration />}>
117− Write the goal and the checks that prove it is done. An intent
118− replaces both the issue and the pull request.
118+ Write the goal and the checks that prove it is done. Think of it
119+ as an issue that can hold any number of competing pull requests.
119120 </FeatureCard>
120121 <FeatureCard title="A fork for every attempt" illustration={<ForkIllustration />}>
121122 Each agent gets its own copy of the repository the moment it
129130 Claude Code and other MCP clients connect to mcp.g1t.sh with one
130131 command. Everything is also a plain REST call at api.g1t.sh.
131132 </FeatureCard>
132− <FeatureCard title="Converge on main" illustration={<ShipIllustration />} soon wide>
133+ <FeatureCard title="Converge on main" illustration={<ShipIllustration />} wide>
133134 However many attempts are in flight, changes reach main one at a
134135 time and in order. An attempt that has fallen behind is told, and
135136 catches up before it lands.
+4−0
2525 4. `record_session` as it goes, so people can see its reasoning.
2626 5. `submit_attempt` with a summary of what changed and why.
2727
28+If shipping reports that `main` has moved, pull `main` from the repository
29+into the fork, push, and the attempt can ship.
30+
2831 ## Tools
2932
3033 | Tool | What it does |
4245 | `read_session` | Read an attempt's recorded session. |
4346 | `submit_attempt` | Mark an attempt finished, with a summary. |
4447 | `abandon_attempt` | Give up on an attempt. |
48+| `ship_attempt` | Land an attempt on `main` and close its intent. Repository owner only. |
4549 | `list_events` | A repository's timeline, newest first. |
4650
4751 Repositories are always given as `owner/name`.
+1−0
6767 | `GET` | `/v1/attempts/{attempt_id}` | An attempt and its intent. |
6868 | `POST` | `/v1/attempts/{attempt_id}/submit` | Finish. Body: `summary`. |
6969 | `POST` | `/v1/attempts/{attempt_id}/abandon` | Give up. |
70+| `POST` | `/v1/attempts/{attempt_id}/ship` | Land it on `main`. Owner only; `409` if `main` has moved. |
7071
7172 Starting an attempt returns the fork's git remote:
7273
+30−7
11 # Concepts
22
3−A pull request assumes one author and one change. g1t assumes many agents
4−working at once, and is built from four ideas.
3+g1t is ordinary git: repositories, commits, branches, clone, push and pull all
4+work as they do anywhere. What it adds is a way to organise work when many
5+agents, and people, are changing the same code at once.
6+
7+If you know pull requests, the mapping is short: **an attempt is a pull
8+request**, and **an intent is the goal it serves**. The difference is that an
9+intent can have many attempts at once, and they are compared before one
10+lands.
511
612 ## Intent
713
8−An intent is a goal stated against a repository. It replaces both the issue
9−("what should happen") and the pull request ("here is a change"), because
10−with agents the two are the same conversation.
14+An intent is a goal stated against a repository. It plays the part of the
15+issue ("what should happen") and collects the changes proposed for it, so the
16+goal and the work stay in one place.
1117
1218 An intent has:
1319
4147 | `shipped` | The attempt was chosen and merged. |
4248 | `abandoned` | The attempt was given up. |
4349
50+## Shipping
51+
52+The owner of a repository ships an attempt to land it. Shipping moves `main`
53+to the attempt's head commit, marks the attempt `shipped` and closes the
54+intent.
55+
56+An attempt can only ship if it contains everything already on `main`. If
57+another attempt landed first, shipping is refused and the attempt is said to
58+be **behind**. Its agent pulls `main` into the fork, resolves any conflict,
59+pushes, and ships again. `main` never loses a commit this way, however many
60+attempts are racing.
61+
4462 ## Session
4563
4664 A session is the record of how an attempt was made: the prompt the agent was
6583 g1t is under active development. These parts of the model are designed but
6684 not available yet:
6785
68−- **Shipping.** Choosing an attempt and merging it into `main` through a
69− landing queue.
86+- **Merging in g1t.** Shipping moves `main` forward to the attempt's head.
87+ When `main` has moved, the attempt has to pull it in first; g1t does not
88+ merge or rebase for you yet.
89+- **Diffs and review.** Seeing an attempt's changes and commenting on them
90+ on the site.
91+- **Pull requests from branches.** Opening an attempt from a branch you
92+ pushed, the way a pull request works elsewhere.
7093 - **Checks.** Running an intent's acceptance checks automatically.
7194 - **Hosted agents.** Starting agents on g1t's own sandboxes. Today you bring
7295 your own agent.
+2−1
11 import { env } from "cloudflare:workers";
22
3−import { identityClient } from "@g1t/contracts";
3+import { identityClient, reposClient } from "@g1t/contracts";
44
55 export const identity = identityClient(env.IDENTITY);
6+export const repos = reposClient(env.REPOS);
+2−2
1−import { env } from "cloudflare:workers";
21 import { Search } from "lucide-react";
32 import { Form } from "react-router";
43
54 import type { Route } from "./+types/explore";
65 import { RepoList } from "../components/repo-list";
6+import { repos } from "../lib/services.server";
77 import { getViewer } from "../lib/session.server";
88
99 export function meta({ loaderData }: Route.MetaArgs) {
1515 /** Serves both /explore and /search?q=. */
1616 export async function loader({ request, context }: Route.LoaderArgs) {
1717 const query = new URL(request.url).searchParams.get("q")?.trim() ?? "";
18− return { query, repos: await env.REPOS.list(getViewer(context), { query }) };
18+ return { query, repos: await repos.list(getViewer(context), { query }) };
1919 }
2020
2121 export default function Explore({ loaderData }: Route.ComponentProps) {
+3−2
1313 Status,
1414 TimeAgo,
1515 } from "../components/ui";
16+import { repos as reposApi } from "../lib/services.server";
1617 import { getViewer } from "../lib/session.server";
1718
1819 export function meta({}: Route.MetaArgs) {
2930 export async function loader({ context }: Route.LoaderArgs) {
3031 const viewer = getViewer(context);
3132 const [repos, attempts] = await Promise.all([
32− env.REPOS.list(viewer, viewer ? { namespace: viewer.username } : {}),
33+ reposApi.list(viewer, viewer ? { namespace: viewer.username } : {}),
3334 env.WORK.listActiveAttempts(viewer),
3435 ]);
3536 // Mission control links to each attempt under its repo.
3637 const attemptRepos = await Promise.all(
37− attempts.map(({ attempt }) => env.REPOS.getById(attempt.repoId, viewer)),
38+ attempts.map(({ attempt }) => reposApi.getById(attempt.repoId, viewer)),
3839 );
3940 return {
4041 viewer,
+2−2
1−import { env } from "cloudflare:workers";
21 import { Form, redirect } from "react-router";
32
43 import type { Route } from "./+types/new";
54 import { Button, ErrorText, Field, Input } from "../components/ui";
5+import { repos } from "../lib/services.server";
66 import { assertSameOrigin, requireUser } from "../lib/session.server";
77
88 export function meta({}: Route.MetaArgs) {
1717 assertSameOrigin(request);
1818 const user = requireUser(context, request);
1919 const form = await request.formData();
20− const result = await env.REPOS.create(user, {
20+ const result = await repos.create(user, {
2121 name: String(form.get("name") ?? ""),
2222 description: String(form.get("description") ?? ""),
2323 isPrivate: form.get("visibility") === "private",
+2−3
1−import { env } from "cloudflare:workers";
21 import { data } from "react-router";
32
43 import type { Route } from "./+types/profile";
54 import { RepoList } from "../components/repo-list";
65 import { Avatar } from "../components/ui";
7−import { identity } from "../lib/services.server";
6+import { identity, repos as reposApi } from "../lib/services.server";
87 import { getViewer } from "../lib/session.server";
98
109 export function meta({ params }: Route.MetaArgs) {
1413 export async function loader({ params, context }: Route.LoaderArgs) {
1514 const [user, repos] = await Promise.all([
1615 identity.userByUsername(params.owner),
17− env.REPOS.list(getViewer(context), { namespace: params.owner }),
16+ reposApi.list(getViewer(context), { namespace: params.owner }),
1817 ]);
1918 if (!user) throw data(null, { status: 404 });
2019 return { user, repos };
+38−9
33 Bot,
44 ChevronRight,
55 GitCommitHorizontal,
6+ Rocket,
67 StickyNote,
78 User,
89 Wrench,
2223 Textarea,
2324 TimeAgo,
2425 } from "../../components/ui";
26+import { repos } from "../../lib/services.server";
2527 import {
2628 assertSameOrigin,
2729 getViewer,
4749 env.WORK.getAttempt(params.id, viewer),
4850 env.WORK.readSession(params.id, viewer),
4951 ]);
50− return { ...unwrap(found), session: unwrap(session), viewer };
52+ const detail = unwrap(found);
53+ const repo = await repos.getById(detail.attempt.repoId, viewer);
54+ return {
55+ ...detail,
56+ session: unwrap(session),
57+ viewer,
58+ // Only the repository's owner can land an attempt.
59+ canShip: repo.ok && repo.value.ownerId === viewer?.id,
60+ defaultBranch: repo.ok ? repo.value.defaultBranch : "main",
61+ };
5162 }
5263
5364 export async function action({ request, params, context }: Route.ActionArgs) {
5465 assertSameOrigin(request);
5566 const user = requireUser(context, request);
5667 const form = await request.formData();
68+ const action = form.get("action");
5769 const result =
58− form.get("action") === "abandon"
59− ? await env.WORK.abandonAttempt(user, params.id)
60− : await env.WORK.submitAttempt(
61− user,
62− params.id,
63− String(form.get("summary") ?? ""),
64− );
70+ action === "ship"
71+ ? await env.WORK.shipAttempt(user, params.id)
72+ : action === "abandon"
73+ ? await env.WORK.abandonAttempt(user, params.id)
74+ : await env.WORK.submitAttempt(
75+ user,
76+ params.id,
77+ String(form.get("summary") ?? ""),
78+ );
6579 return result.ok ? null : { error: result.error.message };
6680 }
6781
139153 actionData,
140154 params,
141155 }: Route.ComponentProps) {
142− const { attempt, intent, session, viewer } = loaderData;
156+ const { attempt, intent, session, viewer, canShip, defaultBranch } = loaderData;
143157 const base = `/${params.owner}/${params.repo}`;
144158 const remote = `https://g1t.sh/${attempt.fork.namespace}/${attempt.fork.name}.git`;
145159 const mine = viewer?.id === attempt.startedBy.id;
212226 </div>
213227
214228 <aside className="space-y-6">
229+ {canShip && active && intent.status === "open" && (
230+ <section className="rounded-xl border border-accent/30 bg-accent/5 p-4">
231+ <h3 className="text-sm font-medium">Ship this attempt</h3>
232+ <p className="mt-1 text-xs text-muted">
233+ Lands its commits on {defaultBranch} and closes the intent.
234+ </p>
235+ <Form method="post" className="mt-3 *:w-full">
236+ <Button variant="accent" type="submit" name="action" value="ship">
237+ <Rocket size={15} />
238+ Ship to {defaultBranch}
239+ </Button>
240+ </Form>
241+ {!(mine && active) && <ErrorText>{actionData?.error}</ErrorText>}
242+ </section>
243+ )}
215244 <section>
216245 <h3 className="text-sm font-medium">Working copy</h3>
217246 <p className="mt-1 text-xs text-muted">
+2−3
1−import { env } from "cloudflare:workers";
2−
31 import type { Route } from "./+types/blob";
42 import { BlobView } from "../../components/repo-view";
53 import { highlight } from "../../lib/highlight.server";
4+import { repos } from "../../lib/services.server";
65 import { getViewer, unwrap } from "../../lib/session.server";
76
87 export function meta({ params }: Route.MetaArgs) {
1211 export async function loader({ params, context }: Route.LoaderArgs) {
1312 const path = { namespace: params.owner, name: params.repo };
1413 const blob = unwrap(
15− await env.REPOS.blob(path, getViewer(context), params.ref, params["*"] ?? ""),
14+ await repos.blob(path, getViewer(context), params.ref, params["*"] ?? ""),
1615 );
1716 return {
1817 blob,
+2−3
1−import { env } from "cloudflare:workers";
2−
31 import type { Route } from "./+types/code";
42 import { TreeView } from "../../components/repo-view";
3+import { repos } from "../../lib/services.server";
54 import { getViewer, unwrap } from "../../lib/session.server";
65
76 export async function loader({ params, context }: Route.LoaderArgs) {
87 const path = { namespace: params.owner, name: params.repo };
9− return unwrap(await env.REPOS.tree(path, getViewer(context), null, ""));
8+ return unwrap(await repos.tree(path, getViewer(context), null, ""));
109 }
1110
1211 export default function Code({ loaderData }: Route.ComponentProps) {
+2−3
1−import { env } from "cloudflare:workers";
2−
31 import type { Route } from "./+types/commits";
42 import { Avatar, EmptyState, TimeAgo } from "../../components/ui";
3+import { repos } from "../../lib/services.server";
54 import { getViewer, unwrap } from "../../lib/session.server";
65
76 const PAGE_SIZE = 50;
1413 const path = { namespace: params.owner, name: params.repo };
1514 return {
1615 commits: unwrap(
17− await env.REPOS.log(path, getViewer(context), null, PAGE_SIZE),
16+ await repos.log(path, getViewer(context), null, PAGE_SIZE),
1817 ),
1918 };
2019 }
+2−1
55
66 import type { Route } from "./+types/layout";
77 import { Pill } from "../../components/ui";
8+import { repos } from "../../lib/services.server";
89 import { getViewer, unwrap } from "../../lib/session.server";
910
1011 export function meta({ params }: Route.MetaArgs) {
1516 const viewer = getViewer(context);
1617 const path = { namespace: params.owner, name: params.repo };
1718 const [repo, intents] = await Promise.all([
18− env.REPOS.get(path, viewer),
19+ repos.get(path, viewer),
1920 env.WORK.listIntents(path, viewer, "open"),
2021 ]);
2122 return {
+2−3
1−import { env } from "cloudflare:workers";
2−
31 import type { Route } from "./+types/tree";
42 import { TreeView } from "../../components/repo-view";
3+import { repos } from "../../lib/services.server";
54 import { getViewer, unwrap } from "../../lib/session.server";
65
76 export function meta({ params }: Route.MetaArgs) {
1211 export async function loader({ params, context }: Route.LoaderArgs) {
1312 const path = { namespace: params.owner, name: params.repo };
1413 return unwrap(
15− await env.REPOS.tree(path, getViewer(context), params.ref, params["*"] ?? ""),
14+ await repos.tree(path, getViewer(context), params.ref, params["*"] ?? ""),
1615 );
1716 }
1817
+2−2
1−import type { EventsApi, ReposApi, ServiceBinding, WorkApi } from "@g1t/contracts";
1+import type { EventsApi, ServiceBinding, WorkApi } from "@g1t/contracts";
22
33 declare global {
44 namespace Cloudflare {
55 interface Env {
66 IDENTITY: ServiceBinding;
77 /** Also serves git over HTTPS through `fetch`. */
8− REPOS: ReposApi & { fetch(request: Request): Promise<Response> };
8+ REPOS: ServiceBinding & { fetch(request: Request): Promise<Response> };
99 WORK: WorkApi;
1010 EVENTS: EventsApi;
1111 }
+45−0
1+//! Events published on the bus. Mirrors `packages/contracts/src/events.ts`.
2+
3+use serde::Serialize;
4+
5+/// What a publisher supplies; the bus fills in the id and time.
6+#[derive(Debug, Serialize)]
7+#[serde(rename_all = "camelCase")]
8+pub struct NewEvent<T: Serialize> {
9+ #[serde(rename = "type")]
10+ pub kind: &'static str,
11+ /// The service that published it.
12+ pub source: &'static str,
13+ /// The repo the event concerns.
14+ pub repo_id: Option<String>,
15+ /// The user or agent that caused it, if any.
16+ pub actor: Option<String>,
17+ pub data: T,
18+}
19+
20+#[derive(Debug, Serialize)]
21+#[serde(rename_all = "camelCase")]
22+pub struct RepoCreated {
23+ pub repo_id: String,
24+ pub namespace: String,
25+ pub name: String,
26+ pub is_private: bool,
27+}
28+
29+#[derive(Debug, Serialize)]
30+#[serde(rename_all = "camelCase")]
31+pub struct RepoForked {
32+ pub repo_id: String,
33+ pub source_repo_id: String,
34+ pub attempt_id: String,
35+}
36+
37+/// `after` is the commit the ref points to once the push has landed.
38+#[derive(Debug, Serialize)]
39+#[serde(rename_all = "camelCase")]
40+pub struct GitPush {
41+ pub repo_id: String,
42+ #[serde(rename = "ref")]
43+ pub git_ref: String,
44+ pub after: String,
45+}
+2−0
44 //! arguments of each of its methods. Services and their callers depend on
55 //! this crate, never on each other's code.
66
7+pub mod events;
78 pub mod identity;
89 mod ids;
910 mod names;
1011 mod outcome;
12+pub mod repos;
1113
1214 pub use ids::new_id;
1315 pub use names::{is_valid_namespace, is_valid_repo_name};
+217−0
1+//! The repos service: repository metadata, contents, forks and git access.
2+//!
3+//! Each `*Args` struct is the argument of the method of the same name,
4+//! served at `POST /rpc/<method>`.
5+
6+use serde::{Deserialize, Serialize};
7+
8+use crate::{User, Viewer};
9+
10+#[derive(Clone, Debug, Serialize, Deserialize)]
11+#[serde(rename_all = "camelCase")]
12+pub struct Repo {
13+ pub id: String,
14+ /// The owning user's (later, workspace's) name: the first URL segment.
15+ pub namespace: String,
16+ pub name: String,
17+ pub description: Option<String>,
18+ pub is_private: bool,
19+ pub owner_id: String,
20+ pub default_branch: String,
21+ /// Set when this repo is an attempt's working copy of another repo.
22+ pub fork_of: Option<String>,
23+ /// Milliseconds since the epoch.
24+ pub created_at: u64,
25+}
26+
27+#[derive(Clone, Debug, Serialize, Deserialize)]
28+pub struct RepoPath {
29+ pub namespace: String,
30+ pub name: String,
31+}
32+
33+#[derive(Clone, Debug, Serialize, Deserialize)]
34+pub struct Signature {
35+ pub name: String,
36+ pub email: String,
37+}
38+
39+#[derive(Clone, Debug, Serialize, Deserialize)]
40+#[serde(rename_all = "camelCase")]
41+pub struct Commit {
42+ pub hash: String,
43+ pub tree_hash: String,
44+ pub message: String,
45+ pub author: Signature,
46+ pub parents: Vec<String>,
47+ /// Milliseconds since the epoch.
48+ pub authored_at: u64,
49+}
50+
51+#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
52+#[serde(rename_all = "lowercase")]
53+pub enum EntryKind {
54+ Tree,
55+ Blob,
56+ Symlink,
57+ Gitlink,
58+ Exec,
59+}
60+
61+#[derive(Clone, Debug, Serialize, Deserialize)]
62+pub struct TreeEntry {
63+ pub name: String,
64+ pub hash: String,
65+ pub kind: EntryKind,
66+}
67+
68+#[derive(Clone, Debug, Serialize, Deserialize)]
69+pub struct Readme {
70+ pub name: String,
71+ /// Null when the file is binary or too large to show.
72+ pub text: Option<String>,
73+}
74+
75+#[derive(Clone, Debug, Serialize, Deserialize)]
76+pub struct TreeView {
77+ pub repo: Repo,
78+ #[serde(rename = "ref")]
79+ pub git_ref: String,
80+ pub path: String,
81+ /// Null when the repo has no commits yet.
82+ pub head: Option<Commit>,
83+ pub entries: Vec<TreeEntry>,
84+ pub readme: Option<Readme>,
85+}
86+
87+#[derive(Clone, Debug, Serialize, Deserialize)]
88+pub struct BlobView {
89+ pub repo: Repo,
90+ #[serde(rename = "ref")]
91+ pub git_ref: String,
92+ pub path: String,
93+ pub size: u64,
94+ /// Null when the file is binary or too large to show.
95+ pub text: Option<String>,
96+}
97+
98+/// A git remote and a short-lived credential for it.
99+#[derive(Clone, Debug, Serialize, Deserialize)]
100+pub struct GitAccess {
101+ pub remote: String,
102+ pub token: String,
103+}
104+
105+#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
106+pub enum GitService {
107+ #[serde(rename = "git-upload-pack")]
108+ UploadPack,
109+ #[serde(rename = "git-receive-pack")]
110+ ReceivePack,
111+}
112+
113+/// The commit `main` points to after an attempt has landed.
114+#[derive(Clone, Debug, Serialize, Deserialize)]
115+pub struct Landed {
116+ pub commit: String,
117+}
118+
119+/// `get`. Returns `Outcome<Repo>`.
120+#[derive(Debug, Serialize, Deserialize)]
121+pub struct GetArgs {
122+ pub path: RepoPath,
123+ pub viewer: Viewer,
124+}
125+
126+/// `get_by_id`. Returns `Outcome<Repo>`.
127+#[derive(Debug, Serialize, Deserialize)]
128+pub struct GetByIdArgs {
129+ pub id: String,
130+ pub viewer: Viewer,
131+}
132+
133+/// `list`: repos the viewer may see, newest first. Returns `Vec<Repo>`.
134+#[derive(Debug, Default, Serialize, Deserialize)]
135+pub struct ListArgs {
136+ pub viewer: Viewer,
137+ #[serde(default)]
138+ pub query: Option<String>,
139+ #[serde(default)]
140+ pub namespace: Option<String>,
141+}
142+
143+/// `create`. Returns `Outcome<Repo>`.
144+#[derive(Debug, Serialize, Deserialize)]
145+#[serde(rename_all = "camelCase")]
146+pub struct CreateArgs {
147+ pub owner: User,
148+ pub name: String,
149+ #[serde(default)]
150+ pub description: Option<String>,
151+ #[serde(default)]
152+ pub is_private: bool,
153+}
154+
155+/// `tree`. Returns `Outcome<TreeView>`.
156+#[derive(Debug, Serialize, Deserialize)]
157+#[serde(rename_all = "camelCase")]
158+pub struct TreeArgs {
159+ pub path: RepoPath,
160+ pub viewer: Viewer,
161+ /// The default branch when absent.
162+ #[serde(default, rename = "ref")]
163+ pub git_ref: Option<String>,
164+ #[serde(default)]
165+ pub tree_path: String,
166+}
167+
168+/// `blob`. Returns `Outcome<BlobView>`.
169+#[derive(Debug, Serialize, Deserialize)]
170+#[serde(rename_all = "camelCase")]
171+pub struct BlobArgs {
172+ pub path: RepoPath,
173+ pub viewer: Viewer,
174+ #[serde(rename = "ref")]
175+ pub git_ref: String,
176+ pub file_path: String,
177+}
178+
179+/// `log`. Returns `Outcome<Vec<Commit>>`.
180+#[derive(Debug, Serialize, Deserialize)]
181+pub struct LogArgs {
182+ pub path: RepoPath,
183+ pub viewer: Viewer,
184+ #[serde(default, rename = "ref")]
185+ pub git_ref: Option<String>,
186+ pub limit: u32,
187+}
188+
189+/// `fork_for_attempt`: a copy-on-write copy of the source repo, hidden from
190+/// listings, for one attempt to work in. Returns `Outcome<Repo>`.
191+#[derive(Debug, Serialize, Deserialize)]
192+#[serde(rename_all = "camelCase")]
193+pub struct ForkArgs {
194+ pub source_id: String,
195+ pub attempt_id: String,
196+ pub actor: User,
197+}
198+
199+/// `git_access`: authorizes a git operation and says where to send it.
200+/// Pushing to a repo that does not exist creates it in the pusher's own
201+/// namespace. Returns `Outcome<GitAccess>`.
202+#[derive(Debug, Serialize, Deserialize)]
203+pub struct GitAccessArgs {
204+ pub path: RepoPath,
205+ pub viewer: Viewer,
206+ pub service: GitService,
207+}
208+
209+/// `land`: moves the default branch of the repo a fork came from to the
210+/// fork's head. Refused with `conflict` when the fork is behind, since
211+/// that would discard commits. Returns `Outcome<Landed>`.
212+#[derive(Debug, Serialize, Deserialize)]
213+#[serde(rename_all = "camelCase")]
214+pub struct LandArgs {
215+ pub fork_id: String,
216+ pub actor: User,
217+}
+94−39
6363 /// Helpers for bindings that workers-rs has no typed wrapper for, such as
6464 /// Artifacts and Email Sending. Values cross the boundary as JSON.
6565 pub mod js {
66+ use std::fmt;
67+
6668 use serde::Serialize;
6769 use serde::de::DeserializeOwned;
68− use worker::js_sys::{Function, JSON, Promise, Reflect};
70+ use worker::js_sys::{Array, Function, JSON, Promise, Reflect};
6971 use worker::wasm_bindgen::{JsCast, JsValue};
7072 use worker::wasm_bindgen_futures::JsFuture;
7173 use worker::{Env, Error, Result};
7274
73− fn error(context: &str, value: JsValue) -> Error {
74− let message = JSON::stringify(&value)
75− .ok()
76− .and_then(|text| text.as_string())
77− .filter(|text| text != "{}")
78− .or_else(|| {
79− Reflect::get(&value, &"message".into())
75+ /// An exception thrown by JavaScript, with its `code` if it had one.
76+ #[derive(Debug)]
77+ pub struct Thrown {
78+ pub code: Option<String>,
79+ pub message: String,
80+ }
81+
82+ impl Thrown {
83+ fn from_value(value: JsValue) -> Self {
84+ let property = |name: &str| {
85+ Reflect::get(&value, &name.into())
8086 .ok()
81− .and_then(|message| message.as_string())
82− })
83− .unwrap_or_else(|| format!("{value:?}"));
84− Error::RustError(format!("{context}: {message}"))
87+ .and_then(|property| property.as_string())
88+ };
89+ Thrown {
90+ code: property("code"),
91+ message: property("message").unwrap_or_else(|| format!("{value:?}")),
92+ }
93+ }
94+
95+ pub fn is(&self, code: &str) -> bool {
96+ self.code.as_deref() == Some(code)
97+ }
8598 }
8699
100+ impl fmt::Display for Thrown {
101+ fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
102+ match &self.code {
103+ Some(code) => write!(f, "{code}: {}", self.message),
104+ None => f.write_str(&self.message),
105+ }
106+ }
107+ }
108+
109+ impl From<Thrown> for Error {
110+ fn from(thrown: Thrown) -> Self {
111+ Error::RustError(thrown.to_string())
112+ }
113+ }
114+
87115 /// The binding called `name`, as a raw JavaScript value.
88116 pub fn binding(env: &Env, name: &str) -> Result<JsValue> {
89− let value = Reflect::get(env.as_ref(), &name.into()).map_err(|e| error(name, e))?;
117+ let value = Reflect::get(env.as_ref(), &name.into()).map_err(Thrown::from_value)?;
90118 if value.is_undefined() {
91119 return Err(Error::RustError(format!(
92120 "binding {name} is not configured"
95123 Ok(value)
96124 }
97125
126+ /// Reads a property of a JavaScript object.
127+ pub fn get(target: &JsValue, name: &str) -> JsValue {
128+ Reflect::get(target, &name.into()).unwrap_or(JsValue::UNDEFINED)
129+ }
130+
98131 pub fn to_js<T: Serialize>(value: &T) -> Result<JsValue> {
99− JSON::parse(&serde_json::to_string(value)?).map_err(|e| error("to_js", e))
132+ Ok(JSON::parse(&serde_json::to_string(value)?).map_err(Thrown::from_value)?)
100133 }
101134
102135 pub fn from_js<T: DeserializeOwned>(value: &JsValue) -> Result<T> {
103− if value.is_undefined() {
104− return Ok(serde_json::from_value(serde_json::Value::Null)?);
105− }
106− let text = JSON::stringify(value)
107− .map_err(|e| error("from_js", e))?
108− .as_string()
109− .unwrap_or_else(|| "null".to_owned());
110− Ok(serde_json::from_str(&text)?)
136+ let text = if value.is_undefined() {
137+ None
138+ } else {
139+ JSON::stringify(value)
140+ .map_err(Thrown::from_value)?
141+ .as_string()
142+ };
143+ let text = text.as_deref().unwrap_or("null");
144+ serde_json::from_str(text).map_err(|error| {
145+ // Say what arrived; a bare serde error is useless in a log.
146+ let seen: String = text.chars().take(300).collect();
147+ Error::RustError(format!(
148+ "unexpected value from JavaScript ({error}): {seen}"
149+ ))
150+ })
111151 }
112152
113− /// Calls `target.method(...args)` and awaits the result if it is a
114− /// promise.
115− pub async fn call(target: &JsValue, method: &str, args: &[JsValue]) -> Result<JsValue> {
116− let function: Function = Reflect::get(target, &method.into())
117− .map_err(|e| error(method, e))?
153+ /// Calls `target[method](...args)` and awaits the result if it is a
154+ /// thenable. `method` may be a name or a symbol.
155+ pub async fn call_key(
156+ target: &JsValue,
157+ method: &JsValue,
158+ args: &[JsValue],
159+ ) -> std::result::Result<JsValue, Thrown> {
160+ let function: Function = Reflect::get(target, method)
161+ .map_err(Thrown::from_value)?
118162 .dyn_into()
119− .map_err(|_| Error::RustError(format!("{method} is not a function")))?;
120− let arguments = worker::js_sys::Array::new();
121− for arg in args {
122− arguments.push(arg);
123− }
124− let returned = function
125− .apply(target, &arguments)
126− .map_err(|e| error(method, e))?;
127− match returned.dyn_into::<Promise>() {
128− Ok(promise) => JsFuture::from(promise).await.map_err(|e| error(method, e)),
129− Err(value) => Ok(value),
130− }
163+ .map_err(|_| Thrown {
164+ code: None,
165+ message: format!("{method:?} is not a function"),
166+ })?;
167+ let arguments: Array = args.iter().collect();
168+ // An RPC stub treats every property access as a remote method, so
169+ // `function.apply(...)` would be sent over the wire as a call to
170+ // "apply". Reflect.apply invokes the function without touching it.
171+ let returned = Reflect::apply(&function, target, &arguments).map_err(Thrown::from_value)?;
172+ // Worker RPC returns its own thenable rather than a Promise, so
173+ // resolve whatever came back instead of testing its type.
174+ JsFuture::from(Promise::resolve(&returned))
175+ .await
176+ .map_err(Thrown::from_value)
177+ }
178+
179+ /// Calls `target.method(...args)`; see [`call_key`].
180+ pub async fn call(
181+ target: &JsValue,
182+ method: &str,
183+ args: &[JsValue],
184+ ) -> std::result::Result<JsValue, Thrown> {
185+ call_key(target, &method.into(), args).await
131186 }
132187 }
+39−33
1717
1818 | Concept | What it is |
1919 | --- | --- |
20−| **Intent** | A goal stated against a repo, with acceptance checks (commands that must pass). Replaces the issue and the pull request. |
20+| **Intent** | A goal stated against a repo, with acceptance checks (commands that must pass). The issue, and the home of every pull request made for it. |
2121 | **Attempt** | One agent's run at an intent, in its own Artifacts fork. Any number run in parallel. |
2222 | **Session** | The agent's full context for an attempt: prompt, messages, tool calls, cost. Stored with the attempt and linked from every commit it produced. |
2323 | **Arena** | The compare view for an intent: every attempt side by side with diff, check results, conflicts against main and against each other, and a reviewer agent's summary. |
3131 agents are still working, and the agents are told.
3232 - **Live lanes.** Watch every attempt progress in real time.
3333
34+## Pull requests are not removed
35+
36+g1t is ordinary git, and the pull request stays. The model extends it rather
37+than replacing it, so an engineer's habits keep working and the agent
38+features are there when wanted.
39+
40+- **An attempt is a pull request.** It has a source (a fork, or a branch
41+ pushed to the repo), a diff, review comments, checks and a merge button.
42+ It is reachable as a pull request, with that name, in the UI and API.
43+- **An intent is the goal above it.** Opening a pull request the familiar way
44+ creates its intent from the title and description, so nobody has to learn
45+ the word to use the product.
46+- **The developer path is unchanged.** Push a branch, open a pull request,
47+ get review, merge.
48+- **The agent path adds to it.** State the intent first, let several
49+ attempts run, compare them, ship one.
50+- **Both paths meet at `main`.** The same landing rules apply to a person's
51+ pull request and an agent's attempt.
52+
3453 ## Converging on main
3554
3655 Twelve intents started together will finish at different times and touch
599618
600619 ## Build order
601620
602−Done: site on g1t.sh; git over HTTPS; accounts and private repos; the four
603−services and the event bus; intents, attempts (a fork each) and session
604−storage, with pages for each. SSH server written, not deployed.
621+Done: site with marketing page and docs; git over HTTPS; accounts with
622+registration, email verification and password reset; intents, attempts and
623+sessions; REST API and MCP server; event bus; shipping an attempt to `main`
624+with a behind check. Identity and repos are in Rust.
605625
606−1. `api.g1t.sh` and `mcp.g1t.sh`; CLI with Claude Code hooks, so an outside
607− agent can claim an intent, push, and record its session.
608−2. OAuth server.
609−3. SSH deployed; registration, email verification, forgot password; search;
610− GitHub import; this repo hosted on g1t.
611−4. Hosted agents in sandboxes behind the runner contract; agent
612− definitions; AI Gateway; push events; checks.
613−5. Merge engine with structural merge; landing queue with speculative
614− checks; projected main; resolve-on-move.
615−6. Arena, diffs, proof bundles, reviewer and adversarial reviewer, risk
616− tiers; work registry with overlap radar,
617− duplicate detection, handoff, ask and wait.
618−7. Adopting outside pushes; protected `main`; approval rules.
619−8. Projects with briefs, planner and dependency graph; mission control with
620− the "needs you" inbox; live session pages and steering.
621−9. Why-blame, signed provenance, forkable sessions, digest and timeline,
622− landing page.
623−10. Workspaces and initiatives; context hub (memory, unified search,
624− Jira and Notion connectors); multi-repo intents.
625−11. Portfolio, standup and ask; checkpoints, stall detection, budgets.
626−12. Moving sessions between local and hosted; installable web app with
627− notifications; preview URLs per attempt.
628−13. Docs view, document intents, templates, explain, living documentation,
629− roles.
630−14. Automations: event bus, triggers, Sentry and generic webhook
631− integrations, write-back.
632−15. Own keys, own endpoints, self-hosted runners; agent leaderboard and
633− routing.
634−16. Large run (100+ agents across many intents), hardening, README, demo.
626+1. Diffs and review on attempts; pull requests from pushed branches, under
627+ that name.
628+2. OAuth server, so MCP clients sign in through the browser with no token
629+ to paste.
630+3. Port work, events and the API to Rust; event storage per the design
631+ above.
632+4. CLI with Claude Code hooks to record sessions automatically.
633+5. Hosted agents in sandboxes; acceptance checks.
634+6. Server-side merge and rebase; landing queue with speculative checks;
635+ resolve-on-move.
636+7. Arena, proof bundles, reviewers, risk tiers; work registry, handoff.
637+8. Projects, mission control, steering; why-blame, digest, timeline.
638+9. Workspaces, context hub, portfolio; automations and integrations.
639+10. SSH; bot protection; own keys, endpoints and runners.
640+11. Large run (100+ agents across many intents), hardening, demo.
635641
636642 Later: code search, mirroring to GitHub, passkeys, SSH
637643 on port 22 without the CLI proxy (needs the Workers inbound TCP private
+1−1
55 "workspaces": ["apps/*", "services/*", "packages/*"],
66 "scripts": {
77 "typecheck": "npm run typecheck --workspaces --if-present",
8− "deploy": "npm run deploy -w @g1t/events -w @g1t/repos -w @g1t/work -w @g1t/api -w @g1t/web"
8+ "deploy": "npm run deploy -w @g1t/events -w @g1t/work -w @g1t/api -w @g1t/web"
99 },
1010 "devDependencies": {
1111 "typescript": "^5.9.3",
+21−0
11 import type { IdentityApi } from "./identity";
2+import type { ReposApi } from "./repos";
23
34 /** A service binding, as far as these clients need it. */
45 export type ServiceBinding = {
5354 removeAccessToken: (user, id) => call("remove_access_token", { user, id }),
5455 };
5556 }
57+
58+export function reposClient(service: ServiceBinding): ReposApi {
59+ const call = <T>(method: string, args: object) => rpc<T>(service, method, args);
60+ return {
61+ get: (path, viewer) => call("get", { path, viewer }),
62+ getById: (id, viewer) => call("get_by_id", { id, viewer }),
63+ list: (viewer, options = {}) => call("list", { viewer, ...options }),
64+ create: (owner, input) => call("create", { owner, ...input }),
65+ tree: (path, viewer, ref, treePath) =>
66+ call("tree", { path, viewer, ref, treePath }),
67+ blob: (path, viewer, ref, filePath) =>
68+ call("blob", { path, viewer, ref, filePath }),
69+ log: (path, viewer, ref, limit) => call("log", { path, viewer, ref, limit }),
70+ forkForAttempt: (sourceId, attemptId, actor) =>
71+ call("fork_for_attempt", { sourceId, attemptId, actor }),
72+ gitAccess: (path, viewer, service) =>
73+ call("git_access", { path, viewer, service }),
74+ land: (forkId, actor) => call("land", { forkId, actor }),
75+ };
76+}
+1−0
1616 "attempt.started": { attemptId: string; intentId: string; repoId: string; agent: string };
1717 "attempt.updated": { attemptId: string; intentId: string; repoId: string; status: string };
1818 "attempt.submitted": { attemptId: string; intentId: string; repoId: string };
19+ "attempt.shipped": { attemptId: string; intentId: string; repoId: string; commit: string };
1920 "session.appended": { attemptId: string; sessionId: string; count: number };
2021 };
2122
+7−0
8585 * repo that does not exist creates it in the pusher's own namespace.
8686 */
8787 gitAccess(path: RepoPath, viewer: Viewer, service: GitService): Promise<Result<GitAccess>>;
88+
89+ /**
90+ * Moves the default branch of the repo a fork came from to the fork's
91+ * head. Refused with "conflict" when the fork is behind, since that would
92+ * discard commits.
93+ */
94+ land(forkId: string, actor: User): Promise<Result<{ commit: string }>>;
8895 }
+5−0
8282 getAttempt(attemptId: string, viewer: Viewer): Promise<Result<{ attempt: Attempt; intent: Intent }>>;
8383 submitAttempt(actor: User, attemptId: string, summary: string): Promise<Result<Attempt>>;
8484 abandonAttempt(actor: User, attemptId: string): Promise<Result<Attempt>>;
85+ /**
86+ * Lands the attempt on the repository's default branch and closes its
87+ * intent as shipped. Only the repository's owner may ship.
88+ */
89+ shipAttempt(actor: User, attemptId: string): Promise<Result<Attempt>>;
8590 /** Attempts in progress that the viewer started, newest first. */
8691 listActiveAttempts(viewer: Viewer): Promise<{ attempt: Attempt; intent: Intent }[]>;
8792
+1−0
1+build/
+16−0
1+[package]
2+name = "g1t-repos"
3+version = "0.1.0"
4+edition.workspace = true
5+license.workspace = true
6+description = "Repository registry, contents, forks, landing and git over HTTPS."
7+
8+[lib]
9+crate-type = ["cdylib"]
10+
11+[dependencies]
12+g1t-contracts.workspace = true
13+g1t-kit.workspace = true
14+serde.workspace = true
15+serde_json.workspace = true
16+worker.workspace = true
+0−16
1−{
2− "name": "@g1t/repos",
3− "version": "0.1.0",
4− "private": true,
5− "type": "module",
6− "license": "MIT",
7− "scripts": {
8− "types": "wrangler types --include-env=false",
9− "typecheck": "wrangler types --include-env=false && tsc -p tsconfig.json",
10− "deploy": "wrangler deploy",
11− "migrate": "wrangler d1 migrations apply DB --remote"
12− },
13− "dependencies": {
14− "@g1t/contracts": "*"
15− }
16−}
+0−96
1−import type { Commit, GitAccess, TreeEntry } from "@g1t/contracts";
2−
3−import type { GitStore } from "./git-store";
4−
5−const TOKEN_TTL_SECONDS = 300;
6−
7−function isCode(error: unknown, code: ArtifactsErrorCode): boolean {
8− return (error as Partial<ArtifactsError> | null)?.code === code;
9−}
10−
11−function toCommit(commit: ArtifactsCommitMetadata): Commit {
12− return {
13− hash: commit.hash,
14− treeHash: commit.treeHash,
15− message: commit.message,
16− author: commit.author,
17− parents: commit.parents,
18− authoredAt: commit.authoredAt * 1000,
19− };
20−}
21−
22−/** {@link GitStore} backed by Cloudflare Artifacts. */
23−export class ArtifactsGitStore implements GitStore {
24− constructor(private readonly artifacts: Artifacts) {}
25−
26− private async use<T>(key: string, fn: (repo: ArtifactsRepo) => Promise<T>): Promise<T> {
27− const repo = await this.artifacts.get(key);
28− try {
29− return await fn(repo);
30− } finally {
31− repo[Symbol.dispose]();
32− }
33− }
34−
35− async create(
36− key: string,
37− options: { description?: string; defaultBranch: string },
38− ): Promise<void> {
39− try {
40− await this.artifacts.create(key, {
41− description: options.description,
42− setDefaultBranch: options.defaultBranch,
43− });
44− } catch (error) {
45− // Left behind by an earlier failed attempt; adopt it.
46− if (!isCode(error, "ALREADY_EXISTS")) throw error;
47− }
48− }
49−
50− async fork(sourceKey: string, targetKey: string): Promise<void> {
51− await this.use(sourceKey, async (repo) => {
52− try {
53− await repo.fork(targetKey, { defaultBranchOnly: true });
54− } catch (error) {
55− if (!isCode(error, "ALREADY_EXISTS")) throw error;
56− }
57− });
58− }
59−
60− async access(key: string, scope: "read" | "write"): Promise<GitAccess> {
61− return this.use(key, async (repo) => {
62− const [info, token] = await Promise.all([
63− repo.info(),
64− repo.createToken(scope, TOKEN_TTL_SECONDS),
65− ]);
66− return { remote: info.remote, token: token.plaintext };
67− });
68− }
69−
70− async log(key: string, ref: string, limit: number): Promise<Commit[]> {
71− return this.use(key, async (repo) =>
72− (await repo.log({ ref, limit })).map(toCommit),
73− );
74− }
75−
76− async readTree(key: string, treeHash: string): Promise<TreeEntry[] | null> {
77− return this.use(key, async (repo) => {
78− const entries = await repo.readTree(treeHash);
79− return (
80− entries?.map((entry) => ({
81− name: entry.name,
82− hash: entry.hash,
83− kind: entry.type,
84− })) ?? null
85− );
86− });
87− }
88−
89− async readBlob(key: string, blobHash: string): Promise<Blob | null> {
90− return this.use(key, (repo) => repo.readBlob(blobHash));
91− }
92−
93− async readFile(key: string, ref: string, path: string): Promise<Blob | null> {
94− return this.use(key, (repo) => repo.readFile({ ref, path }));
95− }
96−}
+0−88
1−import {
2− type GitService,
3− type IdentityApi,
4− type RepoPath,
5− type ReposApi,
6− type Viewer,
7− httpStatus,
8−} from "@g1t/contracts";
9−
10−const GIT_ROUTE =
11− /^\/([^/]+)\/([^/]+?)(?:\.git)?\/(info\/refs|git-upload-pack|git-receive-pack)$/;
12−const FORWARDED_HEADERS = [
13− "accept",
14− "content-encoding",
15− "content-type",
16− "git-protocol",
17− "user-agent",
18−];
19−
20−async function viewerFromBasicAuth(
21− request: Request,
22− identity: IdentityApi,
23−): Promise<Viewer> {
24− const [scheme, encoded] = (request.headers.get("authorization") ?? "").split(" ");
25− if (scheme?.toLowerCase() !== "basic" || !encoded) return null;
26− let decoded: string;
27− try {
28− decoded = atob(encoded);
29− } catch {
30− return null;
31− }
32− const separator = decoded.indexOf(":");
33− if (separator < 0) return null;
34− return identity.userForGitCredentials(
35− decoded.slice(0, separator),
36− decoded.slice(separator + 1),
37− );
38−}
39−
40−/**
41− * Smart HTTP git remote at `/<namespace>/<repo>.git`, proxied to the git
42− * store with a short-lived token. Returns null for requests that are not git.
43− */
44−export async function handleGitHttp(
45− request: Request,
46− identity: IdentityApi,
47− repos: Pick<ReposApi, "gitAccess">,
48− onPush: (path: RepoPath) => void,
49−): Promise<Response | null> {
50− const url = new URL(request.url);
51− const match = GIT_ROUTE.exec(url.pathname);
52− if (!match) return null;
53− const [, namespace, name, endpoint] = match;
54− const service =
55− endpoint === "info/refs" ? url.searchParams.get("service") : endpoint;
56− if (service !== "git-upload-pack" && service !== "git-receive-pack") {
57− return null;
58− }
59−
60− const viewer = await viewerFromBasicAuth(request, identity);
61− const access = await repos.gitAccess(
62− { namespace, name },
63− viewer,
64− service as GitService,
65− );
66− if (!access.ok) {
67− const status = httpStatus(access.error);
68− return new Response(`${access.error.message}\n`, {
69− status,
70− headers: status === 401 ? { "www-authenticate": 'Basic realm="g1t"' } : {},
71− });
72− }
73−
74− const headers = new Headers({ authorization: `Bearer ${access.value.token}` });
75− for (const header of FORWARDED_HEADERS) {
76− const value = request.headers.get(header);
77− if (value) headers.set(header, value);
78− }
79− const response = await fetch(`${access.value.remote}/${endpoint}${url.search}`, {
80− method: request.method,
81− headers,
82− body: request.body,
83− });
84− if (endpoint === "git-receive-pack" && response.ok) {
85− onPush({ namespace, name });
86− }
87− return response;
88−}
+0−23
1−import type { Commit, GitAccess, TreeEntry } from "@g1t/contracts";
2−
3−/**
4− * The storage that actually holds git repositories. The repos service
5− * depends on this port; Artifacts is one adapter for it.
6− *
7− * `key` is the store's own name for a repo.
8− */
9−export interface GitStore {
10− create(key: string, options: { description?: string; defaultBranch: string }): Promise<void>;
11− /** A copy-on-write copy of `sourceKey`. */
12− fork(sourceKey: string, targetKey: string): Promise<void>;
13− /** A remote URL and short-lived credential for git itself. */
14− access(key: string, scope: "read" | "write"): Promise<GitAccess>;
15−
16− /** Newest first along the first-parent chain; empty for an unknown ref. */
17− log(key: string, ref: string, limit: number): Promise<Commit[]>;
18− /** Null when the tree does not exist. */
19− readTree(key: string, treeHash: string): Promise<TreeEntry[] | null>;
20− readBlob(key: string, blobHash: string): Promise<Blob | null>;
21− /** Null when the ref or path does not resolve to a file. */
22− readFile(key: string, ref: string, path: string): Promise<Blob | null>;
23−}
+156−0
1+//! Git over HTTPS: the smart HTTP remote at `/<namespace>/<repo>.git`,
2+//! proxied to the git store with a short-lived token.
3+
4+use g1t_contracts::identity::GitCredentialsArgs;
5+use g1t_contracts::repos::{GitAccess, GitService, RepoPath};
6+use g1t_contracts::{FailureCode, Outcome, Viewer};
7+use worker::js_sys::Uint8Array;
8+use worker::{Fetch, Fetcher, Headers, Method, Request, RequestInit, Response, Result, Url};
9+
10+const ENDPOINTS: [&str; 3] = ["info/refs", "git-upload-pack", "git-receive-pack"];
11+const FORWARDED_HEADERS: [&str; 5] = [
12+ "accept",
13+ "content-encoding",
14+ "content-type",
15+ "git-protocol",
16+ "user-agent",
17+];
18+
19+/// A git request, parsed from its URL.
20+pub struct GitRequest {
21+ pub path: RepoPath,
22+ pub endpoint: &'static str,
23+ pub service: GitService,
24+}
25+
26+/// Parses `/<namespace>/<name>[.git]/<endpoint>`, or returns `None` if the
27+/// request is not git's.
28+pub fn parse(url: &Url) -> Option<GitRequest> {
29+ let path = url.path().strip_prefix('/')?;
30+ let endpoint = ENDPOINTS
31+ .into_iter()
32+ .find(|endpoint| path.ends_with(&format!("/{endpoint}")))?;
33+ let repo = &path[..path.len() - endpoint.len() - 1];
34+ let (namespace, name) = repo.split_once('/')?;
35+ let name = name.strip_suffix(".git").unwrap_or(name);
36+ if namespace.is_empty() || name.is_empty() || name.contains('/') {
37+ return None;
38+ }
39+ let service = if endpoint == "info/refs" {
40+ url.query_pairs()
41+ .find(|(key, _)| key == "service")
42+ .map(|(_, value)| value.into_owned())?
43+ } else {
44+ endpoint.to_owned()
45+ };
46+ let service = match service.as_str() {
47+ "git-upload-pack" => GitService::UploadPack,
48+ "git-receive-pack" => GitService::ReceivePack,
49+ _ => return None,
50+ };
51+ Some(GitRequest {
52+ path: RepoPath {
53+ namespace: namespace.to_owned(),
54+ name: name.to_owned(),
55+ },
56+ endpoint,
57+ service,
58+ })
59+}
60+
61+/// The user named by an HTTP Basic `Authorization` header, as git sends it.
62+pub async fn viewer(request: &Request, identity: &Fetcher) -> Result<Viewer> {
63+ let Some(header) = request.headers().get("authorization")? else {
64+ return Ok(None);
65+ };
66+ let Some((scheme, encoded)) = header.split_once(' ') else {
67+ return Ok(None);
68+ };
69+ if !scheme.eq_ignore_ascii_case("basic") {
70+ return Ok(None);
71+ }
72+ let Some(decoded) = decode_base64(encoded.trim()) else {
73+ return Ok(None);
74+ };
75+ let Some((username, secret)) = decoded.split_once(':') else {
76+ return Ok(None);
77+ };
78+ g1t_kit::call(
79+ identity,
80+ "user_for_git_credentials",
81+ &GitCredentialsArgs {
82+ username: username.to_owned(),
83+ secret: secret.to_owned(),
84+ },
85+ )
86+ .await
87+}
88+
89+/// Standard base64 to a UTF-8 string, or `None` if either step fails.
90+fn decode_base64(input: &str) -> Option<String> {
91+ let mut bytes = Vec::with_capacity(input.len() * 3 / 4);
92+ let mut buffer = 0u32;
93+ let mut bits = 0;
94+ for byte in input.bytes().filter(|byte| *byte != b'=') {
95+ let value = match byte {
96+ b'A'..=b'Z' => byte - b'A',
97+ b'a'..=b'z' => byte - b'a' + 26,
98+ b'0'..=b'9' => byte - b'0' + 52,
99+ b'+' => 62,
100+ b'/' => 63,
101+ _ => return None,
102+ };
103+ buffer = (buffer << 6) | u32::from(value);
104+ bits += 6;
105+ if bits >= 8 {
106+ bits -= 8;
107+ bytes.push((buffer >> bits) as u8);
108+ }
109+ }
110+ String::from_utf8(bytes).ok()
111+}
112+
113+/// The response for a refused git request. Anonymous callers are asked to
114+/// authenticate, which is what makes git prompt for credentials.
115+pub fn refuse<T>(outcome: Outcome<T>) -> Result<Response> {
116+ let Outcome::Fail(failure) = outcome else {
117+ return Response::error("Not found", 404);
118+ };
119+ let mut response = Response::error(failure.message, failure.code.http_status())?;
120+ if failure.code == FailureCode::Unauthenticated {
121+ response
122+ .headers_mut()
123+ .set("www-authenticate", "Basic realm=\"g1t\"")?;
124+ }
125+ Ok(response)
126+}
127+
128+/// Sends the request on to the git store and returns its response as is.
129+pub async fn forward(
130+ mut request: Request,
131+ git: &GitRequest,
132+ access: &GitAccess,
133+) -> Result<Response> {
134+ let headers = Headers::new();
135+ headers.set("authorization", &format!("Bearer {}", access.token))?;
136+ for name in FORWARDED_HEADERS {
137+ if let Some(value) = request.headers().get(name)? {
138+ headers.set(name, &value)?;
139+ }
140+ }
141+ let query = request
142+ .url()?
143+ .query()
144+ .map(|query| format!("?{query}"))
145+ .unwrap_or_default();
146+ let mut init = RequestInit::new();
147+ init.with_method(request.method()).with_headers(headers);
148+ if request.method() == Method::Post {
149+ // Pushes are capped at 100 MB by the platform, so buffering is safe.
150+ let body = request.bytes().await?;
151+ init.with_body(Some(Uint8Array::from(body.as_slice()).into()));
152+ }
153+ let upstream =
154+ Request::new_with_init(&format!("{}/{}{query}", access.remote, git.endpoint), &init)?;
155+ Fetch::Request(upstream).send().await
156+}
+0−312
1−import { WorkerEntrypoint } from "cloudflare:workers";
2−
3−import {
4− type BlobView,
5− type Commit,
6− type CreateRepoInput,
7− type EventsApi,
8− type GitAccess,
9− type GitService,
10− type ServiceBinding,
11− type NewEvent,
12− type Repo,
13− type RepoPath,
14− type ReposApi,
15− type Result,
16− type TreeView,
17− type User,
18− type Viewer,
19− UNVERIFIED,
20− fail,
21− identityClient,
22− isValidNamespace,
23− isValidRepoName,
24− newId,
25− ok,
26−} from "@g1t/contracts";
27−
28−import { ArtifactsGitStore } from "./artifacts-git-store";
29−import { handleGitHttp } from "./git-http";
30−import type { GitStore } from "./git-store";
31−import {
32− RepoRegistry,
33− canRead,
34− canWrite,
35− storeKey,
36−} from "./registry";
37−
38−export interface ReposEnv {
39− DB: D1Database;
40− ARTIFACTS: Artifacts;
41− IDENTITY: ServiceBinding;
42− EVENTS: EventsApi;
43−}
44−
45−/** Namespace that holds every attempt's fork: `attempts/<attempt id>`. */
46−const ATTEMPTS_NAMESPACE = "attempts";
47−const MAX_TEXT_BYTES = 512 * 1024;
48−const README = /^readme(\.(md|markdown|txt))?$/i;
49−const SOURCE = "repos";
50−
51−const NOT_FOUND = fail("not_found", "Repository not found.");
52−
53−/** Decoded text, or null when the blob is too large or looks binary. */
54−async function blobText(blob: Blob): Promise<string | null> {
55− if (blob.size > MAX_TEXT_BYTES) return null;
56− const bytes = new Uint8Array(await blob.arrayBuffer());
57− if (bytes.includes(0)) return null;
58− return new TextDecoder().decode(bytes);
59−}
60−
61−export default class ReposService
62− extends WorkerEntrypoint<ReposEnv>
63− implements ReposApi
64−{
65− private readonly registry = new RepoRegistry(this.env.DB);
66− private readonly store: GitStore = new ArtifactsGitStore(this.env.ARTIFACTS);
67−
68− /** Resolves a repo the viewer may read; private repos look missing. */
69− private async readable(path: RepoPath, viewer: Viewer): Promise<Repo | null> {
70− const repo = await this.registry.byPath(path);
71− return repo && canRead(repo, viewer) ? repo : null;
72− }
73−
74− async get(path: RepoPath, viewer: Viewer): Promise<Result<Repo>> {
75− const repo = await this.readable(path, viewer);
76− return repo ? ok(repo) : NOT_FOUND;
77− }
78−
79− async getById(id: string, viewer: Viewer): Promise<Result<Repo>> {
80− const repo = await this.registry.byId(id);
81− return repo && canRead(repo, viewer) ? ok(repo) : NOT_FOUND;
82− }
83−
84− async list(
85− viewer: Viewer,
86− options: { query?: string; namespace?: string } = {},
87− ): Promise<Repo[]> {
88− return this.registry.list(viewer, options);
89− }
90−
91− async create(owner: User, input: CreateRepoInput): Promise<Result<Repo>> {
92− if (!owner.verified) return UNVERIFIED;
93− const name = input.name.trim().toLowerCase();
94− if (!isValidRepoName(name)) {
95− return fail("invalid", "Use letters, digits, dots, hyphens and underscores only.");
96− }
97− if (!isValidNamespace(owner.username)) {
98− return fail("invalid", "This account cannot own repositories.");
99− }
100− const path = { namespace: owner.username, name };
101− if (await this.registry.byPath(path)) {
102− return fail("conflict", "You already have a repository with that name.");
103− }
104− const repo: Repo = {
105− id: newId("rep"),
106− ...path,
107− description: input.description?.trim() || null,
108− isPrivate: input.isPrivate ?? false,
109− ownerId: owner.id,
110− defaultBranch: "main",
111− forkOf: null,
112− createdAt: Date.now(),
113− };
114− await this.store.create(storeKey(repo), {
115− description: repo.description ?? undefined,
116− defaultBranch: repo.defaultBranch,
117− });
118− await this.registry.insert(repo);
119− await this.publish({
120− type: "repo.created",
121− source: SOURCE,
122− repoId: repo.id,
123− actor: owner.id,
124− data: {
125− repoId: repo.id,
126− namespace: repo.namespace,
127− name: repo.name,
128− isPrivate: repo.isPrivate,
129− },
130− });
131− return ok(repo);
132− }
133−
134− async tree(
135− path: RepoPath,
136− viewer: Viewer,
137− ref: string | null,
138− treePath: string,
139− ): Promise<Result<TreeView>> {
140− const repo = await this.readable(path, viewer);
141− if (!repo) return NOT_FOUND;
142− const key = storeKey(repo);
143− const resolvedRef = ref ?? repo.defaultBranch;
144− const base = { repo, ref: resolvedRef, path: treePath };
145−
146− const [head] = await this.store.log(key, resolvedRef, 1);
147− if (!head) {
148− // An unknown ref is an error; a repo with no commits is just empty.
149− if (ref) return fail("not_found", "No such branch, tag or commit.");
150− return ok({ ...base, head: null, entries: [], readme: null });
151− }
152−
153− let entries = await this.store.readTree(key, head.treeHash);
154− for (const segment of treePath.split("/").filter(Boolean)) {
155− const next = entries?.find(
156− (entry) => entry.name === segment && entry.kind === "tree",
157− );
158− if (!next) return fail("not_found", "No such directory.");
159− entries = await this.store.readTree(key, next.hash);
160− }
161− if (!entries) return fail("not_found", "No such directory.");
162− entries.sort(
163− (a, b) =>
164− Number(b.kind === "tree") - Number(a.kind === "tree") ||
165− a.name.localeCompare(b.name),
166− );
167−
168− const readmeEntry = entries.find(
169− (entry) => entry.kind === "blob" && README.test(entry.name),
170− );
171− const readmeBlob = readmeEntry
172− ? await this.store.readBlob(key, readmeEntry.hash)
173− : null;
174− const readme =
175− readmeEntry && readmeBlob
176− ? { name: readmeEntry.name, text: await blobText(readmeBlob) }
177− : null;
178− return ok({ ...base, head, entries, readme });
179− }
180−
181− async blob(
182− path: RepoPath,
183− viewer: Viewer,
184− ref: string,
185− filePath: string,
186− ): Promise<Result<BlobView>> {
187− const repo = await this.readable(path, viewer);
188− if (!repo) return NOT_FOUND;
189− const blob = filePath
190− ? await this.store.readFile(storeKey(repo), ref, filePath)
191− : null;
192− if (!blob) return fail("not_found", "No such file.");
193− return ok({
194− repo,
195− ref,
196− path: filePath,
197− size: blob.size,
198− text: await blobText(blob),
199− });
200− }
201−
202− async log(
203− path: RepoPath,
204− viewer: Viewer,
205− ref: string | null,
206− limit: number,
207− ): Promise<Result<Commit[]>> {
208− const repo = await this.readable(path, viewer);
209− if (!repo) return NOT_FOUND;
210− return ok(
211− await this.store.log(storeKey(repo), ref ?? repo.defaultBranch, limit),
212− );
213− }
214−
215− async forkForAttempt(
216− sourceId: string,
217− attemptId: string,
218− actor: User,
219− ): Promise<Result<Repo>> {
220− const source = await this.registry.byId(sourceId);
221− if (!source || !canRead(source, actor)) return NOT_FOUND;
222− const fork: Repo = {
223− id: newId("rep"),
224− namespace: ATTEMPTS_NAMESPACE,
225− name: attemptId,
226− description: null,
227− // A fork is exactly as visible as the repo it came from.
228− isPrivate: source.isPrivate,
229− ownerId: actor.id,
230− defaultBranch: source.defaultBranch,
231− forkOf: source.id,
232− createdAt: Date.now(),
233− };
234− await this.store.fork(storeKey(source), storeKey(fork));
235− await this.registry.insert(fork);
236− await this.publish({
237− type: "repo.forked",
238− source: SOURCE,
239− repoId: source.id,
240− actor: actor.id,
241− data: { repoId: fork.id, sourceRepoId: source.id, attemptId },
242− });
243− return ok(fork);
244− }
245−
246− async gitAccess(
247− path: RepoPath,
248− viewer: Viewer,
249− service: GitService,
250− ): Promise<Result<GitAccess>> {
251− const write = service === "git-receive-pack";
252− if (write && viewer && !viewer.verified) return UNVERIFIED;
253− let repo = await this.registry.byPath(path);
254− if (!repo) {
255− // Push to create, in the pusher's own namespace only.
256− if (!write || !viewer || viewer.username !== path.namespace.toLowerCase()) {
257− return this.denied(viewer);
258− }
259− const created = await this.create(viewer, { name: path.name });
260− if (!created.ok) return created;
261− repo = created.value;
262− } else if (write ? !canWrite(repo, viewer) : !canRead(repo, viewer)) {
263− return this.denied(viewer);
264− }
265− return ok(await this.store.access(storeKey(repo), write ? "write" : "read"));
266− }
267−
268− /**
269− * Anonymous callers are asked to authenticate whether or not the repo
270− * exists, so private repos cannot be told apart from missing ones.
271− */
272− private denied(viewer: Viewer): Result<never> {
273− return viewer
274− ? NOT_FOUND
275− : fail("unauthenticated", "Authentication required.");
276− }
277−
278− private async publish(event: NewEvent): Promise<void> {
279− await this.env.EVENTS.publish([event]);
280− }
281−
282− /** Git over HTTPS. */
283− async fetch(request: Request): Promise<Response> {
284− const response = await handleGitHttp(request, identityClient(this.env.IDENTITY), this, (path) =>
285− this.ctx.waitUntil(this.pushed(path)),
286− );
287− return response ?? new Response("Not found\n", { status: 404 });
288− }
289−
290− /**
291− * Publishes `git.push` after a push has gone through the git front end.
292− * Artifacts' own push subscriptions are per repository, which does not fit
293− * a repo per attempt, so the front end reports pushes itself.
294− */
295− private async pushed(path: RepoPath): Promise<void> {
296− const repo = await this.registry.byPath(path);
297− if (!repo) return;
298− const [head] = await this.store.log(storeKey(repo), repo.defaultBranch, 1);
299− if (!head) return;
300− await this.publish({
301− type: "git.push",
302− source: SOURCE,
303− repoId: repo.id,
304− actor: null,
305− data: {
306− repoId: repo.id,
307− ref: `refs/heads/${repo.defaultBranch}`,
308− after: head.hash,
309− },
310− });
311− }
312−}
+150−0
1+//! Landing an attempt: moving a repository's branch forward to a commit
2+//! that so far exists only in a fork.
3+//!
4+//! The Artifacts binding cannot write, so this speaks git's smart HTTP
5+//! protocol directly. It asks the fork for a pack holding exactly the
6+//! objects the target is missing, then pushes that pack to the target
7+//! unchanged. No object is parsed or rebuilt along the way.
8+
9+use g1t_contracts::repos::GitAccess;
10+use worker::js_sys::Uint8Array;
11+use worker::{Error, Fetch, Headers, Method, Request, RequestInit, Result};
12+
13+const ZERO_ID: &str = "0000000000000000000000000000000000000000";
14+const FLUSH: &[u8] = b"0000";
15+/// Side-band channels: pack data, and fatal errors.
16+const PACK_BAND: u8 = 1;
17+const ERROR_BAND: u8 = 3;
18+
19+fn pkt_line(payload: &str) -> Vec<u8> {
20+ format!("{:04x}{payload}", payload.len() + 4).into_bytes()
21+}
22+
23+async fn post(access: &GitAccess, service: &str, body: Vec<u8>) -> Result<Vec<u8>> {
24+ let headers = Headers::new();
25+ headers.set("authorization", &format!("Bearer {}", access.token))?;
26+ headers.set("content-type", &format!("application/x-{service}-request"))?;
27+ headers.set("accept", &format!("application/x-{service}-result"))?;
28+ let mut init = RequestInit::new();
29+ init.with_method(Method::Post)
30+ .with_headers(headers)
31+ .with_body(Some(Uint8Array::from(body.as_slice()).into()));
32+ let request = Request::new_with_init(&format!("{}/{service}", access.remote), &init)?;
33+ let mut response = Fetch::Request(request).send().await?;
34+ let bytes = response.bytes().await?;
35+ if response.status_code() != 200 {
36+ return Err(Error::RustError(format!(
37+ "{service} returned {}: {}",
38+ response.status_code(),
39+ String::from_utf8_lossy(&bytes)
40+ )));
41+ }
42+ Ok(bytes)
43+}
44+
45+/// The payloads of the pkt-lines in `bytes`, and the offset where they stop.
46+fn read_pkt_lines(bytes: &[u8]) -> (Vec<&[u8]>, usize) {
47+ let mut lines = Vec::new();
48+ let mut position = 0;
49+ while position + 4 <= bytes.len() {
50+ let Some(length) = std::str::from_utf8(&bytes[position..position + 4])
51+ .ok()
52+ .and_then(|hex| usize::from_str_radix(hex, 16).ok())
53+ else {
54+ break;
55+ };
56+ if length < 4 {
57+ // Flush, delimiter and response-end packets carry no payload.
58+ position += 4;
59+ continue;
60+ }
61+ let end = (position + length).min(bytes.len());
62+ lines.push(&bytes[position + 4..end]);
63+ position = end;
64+ }
65+ (lines, position)
66+}
67+
68+/// A pack holding everything reachable from `want` that is not reachable
69+/// from `have`.
70+async fn fetch_pack(source: &GitAccess, want: &str, have: Option<&str>) -> Result<Vec<u8>> {
71+ // Side-band framing puts the pack in its own channel, so its exact bytes
72+ // can be recovered. Without it the response ends in a stray flush packet
73+ // that a receiver rejects as junk after the pack.
74+ let mut body = pkt_line(&format!("want {want} side-band-64k\n"));
75+ body.extend_from_slice(FLUSH);
76+ if let Some(have) = have {
77+ body.extend(pkt_line(&format!("have {have}\n")));
78+ }
79+ body.extend(pkt_line("done\n"));
80+
81+ let response = post(source, "git-upload-pack", body).await?;
82+ let (lines, _) = read_pkt_lines(&response);
83+ let mut pack = Vec::new();
84+ for line in lines {
85+ match line.first() {
86+ Some(&PACK_BAND) => pack.extend_from_slice(&line[1..]),
87+ Some(&ERROR_BAND) => {
88+ return Err(Error::RustError(format!(
89+ "the fork refused the fetch: {}",
90+ String::from_utf8_lossy(&line[1..])
91+ )));
92+ }
93+ // ACK and NAK lines, and progress messages.
94+ _ => {}
95+ }
96+ }
97+ if !pack.starts_with(b"PACK") {
98+ return Err(Error::RustError(format!(
99+ "the fork did not send a pack: {}",
100+ String::from_utf8_lossy(&response)
101+ )));
102+ }
103+ Ok(pack)
104+}
105+
106+/// Updates `branch` on the target from `old` to `new`, sending `pack`.
107+/// `Err(reason)` in the inner result means git refused the update, for
108+/// example because the branch is no longer at `old`.
109+async fn push_pack(
110+ target: &GitAccess,
111+ branch: &str,
112+ old: Option<&str>,
113+ new: &str,
114+ pack: Vec<u8>,
115+) -> Result<std::result::Result<(), String>> {
116+ let reference = format!("refs/heads/{branch}");
117+ let mut body = pkt_line(&format!(
118+ "{} {new} {reference}\0 report-status\n",
119+ old.unwrap_or(ZERO_ID)
120+ ));
121+ body.extend_from_slice(FLUSH);
122+ body.extend(pack);
123+
124+ let response = post(target, "git-receive-pack", body).await?;
125+ let (lines, _) = read_pkt_lines(&response);
126+ let lines: Vec<String> = lines
127+ .into_iter()
128+ .map(|line| String::from_utf8_lossy(line).trim_end().to_owned())
129+ .collect();
130+ let unpacked = lines.iter().any(|line| line == "unpack ok");
131+ let updated = lines.iter().any(|line| *line == format!("ok {reference}"));
132+ Ok(if unpacked && updated {
133+ Ok(())
134+ } else {
135+ Err(lines.join("; "))
136+ })
137+}
138+
139+/// Moves `branch` on `target` from `old` to `new`, a commit that exists in
140+/// `source`. The caller must have checked that `new` descends from `old`.
141+pub async fn fast_forward(
142+ source: &GitAccess,
143+ target: &GitAccess,
144+ branch: &str,
145+ old: Option<&str>,
146+ new: &str,
147+) -> Result<std::result::Result<(), String>> {
148+ let pack = fetch_pack(source, new, old).await?;
149+ push_pack(target, branch, old, new, pack).await
150+}
+547−0
1+//! The repos service: repository metadata, contents, forks, landing, and
2+//! git over HTTPS.
3+//!
4+//! Other services reach it over `POST /rpc/<method>`; see
5+//! `g1t_contracts::repos` for the methods and their arguments. Any other
6+//! request is treated as git's smart HTTP protocol.
7+
8+mod git_http;
9+mod land;
10+mod registry;
11+mod store;
12+
13+use g1t_contracts::events::{GitPush, NewEvent, RepoCreated, RepoForked};
14+use g1t_contracts::repos::*;
15+use g1t_contracts::{
16+ FailureCode, Outcome, User, Viewer, is_valid_namespace, is_valid_repo_name, new_id,
17+};
18+use g1t_kit::{args, js, now_ms, reply, rpc_method};
19+use std::collections::{HashMap, HashSet};
20+
21+use serde::Serialize;
22+use worker::wasm_bindgen::JsValue;
23+use worker::{Context, Env, Request, Response, Result, event};
24+
25+use registry::{Registry, can_read, can_write, store_key};
26+use store::{ArtifactsStore, GitRepo, GitStore, Scope};
27+
28+/// Namespace that holds every attempt's fork: `attempts/<attempt id>`.
29+const ATTEMPTS_NAMESPACE: &str = "attempts";
30+const MAX_TEXT_BYTES: usize = 512 * 1024;
31+/// How far back an attempt may have forked and still be landed.
32+const MAX_ANCESTRY: u32 = 1000;
33+const SOURCE: &str = "repos";
34+const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
35+
36+fn not_found<T>() -> Outcome<T> {
37+ Outcome::fail(FailureCode::NotFound, "Repository not found.")
38+}
39+
40+/// Decoded text, or `None` when the file is too large or looks binary.
41+fn text_of(bytes: Vec<u8>) -> Option<String> {
42+ if bytes.len() > MAX_TEXT_BYTES || bytes.contains(&0) {
43+ return None;
44+ }
45+ Some(String::from_utf8_lossy(&bytes).into_owned())
46+}
47+
48+fn is_readme(name: &str) -> bool {
49+ matches!(
50+ name.to_lowercase().as_str(),
51+ "readme" | "readme.md" | "readme.markdown" | "readme.txt"
52+ )
53+}
54+
55+/// Whether `ancestor` is reachable from the newest commit in `history`.
56+///
57+/// `history` is the first-parent chain, which is all the store lists; a fork
58+/// that merged the target branch in has the target's head on a second
59+/// parent, so the walk follows every parent.
60+async fn descends_from<R: GitRepo>(repo: &R, history: &[Commit], ancestor: &str) -> Result<bool> {
61+ let known: HashMap<&str, &[String]> = history
62+ .iter()
63+ .map(|commit| (commit.hash.as_str(), commit.parents.as_slice()))
64+ .collect();
65+ let mut seen = HashSet::new();
66+ let mut queue: Vec<String> = history
67+ .first()
68+ .map(|c| c.hash.clone())
69+ .into_iter()
70+ .collect();
71+ while let Some(hash) = queue.pop() {
72+ if hash == ancestor {
73+ return Ok(true);
74+ }
75+ if !seen.insert(hash.clone()) || seen.len() > MAX_ANCESTRY as usize {
76+ continue;
77+ }
78+ match known.get(hash.as_str()) {
79+ Some(parents) => queue.extend(parents.iter().cloned()),
80+ None => queue.extend(repo.parents(&hash).await?.unwrap_or_default()),
81+ }
82+ }
83+ Ok(false)
84+}
85+
86+struct Repos<S: GitStore> {
87+ registry: Registry,
88+ store: S,
89+ /// The events service, an RPC stub.
90+ events: JsValue,
91+}
92+
93+impl<S: GitStore> Repos<S> {
94+ async fn publish<T: Serialize>(&self, event: NewEvent<T>) -> Result<()> {
95+ js::call(&self.events, "publish", &[js::to_js(&[event])?]).await?;
96+ Ok(())
97+ }
98+
99+ /// Resolves a repo the viewer may read; private repos look missing.
100+ async fn readable(&self, path: &RepoPath, viewer: &Viewer) -> Result<Option<Repo>> {
101+ Ok(self
102+ .registry
103+ .by_path(path)
104+ .await?
105+ .filter(|repo| can_read(repo, viewer)))
106+ }
107+
108+ async fn get(&self, a: GetArgs) -> Result<Outcome<Repo>> {
109+ Ok(self
110+ .readable(&a.path, &a.viewer)
111+ .await?
112+ .map_or_else(not_found, Outcome::Ok))
113+ }
114+
115+ async fn get_by_id(&self, a: GetByIdArgs) -> Result<Outcome<Repo>> {
116+ Ok(self
117+ .registry
118+ .by_id(&a.id)
119+ .await?
120+ .filter(|repo| can_read(repo, &a.viewer))
121+ .map_or_else(not_found, Outcome::Ok))
122+ }
123+
124+ async fn create(&self, a: CreateArgs) -> Result<Outcome<Repo>> {
125+ if !a.owner.verified {
126+ return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
127+ }
128+ let name = a.name.trim().to_lowercase();
129+ if !is_valid_repo_name(&name) {
130+ return Ok(Outcome::fail(
131+ FailureCode::Invalid,
132+ "Use letters, digits, dots, hyphens and underscores only.",
133+ ));
134+ }
135+ if !is_valid_namespace(&a.owner.username) {
136+ return Ok(Outcome::fail(
137+ FailureCode::Invalid,
138+ "This account cannot own repositories.",
139+ ));
140+ }
141+ let path = RepoPath {
142+ namespace: a.owner.username.clone(),
143+ name,
144+ };
145+ if self.registry.by_path(&path).await?.is_some() {
146+ return Ok(Outcome::fail(
147+ FailureCode::Conflict,
148+ "You already have a repository with that name.",
149+ ));
150+ }
151+ let now = now_ms();
152+ let repo = Repo {
153+ id: new_id("rep", now),
154+ namespace: path.namespace,
155+ name: path.name,
156+ description: a
157+ .description
158+ .map(|text| text.trim().to_owned())
159+ .filter(|text| !text.is_empty()),
160+ is_private: a.is_private,
161+ owner_id: a.owner.id.clone(),
162+ default_branch: "main".to_owned(),
163+ fork_of: None,
164+ created_at: now,
165+ };
166+ self.store
167+ .create(
168+ &store_key(&repo),
169+ repo.description.as_deref(),
170+ &repo.default_branch,
171+ )
172+ .await?;
173+ self.registry.insert(&repo).await?;
174+ self.publish(NewEvent {
175+ kind: "repo.created",
176+ source: SOURCE,
177+ repo_id: Some(repo.id.clone()),
178+ actor: Some(a.owner.id),
179+ data: RepoCreated {
180+ repo_id: repo.id.clone(),
181+ namespace: repo.namespace.clone(),
182+ name: repo.name.clone(),
183+ is_private: repo.is_private,
184+ },
185+ })
186+ .await?;
187+ Ok(Outcome::Ok(repo))
188+ }
189+
190+ async fn tree(&self, a: TreeArgs) -> Result<Outcome<TreeView>> {
191+ let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
192+ return Ok(not_found());
193+ };
194+ let git = self.store.open(&store_key(&repo)).await?;
195+ let git_ref = a
196+ .git_ref
197+ .clone()
198+ .unwrap_or_else(|| repo.default_branch.clone());
199+
200+ let Some(head) = git.log(&git_ref, 1).await?.into_iter().next() else {
201+ // An unknown ref is an error; a repo with no commits is just empty.
202+ if a.git_ref.is_some() {
203+ return Ok(Outcome::fail(
204+ FailureCode::NotFound,
205+ "No such branch, tag or commit.",
206+ ));
207+ }
208+ return Ok(Outcome::Ok(TreeView {
209+ repo,
210+ git_ref,
211+ path: a.tree_path,
212+ head: None,
213+ entries: Vec::new(),
214+ readme: None,
215+ }));
216+ };
217+
218+ let no_directory = || Outcome::fail(FailureCode::NotFound, "No such directory.");
219+ let mut entries = git.read_tree(&head.tree_hash).await?;
220+ for segment in a.tree_path.split('/').filter(|segment| !segment.is_empty()) {
221+ let next = entries.as_ref().and_then(|entries| {
222+ entries
223+ .iter()
224+ .find(|entry| entry.name == segment && entry.kind == EntryKind::Tree)
225+ });
226+ let Some(next) = next else {
227+ return Ok(no_directory());
228+ };
229+ entries = git.read_tree(&next.hash).await?;
230+ }
231+ let Some(mut entries) = entries else {
232+ return Ok(no_directory());
233+ };
234+ // Directories first, then by name.
235+ entries.sort_by(|a, b| {
236+ (b.kind == EntryKind::Tree)
237+ .cmp(&(a.kind == EntryKind::Tree))
238+ .then_with(|| a.name.cmp(&b.name))
239+ });
240+
241+ let readme_entry = entries
242+ .iter()
243+ .find(|entry| entry.kind == EntryKind::Blob && is_readme(&entry.name));
244+ let readme = match readme_entry {
245+ Some(entry) => git.read_blob(&entry.hash).await?.map(|bytes| Readme {
246+ name: entry.name.clone(),
247+ text: text_of(bytes),
248+ }),
249+ None => None,
250+ };
251+ Ok(Outcome::Ok(TreeView {
252+ repo,
253+ git_ref,
254+ path: a.tree_path,
255+ head: Some(head),
256+ entries,
257+ readme,
258+ }))
259+ }
260+
261+ async fn blob(&self, a: BlobArgs) -> Result<Outcome<BlobView>> {
262+ let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
263+ return Ok(not_found());
264+ };
265+ let bytes = if a.file_path.is_empty() {
266+ None
267+ } else {
268+ let git = self.store.open(&store_key(&repo)).await?;
269+ git.read_file(&a.git_ref, &a.file_path).await?
270+ };
271+ let Some(bytes) = bytes else {
272+ return Ok(Outcome::fail(FailureCode::NotFound, "No such file."));
273+ };
274+ Ok(Outcome::Ok(BlobView {
275+ repo,
276+ git_ref: a.git_ref,
277+ path: a.file_path,
278+ size: bytes.len() as u64,
279+ text: text_of(bytes),
280+ }))
281+ }
282+
283+ async fn log(&self, a: LogArgs) -> Result<Outcome<Vec<Commit>>> {
284+ let Some(repo) = self.readable(&a.path, &a.viewer).await? else {
285+ return Ok(not_found());
286+ };
287+ let git = self.store.open(&store_key(&repo)).await?;
288+ let git_ref = a.git_ref.unwrap_or_else(|| repo.default_branch.clone());
289+ Ok(Outcome::Ok(git.log(&git_ref, a.limit).await?))
290+ }
291+
292+ async fn fork_for_attempt(&self, a: ForkArgs) -> Result<Outcome<Repo>> {
293+ let viewer = Some(a.actor.clone());
294+ let Some(source) = self
295+ .registry
296+ .by_id(&a.source_id)
297+ .await?
298+ .filter(|repo| can_read(repo, &viewer))
299+ else {
300+ return Ok(not_found());
301+ };
302+ let now = now_ms();
303+ let fork = Repo {
304+ id: new_id("rep", now),
305+ namespace: ATTEMPTS_NAMESPACE.to_owned(),
306+ name: a.attempt_id.clone(),
307+ description: None,
308+ // A fork is exactly as visible as the repo it came from.
309+ is_private: source.is_private,
310+ owner_id: a.actor.id.clone(),
311+ default_branch: source.default_branch.clone(),
312+ fork_of: Some(source.id.clone()),
313+ created_at: now,
314+ };
315+ self.store
316+ .open(&store_key(&source))
317+ .await?
318+ .fork(&store_key(&fork))
319+ .await?;
320+ self.registry.insert(&fork).await?;
321+ self.publish(NewEvent {
322+ kind: "repo.forked",
323+ source: SOURCE,
324+ repo_id: Some(source.id.clone()),
325+ actor: Some(a.actor.id),
326+ data: RepoForked {
327+ repo_id: fork.id.clone(),
328+ source_repo_id: source.id,
329+ attempt_id: a.attempt_id,
330+ },
331+ })
332+ .await?;
333+ Ok(Outcome::Ok(fork))
334+ }
335+
336+ async fn git_access(&self, a: GitAccessArgs) -> Result<Outcome<GitAccess>> {
337+ let write = a.service == GitService::ReceivePack;
338+ // Anonymous callers are asked to authenticate whether or not the repo
339+ // exists, so private repos cannot be told apart from missing ones.
340+ let denied = || match &a.viewer {
341+ Some(_) => not_found(),
342+ None => Outcome::fail(FailureCode::Unauthenticated, "Authentication required."),
343+ };
344+ if let (true, Some(user)) = (write, &a.viewer)
345+ && !user.verified
346+ {
347+ return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
348+ }
349+
350+ let repo = match self.registry.by_path(&a.path).await? {
351+ Some(repo) => {
352+ let allowed = if write {
353+ can_write(&repo, &a.viewer)
354+ } else {
355+ can_read(&repo, &a.viewer)
356+ };
357+ if !allowed {
358+ return Ok(denied());
359+ }
360+ repo
361+ }
362+ None => {
363+ // Push to create, in the pusher's own namespace only.
364+ let owner = a
365+ .viewer
366+ .as_ref()
367+ .filter(|user| write && user.username == a.path.namespace.to_lowercase());
368+ let Some(owner) = owner else {
369+ return Ok(denied());
370+ };
371+ let created = self
372+ .create(CreateArgs {
373+ owner: owner.clone(),
374+ name: a.path.name.clone(),
375+ description: None,
376+ is_private: false,
377+ })
378+ .await?;
379+ match created {
380+ Outcome::Ok(repo) => repo,
381+ Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
382+ }
383+ }
384+ };
385+ let git = self.store.open(&store_key(&repo)).await?;
386+ let scope = if write { Scope::Write } else { Scope::Read };
387+ Ok(Outcome::Ok(git.access(scope).await?))
388+ }
389+
390+ async fn land(&self, a: LandArgs) -> Result<Outcome<Landed>> {
391+ let actor: Viewer = Some(a.actor.clone());
392+ let Some(fork) = self.registry.by_id(&a.fork_id).await? else {
393+ return Ok(not_found());
394+ };
395+ let target = match &fork.fork_of {
396+ Some(id) => self.registry.by_id(id).await?,
397+ None => None,
398+ };
399+ let Some(target) = target.filter(|repo| can_read(repo, &actor)) else {
400+ return Ok(not_found());
401+ };
402+ if !can_write(&target, &actor) {
403+ return Ok(Outcome::fail(
404+ FailureCode::Forbidden,
405+ "Only the repository's owner can land an attempt.",
406+ ));
407+ }
408+ if !a.actor.verified {
409+ return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
410+ }
411+
412+ let branch = &target.default_branch;
413+ let fork_git = self.store.open(&store_key(&fork)).await?;
414+ let target_git = self.store.open(&store_key(&target)).await?;
415+ let history = fork_git.log(branch, MAX_ANCESTRY).await?;
416+ let Some(new) = history.first().map(|commit| commit.hash.clone()) else {
417+ return Ok(Outcome::fail(
418+ FailureCode::Conflict,
419+ "This attempt has no commits to land.",
420+ ));
421+ };
422+ let old = target_git
423+ .log(branch, 1)
424+ .await?
425+ .into_iter()
426+ .next()
427+ .map(|commit| commit.hash);
428+
429+ if old.as_deref() == Some(new.as_str()) {
430+ return Ok(Outcome::Ok(Landed { commit: new }));
431+ }
432+ // Moving the branch to a commit that does not descend from its
433+ // current head would discard whatever landed in between.
434+ if let Some(old) = &old
435+ && !descends_from(&fork_git, &history, old).await?
436+ {
437+ return Ok(Outcome::fail(
438+ FailureCode::Conflict,
439+ format!(
440+ "{branch} has moved since this attempt started. Pull {branch} into the attempt's fork, push, and land again."
441+ ),
442+ ));
443+ }
444+
445+ let source_access = fork_git.access(Scope::Read).await?;
446+ let target_access = target_git.access(Scope::Write).await?;
447+ let pushed =
448+ land::fast_forward(&source_access, &target_access, branch, old.as_deref(), &new)
449+ .await?;
450+ if let Err(reason) = pushed {
451+ // Most often another attempt landed between the check and the push.
452+ return Ok(Outcome::fail(
453+ FailureCode::Conflict,
454+ format!("{branch} could not be updated: {reason}"),
455+ ));
456+ }
457+ self.publish_push(&target, &new, Some(a.actor.id)).await?;
458+ Ok(Outcome::Ok(Landed { commit: new }))
459+ }
460+
461+ async fn publish_push(&self, repo: &Repo, after: &str, actor: Option<String>) -> Result<()> {
462+ self.publish(NewEvent {
463+ kind: "git.push",
464+ source: SOURCE,
465+ repo_id: Some(repo.id.clone()),
466+ actor,
467+ data: GitPush {
468+ repo_id: repo.id.clone(),
469+ git_ref: format!("refs/heads/{}", repo.default_branch),
470+ after: after.to_owned(),
471+ },
472+ })
473+ .await
474+ }
475+
476+ /// Git over HTTPS.
477+ async fn git_http(&self, request: Request, env: &Env) -> Result<Response> {
478+ let Some(git) = git_http::parse(&request.url()?) else {
479+ return Response::error("Not found", 404);
480+ };
481+ let viewer = git_http::viewer(&request, &env.service("IDENTITY")?).await?;
482+ let access = self
483+ .git_access(GitAccessArgs {
484+ path: git.path.clone(),
485+ viewer: viewer.clone(),
486+ service: git.service,
487+ })
488+ .await?;
489+ let access = match access {
490+ Outcome::Ok(access) => access,
491+ refused => return git_http::refuse(refused),
492+ };
493+ let response = git_http::forward(request, &git, &access).await?;
494+
495+ // Artifacts' own push notifications are per repository, which does
496+ // not fit a repo per attempt, so the front end reports pushes itself.
497+ let pushed = git.endpoint == "git-receive-pack" && response.status_code() == 200;
498+ if pushed && let Some(repo) = self.registry.by_path(&git.path).await? {
499+ let head = self
500+ .store
501+ .open(&store_key(&repo))
502+ .await?
503+ .log(&repo.default_branch, 1)
504+ .await?;
505+ if let Some(head) = head.first() {
506+ self.publish_push(&repo, &head.hash, viewer.map(|user: User| user.id))
507+ .await?;
508+ }
509+ }
510+ Ok(response)
511+ }
512+}
513+
514+#[event(fetch)]
515+async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
516+ let repos = Repos {
517+ registry: Registry { db: env.d1("DB")? },
518+ store: ArtifactsStore::new(&env)?,
519+ events: js::binding(&env, "EVENTS")?,
520+ };
521+ let Some(method) = rpc_method(&request) else {
522+ return repos.git_http(request, &env).await;
523+ };
524+ let body: serde_json::Value = request.json().await?;
525+
526+ match method.as_str() {
527+ "get" => reply(&repos.get(args(body)?).await?),
528+ "get_by_id" => reply(&repos.get_by_id(args(body)?).await?),
529+ "list" => {
530+ let a: ListArgs = args(body)?;
531+ reply(
532+ &repos
533+ .registry
534+ .list(&a.viewer, a.query.as_deref(), a.namespace.as_deref())
535+ .await?,
536+ )
537+ }
538+ "create" => reply(&repos.create(args(body)?).await?),
539+ "tree" => reply(&repos.tree(args(body)?).await?),
540+ "blob" => reply(&repos.blob(args(body)?).await?),
541+ "log" => reply(&repos.log(args(body)?).await?),
542+ "fork_for_attempt" => reply(&repos.fork_for_attempt(args(body)?).await?),
543+ "git_access" => reply(&repos.git_access(args(body)?).await?),
544+ "land" => reply(&repos.land(args(body)?).await?),
545+ _ => Response::error("Unknown method", 404),
546+ }
547+}
+148−0
1+//! Repository metadata in D1.
2+
3+use g1t_contracts::Viewer;
4+use g1t_contracts::repos::{Repo, RepoPath};
5+use serde::Deserialize;
6+use worker::wasm_bindgen::JsValue;
7+use worker::{D1Database, Result};
8+
9+#[derive(Deserialize)]
10+struct RepoRow {
11+ id: String,
12+ namespace: String,
13+ name: String,
14+ description: Option<String>,
15+ is_private: u8,
16+ owner_id: String,
17+ default_branch: String,
18+ fork_of: Option<String>,
19+ created_at: u64,
20+}
21+
22+impl From<RepoRow> for Repo {
23+ fn from(row: RepoRow) -> Self {
24+ Repo {
25+ id: row.id,
26+ namespace: row.namespace,
27+ name: row.name,
28+ description: row.description,
29+ is_private: row.is_private != 0,
30+ owner_id: row.owner_id,
31+ default_branch: row.default_branch,
32+ fork_of: row.fork_of,
33+ created_at: row.created_at * 1000,
34+ }
35+ }
36+}
37+
38+/// The key a repo is stored under in the git store.
39+pub fn store_key(repo: &Repo) -> String {
40+ format!("{}--{}", repo.namespace, repo.name)
41+}
42+
43+pub fn can_read(repo: &Repo, viewer: &Viewer) -> bool {
44+ !repo.is_private || can_write(repo, viewer)
45+}
46+
47+pub fn can_write(repo: &Repo, viewer: &Viewer) -> bool {
48+ viewer.as_ref().is_some_and(|user| user.id == repo.owner_id)
49+}
50+
51+fn optional(value: &Option<String>) -> JsValue {
52+ value.as_deref().map_or(JsValue::NULL, JsValue::from)
53+}
54+
55+pub struct Registry {
56+ pub db: D1Database,
57+}
58+
59+impl Registry {
60+ pub async fn by_path(&self, path: &RepoPath) -> Result<Option<Repo>> {
61+ Ok(self
62+ .db
63+ .prepare("SELECT * FROM repos WHERE namespace = ? AND name = ?")
64+ .bind(&[
65+ path.namespace.to_lowercase().into(),
66+ path.name.to_lowercase().into(),
67+ ])?
68+ .first::<RepoRow>(None)
69+ .await?
70+ .map(Repo::from))
71+ }
72+
73+ pub async fn by_id(&self, id: &str) -> Result<Option<Repo>> {
74+ Ok(self
75+ .db
76+ .prepare("SELECT * FROM repos WHERE id = ?")
77+ .bind(&[id.into()])?
78+ .first::<RepoRow>(None)
79+ .await?
80+ .map(Repo::from))
81+ }
82+
83+ /// Repos the viewer may see, newest first. Excludes attempt forks.
84+ pub async fn list(
85+ &self,
86+ viewer: &Viewer,
87+ query: Option<&str>,
88+ namespace: Option<&str>,
89+ ) -> Result<Vec<Repo>> {
90+ let mut conditions = vec!["fork_of IS NULL", "(is_private = 0 OR owner_id = ?)"];
91+ let viewer_id = viewer.as_ref().map_or("", |user| user.id.as_str());
92+ let mut params: Vec<JsValue> = vec![viewer_id.into()];
93+ if let Some(namespace) = namespace {
94+ conditions.push("namespace = ?");
95+ params.push(namespace.to_lowercase().into());
96+ }
97+ if let Some(query) = query.map(str::trim).filter(|query| !query.is_empty()) {
98+ conditions.push("(name LIKE ? ESCAPE '\\' OR description LIKE ? ESCAPE '\\')");
99+ // LIKE wildcards in the query are matched literally.
100+ let escaped: String = query
101+ .chars()
102+ .flat_map(|c| match c {
103+ '\\' | '%' | '_' => vec!['\\', c],
104+ _ => vec![c],
105+ })
106+ .collect();
107+ let pattern = format!("%{escaped}%");
108+ params.push(pattern.as_str().into());
109+ params.push(pattern.into());
110+ }
111+ let sql = format!(
112+ "SELECT * FROM repos WHERE {} ORDER BY created_at DESC, id DESC LIMIT 50",
113+ conditions.join(" AND ")
114+ );
115+ let rows = self
116+ .db
117+ .prepare(sql)
118+ .bind(&params)?
119+ .all()
120+ .await?
121+ .results::<RepoRow>()?;
122+ Ok(rows.into_iter().map(Repo::from).collect())
123+ }
124+
125+ pub async fn insert(&self, repo: &Repo) -> Result<()> {
126+ self.db
127+ .prepare(
128+ "INSERT INTO repos
129+ (id, namespace, name, description, is_private, owner_id,
130+ default_branch, fork_of, created_at)
131+ VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
132+ )
133+ .bind(&[
134+ repo.id.as_str().into(),
135+ repo.namespace.as_str().into(),
136+ repo.name.as_str().into(),
137+ optional(&repo.description),
138+ (repo.is_private as u8).into(),
139+ repo.owner_id.as_str().into(),
140+ repo.default_branch.as_str().into(),
141+ optional(&repo.fork_of),
142+ ((repo.created_at / 1000) as f64).into(),
143+ ])?
144+ .run()
145+ .await?;
146+ Ok(())
147+ }
148+}
+0−108
1−import type { Repo, RepoPath, Viewer } from "@g1t/contracts";
2−
3−type RepoRow = {
4− id: string;
5− namespace: string;
6− name: string;
7− description: string | null;
8− is_private: number;
9− owner_id: string;
10− default_branch: string;
11− fork_of: string | null;
12− created_at: number;
13−};
14−
15−function toRepo(row: RepoRow): Repo {
16− return {
17− id: row.id,
18− namespace: row.namespace,
19− name: row.name,
20− description: row.description,
21− isPrivate: row.is_private === 1,
22− ownerId: row.owner_id,
23− defaultBranch: row.default_branch,
24− forkOf: row.fork_of,
25− createdAt: row.created_at * 1000,
26− };
27−}
28−
29−/** The key a repo is stored under in the git store. */
30−export function storeKey(repo: RepoPath): string {
31− return `${repo.namespace}--${repo.name}`;
32−}
33−
34−export function canRead(repo: Repo, viewer: Viewer): boolean {
35− return !repo.isPrivate || repo.ownerId === viewer?.id;
36−}
37−
38−export function canWrite(repo: Repo, viewer: Viewer): boolean {
39− return repo.ownerId === viewer?.id;
40−}
41−
42−/** Repository metadata in D1. */
43−export class RepoRegistry {
44− constructor(private readonly db: D1Database) {}
45−
46− async byPath(path: RepoPath): Promise<Repo | null> {
47− const row = await this.db
48− .prepare("SELECT * FROM repos WHERE namespace = ? AND name = ?")
49− .bind(path.namespace.toLowerCase(), path.name.toLowerCase())
50− .first<RepoRow>();
51− return row ? toRepo(row) : null;
52− }
53−
54− async byId(id: string): Promise<Repo | null> {
55− const row = await this.db
56− .prepare("SELECT * FROM repos WHERE id = ?")
57− .bind(id)
58− .first<RepoRow>();
59− return row ? toRepo(row) : null;
60− }
61−
62− /** Excludes attempt forks. */
63− async list(
64− viewer: Viewer,
65− options: { query?: string; namespace?: string },
66− ): Promise<Repo[]> {
67− const where = ["fork_of IS NULL", "(is_private = 0 OR owner_id = ?)"];
68− const params: unknown[] = [viewer?.id ?? ""];
69− if (options.namespace) {
70− where.push("namespace = ?");
71− params.push(options.namespace.toLowerCase());
72− }
73− if (options.query?.trim()) {
74− where.push("(name LIKE ? ESCAPE '\\' OR description LIKE ? ESCAPE '\\')");
75− const pattern = `%${options.query.trim().replace(/[\\%_]/g, "\\$&")}%`;
76− params.push(pattern, pattern);
77− }
78− const { results } = await this.db
79− .prepare(
80− `SELECT * FROM repos WHERE ${where.join(" AND ")}
81− ORDER BY created_at DESC, id DESC LIMIT 50`,
82− )
83− .bind(...params)
84− .all<RepoRow>();
85− return results.map(toRepo);
86− }
87−
88− async insert(repo: Repo): Promise<void> {
89− await this.db
90− .prepare(
91− `INSERT INTO repos
92− (id, namespace, name, description, is_private, owner_id, default_branch, fork_of, created_at)
93− VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`,
94− )
95− .bind(
96− repo.id,
97− repo.namespace,
98− repo.name,
99− repo.description,
100− Number(repo.isPrivate),
101− repo.ownerId,
102− repo.defaultBranch,
103− repo.forkOf,
104− Math.floor(repo.createdAt / 1000),
105− )
106− .run();
107− }
108−}
+225−0
1+//! The storage that actually holds git repositories.
2+//!
3+//! The service depends on the [`GitStore`] and [`GitRepo`] ports;
4+//! [`ArtifactsStore`] is the adapter for Cloudflare Artifacts.
5+
6+use g1t_contracts::repos::{Commit, EntryKind, GitAccess, Signature, TreeEntry};
7+use g1t_kit::js;
8+use serde::Deserialize;
9+use worker::js_sys::{Reflect, Uint8Array};
10+use worker::wasm_bindgen::{JsCast, JsValue};
11+use worker::{Env, Result};
12+
13+/// How long a credential handed to git stays valid.
14+const TOKEN_TTL_SECONDS: u32 = 300;
15+
16+#[derive(Clone, Copy)]
17+pub enum Scope {
18+ Read,
19+ Write,
20+}
21+
22+/// A place repositories live. `key` is the store's own name for a repo.
23+#[allow(async_fn_in_trait)]
24+pub trait GitStore {
25+ type Repo: GitRepo;
26+
27+ /// Creates an empty repository. Succeeds if it already exists.
28+ async fn create(
29+ &self,
30+ key: &str,
31+ description: Option<&str>,
32+ default_branch: &str,
33+ ) -> Result<()>;
34+ async fn open(&self, key: &str) -> Result<Self::Repo>;
35+}
36+
37+/// One open repository.
38+#[allow(async_fn_in_trait)]
39+pub trait GitRepo {
40+ /// A remote URL and short-lived credential for git itself.
41+ async fn access(&self, scope: Scope) -> Result<GitAccess>;
42+ /// Newest first along the first-parent chain; empty for an unknown ref.
43+ async fn log(&self, git_ref: &str, limit: u32) -> Result<Vec<Commit>>;
44+ /// The parents of a commit, or `None` if the commit does not exist.
45+ async fn parents(&self, commit_hash: &str) -> Result<Option<Vec<String>>>;
46+ async fn read_tree(&self, tree_hash: &str) -> Result<Option<Vec<TreeEntry>>>;
47+ async fn read_blob(&self, blob_hash: &str) -> Result<Option<Vec<u8>>>;
48+ /// `None` when the ref or path does not resolve to a file.
49+ async fn read_file(&self, git_ref: &str, path: &str) -> Result<Option<Vec<u8>>>;
50+ /// Makes a copy-on-write copy of this repository under `target_key`.
51+ async fn fork(&self, target_key: &str) -> Result<()>;
52+}
53+
54+pub struct ArtifactsStore {
55+ binding: JsValue,
56+}
57+
58+impl ArtifactsStore {
59+ pub fn new(env: &Env) -> Result<Self> {
60+ Ok(Self {
61+ binding: js::binding(env, "ARTIFACTS")?,
62+ })
63+ }
64+}
65+
66+impl GitStore for ArtifactsStore {
67+ type Repo = ArtifactsRepo;
68+
69+ async fn create(
70+ &self,
71+ key: &str,
72+ description: Option<&str>,
73+ default_branch: &str,
74+ ) -> Result<()> {
75+ let options = js::to_js(&serde_json::json!({
76+ "description": description,
77+ "setDefaultBranch": default_branch,
78+ }))?;
79+ match js::call(&self.binding, "create", &[key.into(), options]).await {
80+ // Left behind by an earlier failed attempt; adopt it.
81+ Err(thrown) if !thrown.is("ALREADY_EXISTS") => Err(thrown.into()),
82+ _ => Ok(()),
83+ }
84+ }
85+
86+ async fn open(&self, key: &str) -> Result<ArtifactsRepo> {
87+ Ok(ArtifactsRepo {
88+ handle: js::call(&self.binding, "get", &[key.into()]).await?,
89+ })
90+ }
91+}
92+
93+/// A handle to one Artifacts repository. It is an RPC stub, so it is
94+/// released when dropped.
95+pub struct ArtifactsRepo {
96+ handle: JsValue,
97+}
98+
99+impl Drop for ArtifactsRepo {
100+ fn drop(&mut self) {
101+ let symbol = js::get(&worker::js_sys::global(), "Symbol");
102+ let dispose = js::get(&symbol, "dispose");
103+ if let Ok(function) = Reflect::get(&self.handle, &dispose)
104+ .and_then(|value| value.dyn_into::<worker::js_sys::Function>())
105+ {
106+ let _ = function.call0(&self.handle);
107+ }
108+ }
109+}
110+
111+#[derive(Deserialize)]
112+#[serde(rename_all = "camelCase")]
113+struct RawCommit {
114+ hash: String,
115+ tree_hash: String,
116+ message: String,
117+ author: Signature,
118+ parents: Vec<String>,
119+ /// Seconds since the epoch.
120+ authored_at: u64,
121+}
122+
123+#[derive(Deserialize)]
124+struct RawEntry {
125+ name: String,
126+ hash: String,
127+ #[serde(rename = "type")]
128+ kind: EntryKind,
129+}
130+
131+#[derive(Deserialize)]
132+struct RawInfo {
133+ remote: String,
134+}
135+
136+#[derive(Deserialize)]
137+struct RawToken {
138+ plaintext: String,
139+}
140+
141+/// The bytes of a `Blob`, or `None` for null.
142+async fn blob_bytes(blob: JsValue) -> Result<Option<Vec<u8>>> {
143+ if blob.is_null() || blob.is_undefined() {
144+ return Ok(None);
145+ }
146+ let buffer = js::call(&blob, "arrayBuffer", &[]).await?;
147+ Ok(Some(Uint8Array::new(&buffer).to_vec()))
148+}
149+
150+impl GitRepo for ArtifactsRepo {
151+ async fn access(&self, scope: Scope) -> Result<GitAccess> {
152+ let scope = match scope {
153+ Scope::Read => "read",
154+ Scope::Write => "write",
155+ };
156+ let info: RawInfo = js::from_js(&js::call(&self.handle, "info", &[]).await?)?;
157+ let token: RawToken = js::from_js(
158+ &js::call(
159+ &self.handle,
160+ "createToken",
161+ &[scope.into(), TOKEN_TTL_SECONDS.into()],
162+ )
163+ .await?,
164+ )?;
165+ Ok(GitAccess {
166+ remote: info.remote,
167+ token: token.plaintext,
168+ })
169+ }
170+
171+ async fn log(&self, git_ref: &str, limit: u32) -> Result<Vec<Commit>> {
172+ let options = js::to_js(&serde_json::json!({ "ref": git_ref, "limit": limit }))?;
173+ let commits: Vec<RawCommit> =
174+ js::from_js(&js::call(&self.handle, "log", &[options]).await?)?;
175+ Ok(commits
176+ .into_iter()
177+ .map(|commit| Commit {
178+ hash: commit.hash,
179+ tree_hash: commit.tree_hash,
180+ message: commit.message,
181+ author: commit.author,
182+ parents: commit.parents,
183+ authored_at: commit.authored_at * 1000,
184+ })
185+ .collect())
186+ }
187+
188+ async fn parents(&self, commit_hash: &str) -> Result<Option<Vec<String>>> {
189+ let commit: Option<RawCommit> =
190+ js::from_js(&js::call(&self.handle, "readCommit", &[commit_hash.into()]).await?)?;
191+ Ok(commit.map(|commit| commit.parents))
192+ }
193+
194+ async fn read_tree(&self, tree_hash: &str) -> Result<Option<Vec<TreeEntry>>> {
195+ let entries: Option<Vec<RawEntry>> =
196+ js::from_js(&js::call(&self.handle, "readTree", &[tree_hash.into()]).await?)?;
197+ Ok(entries.map(|entries| {
198+ entries
199+ .into_iter()
200+ .map(|entry| TreeEntry {
201+ name: entry.name,
202+ hash: entry.hash,
203+ kind: entry.kind,
204+ })
205+ .collect()
206+ }))
207+ }
208+
209+ async fn read_blob(&self, blob_hash: &str) -> Result<Option<Vec<u8>>> {
210+ blob_bytes(js::call(&self.handle, "readBlob", &[blob_hash.into()]).await?).await
211+ }
212+
213+ async fn read_file(&self, git_ref: &str, path: &str) -> Result<Option<Vec<u8>>> {
214+ let args = js::to_js(&serde_json::json!({ "ref": git_ref, "path": path }))?;
215+ blob_bytes(js::call(&self.handle, "readFile", &[args]).await?).await
216+ }
217+
218+ async fn fork(&self, target_key: &str) -> Result<()> {
219+ let options = js::to_js(&serde_json::json!({ "defaultBranchOnly": true }))?;
220+ match js::call(&self.handle, "fork", &[target_key.into(), options]).await {
221+ Err(thrown) if !thrown.is("ALREADY_EXISTS") => Err(thrown.into()),
222+ _ => Ok(()),
223+ }
224+ }
225+}
+0−4
1−{
2− "extends": "../../tsconfig.base.json",
3− "include": ["src/**/*", "worker-configuration.d.ts"]
4−}
+2−1
33 "name": "g1t-repos",
44 "account_id": "1e6f2cffa3f445920836e8ebe446bb58",
55 "compatibility_date": "2026-09-26",
6− "main": "./src/index.ts",
6+ "main": "build/index.js",
7+ "build": { "command": "cargo install -q worker-build@0.8.7 && worker-build --release" },
78 "workers_dev": false,
89 "d1_databases": [
910 {
+75−8
1212 type OpenIntentInput,
1313 type RepoPath,
1414 type ReposApi,
15+ type ServiceBinding,
1516 type Result,
1617 type SessionEntry,
1718 type StartAttemptInput,
2223 fail,
2324 newId,
2425 ok,
26+ reposClient,
2527 } from "@g1t/contracts";
2628
2729 import {
3537
3638 export interface WorkEnv {
3739 DB: D1Database;
38− REPOS: ReposApi;
40+ REPOS: ServiceBinding;
3941 EVENTS: EventsApi;
4042 }
4143
5961 return this.env.DB;
6062 }
6163
64+ private get repos(): ReposApi {
65+ return reposClient(this.env.REPOS);
66+ }
67+
6268 private async intentById(id: string): Promise<Intent | null> {
6369 const row = await this.db
6470 .prepare(`SELECT ${INTENT_COLUMNS} FROM intents WHERE id = ?`)
8086 const attempt = await this.attemptById(id);
8187 if (!attempt) return NO_ATTEMPT;
8288 // Reading it must be allowed before "forbidden" may reveal it exists.
83− const repo = await this.env.REPOS.getById(attempt.repoId, actor);
89+ const repo = await this.repos.getById(attempt.repoId, actor);
8490 if (!repo.ok) return NO_ATTEMPT;
8591 if (attempt.startedBy.id !== actor.id) {
8692 return fail("forbidden", "Only the person who started an attempt can change it.");
101107 const title = input.title.trim();
102108 const brief = input.brief.trim();
103109 if (!title) return fail("invalid", "An intent needs a title.");
104− const repo = await this.env.REPOS.get(repoPath, actor);
110+ const repo = await this.repos.get(repoPath, actor);
105111 if (!repo.ok) return repo;
106112 const checks = (input.checks ?? []).map((check) => check.trim()).filter(Boolean);
107113
148154 viewer: Viewer,
149155 status?: IntentStatus,
150156 ): Promise<Result<Intent[]>> {
151− const repo = await this.env.REPOS.get(repoPath, viewer);
157+ const repo = await this.repos.get(repoPath, viewer);
152158 if (!repo.ok) return repo;
153159 const { results } = await this.db
154160 .prepare(
166172 number: number,
167173 viewer: Viewer,
168174 ): Promise<Result<IntentDetail>> {
169− const repo = await this.env.REPOS.get(repoPath, viewer);
175+ const repo = await this.repos.get(repoPath, viewer);
170176 if (!repo.ok) return repo;
171177 const row = await this.db
172178 .prepare(
185191 async withdrawIntent(actor: User, intentId: string): Promise<Result<Intent>> {
186192 const intent = await this.intentById(intentId);
187193 if (!intent) return NO_INTENT;
188− const repo = await this.env.REPOS.getById(intent.repoId, actor);
194+ const repo = await this.repos.getById(intent.repoId, actor);
189195 if (!repo.ok) return NO_INTENT;
190196 if (intent.author.id !== actor.id && repo.value.ownerId !== actor.id) {
191197 return fail("forbidden", "Only the author or the repo owner can withdraw an intent.");
221227 const agent = input.agent.trim() || "agent";
222228
223229 const id = newId("att");
224− const fork = await this.env.REPOS.forkForAttempt(intent.repoId, id, actor);
230+ const fork = await this.repos.forkForAttempt(intent.repoId, id, actor);
225231 if (!fork.ok) return fork.error.code === "not_found" ? NO_INTENT : fork;
226232
227233 const now = Date.now();
272278 ): Promise<Result<{ attempt: Attempt; intent: Intent }>> {
273279 const attempt = await this.attemptById(attemptId);
274280 if (!attempt) return NO_ATTEMPT;
275− const repo = await this.env.REPOS.getById(attempt.repoId, viewer);
281+ const repo = await this.repos.getById(attempt.repoId, viewer);
276282 if (!repo.ok) return NO_ATTEMPT;
277283 return ok({ attempt, intent: (await this.intentById(attempt.intentId))! });
278284 }
322328 return this.setStatus(actor, attemptId, "abandoned", null);
323329 }
324330
331+ async shipAttempt(actor: User, attemptId: string): Promise<Result<Attempt>> {
332+ const row = await this.db
333+ .prepare("SELECT * FROM attempts WHERE id = ?")
334+ .bind(attemptId)
335+ .first<AttemptRow>();
336+ if (!row) return NO_ATTEMPT;
337+ const attempt = toAttempt(row);
338+ // Whether the actor may see and write the repo is decided by repos.
339+ const repo = await this.repos.getById(attempt.repoId, actor);
340+ if (!repo.ok) return NO_ATTEMPT;
341+ if (attempt.status !== "working" && attempt.status !== "submitted") {
342+ return fail("conflict", `This attempt is already ${attempt.status}.`);
343+ }
344+ const intent = (await this.intentById(attempt.intentId))!;
345+ if (intent.status !== "open") {
346+ return fail("conflict", `This intent is already ${intent.status}.`);
347+ }
348+
349+ const landed = await this.repos.land(row.fork_repo_id, actor);
350+ if (!landed.ok) return landed;
351+
352+ const now = Date.now();
353+ await this.db.batch([
354+ this.db
355+ .prepare(
356+ "UPDATE attempts SET status = 'shipped', head_commit = ?, updated_at = ? WHERE id = ?",
357+ )
358+ .bind(landed.value.commit, now, attempt.id),
359+ this.db
360+ .prepare("UPDATE intents SET status = 'shipped' WHERE id = ?")
361+ .bind(intent.id),
362+ ]);
363+ const ids = {
364+ attemptId: attempt.id,
365+ intentId: intent.id,
366+ repoId: attempt.repoId,
367+ };
368+ await this.publish(
369+ {
370+ type: "attempt.shipped",
371+ source: SOURCE,
372+ repoId: attempt.repoId,
373+ actor: actor.id,
374+ data: { ...ids, commit: landed.value.commit },
375+ },
376+ {
377+ type: "intent.closed",
378+ source: SOURCE,
379+ repoId: attempt.repoId,
380+ actor: actor.id,
381+ data: { intentId: intent.id, repoId: attempt.repoId, reason: "shipped" },
382+ },
383+ );
384+ return ok({
385+ ...attempt,
386+ status: "shipped",
387+ headCommit: landed.value.commit,
388+ updatedAt: now,
389+ });
390+ }
391+
325392 async listActiveAttempts(
326393 viewer: Viewer,
327394 ): Promise<{ attempt: Attempt; intent: Intent }[]> {