Events service in Rust, with RFC 3339 times and accurate push events
- events: rewritten in Rust; publishing queues, the queue consumer writes the log and fans out to each subscriber's own queue. Event times are RFC 3339 text, the last place Unix milliseconds were stored - repos and work publish over the same JSON RPC as every other call between services - git.push is published once per branch that moved, naming that branch, instead of always naming the default branch. The branches are read from the push's own commands and checked against the store before reporting - work: a push event updates the pull request for that branch directly
24 files+519−2710/24 viewed
| 825 | 825 | ] | |
| 826 | 826 | ||
| 827 | 827 | [[package]] | |
| 828 | + | name = "g1t-events" | |
| 829 | + | version = "0.1.0" | |
| 830 | + | dependencies = [ | |
| 831 | + | "g1t-contracts", | |
| 832 | + | "g1t-kit", | |
| 833 | + | "serde", | |
| 834 | + | "serde_json", | |
| 835 | + | "worker", | |
| 836 | + | ] | |
| 837 | + | ||
| 838 | + | [[package]] | |
| 828 | 839 | name = "g1t-identity" | |
| 829 | 840 | version = "0.1.0" | |
| 830 | 841 | dependencies = [ |
| 1 | 1 | [workspace] | |
| 2 | 2 | resolver = "3" | |
| 3 | − | members = ["crates/*", "services/identity", "services/repos", "services/work"] | |
| 3 | + | members = ["crates/*", "services/events", "services/identity", "services/repos", "services/work"] | |
| 4 | 4 | ||
| 5 | 5 | [workspace.package] | |
| 6 | 6 | edition = "2024" |
| 64 | 64 | | `services/identity` | Accounts, workspaces, sessions, keys and tokens. Rust. | | |
| 65 | 65 | | `services/repos` | Repository registry, contents, forks, diffs, landing, git over HTTPS. Rust. | | |
| 66 | 66 | | `services/work` | Issues, pull requests, comments and sessions. Rust. | | |
| 67 | − | | `services/events` | The event bus and its log. | | |
| 67 | + | | `services/events` | The event bus and its log. Rust. | | |
| 68 | 68 | | `services/runner` | Starts the sandboxes g1t agents work in. | | |
| 69 | 69 | | `crates/runner` | The program inside a sandbox: runs the agent and reports back. Rust. | | |
| 70 | 70 | | `crates/contracts` | Types and service interfaces for the Rust services. | | |
| 75 | 75 | ||
| 76 | 76 | Each service is its own Worker with its own database. They call each other | |
| 77 | 77 | through service bindings and react to each other through events. Anything | |
| 78 | − | that is not a web UI is written in Rust or on its way there; the events | |
| 79 | − | service and the API are next. | |
| 78 | + | that is not a web UI is written in Rust or on its way there; the API is | |
| 79 | + | next. | |
| 80 | 80 | ||
| 81 | 81 | ## Run your own | |
| 82 | 82 | ||
| 102 | 102 | Deploy everything in dependency order: | |
| 103 | 103 | ||
| 104 | 104 | ```sh | |
| 105 | + | (cd services/events && npx wrangler deploy) | |
| 105 | 106 | (cd services/identity && npx wrangler deploy) | |
| 106 | 107 | (cd services/repos && npx wrangler deploy) | |
| 107 | 108 | (cd services/work && npx wrangler deploy) |
| 4 | 4 | import { | |
| 5 | 5 | type ServiceBinding, | |
| 6 | 6 | type Viewer, | |
| 7 | + | eventsClient, | |
| 7 | 8 | httpStatus, | |
| 8 | 9 | identityClient, | |
| 9 | 10 | reposClient, | |
| 17 | 18 | ||
| 18 | 19 | type Input = Record<string, unknown>; | |
| 19 | 20 | /** The Worker's raw bindings; Rust services are reached through clients. */ | |
| 20 | − | type Bindings = Omit<ApiEnv, "IDENTITY" | "REPOS" | "WORK"> & { | |
| 21 | + | type Bindings = { | |
| 22 | + | EVENTS: ServiceBinding; | |
| 21 | 23 | IDENTITY: ServiceBinding; | |
| 22 | 24 | REPOS: ServiceBinding; | |
| 23 | 25 | WORK: ServiceBinding; | |
| 92 | 94 | app.use(async (c, next) => { | |
| 93 | 95 | const [scheme, token] = (c.req.header("authorization") ?? "").split(" "); | |
| 94 | 96 | const services: ApiEnv = { | |
| 95 | − | ...c.env, | |
| 97 | + | EVENTS: eventsClient(c.env.EVENTS), | |
| 96 | 98 | IDENTITY: identityClient(c.env.IDENTITY), | |
| 97 | 99 | REPOS: reposClient(c.env.REPOS), | |
| 98 | 100 | WORK: workClient(c.env.WORK), |
| 1 | − | import type { EventsApi, RunnerApi, ServiceBinding } from "@g1t/contracts"; | |
| 1 | + | import type { RunnerApi, ServiceBinding } from "@g1t/contracts"; | |
| 2 | 2 | ||
| 3 | 3 | declare global { | |
| 4 | 4 | namespace Cloudflare { | |
| 7 | 7 | /** Also serves git over HTTPS through `fetch`. */ | |
| 8 | 8 | REPOS: ServiceBinding & { fetch(request: Request): Promise<Response> }; | |
| 9 | 9 | WORK: ServiceBinding; | |
| 10 | − | EVENTS: EventsApi; | |
| 11 | 10 | RUNNER: RunnerApi; | |
| 12 | 11 | } | |
| 13 | 12 | } |
| 10 | 10 | { "binding": "IDENTITY", "service": "g1t-identity" }, | |
| 11 | 11 | { "binding": "REPOS", "service": "g1t-repos" }, | |
| 12 | 12 | { "binding": "WORK", "service": "g1t-work" }, | |
| 13 | − | { "binding": "EVENTS", "service": "g1t-events" }, | |
| 14 | 13 | { "binding": "RUNNER", "service": "g1t-runner" } | |
| 15 | 14 | ], | |
| 16 | 15 | "observability": { "enabled": true }, |
| 1 | − | //! Events published on the bus. Mirrors `packages/contracts/src/events.ts`. | |
| 1 | + | //! Events published on the bus, and the events service that carries them. | |
| 2 | + | //! Mirrors `packages/contracts/src/events.ts`. | |
| 2 | 3 | ||
| 3 | 4 | use serde::Serialize; | |
| 4 | 5 | ||
| 34 | 35 | pub pull_id: String, | |
| 35 | 36 | } | |
| 36 | 37 | ||
| 37 | − | /// `after` is the commit the ref points to once the push has landed. | |
| 38 | + | /// One branch moved by a push. `after` is the commit it points to now. | |
| 38 | 39 | #[derive(Debug, Serialize)] | |
| 39 | 40 | #[serde(rename_all = "camelCase")] | |
| 40 | 41 | pub struct GitPush { | |
| 41 | 42 | pub repo_id: String, | |
| 43 | + | /// The full ref, such as `refs/heads/main`. | |
| 42 | 44 | #[serde(rename = "ref")] | |
| 43 | 45 | pub git_ref: String, | |
| 44 | 46 | pub after: String, | |
| 47 | + | /// Whether the ref is the repository's default branch. | |
| 48 | + | pub default_branch: bool, | |
| 45 | 49 | } | |
| 46 | 50 | ||
| 47 | 51 | /// The payload of `issue.opened`, `issue.updated`, `issue.closed` and | |
| 101 | 105 | pub count: u32, | |
| 102 | 106 | } | |
| 103 | 107 | ||
| 104 | − | /// An event as delivered to subscribers. `data` is left as JSON; each | |
| 105 | − | /// subscriber decodes the types it cares about. | |
| 108 | + | /// An event as stored in the log and delivered to subscribers. `data` is | |
| 109 | + | /// left as JSON; each reader decodes the types it cares about. | |
| 110 | + | #[derive(Clone, Debug, Serialize, serde::Deserialize)] | |
| 111 | + | #[serde(rename_all = "camelCase")] | |
| 112 | + | pub struct Event { | |
| 113 | + | /// Sorts by the time the event was published. | |
| 114 | + | pub id: String, | |
| 115 | + | #[serde(rename = "type")] | |
| 116 | + | pub kind: String, | |
| 117 | + | /// The service that published it. | |
| 118 | + | pub source: String, | |
| 119 | + | /// RFC 3339. | |
| 120 | + | pub time: String, | |
| 121 | + | /// The repo the event concerns. | |
| 122 | + | pub repo_id: Option<String>, | |
| 123 | + | /// The user or agent that caused it, if any. | |
| 124 | + | pub actor: Option<String>, | |
| 125 | + | pub data: serde_json::Value, | |
| 126 | + | } | |
| 127 | + | ||
| 128 | + | /// `publish`, as a publisher sends it. Returns nothing. | |
| 129 | + | #[derive(Debug, Serialize)] | |
| 130 | + | pub struct Publish<T: Serialize> { | |
| 131 | + | pub events: Vec<NewEvent<T>>, | |
| 132 | + | } | |
| 133 | + | ||
| 134 | + | /// `publish`, as the events service reads it. | |
| 135 | + | #[derive(Debug, serde::Deserialize)] | |
| 136 | + | pub struct PublishArgs { | |
| 137 | + | pub events: Vec<Published>, | |
| 138 | + | } | |
| 139 | + | ||
| 140 | + | /// A [`NewEvent`] of any type, as received. | |
| 106 | 141 | #[derive(Debug, serde::Deserialize)] | |
| 107 | 142 | #[serde(rename_all = "camelCase")] | |
| 108 | − | pub struct Delivered { | |
| 109 | − | pub id: String, | |
| 143 | + | pub struct Published { | |
| 110 | 144 | #[serde(rename = "type")] | |
| 111 | 145 | pub kind: String, | |
| 112 | − | /// Milliseconds since the epoch, until the bus itself moves to RFC 3339. | |
| 113 | − | pub time: serde_json::Value, | |
| 146 | + | pub source: String, | |
| 147 | + | #[serde(default)] | |
| 114 | 148 | pub repo_id: Option<String>, | |
| 149 | + | #[serde(default)] | |
| 150 | + | pub actor: Option<String>, | |
| 115 | 151 | pub data: serde_json::Value, | |
| 116 | 152 | } | |
| 153 | + | ||
| 154 | + | /// `list`: events from the log, newest first. Returns `Vec<Event>`. | |
| 155 | + | #[derive(Debug, Default, Serialize, serde::Deserialize)] | |
| 156 | + | #[serde(rename_all = "camelCase")] | |
| 157 | + | pub struct ListArgs { | |
| 158 | + | #[serde(default)] | |
| 159 | + | pub repo_id: Option<String>, | |
| 160 | + | /// Only these types; all types when empty. | |
| 161 | + | #[serde(default)] | |
| 162 | + | pub types: Vec<String>, | |
| 163 | + | /// Only events older than this event id. | |
| 164 | + | #[serde(default)] | |
| 165 | + | pub before: Option<String>, | |
| 166 | + | #[serde(default)] | |
| 167 | + | pub limit: Option<u32>, | |
| 168 | + | } |
| 128 | 128 | Reflect::get(target, &name.into()).unwrap_or(JsValue::UNDEFINED) | |
| 129 | 129 | } | |
| 130 | 130 | ||
| 131 | + | /// Sets a property on an object. | |
| 132 | + | pub fn set(target: &JsValue, name: &str, value: &JsValue) { | |
| 133 | + | let _ = Reflect::set(target, &name.into(), value); | |
| 134 | + | } | |
| 135 | + | ||
| 131 | 136 | pub fn to_js<T: Serialize>(value: &T) -> Result<JsValue> { | |
| 132 | 137 | Ok(JSON::parse(&serde_json::to_string(value)?).map_err(Thrown::from_value)?) | |
| 133 | 138 | } |
| 516 | 516 | | `services/identity` | Rust | Worker + D1 | Accounts, workspaces and memberships, sessions, SSH keys, access tokens, device sign-in, OAuth codes and grants | | |
| 517 | 517 | | `services/repos` | Rust | Worker + D1 + Artifacts | Repository registry, contents, forks, diffs, landing, git over HTTPS. Storage sits behind a `GitStore` port with an Artifacts adapter. | | |
| 518 | 518 | | `services/work` | Rust | Worker + D1 | Issues, pull requests, comments, sessions; later a Durable Object per repo for the landing queue and live state | | |
| 519 | − | | `services/events` | TypeScript, moving to Rust | Worker + Queues + D1 | The event bus: durable log, and one queue per subscribing service | | |
| 519 | + | | `services/events` | Rust | Worker + Queues + D1 | The event bus: durable log, and one queue per subscribing service | | |
| 520 | 520 | | `services/runner`, `crates/runner` | TypeScript, Rust | Worker + Containers | Starts a sandbox per g1t agent; the program inside runs the agent harness and reports through the public API | | |
| 521 | 521 | | `apps/web` | TypeScript | Worker | Server-rendered site. Holds no data; calls services over RPC. | | |
| 522 | 522 | | `apps/docs` | TypeScript | Worker (static) | Documentation and the API explorer | | |
| 553 | 553 | ||
| 554 | 554 | The site is TypeScript. Everything behind it is Rust, compiled to | |
| 555 | 555 | WebAssembly for Workers and natively for containers and the CLI. Services | |
| 556 | − | are being ported one at a time; identity, repos and work are done, events | |
| 557 | − | and the API are next. Rust services speak a | |
| 556 | + | are being ported one at a time; identity, repos, work and events are done, | |
| 557 | + | and the API is next. Rust services speak a | |
| 558 | 558 | small JSON protocol over service bindings (`POST /rpc/<method>`), with the | |
| 559 | 559 | types in `crates/contracts`. | |
| 560 | 560 | ||
| 648 | 648 | comments; pull requests in forks or from branches, with diffs and sessions, | |
| 649 | 649 | several per issue; merging with a behind check, which resolves the issue and supersedes | |
| 650 | 650 | the rest; g1t agents in sandboxes with a choice of model; REST API, OpenAPI | |
| 651 | − | and MCP server; event bus. Identity, repos and work are in Rust. | |
| 651 | + | and MCP server; event bus. Identity, repos, work and events are in Rust. | |
| 652 | 652 | ||
| 653 | 653 | 1. Branch protection, and deleting a branch once its pull request merges. | |
| 654 | 654 | 2. Scopes on OAuth grants and access tokens. | |
| 655 | − | 3. Port events and the API to Rust; event storage per the design above. | |
| 655 | + | 3. Port the API to Rust; event storage per the design above. | |
| 656 | 656 | 4. CLI with Claude Code hooks to record sessions automatically. | |
| 657 | 657 | 5. Acceptance checks run in sandboxes; review comments on lines. | |
| 658 | 658 | 6. Server-side merge and rebase; landing queue with speculative checks; |
| 1593 | 1593 | "resolved": "apps/docs", | |
| 1594 | 1594 | "link": true | |
| 1595 | 1595 | }, | |
| 1596 | − | "node_modules/@g1t/events": { | |
| 1597 | − | "resolved": "services/events", | |
| 1598 | − | "link": true | |
| 1599 | − | }, | |
| 1600 | 1596 | "node_modules/@g1t/runner": { | |
| 1601 | 1597 | "resolved": "services/runner", | |
| 1602 | 1598 | "link": true | |
| 9804 | 9800 | "name": "@g1t/theme", | |
| 9805 | 9801 | "version": "0.1.0", | |
| 9806 | 9802 | "license": "MIT" | |
| 9807 | − | }, | |
| 9808 | − | "services/events": { | |
| 9809 | − | "name": "@g1t/events", | |
| 9810 | − | "version": "0.1.0", | |
| 9811 | − | "license": "MIT", | |
| 9812 | − | "dependencies": { | |
| 9813 | − | "@g1t/contracts": "*" | |
| 9814 | − | } | |
| 9815 | 9803 | }, | |
| 9816 | 9804 | "services/runner": { | |
| 9817 | 9805 | "name": "@g1t/runner", |
| 5 | 5 | "workspaces": ["apps/*", "services/*", "packages/*"], | |
| 6 | 6 | "scripts": { | |
| 7 | 7 | "typecheck": "npm run typecheck --workspaces --if-present", | |
| 8 | − | "deploy": "npm run deploy -w @g1t/events -w @g1t/api -w @g1t/web" | |
| 8 | + | "deploy": "npm run deploy -w @g1t/api -w @g1t/web" | |
| 9 | 9 | }, | |
| 10 | 10 | "devDependencies": { | |
| 11 | 11 | "typescript": "^5.9.3", |
| 1 | + | import type { EventsApi } from "./events"; | |
| 1 | 2 | import type { IdentityApi } from "./identity"; | |
| 2 | 3 | import type { ReposApi } from "./repos"; | |
| 3 | 4 | import type { WorkApi } from "./work"; | |
| 127 | 128 | call("read_session", { repo, number, viewer, afterSeq }), | |
| 128 | 129 | }; | |
| 129 | 130 | } | |
| 131 | + | ||
| 132 | + | export function eventsClient(service: ServiceBinding): EventsApi { | |
| 133 | + | const call = <T>(method: string, args: object) => rpc<T>(service, method, args); | |
| 134 | + | return { | |
| 135 | + | publish: (events) => call("publish", { events }), | |
| 136 | + | list: (query) => call("list", query), | |
| 137 | + | }; | |
| 138 | + | } |
| 9 | 9 | export type EventPayloads = { | |
| 10 | 10 | "repo.created": { repoId: string; namespace: string; name: string; isPrivate: boolean }; | |
| 11 | 11 | "repo.forked": { repoId: string; sourceRepoId: string; pullId: string }; | |
| 12 | − | /** `after` is the commit the ref points to once the push has landed. */ | |
| 13 | − | "git.push": { repoId: string; ref: string; after: string }; | |
| 12 | + | /** | |
| 13 | + | * One branch moved by a push. `ref` is the full ref, `after` the commit it | |
| 14 | + | * points to now, and `defaultBranch` whether it is the default branch. | |
| 15 | + | */ | |
| 16 | + | "git.push": { repoId: string; ref: string; after: string; defaultBranch: boolean }; | |
| 14 | 17 | "issue.opened": { issueId: string; repoId: string; number: number; title: string }; | |
| 15 | 18 | "issue.updated": { issueId: string; repoId: string; number: number }; | |
| 16 | 19 | /** `resolvedBy` is the number of the pull request whose merge closed it. */ | |
| 40 | 43 | type: K; | |
| 41 | 44 | /** The service that published it. */ | |
| 42 | 45 | source: string; | |
| 43 | − | /** Milliseconds since the epoch. */ | |
| 44 | − | time: number; | |
| 46 | + | /** RFC 3339. */ | |
| 47 | + | time: string; | |
| 45 | 48 | /** The repo the event concerns, used to scope timelines and deliveries. */ | |
| 46 | 49 | repoId: string | null; | |
| 47 | 50 | /** The user or agent that caused it, if any. */ | |
| 68 | 71 | publish(events: NewEvent[]): Promise<void>; | |
| 69 | 72 | /** Newest first. */ | |
| 70 | 73 | list(query: EventQuery): Promise<G1tEvent[]>; | |
| 71 | − | } | |
| 72 | − | ||
| 73 | − | /** Implemented by services that consume events from the bus. */ | |
| 74 | − | export interface EventSubscriber { | |
| 75 | − | onEvents(events: G1tEvent[]): Promise<void>; | |
| 76 | 74 | } |
| 1 | + | [package] | |
| 2 | + | name = "g1t-events" | |
| 3 | + | version = "0.1.0" | |
| 4 | + | edition.workspace = true | |
| 5 | + | license.workspace = true | |
| 6 | + | description = "The event bus and its log." | |
| 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 |
| 4 | 4 | id TEXT PRIMARY KEY, | |
| 5 | 5 | type TEXT NOT NULL, | |
| 6 | 6 | source TEXT NOT NULL, | |
| 7 | − | time INTEGER NOT NULL, | |
| 7 | + | -- RFC 3339 UTC. | |
| 8 | + | time TEXT NOT NULL, | |
| 8 | 9 | repo_id TEXT, | |
| 9 | 10 | actor TEXT, | |
| 11 | + | -- JSON. | |
| 10 | 12 | data TEXT NOT NULL | |
| 11 | 13 | ); | |
| 12 | 14 | CREATE INDEX events_repo ON events (repo_id, id); |
| 1 | − | { | |
| 2 | − | "name": "@g1t/events", | |
| 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 | − | } |
| 1 | − | import { WorkerEntrypoint } from "cloudflare:workers"; | |
| 2 | − | ||
| 3 | − | import { | |
| 4 | − | type EventQuery, | |
| 5 | − | type EventsApi, | |
| 6 | − | type G1tEvent, | |
| 7 | − | type NewEvent, | |
| 8 | − | newId, | |
| 9 | − | } from "@g1t/contracts"; | |
| 10 | − | ||
| 11 | − | export interface EventsEnv { | |
| 12 | − | DB: D1Database; | |
| 13 | − | BUS: Queue<G1tEvent>; | |
| 14 | − | /** Every `SUBSCRIBER_*` binding is a queue that receives all events. */ | |
| 15 | − | [subscriber: `SUBSCRIBER_${string}`]: Queue<G1tEvent>; | |
| 16 | − | } | |
| 17 | − | ||
| 18 | − | type EventRow = { | |
| 19 | − | id: string; | |
| 20 | − | type: string; | |
| 21 | − | source: string; | |
| 22 | − | time: number; | |
| 23 | − | repo_id: string | null; | |
| 24 | − | actor: string | null; | |
| 25 | − | data: string; | |
| 26 | − | }; | |
| 27 | − | ||
| 28 | − | const DEFAULT_PAGE = 50; | |
| 29 | − | const MAX_PAGE = 200; | |
| 30 | − | ||
| 31 | − | function toEvent(row: EventRow): G1tEvent { | |
| 32 | − | return { | |
| 33 | − | id: row.id, | |
| 34 | − | type: row.type, | |
| 35 | − | source: row.source, | |
| 36 | − | time: row.time, | |
| 37 | − | repoId: row.repo_id, | |
| 38 | − | actor: row.actor, | |
| 39 | − | data: JSON.parse(row.data), | |
| 40 | − | } as G1tEvent; | |
| 41 | − | } | |
| 42 | − | ||
| 43 | − | /** | |
| 44 | − | * The event bus. Publishing enqueues; the queue consumer below writes the | |
| 45 | − | * durable log and fans each batch out to every subscriber's own queue. | |
| 46 | − | */ | |
| 47 | − | export default class EventsService | |
| 48 | − | extends WorkerEntrypoint<EventsEnv> | |
| 49 | − | implements EventsApi | |
| 50 | − | { | |
| 51 | − | async publish(events: NewEvent[]): Promise<void> { | |
| 52 | − | if (events.length === 0) return; | |
| 53 | − | const now = Date.now(); | |
| 54 | − | await this.env.BUS.sendBatch( | |
| 55 | − | events.map((event) => ({ | |
| 56 | − | body: { ...event, id: newId("evt", now), time: now } as G1tEvent, | |
| 57 | − | })), | |
| 58 | − | ); | |
| 59 | − | } | |
| 60 | − | ||
| 61 | − | /** Callers must have checked that the viewer may see `query.repoId`. */ | |
| 62 | − | async list(query: EventQuery): Promise<G1tEvent[]> { | |
| 63 | − | const where: string[] = []; | |
| 64 | − | const params: unknown[] = []; | |
| 65 | − | if (query.repoId) { | |
| 66 | − | where.push("repo_id = ?"); | |
| 67 | − | params.push(query.repoId); | |
| 68 | − | } | |
| 69 | − | if (query.types?.length) { | |
| 70 | − | where.push(`type IN (${query.types.map(() => "?").join(", ")})`); | |
| 71 | − | params.push(...query.types); | |
| 72 | − | } | |
| 73 | − | if (query.before) { | |
| 74 | − | where.push("id < ?"); | |
| 75 | − | params.push(query.before); | |
| 76 | − | } | |
| 77 | − | const limit = Math.min(query.limit ?? DEFAULT_PAGE, MAX_PAGE); | |
| 78 | − | const { results } = await this.env.DB.prepare( | |
| 79 | − | `SELECT * FROM events ${where.length ? `WHERE ${where.join(" AND ")}` : ""} | |
| 80 | − | ORDER BY id DESC LIMIT ?`, | |
| 81 | − | ) | |
| 82 | − | .bind(...params, limit) | |
| 83 | − | .all<EventRow>(); | |
| 84 | − | return results.map(toEvent); | |
| 85 | − | } | |
| 86 | − | ||
| 87 | − | async queue(batch: MessageBatch<G1tEvent>): Promise<void> { | |
| 88 | − | const events = batch.messages.map((message) => message.body); | |
| 89 | − | const insert = this.env.DB.prepare( | |
| 90 | − | // Redelivered batches must not duplicate log rows. | |
| 91 | − | `INSERT OR IGNORE INTO events (id, type, source, time, repo_id, actor, data) | |
| 92 | − | VALUES (?, ?, ?, ?, ?, ?, ?)`, | |
| 93 | − | ); | |
| 94 | − | await this.env.DB.batch( | |
| 95 | − | events.map((event) => | |
| 96 | − | insert.bind( | |
| 97 | − | event.id, | |
| 98 | − | event.type, | |
| 99 | − | event.source, | |
| 100 | − | event.time, | |
| 101 | − | event.repoId, | |
| 102 | − | event.actor, | |
| 103 | − | JSON.stringify(event.data), | |
| 104 | − | ), | |
| 105 | − | ), | |
| 106 | − | ); | |
| 107 | − | ||
| 108 | − | const subscribers = Object.entries(this.env) | |
| 109 | − | .filter(([name]) => name.startsWith("SUBSCRIBER_")) | |
| 110 | − | .map(([, queue]) => queue as Queue<G1tEvent>); | |
| 111 | − | await Promise.all( | |
| 112 | − | subscribers.map((queue) => | |
| 113 | − | queue.sendBatch(events.map((event) => ({ body: event }))), | |
| 114 | − | ), | |
| 115 | − | ); | |
| 116 | − | } | |
| 117 | − | } |
| 1 | + | //! The events service: the bus every state change in g1t is published on, | |
| 2 | + | //! and its durable log. | |
| 3 | + | //! | |
| 4 | + | //! Publishing puts events on a queue and returns. The queue consumer writes | |
| 5 | + | //! them to the log and passes each batch on to every subscriber's own | |
| 6 | + | //! queue, so a slow or failing subscriber holds up nobody else. | |
| 7 | + | //! | |
| 8 | + | //! Other services reach it over `POST /rpc/<method>`; see | |
| 9 | + | //! `g1t_contracts::events` for the methods and their arguments. | |
| 10 | + | ||
| 11 | + | use g1t_contracts::events::{Event, ListArgs, PublishArgs}; | |
| 12 | + | use g1t_contracts::new_id; | |
| 13 | + | use g1t_contracts::time::rfc3339; | |
| 14 | + | use g1t_kit::{args, js, now_ms, reply, rpc_method}; | |
| 15 | + | use serde::Deserialize; | |
| 16 | + | use worker::js_sys::{Array, Object}; | |
| 17 | + | use worker::wasm_bindgen::{JsCast, JsValue}; | |
| 18 | + | use worker::{Context, D1Database, Env, MessageBatch, Request, Response, Result, event}; | |
| 19 | + | ||
| 20 | + | const DEFAULT_PAGE: u32 = 50; | |
| 21 | + | const MAX_PAGE: u32 = 200; | |
| 22 | + | /// Every binding whose name starts with this is a queue that receives all | |
| 23 | + | /// events: one per subscribing service. | |
| 24 | + | const SUBSCRIBER_PREFIX: &str = "SUBSCRIBER_"; | |
| 25 | + | ||
| 26 | + | #[derive(Deserialize)] | |
| 27 | + | struct EventRow { | |
| 28 | + | id: String, | |
| 29 | + | #[serde(rename = "type")] | |
| 30 | + | kind: String, | |
| 31 | + | source: String, | |
| 32 | + | time: String, | |
| 33 | + | repo_id: Option<String>, | |
| 34 | + | actor: Option<String>, | |
| 35 | + | /// JSON. | |
| 36 | + | data: String, | |
| 37 | + | } | |
| 38 | + | ||
| 39 | + | impl From<EventRow> for Event { | |
| 40 | + | fn from(row: EventRow) -> Self { | |
| 41 | + | Event { | |
| 42 | + | id: row.id, | |
| 43 | + | kind: row.kind, | |
| 44 | + | source: row.source, | |
| 45 | + | time: row.time, | |
| 46 | + | repo_id: row.repo_id, | |
| 47 | + | actor: row.actor, | |
| 48 | + | data: serde_json::from_str(&row.data).unwrap_or_default(), | |
| 49 | + | } | |
| 50 | + | } | |
| 51 | + | } | |
| 52 | + | ||
| 53 | + | fn optional(value: &Option<String>) -> JsValue { | |
| 54 | + | value.as_deref().map_or(JsValue::NULL, JsValue::from) | |
| 55 | + | } | |
| 56 | + | ||
| 57 | + | /// Sends `events` to a queue binding, one message each. | |
| 58 | + | async fn send(queue: &JsValue, events: &[Event]) -> Result<()> { | |
| 59 | + | let messages = Array::new(); | |
| 60 | + | for event in events { | |
| 61 | + | let message = Object::new(); | |
| 62 | + | js::set(&message, "body", &js::to_js(event)?); | |
| 63 | + | messages.push(&message); | |
| 64 | + | } | |
| 65 | + | js::call(queue, "sendBatch", &[messages.into()]).await?; | |
| 66 | + | Ok(()) | |
| 67 | + | } | |
| 68 | + | ||
| 69 | + | struct Events { | |
| 70 | + | db: D1Database, | |
| 71 | + | env: Env, | |
| 72 | + | } | |
| 73 | + | ||
| 74 | + | impl Events { | |
| 75 | + | /// Assigns each event its id and time and puts it on the bus. | |
| 76 | + | async fn publish(&self, a: PublishArgs) -> Result<()> { | |
| 77 | + | if a.events.is_empty() { | |
| 78 | + | return Ok(()); | |
| 79 | + | } | |
| 80 | + | let now = now_ms(); | |
| 81 | + | let events: Vec<Event> = a | |
| 82 | + | .events | |
| 83 | + | .into_iter() | |
| 84 | + | .map(|event| Event { | |
| 85 | + | id: new_id("evt", now), | |
| 86 | + | kind: event.kind, | |
| 87 | + | source: event.source, | |
| 88 | + | time: rfc3339(now), | |
| 89 | + | repo_id: event.repo_id, | |
| 90 | + | actor: event.actor, | |
| 91 | + | data: event.data, | |
| 92 | + | }) | |
| 93 | + | .collect(); | |
| 94 | + | send(&js::binding(&self.env, "BUS")?, &events).await | |
| 95 | + | } | |
| 96 | + | ||
| 97 | + | /// Newest first. Callers must have checked that the viewer may see the | |
| 98 | + | /// repository asked about. | |
| 99 | + | async fn list(&self, a: ListArgs) -> Result<Vec<Event>> { | |
| 100 | + | let mut conditions = Vec::new(); | |
| 101 | + | let mut values: Vec<JsValue> = Vec::new(); | |
| 102 | + | if let Some(repo_id) = &a.repo_id { | |
| 103 | + | conditions.push("repo_id = ?".to_owned()); | |
| 104 | + | values.push(repo_id.as_str().into()); | |
| 105 | + | } | |
| 106 | + | if !a.types.is_empty() { | |
| 107 | + | let marks = vec!["?"; a.types.len()].join(", "); | |
| 108 | + | conditions.push(format!("type IN ({marks})")); | |
| 109 | + | values.extend(a.types.iter().map(|kind| JsValue::from(kind.as_str()))); | |
| 110 | + | } | |
| 111 | + | if let Some(before) = &a.before { | |
| 112 | + | conditions.push("id < ?".to_owned()); | |
| 113 | + | values.push(before.as_str().into()); | |
| 114 | + | } | |
| 115 | + | let filter = if conditions.is_empty() { | |
| 116 | + | String::new() | |
| 117 | + | } else { | |
| 118 | + | format!("WHERE {}", conditions.join(" AND ")) | |
| 119 | + | }; | |
| 120 | + | values.push(a.limit.unwrap_or(DEFAULT_PAGE).min(MAX_PAGE).into()); | |
| 121 | + | let rows = self | |
| 122 | + | .db | |
| 123 | + | .prepare(format!( | |
| 124 | + | "SELECT * FROM events {filter} ORDER BY id DESC LIMIT ?" | |
| 125 | + | )) | |
| 126 | + | .bind(&values)? | |
| 127 | + | .all() | |
| 128 | + | .await? | |
| 129 | + | .results::<EventRow>()?; | |
| 130 | + | Ok(rows.into_iter().map(Event::from).collect()) | |
| 131 | + | } | |
| 132 | + | ||
| 133 | + | /// Writes a batch from the bus to the log, then hands it to every | |
| 134 | + | /// subscriber. | |
| 135 | + | async fn deliver(&self, events: &[Event]) -> Result<()> { | |
| 136 | + | let mut statements = Vec::with_capacity(events.len()); | |
| 137 | + | for event in events { | |
| 138 | + | statements.push( | |
| 139 | + | self.db | |
| 140 | + | .prepare( | |
| 141 | + | // Redelivered batches must not duplicate log rows. | |
| 142 | + | "INSERT OR IGNORE INTO events (id, type, source, time, repo_id, actor, data) | |
| 143 | + | VALUES (?, ?, ?, ?, ?, ?, ?)", | |
| 144 | + | ) | |
| 145 | + | .bind(&[ | |
| 146 | + | event.id.as_str().into(), | |
| 147 | + | event.kind.as_str().into(), | |
| 148 | + | event.source.as_str().into(), | |
| 149 | + | event.time.as_str().into(), | |
| 150 | + | optional(&event.repo_id), | |
| 151 | + | optional(&event.actor), | |
| 152 | + | serde_json::to_string(&event.data)?.into(), | |
| 153 | + | ])?, | |
| 154 | + | ); | |
| 155 | + | } | |
| 156 | + | self.db.batch(statements).await?; | |
| 157 | + | ||
| 158 | + | let bindings: &JsValue = self.env.as_ref(); | |
| 159 | + | for name in Object::keys(bindings.unchecked_ref::<Object>()).iter() { | |
| 160 | + | let Some(name) = name | |
| 161 | + | .as_string() | |
| 162 | + | .filter(|name| name.starts_with(SUBSCRIBER_PREFIX)) | |
| 163 | + | else { | |
| 164 | + | continue; | |
| 165 | + | }; | |
| 166 | + | send(&js::binding(&self.env, &name)?, events).await?; | |
| 167 | + | } | |
| 168 | + | Ok(()) | |
| 169 | + | } | |
| 170 | + | } | |
| 171 | + | ||
| 172 | + | #[event(fetch)] | |
| 173 | + | async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> { | |
| 174 | + | let Some(method) = rpc_method(&request) else { | |
| 175 | + | return Response::error("Not found", 404); | |
| 176 | + | }; | |
| 177 | + | let body: serde_json::Value = request.json().await?; | |
| 178 | + | let events = Events { | |
| 179 | + | db: env.d1("DB")?, | |
| 180 | + | env, | |
| 181 | + | }; | |
| 182 | + | match method.as_str() { | |
| 183 | + | "publish" => reply(&events.publish(args(body)?).await?), | |
| 184 | + | "list" => reply(&events.list(args(body)?).await?), | |
| 185 | + | _ => Response::error("Unknown method", 404), | |
| 186 | + | } | |
| 187 | + | } | |
| 188 | + | ||
| 189 | + | /// Events from the bus. A batch that fails is retried whole, which is safe | |
| 190 | + | /// because both the log and subscribers ignore an event they have seen. | |
| 191 | + | #[event(queue)] | |
| 192 | + | async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> { | |
| 193 | + | let events = Events { | |
| 194 | + | db: env.d1("DB")?, | |
| 195 | + | env, | |
| 196 | + | }; | |
| 197 | + | let delivered: Vec<Event> = batch | |
| 198 | + | .messages()? | |
| 199 | + | .into_iter() | |
| 200 | + | .map(|message| message.into_body()) | |
| 201 | + | .collect(); | |
| 202 | + | events.deliver(&delivered).await?; | |
| 203 | + | batch.ack_all(); | |
| 204 | + | Ok(()) | |
| 205 | + | } |
| 1 | − | { | |
| 2 | − | "extends": "../../tsconfig.base.json", | |
| 3 | − | "include": ["src/**/*", "worker-configuration.d.ts"] | |
| 4 | − | } |
| 3 | 3 | "name": "g1t-events", | |
| 4 | 4 | "account_id": "1e6f2cffa3f445920836e8ebe446bb58", | |
| 5 | 5 | "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" }, | |
| 7 | 8 | "workers_dev": false, | |
| 8 | 9 | "d1_databases": [ | |
| 9 | 10 | { |
| 125 | 125 | Ok(response) | |
| 126 | 126 | } | |
| 127 | 127 | ||
| 128 | + | const ZERO_ID: &str = "0000000000000000000000000000000000000000"; | |
| 129 | + | const HEADS: &str = "refs/heads/"; | |
| 130 | + | ||
| 131 | + | /// The branches a push asks to move, as `(branch, new commit)`, read from | |
| 132 | + | /// the commands at the start of a receive-pack request. Deletions and refs | |
| 133 | + | /// that are not branches are left out. | |
| 134 | + | fn pushed_branches(body: &[u8]) -> Vec<(String, String)> { | |
| 135 | + | let mut branches = Vec::new(); | |
| 136 | + | let mut position = 0; | |
| 137 | + | // Commands are pkt-lines; a flush packet ends them and the pack follows. | |
| 138 | + | while let Some(length) = body | |
| 139 | + | .get(position..position + 4) | |
| 140 | + | .and_then(|hex| std::str::from_utf8(hex).ok()) | |
| 141 | + | .and_then(|hex| usize::from_str_radix(hex, 16).ok()) | |
| 142 | + | { | |
| 143 | + | if length < 4 || position + length > body.len() { | |
| 144 | + | break; | |
| 145 | + | } | |
| 146 | + | let line = &body[position + 4..position + length]; | |
| 147 | + | position += length; | |
| 148 | + | // `<old> <new> <ref>`, and on the first command a NUL then capabilities. | |
| 149 | + | let line = line.split(|byte| *byte == 0).next().unwrap_or_default(); | |
| 150 | + | let Ok(line) = std::str::from_utf8(line) else { | |
| 151 | + | continue; | |
| 152 | + | }; | |
| 153 | + | let mut parts = line.trim_end().splitn(3, ' '); | |
| 154 | + | let (Some(_old), Some(new), Some(name)) = (parts.next(), parts.next(), parts.next()) else { | |
| 155 | + | continue; | |
| 156 | + | }; | |
| 157 | + | if let Some(branch) = name.strip_prefix(HEADS) | |
| 158 | + | && new != ZERO_ID | |
| 159 | + | { | |
| 160 | + | branches.push((branch.to_owned(), new.to_owned())); | |
| 161 | + | } | |
| 162 | + | } | |
| 163 | + | branches | |
| 164 | + | } | |
| 165 | + | ||
| 166 | + | /// The git store's answer, and what the request asked it to change. | |
| 167 | + | pub struct Forwarded { | |
| 168 | + | pub response: Response, | |
| 169 | + | /// For a push: the branches it asks to move and the commits to move | |
| 170 | + | /// them to. Whether each moved is for the caller to confirm. | |
| 171 | + | pub pushed: Vec<(String, String)>, | |
| 172 | + | } | |
| 173 | + | ||
| 128 | 174 | /// Sends the request on to the git store and returns its response as is. | |
| 129 | 175 | pub async fn forward( | |
| 130 | 176 | mut request: Request, | |
| 131 | 177 | git: &GitRequest, | |
| 132 | 178 | access: &GitAccess, | |
| 133 | − | ) -> Result<Response> { | |
| 179 | + | ) -> Result<Forwarded> { | |
| 134 | 180 | let headers = Headers::new(); | |
| 135 | 181 | headers.set("authorization", &format!("Bearer {}", access.token))?; | |
| 136 | 182 | for name in FORWARDED_HEADERS { | |
| 145 | 191 | .unwrap_or_default(); | |
| 146 | 192 | let mut init = RequestInit::new(); | |
| 147 | 193 | init.with_method(request.method()).with_headers(headers); | |
| 194 | + | let mut pushed = Vec::new(); | |
| 148 | 195 | if request.method() == Method::Post { | |
| 149 | 196 | // Pushes are capped at 100 MB by the platform, so buffering is safe. | |
| 150 | 197 | let body = request.bytes().await?; | |
| 198 | + | if git.endpoint == "git-receive-pack" { | |
| 199 | + | pushed = pushed_branches(&body); | |
| 200 | + | } | |
| 151 | 201 | init.with_body(Some(Uint8Array::from(body.as_slice()).into())); | |
| 152 | 202 | } | |
| 153 | 203 | let upstream = | |
| 154 | 204 | Request::new_with_init(&format!("{}/{}{query}", access.remote, git.endpoint), &init)?; | |
| 155 | − | Fetch::Request(upstream).send().await | |
| 205 | + | Ok(Forwarded { | |
| 206 | + | response: Fetch::Request(upstream).send().await?, | |
| 207 | + | pushed, | |
| 208 | + | }) | |
| 209 | + | } | |
| 210 | + | ||
| 211 | + | #[cfg(test)] | |
| 212 | + | mod tests { | |
| 213 | + | use super::{ZERO_ID, pushed_branches}; | |
| 214 | + | ||
| 215 | + | fn pkt(payload: &str) -> Vec<u8> { | |
| 216 | + | format!("{:04x}{payload}", payload.len() + 4).into_bytes() | |
| 217 | + | } | |
| 218 | + | ||
| 219 | + | #[test] | |
| 220 | + | fn pushed_branches_are_read_from_the_commands() { | |
| 221 | + | let old = "c71546fcd893ef8b0f57388b65e620d759705dda"; | |
| 222 | + | let new = "4807077b296e6edbf410d55e72749d3e1170c291"; | |
| 223 | + | let body = [ | |
| 224 | + | pkt(&format!( | |
| 225 | + | "{old} {new} refs/heads/main\0 report-status side-band-64k\n" | |
| 226 | + | )), | |
| 227 | + | pkt(&format!("{ZERO_ID} {new} refs/heads/feature/x\n")), | |
| 228 | + | pkt(&format!("{old} {ZERO_ID} refs/heads/gone\n")), | |
| 229 | + | pkt(&format!("{ZERO_ID} {new} refs/tags/v1\n")), | |
| 230 | + | b"0000".to_vec(), | |
| 231 | + | b"PACK\0\0\0\x02\0\0\0\0".to_vec(), | |
| 232 | + | ] | |
| 233 | + | .concat(); | |
| 234 | + | assert_eq!( | |
| 235 | + | pushed_branches(&body), | |
| 236 | + | [ | |
| 237 | + | ("main".to_owned(), new.to_owned()), | |
| 238 | + | ("feature/x".to_owned(), new.to_owned()), | |
| 239 | + | ] | |
| 240 | + | ); | |
| 241 | + | } | |
| 242 | + | ||
| 243 | + | #[test] | |
| 244 | + | fn a_fetch_request_names_no_branches() { | |
| 245 | + | assert!( | |
| 246 | + | pushed_branches(b"0032want c71546fcd893ef8b0f57388b65e620d759705dda\n0000").is_empty() | |
| 247 | + | ); | |
| 248 | + | } | |
| 156 | 249 | } |
| 12 | 12 | mod registry; | |
| 13 | 13 | mod store; | |
| 14 | 14 | ||
| 15 | − | use g1t_contracts::events::{GitPush, NewEvent, RepoCreated, RepoForked}; | |
| 15 | + | use g1t_contracts::events::{GitPush, NewEvent, Publish, RepoCreated, RepoForked}; | |
| 16 | 16 | use g1t_contracts::repos::*; | |
| 17 | 17 | use g1t_contracts::time::rfc3339; | |
| 18 | 18 | use g1t_contracts::{FailureCode, Outcome, User, Viewer, is_valid_repo_name, new_id}; | |
| 19 | − | use g1t_kit::{args, js, now_ms, reply, rpc_method}; | |
| 19 | + | use g1t_kit::{args, now_ms, reply, rpc_method}; | |
| 20 | 20 | use std::collections::{HashMap, HashSet, VecDeque}; | |
| 21 | 21 | ||
| 22 | 22 | use serde::Serialize; | |
| 23 | − | use worker::wasm_bindgen::JsValue; | |
| 24 | − | use worker::{Context, Env, Request, Response, Result, event}; | |
| 23 | + | use worker::{Context, Env, Fetcher, Request, Response, Result, event}; | |
| 25 | 24 | ||
| 26 | 25 | use registry::{Registry, can_read, can_write, store_key}; | |
| 27 | 26 | use store::{ArtifactsStore, GitRepo, GitStore, Scope}; | |
| 119 | 118 | struct Repos<S: GitStore> { | |
| 120 | 119 | registry: Registry, | |
| 121 | 120 | store: S, | |
| 122 | − | /// The events service, an RPC stub. | |
| 123 | − | events: JsValue, | |
| 121 | + | events: Fetcher, | |
| 124 | 122 | } | |
| 125 | 123 | ||
| 126 | 124 | impl<S: GitStore> Repos<S> { | |
| 127 | 125 | async fn publish<T: Serialize>(&self, event: NewEvent<T>) -> Result<()> { | |
| 128 | − | js::call(&self.events, "publish", &[js::to_js(&[event])?]).await?; | |
| 129 | − | Ok(()) | |
| 126 | + | g1t_kit::call( | |
| 127 | + | &self.events, | |
| 128 | + | "publish", | |
| 129 | + | &Publish { | |
| 130 | + | events: vec![event], | |
| 131 | + | }, | |
| 132 | + | ) | |
| 133 | + | .await | |
| 130 | 134 | } | |
| 131 | 135 | ||
| 132 | 136 | /// Resolves a repo the viewer may read; private repos look missing. | |
| 542 | 546 | format!("{branch} could not be updated: {reason}"), | |
| 543 | 547 | )); | |
| 544 | 548 | } | |
| 545 | − | self.publish_push(&target, &new, Some(a.actor.id)).await?; | |
| 549 | + | self.publish_push(&target, branch, &new, Some(a.actor.id)) | |
| 550 | + | .await?; | |
| 546 | 551 | Ok(Outcome::Ok(Landed { | |
| 547 | 552 | commit: new, | |
| 548 | 553 | previous: old, | |
| 610 | 615 | })) | |
| 611 | 616 | } | |
| 612 | 617 | ||
| 613 | − | async fn publish_push(&self, repo: &Repo, after: &str, actor: Option<String>) -> Result<()> { | |
| 618 | + | /// Reports that `branch` of `repo` now points to `after`. | |
| 619 | + | async fn publish_push( | |
| 620 | + | &self, | |
| 621 | + | repo: &Repo, | |
| 622 | + | branch: &str, | |
| 623 | + | after: &str, | |
| 624 | + | actor: Option<String>, | |
| 625 | + | ) -> Result<()> { | |
| 614 | 626 | self.publish(NewEvent { | |
| 615 | 627 | kind: "git.push", | |
| 616 | 628 | source: SOURCE, | |
| 618 | 630 | actor, | |
| 619 | 631 | data: GitPush { | |
| 620 | 632 | repo_id: repo.id.clone(), | |
| 621 | − | git_ref: format!("refs/heads/{}", repo.default_branch), | |
| 633 | + | git_ref: format!("refs/heads/{branch}"), | |
| 622 | 634 | after: after.to_owned(), | |
| 635 | + | default_branch: branch == repo.default_branch, | |
| 623 | 636 | }, | |
| 624 | 637 | }) | |
| 625 | 638 | .await | |
| 642 | 655 | Outcome::Ok(access) => access, | |
| 643 | 656 | refused => return git_http::refuse(refused), | |
| 644 | 657 | }; | |
| 645 | − | let response = git_http::forward(request, &git, &access).await?; | |
| 658 | + | let forwarded = git_http::forward(request, &git, &access).await?; | |
| 646 | 659 | ||
| 647 | 660 | // Artifacts' own push notifications are per repository, which does | |
| 648 | − | // not fit a repo per pull request, so the front end reports pushes itself. | |
| 649 | − | let pushed = git.endpoint == "git-receive-pack" && response.status_code() == 200; | |
| 650 | − | if pushed && let Some(repo) = self.registry.by_path(&git.path).await? { | |
| 651 | − | let head = self | |
| 652 | − | .store | |
| 653 | − | .open(&store_key(&repo)) | |
| 654 | − | .await? | |
| 655 | − | .log(&repo.default_branch, 1) | |
| 656 | − | .await?; | |
| 657 | − | if let Some(head) = head.first() { | |
| 658 | − | self.publish_push(&repo, &head.hash, viewer.map(|user: User| user.id)) | |
| 659 | − | .await?; | |
| 661 | + | // not fit a repo per pull request, so the front end reports pushes | |
| 662 | + | // itself: one event for each branch that moved. | |
| 663 | + | let accepted = forwarded.response.status_code() == 200 && !forwarded.pushed.is_empty(); | |
| 664 | + | if accepted && let Some(repo) = self.registry.by_path(&git.path).await? { | |
| 665 | + | let stored = self.store.open(&store_key(&repo)).await?; | |
| 666 | + | let actor = viewer.map(|user: User| user.id); | |
| 667 | + | for (branch, pushed) in &forwarded.pushed { | |
| 668 | + | // The store can refuse one ref and accept another, so each | |
| 669 | + | // is checked against where the branch actually is. | |
| 670 | + | let head = stored.log(branch, 1).await?; | |
| 671 | + | if head.first().is_some_and(|commit| commit.hash == *pushed) { | |
| 672 | + | self.publish_push(&repo, branch, pushed, actor.clone()) | |
| 673 | + | .await?; | |
| 674 | + | } | |
| 660 | 675 | } | |
| 661 | 676 | } | |
| 662 | − | Ok(response) | |
| 677 | + | Ok(forwarded.response) | |
| 663 | 678 | } | |
| 664 | 679 | } | |
| 665 | 680 | ||
| 668 | 683 | let repos = Repos { | |
| 669 | 684 | registry: Registry { db: env.d1("DB")? }, | |
| 670 | 685 | store: ArtifactsStore::new(&env)?, | |
| 671 | − | events: js::binding(&env, "EVENTS")?, | |
| 686 | + | events: env.service("EVENTS")?, | |
| 672 | 687 | }; | |
| 673 | 688 | let Some(method) = rpc_method(&request) else { | |
| 674 | 689 | return repos.git_http(request, &env).await; |
| 7 | 7 | mod rows; | |
| 8 | 8 | ||
| 9 | 9 | use g1t_contracts::events::{ | |
| 10 | − | CommentCreated, Delivered, IssueEvent, NewEvent, PullEvent, SessionAppended, | |
| 10 | + | CommentCreated, Event, IssueEvent, NewEvent, Publish, PullEvent, SessionAppended, | |
| 11 | 11 | }; | |
| 12 | 12 | use g1t_contracts::repos::{ForkArgs, GetArgs, HeadArgs, LandArgs, Landed, Repo, RepoPath}; | |
| 13 | 13 | use g1t_contracts::time::rfc3339; | |
| 14 | 14 | use g1t_contracts::work::*; | |
| 15 | 15 | use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id}; | |
| 16 | − | use g1t_kit::{args, js, now_ms, reply, rpc_method}; | |
| 16 | + | use g1t_kit::{args, now_ms, reply, rpc_method}; | |
| 17 | 17 | use serde::Serialize; | |
| 18 | 18 | use worker::wasm_bindgen::JsValue; | |
| 19 | 19 | use worker::{ | |
| 20 | 20 | Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event, | |
| 21 | 21 | }; | |
| 22 | 22 | ||
| 23 | − | use rows::{BranchRow, CommentRow, IssueRow, NumberRow, PullRow, SessionRow, ValueRow}; | |
| 23 | + | use rows::{CommentRow, IssueRow, NumberRow, PullRow, SessionRow, ValueRow}; | |
| 24 | 24 | ||
| 25 | 25 | const SOURCE: &str = "work"; | |
| 26 | 26 | const MAX_ENTRY_BATCH: usize = 200; | |
| 84 | 84 | struct Work { | |
| 85 | 85 | db: D1Database, | |
| 86 | 86 | repos: Fetcher, | |
| 87 | − | /// The events service, an RPC stub. | |
| 88 | − | events: JsValue, | |
| 87 | + | events: Fetcher, | |
| 89 | 88 | } | |
| 90 | 89 | ||
| 91 | 90 | impl Work { | |
| 103 | 102 | actor: Some(actor.id.clone()), | |
| 104 | 103 | data, | |
| 105 | 104 | }; | |
| 106 | − | js::call(&self.events, "publish", &[js::to_js(&[event])?]).await?; | |
| 107 | − | Ok(()) | |
| 105 | + | g1t_kit::call( | |
| 106 | + | &self.events, | |
| 107 | + | "publish", | |
| 108 | + | &Publish { | |
| 109 | + | events: vec![event], | |
| 110 | + | }, | |
| 111 | + | ) | |
| 112 | + | .await | |
| 108 | 113 | } | |
| 109 | 114 | ||
| 110 | 115 | /// The repository, if the viewer may see it. Whether they may is | |
| 1098 | 1103 | )) | |
| 1099 | 1104 | } | |
| 1100 | 1105 | ||
| 1101 | − | /// A push moves the head of the pull requests it concerns: the one | |
| 1102 | − | /// whose fork was pushed to, or those from branches of the repository | |
| 1103 | − | /// that was. | |
| 1104 | − | async fn on_event(&self, event: &Delivered) -> Result<()> { | |
| 1106 | + | /// A push moves the head of the pull request it concerns: the one whose | |
| 1107 | + | /// fork was pushed to, or the one opened from the branch that moved. | |
| 1108 | + | async fn on_event(&self, event: &Event) -> Result<()> { | |
| 1105 | 1109 | if event.kind != "git.push" { | |
| 1106 | 1110 | return Ok(()); | |
| 1107 | 1111 | } | |
| 1108 | − | let (Some(repo_id), Some(after)) = (event.repo_id.as_deref(), event.data["after"].as_str()) | |
| 1109 | − | else { | |
| 1112 | + | let (Some(repo_id), Some(after), Some(git_ref)) = ( | |
| 1113 | + | event.repo_id.as_deref(), | |
| 1114 | + | event.data["after"].as_str(), | |
| 1115 | + | event.data["ref"].as_str(), | |
| 1116 | + | ) else { | |
| 1110 | 1117 | return Ok(()); | |
| 1111 | 1118 | }; | |
| 1112 | 1119 | let now = rfc3339(now_ms()); | |
| 1113 | − | self.db | |
| 1114 | − | .prepare( | |
| 1115 | − | "UPDATE pulls SET head_commit = ?, updated_at = ? | |
| 1116 | − | WHERE fork_repo_id = ? AND status IN ('draft', 'open')", | |
| 1117 | − | ) | |
| 1118 | − | .bind(&[after.into(), now.as_str().into(), repo_id.into()])? | |
| 1119 | − | .run() | |
| 1120 | − | .await?; | |
| 1121 | − | ||
| 1122 | − | // The event does not say which branch moved, so each open pull | |
| 1123 | − | // request from a branch of this repository is checked. | |
| 1124 | − | let from_branches = self | |
| 1125 | − | .db | |
| 1126 | − | .prepare( | |
| 1127 | − | "SELECT id, source_branch, head_commit FROM pulls | |
| 1128 | − | WHERE repo_id = ? AND source_branch IS NOT NULL AND status IN ('draft', 'open')", | |
| 1129 | − | ) | |
| 1130 | − | .bind(&[repo_id.into()])? | |
| 1131 | − | .all() | |
| 1132 | − | .await? | |
| 1133 | − | .results::<BranchRow>()?; | |
| 1134 | − | for row in from_branches { | |
| 1135 | − | let head: Option<String> = g1t_kit::call( | |
| 1136 | − | &self.repos, | |
| 1137 | − | "head", | |
| 1138 | − | &HeadArgs { | |
| 1139 | − | repo_id: repo_id.to_owned(), | |
| 1140 | − | branch: row.source_branch, | |
| 1141 | − | }, | |
| 1142 | − | ) | |
| 1143 | − | .await?; | |
| 1144 | − | if let Some(head) = head.filter(|head| Some(head) != row.head_commit.as_ref()) { | |
| 1120 | + | let active = "status IN ('draft', 'open')"; | |
| 1121 | + | let mut statements = Vec::new(); | |
| 1122 | + | // A fork carries its pull request on its default branch. | |
| 1123 | + | if event.data["defaultBranch"].as_bool() == Some(true) { | |
| 1124 | + | statements.push( | |
| 1125 | + | self.db | |
| 1126 | + | .prepare(format!( | |
| 1127 | + | "UPDATE pulls SET head_commit = ?, updated_at = ? | |
| 1128 | + | WHERE fork_repo_id = ? AND {active}" | |
| 1129 | + | )) | |
| 1130 | + | .bind(&[after.into(), now.as_str().into(), repo_id.into()])?, | |
| 1131 | + | ); | |
| 1132 | + | } | |
| 1133 | + | if let Some(branch) = git_ref.strip_prefix("refs/heads/") { | |
| 1134 | + | statements.push( | |
| 1145 | 1135 | self.db | |
| 1146 | − | .prepare("UPDATE pulls SET head_commit = ?, updated_at = ? WHERE id = ?") | |
| 1147 | − | .bind(&[head.into(), now.as_str().into(), row.id.into()])? | |
| 1148 | − | .run() | |
| 1149 | − | .await?; | |
| 1150 | − | } | |
| 1136 | + | .prepare(format!( | |
| 1137 | + | "UPDATE pulls SET head_commit = ?, updated_at = ? | |
| 1138 | + | WHERE repo_id = ? AND source_branch = ? AND {active}" | |
| 1139 | + | )) | |
| 1140 | + | .bind(&[ | |
| 1141 | + | after.into(), | |
| 1142 | + | now.as_str().into(), | |
| 1143 | + | repo_id.into(), | |
| 1144 | + | branch.into(), | |
| 1145 | + | ])?, | |
| 1146 | + | ); | |
| 1151 | 1147 | } | |
| 1148 | + | self.db.batch(statements).await?; | |
| 1152 | 1149 | Ok(()) | |
| 1153 | 1150 | } | |
| 1154 | 1151 | } | |
| 1157 | 1154 | Ok(Work { | |
| 1158 | 1155 | db: env.d1("DB")?, | |
| 1159 | 1156 | repos: env.service("REPOS")?, | |
| 1160 | − | events: js::binding(env, "EVENTS")?, | |
| 1157 | + | events: env.service("EVENTS")?, | |
| 1161 | 1158 | }) | |
| 1162 | 1159 | } | |
| 1163 | 1160 | ||
| 1194 | 1191 | ||
| 1195 | 1192 | /// Events from the bus, delivered on this service's own queue. | |
| 1196 | 1193 | #[event(queue)] | |
| 1197 | − | async fn queue(batch: MessageBatch<Delivered>, env: Env, _ctx: Context) -> Result<()> { | |
| 1194 | + | async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> { | |
| 1198 | 1195 | let work = service(&env)?; | |
| 1199 | 1196 | for message in batch.messages()? { | |
| 1200 | 1197 | work.on_event(message.body()).await?; |
| 160 | 160 | } | |
| 161 | 161 | } | |
| 162 | 162 | ||
| 163 | − | /// An open pull request made from a branch of its repository. | |
| 164 | − | #[derive(Deserialize)] | |
| 165 | − | pub struct BranchRow { | |
| 166 | − | pub id: String, | |
| 167 | − | pub source_branch: String, | |
| 168 | − | pub head_commit: Option<String>, | |
| 169 | − | } | |
| 170 | − | ||
| 171 | 163 | /// A single number selected as `n`. | |
| 172 | 164 | #[derive(Deserialize)] | |
| 173 | 165 | pub struct NumberRow { |