Commit

Work service in Rust, with RFC 3339 timestamps

- services/work rewritten in Rust: intents, attempts, sessions, shipping, and the queue consumer that moves an attempt's head on a push - work database migrated to text timestamps - site, API and runner call it through a typed client - identity migration corrected so rebuilding tables cannot cascade-delete rows that reference the old users table

syntaqxcommitted Parent82eeeb9Browse files
30 files+1312−6670/30 viewed
+11−0
889889 ]
890890
891891 [[package]]
892+name = "g1t-work"
893+version = "0.1.0"
894+dependencies = [
895+ "g1t-contracts",
896+ "g1t-kit",
897+ "serde",
898+ "serde_json",
899+ "worker",
900+]
901+
902+[[package]]
892903 name = "generic-array"
893904 version = "0.14.7"
894905 source = "registry+https://github.com/rust-lang/crates.io-index"
+1−1
11 [workspace]
22 resolver = "3"
3−members = ["crates/*", "services/identity", "services/repos"]
3+members = ["crates/*", "services/identity", "services/repos", "services/work"]
44
55 [workspace.package]
66 edition = "2024"
+4−1
77 httpStatus,
88 identityClient,
99 reposClient,
10+ workClient,
1011 } from "@g1t/contracts";
1112
1213 import { handleMcp } from "./mcp";
1516
1617 type Input = Record<string, unknown>;
1718 /** The Worker's raw bindings; Rust services are reached through clients. */
18−type Bindings = Omit<ApiEnv, "IDENTITY" | "REPOS"> & {
19+type Bindings = Omit<ApiEnv, "IDENTITY" | "REPOS" | "WORK"> & {
1920 IDENTITY: ServiceBinding;
2021 REPOS: ServiceBinding;
22+ WORK: ServiceBinding;
2123 };
2224 type App = { Bindings: Bindings; Variables: { viewer: Viewer; services: ApiEnv } };
2325
7678 ...c.env,
7779 IDENTITY: identityClient(c.env.IDENTITY),
7880 REPOS: reposClient(c.env.REPOS),
81+ WORK: workClient(c.env.WORK),
7982 };
8083 c.set("services", services);
8184 let viewer: Viewer = null;
+2−1
11 import { env } from "cloudflare:workers";
22
3−import { identityClient, reposClient } from "@g1t/contracts";
3+import { identityClient, reposClient, workClient } from "@g1t/contracts";
44
55 export const identity = identityClient(env.IDENTITY);
66 export const repos = reposClient(env.REPOS);
7+export const work = workClient(env.WORK);
+2−3
1−import { env } from "cloudflare:workers";
21 import { ArrowRight, Plus } from "lucide-react";
32 import { Link } from "react-router";
43
1312 Status,
1413 TimeAgo,
1514 } from "../components/ui";
16−import { repos as reposApi } from "../lib/services.server";
15+import { repos as reposApi, work } from "../lib/services.server";
1716 import { getViewer } from "../lib/session.server";
1817
1918 export function meta({}: Route.MetaArgs) {
3130 const viewer = getViewer(context);
3231 const [repos, attempts] = await Promise.all([
3332 reposApi.list(viewer, viewer ? { namespace: viewer.username } : {}),
34− env.WORK.listActiveAttempts(viewer),
33+ work.listActiveAttempts(viewer),
3534 ]);
3635 // Mission control links to each attempt under its repo.
3736 const attemptRepos = await Promise.all(
+6−7
1−import { env } from "cloudflare:workers";
21 import {
32 Bot,
43 ChevronRight,
2726 Textarea,
2827 TimeAgo,
2928 } from "../../components/ui";
30−import { repos } from "../../lib/services.server";
29+import { repos, work } from "../../lib/services.server";
3130 import {
3231 assertSameOrigin,
3332 getViewer,
5150 export async function loader({ params, context, request }: Route.LoaderArgs) {
5251 const viewer = getViewer(context);
5352 const [found, session] = await Promise.all([
54− env.WORK.getAttempt(params.id, viewer),
55− env.WORK.readSession(params.id, viewer),
53+ work.getAttempt(params.id, viewer),
54+ work.readSession(params.id, viewer),
5655 ]);
5756 const detail = unwrap(found);
5857 const repo = await repos.getById(detail.attempt.repoId, viewer);
8382 const action = form.get("action");
8483 const result =
8584 action === "ship"
86− ? await env.WORK.shipAttempt(user, params.id)
85+ ? await work.shipAttempt(user, params.id)
8786 : action === "abandon"
88− ? await env.WORK.abandonAttempt(user, params.id)
89− : await env.WORK.submitAttempt(
87+ ? await work.abandonAttempt(user, params.id)
88+ : await work.submitAttempt(
9089 user,
9190 params.id,
9291 String(form.get("summary") ?? ""),
+2−2
1−import { env } from "cloudflare:workers";
21 import { Form, redirect } from "react-router";
32
43 import type { Route } from "./+types/intent-new";
54 import { Button, ErrorText, Field, Input, Textarea } from "../../components/ui";
5+import { work } from "../../lib/services.server";
66 import { assertSameOrigin, requireUser } from "../../lib/session.server";
77
88 export function loader({ request, context }: Route.LoaderArgs) {
1414 assertSameOrigin(request);
1515 const user = requireUser(context, request);
1616 const form = await request.formData();
17− const result = await env.WORK.openIntent(
17+ const result = await work.openIntent(
1818 user,
1919 { namespace: params.owner, name: params.repo },
2020 {
+4−3
2222 Textarea,
2323 TimeAgo,
2424 } from "../../components/ui";
25+import { work } from "../../lib/services.server";
2526 import {
2627 assertSameOrigin,
2728 getViewer,
3940 export async function loader({ params, context }: Route.LoaderArgs) {
4041 const viewer = getViewer(context);
4142 const detail = unwrap(
42− await env.WORK.getIntent(
43+ await work.getIntent(
4344 { namespace: params.owner, name: params.repo },
4445 Number(params.number),
4546 viewer,
6768 return result.ok ? null : { error: result.error.message };
6869 }
6970 if (form.get("action") === "withdraw") {
70− const result = await env.WORK.withdrawIntent(user, intentId);
71+ const result = await work.withdrawIntent(user, intentId);
7172 return result.ok ? null : { error: result.error.message };
7273 }
73− const result = await env.WORK.startAttempt(user, intentId, {
74+ const result = await work.startAttempt(user, intentId, {
7475 agent: String(form.get("agent") ?? ""),
7576 runtime: "external",
7677 });
+2−2
1−import { env } from "cloudflare:workers";
21 import { CircleDot, GitMerge, Plus, XCircle } from "lucide-react";
32 import { Link } from "react-router";
43
65
76 import type { Route } from "./+types/intents";
87 import { ButtonLink, EmptyState, TimeAgo } from "../../components/ui";
8+import { work } from "../../lib/services.server";
99 import { getViewer, unwrap } from "../../lib/session.server";
1010
1111 export function meta({ params }: Route.MetaArgs) {
1515 export async function loader({ params, context }: Route.LoaderArgs) {
1616 const path = { namespace: params.owner, name: params.repo };
1717 return {
18− intents: unwrap(await env.WORK.listIntents(path, getViewer(context))),
18+ intents: unwrap(await work.listIntents(path, getViewer(context))),
1919 };
2020 }
2121
+2−3
1−import { env } from "cloudflare:workers";
21 import { BookMarked, Code2, History, Lock, Target } from "lucide-react";
32 import type { ReactNode } from "react";
43 import { Link, NavLink, Outlet } from "react-router";
54
65 import type { Route } from "./+types/layout";
76 import { Pill } from "../../components/ui";
8−import { repos } from "../../lib/services.server";
7+import { repos, work } from "../../lib/services.server";
98 import { getViewer, unwrap } from "../../lib/session.server";
109
1110 export function meta({ params }: Route.MetaArgs) {
1716 const path = { namespace: params.owner, name: params.repo };
1817 const [repo, intents] = await Promise.all([
1918 repos.get(path, viewer),
20− env.WORK.listIntents(path, viewer, "open"),
19+ work.listIntents(path, viewer, "open"),
2120 ]);
2221 return {
2322 repo: unwrap(repo),
+2−2
1−import type { EventsApi, RunnerApi, ServiceBinding, WorkApi } from "@g1t/contracts";
1+import type { EventsApi, RunnerApi, ServiceBinding } from "@g1t/contracts";
22
33 declare global {
44 namespace Cloudflare {
66 IDENTITY: ServiceBinding;
77 /** Also serves git over HTTPS through `fetch`. */
88 REPOS: ServiceBinding & { fetch(request: Request): Promise<Response> };
9− WORK: WorkApi;
9+ WORK: ServiceBinding;
1010 EVENTS: EventsApi;
1111 RUNNER: RunnerApi;
1212 }
+56−0
4343 pub git_ref: String,
4444 pub after: String,
4545 }
46+
47+#[derive(Debug, Serialize)]
48+#[serde(rename_all = "camelCase")]
49+pub struct IntentOpened {
50+ pub intent_id: String,
51+ pub repo_id: String,
52+ pub number: u32,
53+ pub title: String,
54+}
55+
56+#[derive(Debug, Serialize)]
57+#[serde(rename_all = "camelCase")]
58+pub struct IntentClosed {
59+ pub intent_id: String,
60+ pub repo_id: String,
61+ /// `shipped` or `withdrawn`.
62+ pub reason: &'static str,
63+}
64+
65+/// The payload of `attempt.started`, `attempt.updated`, `attempt.submitted`
66+/// and `attempt.shipped`; each uses the fields that apply to it.
67+#[derive(Debug, Default, Serialize)]
68+#[serde(rename_all = "camelCase")]
69+pub struct AttemptEvent {
70+ pub attempt_id: String,
71+ pub intent_id: String,
72+ pub repo_id: String,
73+ #[serde(skip_serializing_if = "Option::is_none")]
74+ pub agent: Option<String>,
75+ #[serde(skip_serializing_if = "Option::is_none")]
76+ pub status: Option<&'static str>,
77+ #[serde(skip_serializing_if = "Option::is_none")]
78+ pub commit: Option<String>,
79+}
80+
81+#[derive(Debug, Serialize)]
82+#[serde(rename_all = "camelCase")]
83+pub struct SessionAppended {
84+ pub attempt_id: String,
85+ pub session_id: String,
86+ pub count: u32,
87+}
88+
89+/// An event as delivered to subscribers. `data` is left as JSON; each
90+/// subscriber decodes the types it cares about.
91+#[derive(Debug, serde::Deserialize)]
92+#[serde(rename_all = "camelCase")]
93+pub struct Delivered {
94+ pub id: String,
95+ #[serde(rename = "type")]
96+ pub kind: String,
97+ /// Milliseconds since the epoch, until the bus itself moves to RFC 3339.
98+ pub time: serde_json::Value,
99+ pub repo_id: Option<String>,
100+ pub data: serde_json::Value,
101+}
+1−0
1111 mod outcome;
1212 pub mod repos;
1313 pub mod time;
14+pub mod work;
1415
1516 pub use ids::new_id;
1617 pub use names::{is_valid_namespace, is_valid_repo_name};
+238−0
1+//! The work service: intents, attempts and sessions.
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::repos::RepoPath;
9+use crate::{User, Viewer};
10+
11+#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
12+#[serde(rename_all = "lowercase")]
13+pub enum IntentStatus {
14+ Open,
15+ Shipped,
16+ Withdrawn,
17+}
18+
19+/// A goal stated against a repo: the issue, and the home of every attempt
20+/// made for it.
21+#[derive(Clone, Debug, Serialize, Deserialize)]
22+#[serde(rename_all = "camelCase")]
23+pub struct Intent {
24+ pub id: String,
25+ pub repo_id: String,
26+ /// Sequential per repo, shown as `#12`.
27+ pub number: u32,
28+ pub title: String,
29+ /// The goal in prose: what an agent is given to work from.
30+ pub brief: String,
31+ /// Commands that must pass for an attempt to be accepted.
32+ pub checks: Vec<String>,
33+ pub status: IntentStatus,
34+ pub author: User,
35+ /// RFC 3339.
36+ pub created_at: String,
37+ pub attempt_count: u32,
38+}
39+
40+#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
41+#[serde(rename_all = "lowercase")]
42+pub enum AttemptStatus {
43+ Working,
44+ Submitted,
45+ Shipped,
46+ Abandoned,
47+}
48+
49+impl AttemptStatus {
50+ pub fn as_str(self) -> &'static str {
51+ match self {
52+ AttemptStatus::Working => "working",
53+ AttemptStatus::Submitted => "submitted",
54+ AttemptStatus::Shipped => "shipped",
55+ AttemptStatus::Abandoned => "abandoned",
56+ }
57+ }
58+
59+ /// Whether the attempt can still be changed or shipped.
60+ pub fn is_active(self) -> bool {
61+ matches!(self, AttemptStatus::Working | AttemptStatus::Submitted)
62+ }
63+}
64+
65+/// Where the agent runs: on g1t's sandboxes, or in someone's own session.
66+#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
67+#[serde(rename_all = "lowercase")]
68+pub enum AttemptRuntime {
69+ Hosted,
70+ External,
71+}
72+
73+/// One agent's or person's run at an intent, in its own fork.
74+#[derive(Clone, Debug, Serialize, Deserialize)]
75+#[serde(rename_all = "camelCase")]
76+pub struct Attempt {
77+ pub id: String,
78+ pub intent_id: String,
79+ pub repo_id: String,
80+ /// Sequential per intent.
81+ pub number: u32,
82+ /// A label for the agent doing the work, e.g. `claude-code`.
83+ pub agent: String,
84+ pub runtime: AttemptRuntime,
85+ pub status: AttemptStatus,
86+ /// The agent's own account of what it did, set on submit.
87+ pub summary: Option<String>,
88+ pub fork: RepoPath,
89+ /// The fork's repository id.
90+ pub fork_repo_id: String,
91+ pub head_commit: Option<String>,
92+ /// For a shipped attempt, what the branch pointed to before it landed.
93+ /// Comparing against it shows what the attempt changed.
94+ pub landed_base: Option<String>,
95+ pub started_by: User,
96+ /// RFC 3339.
97+ pub created_at: String,
98+ /// RFC 3339.
99+ pub updated_at: String,
100+}
101+
102+#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
103+#[serde(rename_all = "snake_case")]
104+pub enum SessionEntryKind {
105+ Prompt,
106+ Message,
107+ ToolCall,
108+ ToolResult,
109+ Note,
110+}
111+
112+/// One step of an agent's session: the "why" behind an attempt's commits.
113+#[derive(Clone, Debug, Serialize, Deserialize)]
114+pub struct SessionEntry {
115+ pub seq: u32,
116+ pub kind: SessionEntryKind,
117+ pub text: String,
118+ /// For tool calls and results.
119+ pub tool: Option<String>,
120+ /// The fork's head commit when this entry was recorded, if known.
121+ pub commit: Option<String>,
122+ /// RFC 3339.
123+ pub at: String,
124+}
125+
126+#[derive(Clone, Debug, Serialize, Deserialize)]
127+pub struct NewSessionEntry {
128+ pub kind: SessionEntryKind,
129+ pub text: String,
130+ #[serde(default)]
131+ pub tool: Option<String>,
132+ #[serde(default)]
133+ pub commit: Option<String>,
134+}
135+
136+#[derive(Clone, Debug, Serialize, Deserialize)]
137+pub struct IntentDetail {
138+ pub intent: Intent,
139+ pub attempts: Vec<Attempt>,
140+}
141+
142+#[derive(Clone, Debug, Serialize, Deserialize)]
143+pub struct AttemptDetail {
144+ pub attempt: Attempt,
145+ pub intent: Intent,
146+}
147+
148+/// `open_intent`. Returns `Outcome<Intent>`.
149+#[derive(Debug, Serialize, Deserialize)]
150+pub struct OpenIntentArgs {
151+ pub actor: User,
152+ pub repo: RepoPath,
153+ pub title: String,
154+ #[serde(default)]
155+ pub brief: String,
156+ #[serde(default)]
157+ pub checks: Vec<String>,
158+}
159+
160+/// `list_intents`. Returns `Outcome<Vec<Intent>>`.
161+#[derive(Debug, Serialize, Deserialize)]
162+pub struct ListIntentsArgs {
163+ pub repo: RepoPath,
164+ pub viewer: Viewer,
165+ #[serde(default)]
166+ pub status: Option<IntentStatus>,
167+}
168+
169+/// `get_intent`. Returns `Outcome<IntentDetail>`.
170+#[derive(Debug, Serialize, Deserialize)]
171+pub struct GetIntentArgs {
172+ pub repo: RepoPath,
173+ pub number: u32,
174+ pub viewer: Viewer,
175+}
176+
177+/// `withdraw_intent`. Returns `Outcome<Intent>`.
178+#[derive(Debug, Serialize, Deserialize)]
179+#[serde(rename_all = "camelCase")]
180+pub struct IntentActionArgs {
181+ pub actor: User,
182+ pub intent_id: String,
183+}
184+
185+/// `start_attempt`: forks the repo for the agent and returns the attempt to
186+/// push to. Returns `Outcome<Attempt>`.
187+#[derive(Debug, Serialize, Deserialize)]
188+#[serde(rename_all = "camelCase")]
189+pub struct StartAttemptArgs {
190+ pub actor: User,
191+ pub intent_id: String,
192+ pub agent: String,
193+ pub runtime: AttemptRuntime,
194+}
195+
196+/// `get_attempt` and `read_session`.
197+#[derive(Debug, Serialize, Deserialize)]
198+#[serde(rename_all = "camelCase")]
199+pub struct AttemptViewArgs {
200+ pub attempt_id: String,
201+ pub viewer: Viewer,
202+ /// For `read_session`: only entries after this sequence number.
203+ #[serde(default)]
204+ pub after_seq: u32,
205+}
206+
207+/// `submit_attempt`, `abandon_attempt` and `ship_attempt`.
208+/// Each returns `Outcome<Attempt>`.
209+#[derive(Debug, Serialize, Deserialize)]
210+#[serde(rename_all = "camelCase")]
211+pub struct AttemptActionArgs {
212+ pub actor: User,
213+ pub attempt_id: String,
214+ /// For `submit_attempt`: what changed and why.
215+ #[serde(default)]
216+ pub summary: String,
217+}
218+
219+/// `list_active_attempts`: attempts in progress that the viewer started,
220+/// newest first. Returns `Vec<AttemptDetail>`.
221+#[derive(Debug, Serialize, Deserialize)]
222+pub struct ViewerArgs {
223+ pub viewer: Viewer,
224+}
225+
226+/// `append_session`. Returns `Outcome<Appended>`.
227+#[derive(Debug, Serialize, Deserialize)]
228+#[serde(rename_all = "camelCase")]
229+pub struct AppendSessionArgs {
230+ pub actor: User,
231+ pub attempt_id: String,
232+ pub entries: Vec<NewSessionEntry>,
233+}
234+
235+#[derive(Debug, Serialize, Deserialize)]
236+pub struct Appended {
237+ pub count: u32,
238+}
+0−12
16091609 "resolved": "apps/web",
16101610 "link": true
16111611 },
1612− "node_modules/@g1t/work": {
1613− "resolved": "services/work",
1614− "link": true
1615− },
16161612 "node_modules/@img/colour": {
16171613 "version": "1.1.0",
16181614 "resolved": "https://registry.npmjs.org/@img/colour/-/colour-1.1.0.tgz",
98239819 "license": "MIT",
98249820 "dependencies": {
98259821 "@cloudflare/containers": "^0.3.7",
9826− "@g1t/contracts": "*"
9827− }
9828− },
9829− "services/work": {
9830− "name": "@g1t/work",
9831− "version": "0.1.0",
9832− "license": "MIT",
9833− "dependencies": {
98349822 "@g1t/contracts": "*"
98359823 }
98369824 }
+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/work -w @g1t/api -w @g1t/web"
8+ "deploy": "npm run deploy -w @g1t/events -w @g1t/api -w @g1t/web"
99 },
1010 "devDependencies": {
1111 "typescript": "^5.9.3",
+23−0
11 import type { IdentityApi } from "./identity";
22 import type { ReposApi } from "./repos";
3+import type { WorkApi } from "./work";
34
45 /** A service binding, as far as these clients need it. */
56 export type ServiceBinding = {
8182 compare: (repoId, viewer, base) => call("compare", { repoId, viewer, base }),
8283 };
8384 }
85+
86+export function workClient(service: ServiceBinding): WorkApi {
87+ const call = <T>(method: string, args: object) => rpc<T>(service, method, args);
88+ return {
89+ openIntent: (actor, repo, input) => call("open_intent", { actor, repo, ...input }),
90+ listIntents: (repo, viewer, status) => call("list_intents", { repo, viewer, status }),
91+ getIntent: (repo, number, viewer) => call("get_intent", { repo, number, viewer }),
92+ withdrawIntent: (actor, intentId) => call("withdraw_intent", { actor, intentId }),
93+ startAttempt: (actor, intentId, input) =>
94+ call("start_attempt", { actor, intentId, ...input }),
95+ getAttempt: (attemptId, viewer) => call("get_attempt", { attemptId, viewer }),
96+ submitAttempt: (actor, attemptId, summary) =>
97+ call("submit_attempt", { actor, attemptId, summary }),
98+ abandonAttempt: (actor, attemptId) => call("abandon_attempt", { actor, attemptId }),
99+ shipAttempt: (actor, attemptId) => call("ship_attempt", { actor, attemptId }),
100+ listActiveAttempts: (viewer) => call("list_active_attempts", { viewer }),
101+ appendSession: (actor, attemptId, entries) =>
102+ call("append_session", { actor, attemptId, entries }),
103+ readSession: (attemptId, viewer, afterSeq = 0) =>
104+ call("read_session", { attemptId, viewer, afterSeq }),
105+ };
106+}
+10−7
1−import type { EventSubscriber } from "./events";
21 import type { User, Viewer } from "./identity";
32 import type { RepoPath } from "./repos";
43 import type { Result } from "./result";
1817 checks: string[];
1918 status: IntentStatus;
2019 author: User;
21− createdAt: number;
20+ /** RFC 3339. */
21+ createdAt: string;
2222 attemptCount: number;
2323 };
2424
5050 */
5151 landedBase: string | null;
5252 startedBy: User;
53− createdAt: number;
54− updatedAt: number;
53+ /** RFC 3339. */
54+ createdAt: string;
55+ /** RFC 3339. */
56+ updatedAt: string;
5557 };
5658
5759 export type SessionEntryKind = "prompt" | "message" | "tool_call" | "tool_result" | "note";
6567 tool: string | null;
6668 /** The fork's head commit when this entry was recorded, if known. */
6769 commit: string | null;
68− at: number;
70+ /** RFC 3339. */
71+ at: string;
6972 };
7073
7174 export type NewSessionEntry = Pick<SessionEntry, "kind" | "text"> &
72− Partial<Pick<SessionEntry, "tool" | "commit" | "at">>;
75+ Partial<Pick<SessionEntry, "tool" | "commit">>;
7376
7477 export type IntentDetail = { intent: Intent; attempts: Attempt[] };
7578
7881 export type StartAttemptInput = { agent: string; runtime: AttemptRuntime };
7982
8083 /** Intents, attempts and sessions. */
81−export interface WorkApi extends EventSubscriber {
84+export interface WorkApi {
8285 openIntent(actor: User, repo: RepoPath, input: OpenIntentInput): Promise<Result<Intent>>;
8386 listIntents(repo: RepoPath, viewer: Viewer, status?: IntentStatus): Promise<Result<Intent[]>>;
8487 getIntent(repo: RepoPath, number: number, viewer: Viewer): Promise<Result<IntentDetail>>;
+8−4
11 -- Every timestamp becomes RFC 3339 UTC text (2026-10-02T05:16:19.000Z)
22 -- instead of Unix seconds. SQLite cannot change a column's type, so each
33 -- table is rebuilt and its rows converted.
4+--
5+-- The new child tables must reference users_new, not users: dropping the old
6+-- users table deletes its rows, and that cascades to anything still pointing
7+-- at it. The renames at the end carry the references along.
48 PRAGMA defer_foreign_keys = on;
59
610 CREATE TABLE users_new (
2024 -- id is the SHA-256 of the session token.
2125 CREATE TABLE sessions_new (
2226 id TEXT PRIMARY KEY,
23− user_id TEXT NOT NULL REFERENCES users (id) ON DELETE CASCADE,
27+ user_id TEXT NOT NULL REFERENCES users_new (id) ON DELETE CASCADE,
2428 expires_at TEXT NOT NULL
2529 );
2630 INSERT INTO sessions_new (id, user_id, expires_at)
2832
2933 CREATE TABLE access_tokens_new (
3034 id TEXT PRIMARY KEY,
31− user_id TEXT NOT NULL REFERENCES users (id) ON DELETE CASCADE,
35+ user_id TEXT NOT NULL REFERENCES users_new (id) ON DELETE CASCADE,
3236 name TEXT NOT NULL,
3337 token_hash TEXT NOT NULL UNIQUE,
3438 created_at TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ', 'now')),
4347
4448 CREATE TABLE ssh_keys_new (
4549 id TEXT PRIMARY KEY,
46− user_id TEXT NOT NULL REFERENCES users (id) ON DELETE CASCADE,
50+ user_id TEXT NOT NULL REFERENCES users_new (id) ON DELETE CASCADE,
4751 title TEXT NOT NULL,
4852 public_key TEXT NOT NULL,
4953 fingerprint TEXT NOT NULL UNIQUE,
5761 -- One-time links sent by email. id is the SHA-256 of the token in the link.
5862 CREATE TABLE email_tokens_new (
5963 id TEXT PRIMARY KEY,
60− user_id TEXT NOT NULL REFERENCES users (id) ON DELETE CASCADE,
64+ user_id TEXT NOT NULL REFERENCES users_new (id) ON DELETE CASCADE,
6165 -- 'verify' or 'reset'
6266 kind TEXT NOT NULL,
6367 expires_at TEXT NOT NULL
+6−5
1111 type ServiceBinding,
1212 type User,
1313 type Viewer,
14− type WorkApi,
1514 fail,
1615 identityClient,
1716 ok,
17+ workClient,
1818 } from "@g1t/contracts";
1919
2020 export interface RunnerEnv {
2121 SANDBOX: DurableObjectNamespace<AttemptSandbox>;
2222 IDENTITY: ServiceBinding;
23− WORK: WorkApi;
23+ WORK: ServiceBinding;
2424 /** Secret. The model key the hosted agent runs on. */
2525 ANTHROPIC_API_KEY?: string;
2626 /**
7676 // sandbox that was killed before it could; abandoning twice is refused
7777 // harmlessly.
7878 const run = await this.ctx.storage.get<Pick<RunRequest, "actor" | "attemptId">>("run");
79− if (run) await this.env.WORK.abandonAttempt(run.actor, run.attemptId);
79+ if (run) await workClient(this.env.WORK).abandonAttempt(run.actor, run.attemptId);
8080 }
8181 }
8282
163163 if (!model) return fail("invalid", "That model is not available.");
164164 const count = Math.min(Math.max(Math.trunc(input.count) || 1, 1), MAX_AGENTS_PER_RUN);
165165 const identity = identityClient(this.env.IDENTITY);
166+ const work = workClient(this.env.WORK);
166167
167168 const attempts: Attempt[] = [];
168169 for (let i = 0; i < count; i++) {
169− const started = await this.env.WORK.startAttempt(actor, intentId, {
170+ const started = await work.startAttempt(actor, intentId, {
170171 agent: AGENT,
171172 runtime: "hosted",
172173 });
175176 const attempt = started.value;
176177 attempts.push(attempt);
177178
178− const found = await this.env.WORK.getAttempt(attempt.id, actor);
179+ const found = await work.getAttempt(attempt.id, actor);
179180 if (!found.ok) return found;
180181 // The sandbox acts as the person who started it, through a token
181182 // that only lives as long as a run can.
+1−0
1+build/
+16−0
1+[package]
2+name = "g1t-work"
3+version = "0.1.0"
4+edition.workspace = true
5+license.workspace = true
6+description = "Intents, attempts and sessions."
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
+80−0
1+-- Every timestamp becomes RFC 3339 UTC text instead of Unix milliseconds.
2+-- SQLite cannot change a column's type, so each table is rebuilt.
3+--
4+-- The new tables reference each other's new names; the renames at the end
5+-- carry those references along.
6+PRAGMA defer_foreign_keys = on;
7+
8+CREATE TABLE intents_new (
9+ id TEXT PRIMARY KEY,
10+ repo_id TEXT NOT NULL,
11+ number INTEGER NOT NULL,
12+ title TEXT NOT NULL,
13+ brief TEXT NOT NULL,
14+ -- JSON array of commands.
15+ checks TEXT NOT NULL DEFAULT '[]',
16+ status TEXT NOT NULL DEFAULT 'open',
17+ author_id TEXT NOT NULL,
18+ author_name TEXT NOT NULL,
19+ created_at TEXT NOT NULL,
20+ UNIQUE (repo_id, number)
21+);
22+INSERT INTO intents_new
23+SELECT id, repo_id, number, title, brief, checks, status, author_id, author_name,
24+ strftime('%Y-%m-%dT%H:%M:%fZ', created_at / 1000.0, 'unixepoch')
25+FROM intents;
26+
27+CREATE TABLE attempts_new (
28+ id TEXT PRIMARY KEY,
29+ intent_id TEXT NOT NULL REFERENCES intents_new (id),
30+ repo_id TEXT NOT NULL,
31+ number INTEGER NOT NULL,
32+ agent TEXT NOT NULL,
33+ runtime TEXT NOT NULL,
34+ status TEXT NOT NULL DEFAULT 'working',
35+ summary TEXT,
36+ fork_repo_id TEXT NOT NULL UNIQUE,
37+ fork_namespace TEXT NOT NULL,
38+ fork_name TEXT NOT NULL,
39+ head_commit TEXT,
40+ -- What the branch pointed to before a shipped attempt landed.
41+ landed_base TEXT,
42+ started_by_id TEXT NOT NULL,
43+ started_by_name TEXT NOT NULL,
44+ created_at TEXT NOT NULL,
45+ updated_at TEXT NOT NULL,
46+ UNIQUE (intent_id, number)
47+);
48+INSERT INTO attempts_new
49+SELECT id, intent_id, repo_id, number, agent, runtime, status, summary,
50+ fork_repo_id, fork_namespace, fork_name, head_commit, landed_base,
51+ started_by_id, started_by_name,
52+ strftime('%Y-%m-%dT%H:%M:%fZ', created_at / 1000.0, 'unixepoch'),
53+ strftime('%Y-%m-%dT%H:%M:%fZ', updated_at / 1000.0, 'unixepoch')
54+FROM attempts;
55+
56+CREATE TABLE session_entries_new (
57+ attempt_id TEXT NOT NULL REFERENCES attempts_new (id),
58+ seq INTEGER NOT NULL,
59+ kind TEXT NOT NULL,
60+ text TEXT NOT NULL,
61+ tool TEXT,
62+ -- The fork's head when the entry was recorded: links reasoning to code.
63+ "commit" TEXT,
64+ at TEXT NOT NULL,
65+ PRIMARY KEY (attempt_id, seq)
66+);
67+INSERT INTO session_entries_new
68+SELECT attempt_id, seq, kind, text, tool, "commit",
69+ strftime('%Y-%m-%dT%H:%M:%fZ', at / 1000.0, 'unixepoch')
70+FROM session_entries;
71+
72+DROP TABLE session_entries;
73+DROP TABLE attempts;
74+DROP TABLE intents;
75+
76+ALTER TABLE intents_new RENAME TO intents;
77+ALTER TABLE attempts_new RENAME TO attempts;
78+ALTER TABLE session_entries_new RENAME TO session_entries;
79+
80+CREATE INDEX attempts_starter ON attempts (started_by_id, status);
+0−16
1−{
2− "name": "@g1t/work",
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−493
1−import { WorkerEntrypoint } from "cloudflare:workers";
2−
3−import {
4− type Attempt,
5− type EventsApi,
6− type G1tEvent,
7− type Intent,
8− type IntentDetail,
9− type IntentStatus,
10− type NewEvent,
11− type NewSessionEntry,
12− type OpenIntentInput,
13− type RepoPath,
14− type ReposApi,
15− type ServiceBinding,
16− type Result,
17− type SessionEntry,
18− type StartAttemptInput,
19− type User,
20− type Viewer,
21− type WorkApi,
22− UNVERIFIED,
23− fail,
24− newId,
25− ok,
26− reposClient,
27−} from "@g1t/contracts";
28−
29−import {
30− type AttemptRow,
31− type IntentRow,
32− type SessionRow,
33− toAttempt,
34− toIntent,
35− toSessionEntry,
36−} from "./rows";
37−
38−export interface WorkEnv {
39− DB: D1Database;
40− REPOS: ServiceBinding;
41− EVENTS: EventsApi;
42−}
43−
44−const SOURCE = "work";
45−const MAX_ENTRY_BATCH = 200;
46−const MAX_ENTRY_CHARS = 64_000;
47−const SESSION_PAGE = 500;
48−
49−const INTENT_COLUMNS = `intents.*,
50− (SELECT count(*) FROM attempts WHERE attempts.intent_id = intents.id) AS attempt_count`;
51−
52−const NO_INTENT = fail("not_found", "Intent not found.");
53−const NO_ATTEMPT = fail("not_found", "Attempt not found.");
54−const SIGN_IN = fail("unauthenticated", "Sign in to do that.");
55−
56−export default class WorkService
57− extends WorkerEntrypoint<WorkEnv>
58− implements WorkApi
59−{
60− private get db(): D1Database {
61− return this.env.DB;
62− }
63−
64− private get repos(): ReposApi {
65− return reposClient(this.env.REPOS);
66− }
67−
68− private async intentById(id: string): Promise<Intent | null> {
69− const row = await this.db
70− .prepare(`SELECT ${INTENT_COLUMNS} FROM intents WHERE id = ?`)
71− .bind(id)
72− .first<IntentRow>();
73− return row ? toIntent(row) : null;
74− }
75−
76− private async attemptById(id: string): Promise<Attempt | null> {
77− const row = await this.db
78− .prepare("SELECT * FROM attempts WHERE id = ?")
79− .bind(id)
80− .first<AttemptRow>();
81− return row ? toAttempt(row) : null;
82− }
83−
84− /** The attempt, if `actor` is the one running it. */
85− private async ownAttempt(actor: User, id: string): Promise<Result<Attempt>> {
86− const attempt = await this.attemptById(id);
87− if (!attempt) return NO_ATTEMPT;
88− // Reading it must be allowed before "forbidden" may reveal it exists.
89− const repo = await this.repos.getById(attempt.repoId, actor);
90− if (!repo.ok) return NO_ATTEMPT;
91− if (attempt.startedBy.id !== actor.id) {
92− return fail("forbidden", "Only the person who started an attempt can change it.");
93− }
94− return ok(attempt);
95− }
96−
97− private publish(...events: NewEvent[]): Promise<void> {
98− return this.env.EVENTS.publish(events);
99− }
100−
101− async openIntent(
102− actor: User,
103− repoPath: RepoPath,
104− input: OpenIntentInput,
105− ): Promise<Result<Intent>> {
106− if (!actor.verified) return UNVERIFIED;
107− const title = input.title.trim();
108− const brief = input.brief.trim();
109− if (!title) return fail("invalid", "An intent needs a title.");
110− const repo = await this.repos.get(repoPath, actor);
111− if (!repo.ok) return repo;
112− const checks = (input.checks ?? []).map((check) => check.trim()).filter(Boolean);
113−
114− const id = newId("int");
115− // Numbering and insert are one statement, so concurrent opens on the
116− // same repo cannot take the same number.
117− await this.db
118− .prepare(
119− `INSERT INTO intents
120− (id, repo_id, number, title, brief, checks, author_id, author_name, created_at)
121− SELECT ?, ?, COALESCE(MAX(number), 0) + 1, ?, ?, ?, ?, ?, ?
122− FROM intents WHERE repo_id = ?`,
123− )
124− .bind(
125− id,
126− repo.value.id,
127− title,
128− brief,
129− JSON.stringify(checks),
130− actor.id,
131− actor.username,
132− Date.now(),
133− repo.value.id,
134− )
135− .run();
136− const intent = (await this.intentById(id))!;
137− await this.publish({
138− type: "intent.opened",
139− source: SOURCE,
140− repoId: intent.repoId,
141− actor: actor.id,
142− data: {
143− intentId: intent.id,
144− repoId: intent.repoId,
145− number: intent.number,
146− title: intent.title,
147− },
148− });
149− return ok(intent);
150− }
151−
152− async listIntents(
153− repoPath: RepoPath,
154− viewer: Viewer,
155− status?: IntentStatus,
156− ): Promise<Result<Intent[]>> {
157− const repo = await this.repos.get(repoPath, viewer);
158− if (!repo.ok) return repo;
159− const { results } = await this.db
160− .prepare(
161− `SELECT ${INTENT_COLUMNS} FROM intents
162− WHERE repo_id = ? AND (? IS NULL OR status = ?)
163− ORDER BY number DESC LIMIT 100`,
164− )
165− .bind(repo.value.id, status ?? null, status ?? null)
166− .all<IntentRow>();
167− return ok(results.map(toIntent));
168− }
169−
170− async getIntent(
171− repoPath: RepoPath,
172− number: number,
173− viewer: Viewer,
174− ): Promise<Result<IntentDetail>> {
175− const repo = await this.repos.get(repoPath, viewer);
176− if (!repo.ok) return repo;
177− const row = await this.db
178− .prepare(
179− `SELECT ${INTENT_COLUMNS} FROM intents WHERE repo_id = ? AND number = ?`,
180− )
181− .bind(repo.value.id, number)
182− .first<IntentRow>();
183− if (!row) return NO_INTENT;
184− const attempts = await this.db
185− .prepare("SELECT * FROM attempts WHERE intent_id = ? ORDER BY number")
186− .bind(row.id)
187− .all<AttemptRow>();
188− return ok({ intent: toIntent(row), attempts: attempts.results.map(toAttempt) });
189− }
190−
191− async withdrawIntent(actor: User, intentId: string): Promise<Result<Intent>> {
192− const intent = await this.intentById(intentId);
193− if (!intent) return NO_INTENT;
194− const repo = await this.repos.getById(intent.repoId, actor);
195− if (!repo.ok) return NO_INTENT;
196− if (intent.author.id !== actor.id && repo.value.ownerId !== actor.id) {
197− return fail("forbidden", "Only the author or the repo owner can withdraw an intent.");
198− }
199− if (intent.status !== "open") {
200− return fail("conflict", `This intent is already ${intent.status}.`);
201− }
202− await this.db
203− .prepare("UPDATE intents SET status = 'withdrawn' WHERE id = ?")
204− .bind(intent.id)
205− .run();
206− await this.publish({
207− type: "intent.closed",
208− source: SOURCE,
209− repoId: intent.repoId,
210− actor: actor.id,
211− data: { intentId: intent.id, repoId: intent.repoId, reason: "withdrawn" },
212− });
213− return ok({ ...intent, status: "withdrawn" });
214− }
215−
216− async startAttempt(
217− actor: User,
218− intentId: string,
219− input: StartAttemptInput,
220− ): Promise<Result<Attempt>> {
221− if (!actor.verified) return UNVERIFIED;
222− const intent = await this.intentById(intentId);
223− if (!intent) return NO_INTENT;
224− if (intent.status !== "open") {
225− return fail("conflict", `This intent is already ${intent.status}.`);
226− }
227− const agent = input.agent.trim() || "agent";
228−
229− const id = newId("att");
230− const fork = await this.repos.forkForAttempt(intent.repoId, id, actor);
231− if (!fork.ok) return fork.error.code === "not_found" ? NO_INTENT : fork;
232−
233− const now = Date.now();
234− await this.db
235− .prepare(
236− `INSERT INTO attempts
237− (id, intent_id, repo_id, number, agent, runtime, fork_repo_id,
238− fork_namespace, fork_name, started_by_id, started_by_name,
239− created_at, updated_at)
240− SELECT ?, ?, ?, COALESCE(MAX(number), 0) + 1, ?, ?, ?, ?, ?, ?, ?, ?, ?
241− FROM attempts WHERE intent_id = ?`,
242− )
243− .bind(
244− id,
245− intent.id,
246− intent.repoId,
247− agent,
248− input.runtime,
249− fork.value.id,
250− fork.value.namespace,
251− fork.value.name,
252− actor.id,
253− actor.username,
254− now,
255− now,
256− intent.id,
257− )
258− .run();
259− const attempt = (await this.attemptById(id))!;
260− await this.publish({
261− type: "attempt.started",
262− source: SOURCE,
263− repoId: intent.repoId,
264− actor: actor.id,
265− data: {
266− attemptId: attempt.id,
267− intentId: intent.id,
268− repoId: intent.repoId,
269− agent: attempt.agent,
270− },
271− });
272− return ok(attempt);
273− }
274−
275− async getAttempt(
276− attemptId: string,
277− viewer: Viewer,
278− ): Promise<Result<{ attempt: Attempt; intent: Intent }>> {
279− const attempt = await this.attemptById(attemptId);
280− if (!attempt) return NO_ATTEMPT;
281− const repo = await this.repos.getById(attempt.repoId, viewer);
282− if (!repo.ok) return NO_ATTEMPT;
283− return ok({ attempt, intent: (await this.intentById(attempt.intentId))! });
284− }
285−
286− private async setStatus(
287− actor: User,
288− attemptId: string,
289− status: "submitted" | "abandoned",
290− summary: string | null,
291− ): Promise<Result<Attempt>> {
292− const own = await this.ownAttempt(actor, attemptId);
293− if (!own.ok) return own;
294− const attempt = own.value;
295− if (attempt.status !== "working" && attempt.status !== "submitted") {
296− return fail("conflict", `This attempt is already ${attempt.status}.`);
297− }
298− const now = Date.now();
299− await this.db
300− .prepare(
301− "UPDATE attempts SET status = ?, summary = COALESCE(?, summary), updated_at = ? WHERE id = ?",
302− )
303− .bind(status, summary, now, attempt.id)
304− .run();
305− const ids = {
306− attemptId: attempt.id,
307− intentId: attempt.intentId,
308− repoId: attempt.repoId,
309− };
310− await this.publish(
311− status === "submitted"
312− ? { type: "attempt.submitted", source: SOURCE, repoId: attempt.repoId, actor: actor.id, data: ids }
313− : { type: "attempt.updated", source: SOURCE, repoId: attempt.repoId, actor: actor.id, data: { ...ids, status } },
314− );
315− return ok({
316− ...attempt,
317− status,
318− summary: summary ?? attempt.summary,
319− updatedAt: now,
320− });
321− }
322−
323− submitAttempt(actor: User, attemptId: string, summary: string): Promise<Result<Attempt>> {
324− return this.setStatus(actor, attemptId, "submitted", summary.trim() || null);
325− }
326−
327− abandonAttempt(actor: User, attemptId: string): Promise<Result<Attempt>> {
328− return this.setStatus(actor, attemptId, "abandoned", null);
329− }
330−
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 = ?, landed_base = ?, updated_at = ? WHERE id = ?",
357− )
358− .bind(landed.value.commit, landed.value.previous, 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− landedBase: landed.value.previous,
389− updatedAt: now,
390− });
391− }
392−
393− async listActiveAttempts(
394− viewer: Viewer,
395− ): Promise<{ attempt: Attempt; intent: Intent }[]> {
396− if (!viewer) return [];
397− const { results } = await this.db
398− .prepare(
399− `SELECT * FROM attempts
400− WHERE started_by_id = ? AND status IN ('working', 'submitted')
401− ORDER BY updated_at DESC LIMIT 50`,
402− )
403− .bind(viewer.id)
404− .all<AttemptRow>();
405− return Promise.all(
406− results.map(async (row) => ({
407− attempt: toAttempt(row),
408− intent: (await this.intentById(row.intent_id))!,
409− })),
410− );
411− }
412−
413− async appendSession(
414− actor: User,
415− attemptId: string,
416− entries: NewSessionEntry[],
417− ): Promise<Result<{ count: number }>> {
418− if (!actor) return SIGN_IN;
419− if (entries.length === 0) return ok({ count: 0 });
420− if (entries.length > MAX_ENTRY_BATCH) {
421− return fail("invalid", `Send at most ${MAX_ENTRY_BATCH} entries at a time.`);
422− }
423− const own = await this.ownAttempt(actor, attemptId);
424− if (!own.ok) return own;
425− const attempt = own.value;
426−
427− const now = Date.now();
428− // Each insert takes the next sequence number itself, so two writers
429− // appending at once cannot collide.
430− const insert = this.db.prepare(
431− `INSERT INTO session_entries (attempt_id, seq, kind, text, tool, "commit", at)
432− SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
433− FROM session_entries WHERE attempt_id = ?`,
434− );
435− await this.db.batch([
436− ...entries.map((entry) =>
437− insert.bind(
438− attempt.id,
439− entry.kind,
440− entry.text.slice(0, MAX_ENTRY_CHARS),
441− entry.tool ?? null,
442− entry.commit ?? attempt.headCommit,
443− entry.at ?? now,
444− attempt.id,
445− ),
446− ),
447− this.db
448− .prepare("UPDATE attempts SET updated_at = ? WHERE id = ?")
449− .bind(now, attempt.id),
450− ]);
451− await this.publish({
452− type: "session.appended",
453− source: SOURCE,
454− repoId: attempt.repoId,
455− actor: actor.id,
456− data: { attemptId: attempt.id, sessionId: attempt.id, count: entries.length },
457− });
458− return ok({ count: entries.length });
459− }
460−
461− async readSession(
462− attemptId: string,
463− viewer: Viewer,
464− afterSeq = 0,
465− ): Promise<Result<SessionEntry[]>> {
466− const found = await this.getAttempt(attemptId, viewer);
467− if (!found.ok) return found;
468− const { results } = await this.db
469− .prepare(
470− "SELECT * FROM session_entries WHERE attempt_id = ? AND seq > ? ORDER BY seq LIMIT ?",
471− )
472− .bind(attemptId, afterSeq, SESSION_PAGE)
473− .all<SessionRow>();
474− return ok(results.map(toSessionEntry));
475− }
476−
477− /** A push to an attempt's fork moves that attempt's head. */
478− async onEvents(events: G1tEvent[]): Promise<void> {
479− for (const event of events) {
480− if (event.type !== "git.push") continue;
481− await this.db
482− .prepare(
483− "UPDATE attempts SET head_commit = ?, updated_at = ? WHERE fork_repo_id = ?",
484− )
485− .bind(event.data.after, event.time, event.data.repoId)
486− .run();
487− }
488− }
489−
490− async queue(batch: MessageBatch<G1tEvent>): Promise<void> {
491− await this.onEvents(batch.messages.map((message) => message.body));
492− }
493−}
+714−0
1+//! The work service: intents, attempts and sessions.
2+//!
3+//! Other services reach it over `POST /rpc/<method>`; see
4+//! `g1t_contracts::work` for the methods and their arguments. It also
5+//! consumes its queue of events from the bus.
6+
7+mod rows;
8+
9+use g1t_contracts::events::{
10+ AttemptEvent, Delivered, IntentClosed, IntentOpened, NewEvent, SessionAppended,
11+};
12+use g1t_contracts::repos::{ForkArgs, GetArgs, GetByIdArgs, LandArgs, Landed, Repo};
13+use g1t_contracts::time::rfc3339;
14+use g1t_contracts::work::*;
15+use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
16+use g1t_kit::{args, js, now_ms, reply, rpc_method};
17+use serde::Serialize;
18+use worker::wasm_bindgen::JsValue;
19+use worker::{
20+ Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event,
21+};
22+
23+use rows::{AttemptRow, IntentRow, SessionRow};
24+
25+const SOURCE: &str = "work";
26+const MAX_ENTRY_BATCH: usize = 200;
27+const MAX_ENTRY_CHARS: usize = 64_000;
28+const SESSION_PAGE: u32 = 500;
29+const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
30+
31+const INTENT_COLUMNS: &str = "intents.*,
32+ (SELECT count(*) FROM attempts WHERE attempts.intent_id = intents.id) AS attempt_count";
33+
34+fn no_intent<T>() -> Outcome<T> {
35+ Outcome::fail(FailureCode::NotFound, "Intent not found.")
36+}
37+
38+fn no_attempt<T>() -> Outcome<T> {
39+ Outcome::fail(FailureCode::NotFound, "Attempt not found.")
40+}
41+
42+fn optional(value: &Option<String>) -> JsValue {
43+ value.as_deref().map_or(JsValue::NULL, JsValue::from)
44+}
45+
46+struct Work {
47+ db: D1Database,
48+ repos: Fetcher,
49+ /// The events service, an RPC stub.
50+ events: JsValue,
51+}
52+
53+impl Work {
54+ async fn publish<T: Serialize>(&self, event: NewEvent<T>) -> Result<()> {
55+ js::call(&self.events, "publish", &[js::to_js(&[event])?]).await?;
56+ Ok(())
57+ }
58+
59+ async fn repo_by_path(
60+ &self,
61+ path: &g1t_contracts::repos::RepoPath,
62+ viewer: &Viewer,
63+ ) -> Result<Outcome<Repo>> {
64+ g1t_kit::call(
65+ &self.repos,
66+ "get",
67+ &GetArgs {
68+ path: path.clone(),
69+ viewer: viewer.clone(),
70+ },
71+ )
72+ .await
73+ }
74+
75+ async fn repo_by_id(&self, id: &str, viewer: &Viewer) -> Result<Outcome<Repo>> {
76+ g1t_kit::call(
77+ &self.repos,
78+ "get_by_id",
79+ &GetByIdArgs {
80+ id: id.to_owned(),
81+ viewer: viewer.clone(),
82+ },
83+ )
84+ .await
85+ }
86+
87+ async fn intent_by_id(&self, id: &str) -> Result<Option<Intent>> {
88+ Ok(self
89+ .db
90+ .prepare(format!("SELECT {INTENT_COLUMNS} FROM intents WHERE id = ?"))
91+ .bind(&[id.into()])?
92+ .first::<IntentRow>(None)
93+ .await?
94+ .map(Intent::from))
95+ }
96+
97+ async fn attempt_by_id(&self, id: &str) -> Result<Option<Attempt>> {
98+ Ok(self
99+ .db
100+ .prepare("SELECT * FROM attempts WHERE id = ?")
101+ .bind(&[id.into()])?
102+ .first::<AttemptRow>(None)
103+ .await?
104+ .map(Attempt::from))
105+ }
106+
107+ /// The attempt, if `actor` is the one running it.
108+ async fn own_attempt(&self, actor: &User, id: &str) -> Result<Outcome<Attempt>> {
109+ let Some(attempt) = self.attempt_by_id(id).await? else {
110+ return Ok(no_attempt());
111+ };
112+ // Reading it must be allowed before "forbidden" may reveal it exists.
113+ let viewer = Some(actor.clone());
114+ if let Outcome::Fail(_) = self.repo_by_id(&attempt.repo_id, &viewer).await? {
115+ return Ok(no_attempt());
116+ }
117+ if attempt.started_by.id != actor.id {
118+ return Ok(Outcome::fail(
119+ FailureCode::Forbidden,
120+ "Only the person who started an attempt can change it.",
121+ ));
122+ }
123+ Ok(Outcome::Ok(attempt))
124+ }
125+
126+ fn attempt_event(attempt: &Attempt) -> AttemptEvent {
127+ AttemptEvent {
128+ attempt_id: attempt.id.clone(),
129+ intent_id: attempt.intent_id.clone(),
130+ repo_id: attempt.repo_id.clone(),
131+ ..AttemptEvent::default()
132+ }
133+ }
134+
135+ async fn open_intent(&self, a: OpenIntentArgs) -> Result<Outcome<Intent>> {
136+ if !a.actor.verified {
137+ return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
138+ }
139+ let title = a.title.trim();
140+ if title.is_empty() {
141+ return Ok(Outcome::fail(
142+ FailureCode::Invalid,
143+ "An intent needs a title.",
144+ ));
145+ }
146+ let repo = match self.repo_by_path(&a.repo, &Some(a.actor.clone())).await? {
147+ Outcome::Ok(repo) => repo,
148+ Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
149+ };
150+ let checks: Vec<&str> = a
151+ .checks
152+ .iter()
153+ .map(|check| check.trim())
154+ .filter(|check| !check.is_empty())
155+ .collect();
156+
157+ let now = now_ms();
158+ let id = new_id("int", now);
159+ // Numbering and insert are one statement, so concurrent opens on the
160+ // same repo cannot take the same number.
161+ self.db
162+ .prepare(
163+ "INSERT INTO intents
164+ (id, repo_id, number, title, brief, checks, author_id, author_name, created_at)
165+ SELECT ?, ?, COALESCE(MAX(number), 0) + 1, ?, ?, ?, ?, ?, ?
166+ FROM intents WHERE repo_id = ?",
167+ )
168+ .bind(&[
169+ id.as_str().into(),
170+ repo.id.as_str().into(),
171+ title.into(),
172+ a.brief.trim().into(),
173+ serde_json::to_string(&checks)?.into(),
174+ a.actor.id.as_str().into(),
175+ a.actor.username.as_str().into(),
176+ rfc3339(now).into(),
177+ repo.id.as_str().into(),
178+ ])?
179+ .run()
180+ .await?;
181+ let Some(intent) = self.intent_by_id(&id).await? else {
182+ return Ok(no_intent());
183+ };
184+ self.publish(NewEvent {
185+ kind: "intent.opened",
186+ source: SOURCE,
187+ repo_id: Some(intent.repo_id.clone()),
188+ actor: Some(a.actor.id),
189+ data: IntentOpened {
190+ intent_id: intent.id.clone(),
191+ repo_id: intent.repo_id.clone(),
192+ number: intent.number,
193+ title: intent.title.clone(),
194+ },
195+ })
196+ .await?;
197+ Ok(Outcome::Ok(intent))
198+ }
199+
200+ async fn list_intents(&self, a: ListIntentsArgs) -> Result<Outcome<Vec<Intent>>> {
201+ let repo = match self.repo_by_path(&a.repo, &a.viewer).await? {
202+ Outcome::Ok(repo) => repo,
203+ Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
204+ };
205+ let status = a
206+ .status
207+ .and_then(|status| serde_json::to_value(status).ok())
208+ .and_then(|value| value.as_str().map(str::to_owned));
209+ let rows = self
210+ .db
211+ .prepare(format!(
212+ "SELECT {INTENT_COLUMNS} FROM intents
213+ WHERE repo_id = ? AND (? IS NULL OR status = ?)
214+ ORDER BY number DESC LIMIT 100"
215+ ))
216+ .bind(&[repo.id.into(), optional(&status), optional(&status)])?
217+ .all()
218+ .await?
219+ .results::<IntentRow>()?;
220+ Ok(Outcome::Ok(rows.into_iter().map(Intent::from).collect()))
221+ }
222+
223+ async fn get_intent(&self, a: GetIntentArgs) -> Result<Outcome<IntentDetail>> {
224+ let repo = match self.repo_by_path(&a.repo, &a.viewer).await? {
225+ Outcome::Ok(repo) => repo,
226+ Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
227+ };
228+ let row = self
229+ .db
230+ .prepare(format!(
231+ "SELECT {INTENT_COLUMNS} FROM intents WHERE repo_id = ? AND number = ?"
232+ ))
233+ .bind(&[repo.id.into(), a.number.into()])?
234+ .first::<IntentRow>(None)
235+ .await?;
236+ let Some(intent) = row.map(Intent::from) else {
237+ return Ok(no_intent());
238+ };
239+ let attempts = self
240+ .db
241+ .prepare("SELECT * FROM attempts WHERE intent_id = ? ORDER BY number")
242+ .bind(&[intent.id.as_str().into()])?
243+ .all()
244+ .await?
245+ .results::<AttemptRow>()?;
246+ Ok(Outcome::Ok(IntentDetail {
247+ intent,
248+ attempts: attempts.into_iter().map(Attempt::from).collect(),
249+ }))
250+ }
251+
252+ async fn withdraw_intent(&self, a: IntentActionArgs) -> Result<Outcome<Intent>> {
253+ let Some(mut intent) = self.intent_by_id(&a.intent_id).await? else {
254+ return Ok(no_intent());
255+ };
256+ let repo = match self
257+ .repo_by_id(&intent.repo_id, &Some(a.actor.clone()))
258+ .await?
259+ {
260+ Outcome::Ok(repo) => repo,
261+ Outcome::Fail(_) => return Ok(no_intent()),
262+ };
263+ if intent.author.id != a.actor.id && repo.owner_id != a.actor.id {
264+ return Ok(Outcome::fail(
265+ FailureCode::Forbidden,
266+ "Only the author or the repo owner can withdraw an intent.",
267+ ));
268+ }
269+ if intent.status != IntentStatus::Open {
270+ return Ok(Outcome::fail(
271+ FailureCode::Conflict,
272+ "This intent is already closed.",
273+ ));
274+ }
275+ self.db
276+ .prepare("UPDATE intents SET status = 'withdrawn' WHERE id = ?")
277+ .bind(&[intent.id.as_str().into()])?
278+ .run()
279+ .await?;
280+ self.publish(NewEvent {
281+ kind: "intent.closed",
282+ source: SOURCE,
283+ repo_id: Some(intent.repo_id.clone()),
284+ actor: Some(a.actor.id),
285+ data: IntentClosed {
286+ intent_id: intent.id.clone(),
287+ repo_id: intent.repo_id.clone(),
288+ reason: "withdrawn",
289+ },
290+ })
291+ .await?;
292+ intent.status = IntentStatus::Withdrawn;
293+ Ok(Outcome::Ok(intent))
294+ }
295+
296+ async fn start_attempt(&self, a: StartAttemptArgs) -> Result<Outcome<Attempt>> {
297+ if !a.actor.verified {
298+ return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
299+ }
300+ let Some(intent) = self.intent_by_id(&a.intent_id).await? else {
301+ return Ok(no_intent());
302+ };
303+ if intent.status != IntentStatus::Open {
304+ return Ok(Outcome::fail(
305+ FailureCode::Conflict,
306+ "This intent is already closed.",
307+ ));
308+ }
309+ let agent = match a.agent.trim() {
310+ "" => "agent",
311+ agent => agent,
312+ };
313+ let runtime = match a.runtime {
314+ AttemptRuntime::Hosted => "hosted",
315+ AttemptRuntime::External => "external",
316+ };
317+
318+ let now = now_ms();
319+ let id = new_id("att", now);
320+ let fork: Outcome<Repo> = g1t_kit::call(
321+ &self.repos,
322+ "fork_for_attempt",
323+ &ForkArgs {
324+ source_id: intent.repo_id.clone(),
325+ attempt_id: id.clone(),
326+ actor: a.actor.clone(),
327+ },
328+ )
329+ .await?;
330+ let fork = match fork {
331+ Outcome::Ok(fork) => fork,
332+ // The repo is not visible to this actor, so neither is the intent.
333+ Outcome::Fail(failure) if failure.code == FailureCode::NotFound => {
334+ return Ok(no_intent());
335+ }
336+ Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
337+ };
338+
339+ let timestamp = rfc3339(now);
340+ self.db
341+ .prepare(
342+ "INSERT INTO attempts
343+ (id, intent_id, repo_id, number, agent, runtime, fork_repo_id,
344+ fork_namespace, fork_name, started_by_id, started_by_name,
345+ created_at, updated_at)
346+ SELECT ?, ?, ?, COALESCE(MAX(number), 0) + 1, ?, ?, ?, ?, ?, ?, ?, ?, ?
347+ FROM attempts WHERE intent_id = ?",
348+ )
349+ .bind(&[
350+ id.as_str().into(),
351+ intent.id.as_str().into(),
352+ intent.repo_id.as_str().into(),
353+ agent.into(),
354+ runtime.into(),
355+ fork.id.into(),
356+ fork.namespace.into(),
357+ fork.name.into(),
358+ a.actor.id.as_str().into(),
359+ a.actor.username.as_str().into(),
360+ timestamp.as_str().into(),
361+ timestamp.as_str().into(),
362+ intent.id.as_str().into(),
363+ ])?
364+ .run()
365+ .await?;
366+ let Some(attempt) = self.attempt_by_id(&id).await? else {
367+ return Ok(no_attempt());
368+ };
369+ self.publish(NewEvent {
370+ kind: "attempt.started",
371+ source: SOURCE,
372+ repo_id: Some(attempt.repo_id.clone()),
373+ actor: Some(a.actor.id),
374+ data: AttemptEvent {
375+ agent: Some(attempt.agent.clone()),
376+ ..Self::attempt_event(&attempt)
377+ },
378+ })
379+ .await?;
380+ Ok(Outcome::Ok(attempt))
381+ }
382+
383+ async fn get_attempt(&self, a: AttemptViewArgs) -> Result<Outcome<AttemptDetail>> {
384+ let Some(attempt) = self.attempt_by_id(&a.attempt_id).await? else {
385+ return Ok(no_attempt());
386+ };
387+ if let Outcome::Fail(_) = self.repo_by_id(&attempt.repo_id, &a.viewer).await? {
388+ return Ok(no_attempt());
389+ }
390+ let Some(intent) = self.intent_by_id(&attempt.intent_id).await? else {
391+ return Ok(no_attempt());
392+ };
393+ Ok(Outcome::Ok(AttemptDetail { attempt, intent }))
394+ }
395+
396+ /// Submits or abandons an attempt on behalf of the one running it.
397+ async fn close_attempt(
398+ &self,
399+ a: AttemptActionArgs,
400+ status: AttemptStatus,
401+ ) -> Result<Outcome<Attempt>> {
402+ let mut attempt = match self.own_attempt(&a.actor, &a.attempt_id).await? {
403+ Outcome::Ok(attempt) => attempt,
404+ failed => return Ok(failed),
405+ };
406+ if !attempt.status.is_active() {
407+ return Ok(Outcome::fail(
408+ FailureCode::Conflict,
409+ format!("This attempt is already {}.", attempt.status.as_str()),
410+ ));
411+ }
412+ let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
413+ let now = rfc3339(now_ms());
414+ self.db
415+ .prepare(
416+ "UPDATE attempts SET status = ?, summary = COALESCE(?, summary), updated_at = ?
417+ WHERE id = ?",
418+ )
419+ .bind(&[
420+ status.as_str().into(),
421+ optional(&summary),
422+ now.as_str().into(),
423+ attempt.id.as_str().into(),
424+ ])?
425+ .run()
426+ .await?;
427+ let submitted = status == AttemptStatus::Submitted;
428+ self.publish(NewEvent {
429+ kind: if submitted {
430+ "attempt.submitted"
431+ } else {
432+ "attempt.updated"
433+ },
434+ source: SOURCE,
435+ repo_id: Some(attempt.repo_id.clone()),
436+ actor: Some(a.actor.id),
437+ data: AttemptEvent {
438+ status: (!submitted).then_some(status.as_str()),
439+ ..Self::attempt_event(&attempt)
440+ },
441+ })
442+ .await?;
443+ attempt.status = status;
444+ attempt.summary = summary.or(attempt.summary);
445+ attempt.updated_at = now;
446+ Ok(Outcome::Ok(attempt))
447+ }
448+
449+ async fn ship_attempt(&self, a: AttemptActionArgs) -> Result<Outcome<Attempt>> {
450+ let Some(mut attempt) = self.attempt_by_id(&a.attempt_id).await? else {
451+ return Ok(no_attempt());
452+ };
453+ // Whether the actor may see and write the repo is decided by repos.
454+ let viewer = Some(a.actor.clone());
455+ if let Outcome::Fail(_) = self.repo_by_id(&attempt.repo_id, &viewer).await? {
456+ return Ok(no_attempt());
457+ }
458+ if !attempt.status.is_active() {
459+ return Ok(Outcome::fail(
460+ FailureCode::Conflict,
461+ format!("This attempt is already {}.", attempt.status.as_str()),
462+ ));
463+ }
464+ let Some(intent) = self.intent_by_id(&attempt.intent_id).await? else {
465+ return Ok(no_intent());
466+ };
467+ if intent.status != IntentStatus::Open {
468+ return Ok(Outcome::fail(
469+ FailureCode::Conflict,
470+ "This intent is already closed.",
471+ ));
472+ }
473+
474+ let landed: Outcome<Landed> = g1t_kit::call(
475+ &self.repos,
476+ "land",
477+ &LandArgs {
478+ fork_id: attempt.fork_repo_id.clone(),
479+ actor: a.actor.clone(),
480+ },
481+ )
482+ .await?;
483+ let landed = match landed {
484+ Outcome::Ok(landed) => landed,
485+ Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
486+ };
487+
488+ let now = rfc3339(now_ms());
489+ self.db
490+ .batch(vec![
491+ self.db
492+ .prepare(
493+ "UPDATE attempts
494+ SET status = 'shipped', head_commit = ?, landed_base = ?, updated_at = ?
495+ WHERE id = ?",
496+ )
497+ .bind(&[
498+ landed.commit.as_str().into(),
499+ optional(&landed.previous),
500+ now.as_str().into(),
501+ attempt.id.as_str().into(),
502+ ])?,
503+ self.db
504+ .prepare("UPDATE intents SET status = 'shipped' WHERE id = ?")
505+ .bind(&[intent.id.as_str().into()])?,
506+ ])
507+ .await?;
508+ self.publish(NewEvent {
509+ kind: "attempt.shipped",
510+ source: SOURCE,
511+ repo_id: Some(attempt.repo_id.clone()),
512+ actor: Some(a.actor.id.clone()),
513+ data: AttemptEvent {
514+ commit: Some(landed.commit.clone()),
515+ ..Self::attempt_event(&attempt)
516+ },
517+ })
518+ .await?;
519+ self.publish(NewEvent {
520+ kind: "intent.closed",
521+ source: SOURCE,
522+ repo_id: Some(attempt.repo_id.clone()),
523+ actor: Some(a.actor.id),
524+ data: IntentClosed {
525+ intent_id: intent.id,
526+ repo_id: attempt.repo_id.clone(),
527+ reason: "shipped",
528+ },
529+ })
530+ .await?;
531+
532+ attempt.status = AttemptStatus::Shipped;
533+ attempt.head_commit = Some(landed.commit);
534+ attempt.landed_base = landed.previous;
535+ attempt.updated_at = now;
536+ Ok(Outcome::Ok(attempt))
537+ }
538+
539+ async fn list_active_attempts(&self, a: ViewerArgs) -> Result<Vec<AttemptDetail>> {
540+ let Some(viewer) = a.viewer else {
541+ return Ok(Vec::new());
542+ };
543+ let rows = self
544+ .db
545+ .prepare(
546+ "SELECT * FROM attempts
547+ WHERE started_by_id = ? AND status IN ('working', 'submitted')
548+ ORDER BY updated_at DESC LIMIT 50",
549+ )
550+ .bind(&[viewer.id.into()])?
551+ .all()
552+ .await?
553+ .results::<AttemptRow>()?;
554+ let mut active = Vec::with_capacity(rows.len());
555+ for attempt in rows.into_iter().map(Attempt::from) {
556+ if let Some(intent) = self.intent_by_id(&attempt.intent_id).await? {
557+ active.push(AttemptDetail { attempt, intent });
558+ }
559+ }
560+ Ok(active)
561+ }
562+
563+ async fn append_session(&self, a: AppendSessionArgs) -> Result<Outcome<Appended>> {
564+ if a.entries.is_empty() {
565+ return Ok(Outcome::Ok(Appended { count: 0 }));
566+ }
567+ if a.entries.len() > MAX_ENTRY_BATCH {
568+ return Ok(Outcome::fail(
569+ FailureCode::Invalid,
570+ format!("Send at most {MAX_ENTRY_BATCH} entries at a time."),
571+ ));
572+ }
573+ let attempt = match self.own_attempt(&a.actor, &a.attempt_id).await? {
574+ Outcome::Ok(attempt) => attempt,
575+ Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
576+ };
577+
578+ let now = rfc3339(now_ms());
579+ let count = a.entries.len() as u32;
580+ let mut statements = Vec::with_capacity(a.entries.len() + 1);
581+ for entry in a.entries {
582+ let kind = serde_json::to_value(entry.kind)?;
583+ let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
584+ // Each insert takes the next sequence number itself, so two
585+ // writers appending at once cannot collide.
586+ statements.push(
587+ self.db
588+ .prepare(
589+ "INSERT INTO session_entries (attempt_id, seq, kind, text, tool, \"commit\", at)
590+ SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
591+ FROM session_entries WHERE attempt_id = ?",
592+ )
593+ .bind(&[
594+ attempt.id.as_str().into(),
595+ kind.as_str().unwrap_or("note").into(),
596+ text.into(),
597+ optional(&entry.tool),
598+ optional(&entry.commit.or_else(|| attempt.head_commit.clone())),
599+ now.as_str().into(),
600+ attempt.id.as_str().into(),
601+ ])?,
602+ );
603+ }
604+ statements.push(
605+ self.db
606+ .prepare("UPDATE attempts SET updated_at = ? WHERE id = ?")
607+ .bind(&[now.as_str().into(), attempt.id.as_str().into()])?,
608+ );
609+ self.db.batch(statements).await?;
610+ self.publish(NewEvent {
611+ kind: "session.appended",
612+ source: SOURCE,
613+ repo_id: Some(attempt.repo_id.clone()),
614+ actor: Some(a.actor.id),
615+ data: SessionAppended {
616+ attempt_id: attempt.id.clone(),
617+ session_id: attempt.id,
618+ count,
619+ },
620+ })
621+ .await?;
622+ Ok(Outcome::Ok(Appended { count }))
623+ }
624+
625+ async fn read_session(&self, a: AttemptViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
626+ let after_seq = a.after_seq;
627+ let attempt_id = a.attempt_id.clone();
628+ if let Outcome::Fail(failure) = self.get_attempt(a).await? {
629+ return Ok(Outcome::Fail(failure));
630+ }
631+ let rows = self
632+ .db
633+ .prepare(
634+ "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
635+ WHERE attempt_id = ? AND seq > ? ORDER BY seq LIMIT ?",
636+ )
637+ .bind(&[attempt_id.into(), after_seq.into(), SESSION_PAGE.into()])?
638+ .all()
639+ .await?
640+ .results::<SessionRow>()?;
641+ Ok(Outcome::Ok(
642+ rows.into_iter().map(SessionEntry::from).collect(),
643+ ))
644+ }
645+
646+ /// A push to an attempt's fork moves that attempt's head.
647+ async fn on_event(&self, event: &Delivered) -> Result<()> {
648+ if event.kind != "git.push" {
649+ return Ok(());
650+ }
651+ let (Some(repo_id), Some(after)) = (event.repo_id.as_deref(), event.data["after"].as_str())
652+ else {
653+ return Ok(());
654+ };
655+ self.db
656+ .prepare("UPDATE attempts SET head_commit = ?, updated_at = ? WHERE fork_repo_id = ?")
657+ .bind(&[after.into(), rfc3339(now_ms()).into(), repo_id.into()])?
658+ .run()
659+ .await?;
660+ Ok(())
661+ }
662+}
663+
664+fn service(env: &Env) -> Result<Work> {
665+ Ok(Work {
666+ db: env.d1("DB")?,
667+ repos: env.service("REPOS")?,
668+ events: js::binding(env, "EVENTS")?,
669+ })
670+}
671+
672+#[event(fetch)]
673+async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
674+ let Some(method) = rpc_method(&request) else {
675+ return Response::error("Not found", 404);
676+ };
677+ let body: serde_json::Value = request.json().await?;
678+ let work = service(&env)?;
679+
680+ match method.as_str() {
681+ "open_intent" => reply(&work.open_intent(args(body)?).await?),
682+ "list_intents" => reply(&work.list_intents(args(body)?).await?),
683+ "get_intent" => reply(&work.get_intent(args(body)?).await?),
684+ "withdraw_intent" => reply(&work.withdraw_intent(args(body)?).await?),
685+ "start_attempt" => reply(&work.start_attempt(args(body)?).await?),
686+ "get_attempt" => reply(&work.get_attempt(args(body)?).await?),
687+ "submit_attempt" => reply(
688+ &work
689+ .close_attempt(args(body)?, AttemptStatus::Submitted)
690+ .await?,
691+ ),
692+ "abandon_attempt" => reply(
693+ &work
694+ .close_attempt(args(body)?, AttemptStatus::Abandoned)
695+ .await?,
696+ ),
697+ "ship_attempt" => reply(&work.ship_attempt(args(body)?).await?),
698+ "list_active_attempts" => reply(&work.list_active_attempts(args(body)?).await?),
699+ "append_session" => reply(&work.append_session(args(body)?).await?),
700+ "read_session" => reply(&work.read_session(args(body)?).await?),
701+ _ => Response::error("Unknown method", 404),
702+ }
703+}
704+
705+/// Events from the bus, delivered on this service's own queue.
706+#[event(queue)]
707+async fn queue(batch: MessageBatch<Delivered>, env: Env, _ctx: Context) -> Result<()> {
708+ let work = service(&env)?;
709+ for message in batch.messages()? {
710+ work.on_event(message.body()).await?;
711+ message.ack();
712+ }
713+ Ok(())
714+}
+118−0
1+//! Rows as they come out of D1, and their conversion to contract types.
2+
3+use g1t_contracts::User;
4+use g1t_contracts::repos::RepoPath;
5+use g1t_contracts::work::{
6+ Attempt, AttemptRuntime, AttemptStatus, Intent, IntentStatus, SessionEntry, SessionEntryKind,
7+};
8+use serde::Deserialize;
9+
10+#[derive(Deserialize)]
11+pub struct IntentRow {
12+ pub id: String,
13+ pub repo_id: String,
14+ pub number: u32,
15+ pub title: String,
16+ pub brief: String,
17+ /// JSON array of commands.
18+ pub checks: String,
19+ pub status: IntentStatus,
20+ pub author_id: String,
21+ pub author_name: String,
22+ pub created_at: String,
23+ pub attempt_count: u32,
24+}
25+
26+impl From<IntentRow> for Intent {
27+ fn from(row: IntentRow) -> Self {
28+ Intent {
29+ id: row.id,
30+ repo_id: row.repo_id,
31+ number: row.number,
32+ title: row.title,
33+ brief: row.brief,
34+ checks: serde_json::from_str(&row.checks).unwrap_or_default(),
35+ status: row.status,
36+ author: User {
37+ id: row.author_id,
38+ username: row.author_name,
39+ verified: false,
40+ },
41+ created_at: row.created_at,
42+ attempt_count: row.attempt_count,
43+ }
44+ }
45+}
46+
47+#[derive(Deserialize)]
48+pub struct AttemptRow {
49+ pub id: String,
50+ pub intent_id: String,
51+ pub repo_id: String,
52+ pub number: u32,
53+ pub agent: String,
54+ pub runtime: AttemptRuntime,
55+ pub status: AttemptStatus,
56+ pub summary: Option<String>,
57+ pub fork_repo_id: String,
58+ pub fork_namespace: String,
59+ pub fork_name: String,
60+ pub head_commit: Option<String>,
61+ pub landed_base: Option<String>,
62+ pub started_by_id: String,
63+ pub started_by_name: String,
64+ pub created_at: String,
65+ pub updated_at: String,
66+}
67+
68+impl From<AttemptRow> for Attempt {
69+ fn from(row: AttemptRow) -> Self {
70+ Attempt {
71+ id: row.id,
72+ intent_id: row.intent_id,
73+ repo_id: row.repo_id,
74+ number: row.number,
75+ agent: row.agent,
76+ runtime: row.runtime,
77+ status: row.status,
78+ summary: row.summary,
79+ fork: RepoPath {
80+ namespace: row.fork_namespace,
81+ name: row.fork_name,
82+ },
83+ fork_repo_id: row.fork_repo_id,
84+ head_commit: row.head_commit,
85+ landed_base: row.landed_base,
86+ started_by: User {
87+ id: row.started_by_id,
88+ username: row.started_by_name,
89+ verified: false,
90+ },
91+ created_at: row.created_at,
92+ updated_at: row.updated_at,
93+ }
94+ }
95+}
96+
97+#[derive(Deserialize)]
98+pub struct SessionRow {
99+ pub seq: u32,
100+ pub kind: SessionEntryKind,
101+ pub text: String,
102+ pub tool: Option<String>,
103+ pub commit: Option<String>,
104+ pub at: String,
105+}
106+
107+impl From<SessionRow> for SessionEntry {
108+ fn from(row: SessionRow) -> Self {
109+ SessionEntry {
110+ seq: row.seq,
111+ kind: row.kind,
112+ text: row.text,
113+ tool: row.tool,
114+ commit: row.commit,
115+ at: row.at,
116+ }
117+ }
118+}
+0−99
1−import type {
2− Attempt,
3− AttemptRuntime,
4− AttemptStatus,
5− Intent,
6− IntentStatus,
7− SessionEntry,
8− SessionEntryKind,
9−} from "@g1t/contracts";
10−
11−export type IntentRow = {
12− id: string;
13− repo_id: string;
14− number: number;
15− title: string;
16− brief: string;
17− checks: string;
18− status: IntentStatus;
19− author_id: string;
20− author_name: string;
21− created_at: number;
22− attempt_count: number;
23−};
24−
25−export type AttemptRow = {
26− id: string;
27− intent_id: string;
28− repo_id: string;
29− number: number;
30− agent: string;
31− runtime: AttemptRuntime;
32− status: AttemptStatus;
33− summary: string | null;
34− fork_repo_id: string;
35− fork_namespace: string;
36− fork_name: string;
37− head_commit: string | null;
38− landed_base: string | null;
39− started_by_id: string;
40− started_by_name: string;
41− created_at: number;
42− updated_at: number;
43−};
44−
45−export type SessionRow = {
46− attempt_id: string;
47− seq: number;
48− kind: SessionEntryKind;
49− text: string;
50− tool: string | null;
51− commit: string | null;
52− at: number;
53−};
54−
55−export function toIntent(row: IntentRow): Intent {
56− return {
57− id: row.id,
58− repoId: row.repo_id,
59− number: row.number,
60− title: row.title,
61− brief: row.brief,
62− checks: JSON.parse(row.checks),
63− status: row.status,
64− author: { id: row.author_id, username: row.author_name },
65− createdAt: row.created_at,
66− attemptCount: row.attempt_count,
67− };
68−}
69−
70−export function toAttempt(row: AttemptRow): Attempt {
71− return {
72− id: row.id,
73− intentId: row.intent_id,
74− repoId: row.repo_id,
75− number: row.number,
76− agent: row.agent,
77− runtime: row.runtime,
78− status: row.status,
79− summary: row.summary,
80− fork: { namespace: row.fork_namespace, name: row.fork_name },
81− forkRepoId: row.fork_repo_id,
82− headCommit: row.head_commit,
83− landedBase: row.landed_base,
84− startedBy: { id: row.started_by_id, username: row.started_by_name },
85− createdAt: row.created_at,
86− updatedAt: row.updated_at,
87− };
88−}
89−
90−export function toSessionEntry(row: SessionRow): SessionEntry {
91− return {
92− seq: row.seq,
93− kind: row.kind,
94− text: row.text,
95− tool: row.tool,
96− commit: row.commit,
97− at: row.at,
98− };
99−}
+0−4
1−{
2− "extends": "../../tsconfig.base.json",
3− "include": ["src/**/*", "worker-configuration.d.ts"]
4−}
+2−1
33 "name": "g1t-work",
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 {