g1t/services/work/src/index.ts

422 lines12,783 bytesCodeBlame
1import { WorkerEntrypoint } from "cloudflare:workers";
2
3import {
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 Result,
16 type SessionEntry,
17 type StartAttemptInput,
18 type User,
19 type Viewer,
20 type WorkApi,
21 fail,
22 newId,
23 ok,
24} from "@g1t/contracts";
25
26import {
27 type AttemptRow,
28 type IntentRow,
29 type SessionRow,
30 toAttempt,
31 toIntent,
32 toSessionEntry,
33} from "./rows";
34
35export interface WorkEnv {
36 DB: D1Database;
37 REPOS: ReposApi;
38 EVENTS: EventsApi;
39}
40
41const SOURCE = "work";
42const MAX_ENTRY_BATCH = 200;
43const MAX_ENTRY_CHARS = 64_000;
44const SESSION_PAGE = 500;
45
46const INTENT_COLUMNS = `intents.*,
47 (SELECT count(*) FROM attempts WHERE attempts.intent_id = intents.id) AS attempt_count`;
48
49const NO_INTENT = fail("not_found", "Intent not found.");
50const NO_ATTEMPT = fail("not_found", "Attempt not found.");
51const SIGN_IN = fail("unauthenticated", "Sign in to do that.");
52
53export default class WorkService
54 extends WorkerEntrypoint<WorkEnv>
55 implements WorkApi
56{
57 private get db(): D1Database {
58 return this.env.DB;
59 }
60
61 private async intentById(id: string): Promise<Intent | null> {
62 const row = await this.db
63 .prepare(`SELECT ${INTENT_COLUMNS} FROM intents WHERE id = ?`)
64 .bind(id)
65 .first<IntentRow>();
66 return row ? toIntent(row) : null;
67 }
68
69 private async attemptById(id: string): Promise<Attempt | null> {
70 const row = await this.db
71 .prepare("SELECT * FROM attempts WHERE id = ?")
72 .bind(id)
73 .first<AttemptRow>();
74 return row ? toAttempt(row) : null;
75 }
76
77 /** The attempt, if `actor` is the one running it. */
78 private async ownAttempt(actor: User, id: string): Promise<Result<Attempt>> {
79 const attempt = await this.attemptById(id);
80 if (!attempt) return NO_ATTEMPT;
81 // Reading it must be allowed before "forbidden" may reveal it exists.
82 const repo = await this.env.REPOS.getById(attempt.repoId, actor);
83 if (!repo.ok) return NO_ATTEMPT;
84 if (attempt.startedBy.id !== actor.id) {
85 return fail("forbidden", "Only the person who started an attempt can change it.");
86 }
87 return ok(attempt);
88 }
89
90 private publish(...events: NewEvent[]): Promise<void> {
91 return this.env.EVENTS.publish(events);
92 }
93
94 async openIntent(
95 actor: User,
96 repoPath: RepoPath,
97 input: OpenIntentInput,
98 ): Promise<Result<Intent>> {
99 const title = input.title.trim();
100 const brief = input.brief.trim();
101 if (!title) return fail("invalid", "An intent needs a title.");
102 const repo = await this.env.REPOS.get(repoPath, actor);
103 if (!repo.ok) return repo;
104 const checks = (input.checks ?? []).map((check) => check.trim()).filter(Boolean);
105
106 const id = newId("int");
107 // Numbering and insert are one statement, so concurrent opens on the
108 // same repo cannot take the same number.
109 await this.db
110 .prepare(
111 `INSERT INTO intents
112 (id, repo_id, number, title, brief, checks, author_id, author_name, created_at)
113 SELECT ?, ?, COALESCE(MAX(number), 0) + 1, ?, ?, ?, ?, ?, ?
114 FROM intents WHERE repo_id = ?`,
115 )
116 .bind(
117 id,
118 repo.value.id,
119 title,
120 brief,
121 JSON.stringify(checks),
122 actor.id,
123 actor.username,
124 Date.now(),
125 repo.value.id,
126 )
127 .run();
128 const intent = (await this.intentById(id))!;
129 await this.publish({
130 type: "intent.opened",
131 source: SOURCE,
132 repoId: intent.repoId,
133 actor: actor.id,
134 data: {
135 intentId: intent.id,
136 repoId: intent.repoId,
137 number: intent.number,
138 title: intent.title,
139 },
140 });
141 return ok(intent);
142 }
143
144 async listIntents(
145 repoPath: RepoPath,
146 viewer: Viewer,
147 status?: IntentStatus,
148 ): Promise<Result<Intent[]>> {
149 const repo = await this.env.REPOS.get(repoPath, viewer);
150 if (!repo.ok) return repo;
151 const { results } = await this.db
152 .prepare(
153 `SELECT ${INTENT_COLUMNS} FROM intents
154 WHERE repo_id = ? AND (? IS NULL OR status = ?)
155 ORDER BY number DESC LIMIT 100`,
156 )
157 .bind(repo.value.id, status ?? null, status ?? null)
158 .all<IntentRow>();
159 return ok(results.map(toIntent));
160 }
161
162 async getIntent(
163 repoPath: RepoPath,
164 number: number,
165 viewer: Viewer,
166 ): Promise<Result<IntentDetail>> {
167 const repo = await this.env.REPOS.get(repoPath, viewer);
168 if (!repo.ok) return repo;
169 const row = await this.db
170 .prepare(
171 `SELECT ${INTENT_COLUMNS} FROM intents WHERE repo_id = ? AND number = ?`,
172 )
173 .bind(repo.value.id, number)
174 .first<IntentRow>();
175 if (!row) return NO_INTENT;
176 const attempts = await this.db
177 .prepare("SELECT * FROM attempts WHERE intent_id = ? ORDER BY number")
178 .bind(row.id)
179 .all<AttemptRow>();
180 return ok({ intent: toIntent(row), attempts: attempts.results.map(toAttempt) });
181 }
182
183 async withdrawIntent(actor: User, intentId: string): Promise<Result<Intent>> {
184 const intent = await this.intentById(intentId);
185 if (!intent) return NO_INTENT;
186 const repo = await this.env.REPOS.getById(intent.repoId, actor);
187 if (!repo.ok) return NO_INTENT;
188 if (intent.author.id !== actor.id && repo.value.ownerId !== actor.id) {
189 return fail("forbidden", "Only the author or the repo owner can withdraw an intent.");
190 }
191 if (intent.status !== "open") {
192 return fail("conflict", `This intent is already ${intent.status}.`);
193 }
194 await this.db
195 .prepare("UPDATE intents SET status = 'withdrawn' WHERE id = ?")
196 .bind(intent.id)
197 .run();
198 await this.publish({
199 type: "intent.closed",
200 source: SOURCE,
201 repoId: intent.repoId,
202 actor: actor.id,
203 data: { intentId: intent.id, repoId: intent.repoId, reason: "withdrawn" },
204 });
205 return ok({ ...intent, status: "withdrawn" });
206 }
207
208 async startAttempt(
209 actor: User,
210 intentId: string,
211 input: StartAttemptInput,
212 ): Promise<Result<Attempt>> {
213 const intent = await this.intentById(intentId);
214 if (!intent) return NO_INTENT;
215 if (intent.status !== "open") {
216 return fail("conflict", `This intent is already ${intent.status}.`);
217 }
218 const agent = input.agent.trim() || "agent";
219
220 const id = newId("att");
221 const fork = await this.env.REPOS.forkForAttempt(intent.repoId, id, actor);
222 if (!fork.ok) return fork.error.code === "not_found" ? NO_INTENT : fork;
223
224 const now = Date.now();
225 await this.db
226 .prepare(
227 `INSERT INTO attempts
228 (id, intent_id, repo_id, number, agent, runtime, fork_repo_id,
229 fork_namespace, fork_name, started_by_id, started_by_name,
230 created_at, updated_at)
231 SELECT ?, ?, ?, COALESCE(MAX(number), 0) + 1, ?, ?, ?, ?, ?, ?, ?, ?, ?
232 FROM attempts WHERE intent_id = ?`,
233 )
234 .bind(
235 id,
236 intent.id,
237 intent.repoId,
238 agent,
239 input.runtime,
240 fork.value.id,
241 fork.value.namespace,
242 fork.value.name,
243 actor.id,
244 actor.username,
245 now,
246 now,
247 intent.id,
248 )
249 .run();
250 const attempt = (await this.attemptById(id))!;
251 await this.publish({
252 type: "attempt.started",
253 source: SOURCE,
254 repoId: intent.repoId,
255 actor: actor.id,
256 data: {
257 attemptId: attempt.id,
258 intentId: intent.id,
259 repoId: intent.repoId,
260 agent: attempt.agent,
261 },
262 });
263 return ok(attempt);
264 }
265
266 async getAttempt(
267 attemptId: string,
268 viewer: Viewer,
269 ): Promise<Result<{ attempt: Attempt; intent: Intent }>> {
270 const attempt = await this.attemptById(attemptId);
271 if (!attempt) return NO_ATTEMPT;
272 const repo = await this.env.REPOS.getById(attempt.repoId, viewer);
273 if (!repo.ok) return NO_ATTEMPT;
274 return ok({ attempt, intent: (await this.intentById(attempt.intentId))! });
275 }
276
277 private async setStatus(
278 actor: User,
279 attemptId: string,
280 status: "submitted" | "abandoned",
281 summary: string | null,
282 ): Promise<Result<Attempt>> {
283 const own = await this.ownAttempt(actor, attemptId);
284 if (!own.ok) return own;
285 const attempt = own.value;
286 if (attempt.status !== "working" && attempt.status !== "submitted") {
287 return fail("conflict", `This attempt is already ${attempt.status}.`);
288 }
289 const now = Date.now();
290 await this.db
291 .prepare(
292 "UPDATE attempts SET status = ?, summary = COALESCE(?, summary), updated_at = ? WHERE id = ?",
293 )
294 .bind(status, summary, now, attempt.id)
295 .run();
296 const ids = {
297 attemptId: attempt.id,
298 intentId: attempt.intentId,
299 repoId: attempt.repoId,
300 };
301 await this.publish(
302 status === "submitted"
303 ? { type: "attempt.submitted", source: SOURCE, repoId: attempt.repoId, actor: actor.id, data: ids }
304 : { type: "attempt.updated", source: SOURCE, repoId: attempt.repoId, actor: actor.id, data: { ...ids, status } },
305 );
306 return ok({
307 ...attempt,
308 status,
309 summary: summary ?? attempt.summary,
310 updatedAt: now,
311 });
312 }
313
314 submitAttempt(actor: User, attemptId: string, summary: string): Promise<Result<Attempt>> {
315 return this.setStatus(actor, attemptId, "submitted", summary.trim() || null);
316 }
317
318 abandonAttempt(actor: User, attemptId: string): Promise<Result<Attempt>> {
319 return this.setStatus(actor, attemptId, "abandoned", null);
320 }
321
322 async listActiveAttempts(
323 viewer: Viewer,
324 ): Promise<{ attempt: Attempt; intent: Intent }[]> {
325 if (!viewer) return [];
326 const { results } = await this.db
327 .prepare(
328 `SELECT * FROM attempts
329 WHERE started_by_id = ? AND status IN ('working', 'submitted')
330 ORDER BY updated_at DESC LIMIT 50`,
331 )
332 .bind(viewer.id)
333 .all<AttemptRow>();
334 return Promise.all(
335 results.map(async (row) => ({
336 attempt: toAttempt(row),
337 intent: (await this.intentById(row.intent_id))!,
338 })),
339 );
340 }
341
342 async appendSession(
343 actor: User,
344 attemptId: string,
345 entries: NewSessionEntry[],
346 ): Promise<Result<{ count: number }>> {
347 if (!actor) return SIGN_IN;
348 if (entries.length === 0) return ok({ count: 0 });
349 if (entries.length > MAX_ENTRY_BATCH) {
350 return fail("invalid", `Send at most ${MAX_ENTRY_BATCH} entries at a time.`);
351 }
352 const own = await this.ownAttempt(actor, attemptId);
353 if (!own.ok) return own;
354 const attempt = own.value;
355
356 const now = Date.now();
357 // Each insert takes the next sequence number itself, so two writers
358 // appending at once cannot collide.
359 const insert = this.db.prepare(
360 `INSERT INTO session_entries (attempt_id, seq, kind, text, tool, "commit", at)
361 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
362 FROM session_entries WHERE attempt_id = ?`,
363 );
364 await this.db.batch([
365 ...entries.map((entry) =>
366 insert.bind(
367 attempt.id,
368 entry.kind,
369 entry.text.slice(0, MAX_ENTRY_CHARS),
370 entry.tool ?? null,
371 entry.commit ?? attempt.headCommit,
372 entry.at ?? now,
373 attempt.id,
374 ),
375 ),
376 this.db
377 .prepare("UPDATE attempts SET updated_at = ? WHERE id = ?")
378 .bind(now, attempt.id),
379 ]);
380 await this.publish({
381 type: "session.appended",
382 source: SOURCE,
383 repoId: attempt.repoId,
384 actor: actor.id,
385 data: { attemptId: attempt.id, sessionId: attempt.id, count: entries.length },
386 });
387 return ok({ count: entries.length });
388 }
389
390 async readSession(
391 attemptId: string,
392 viewer: Viewer,
393 afterSeq = 0,
394 ): Promise<Result<SessionEntry[]>> {
395 const found = await this.getAttempt(attemptId, viewer);
396 if (!found.ok) return found;
397 const { results } = await this.db
398 .prepare(
399 "SELECT * FROM session_entries WHERE attempt_id = ? AND seq > ? ORDER BY seq LIMIT ?",
400 )
401 .bind(attemptId, afterSeq, SESSION_PAGE)
402 .all<SessionRow>();
403 return ok(results.map(toSessionEntry));
404 }
405
406 /** A push to an attempt's fork moves that attempt's head. */
407 async onEvents(events: G1tEvent[]): Promise<void> {
408 for (const event of events) {
409 if (event.type !== "git.push") continue;
410 await this.db
411 .prepare(
412 "UPDATE attempts SET head_commit = ?, updated_at = ? WHERE fork_repo_id = ?",
413 )
414 .bind(event.data.after, event.time, event.data.repoId)
415 .run();
416 }
417 }
418
419 async queue(batch: MessageBatch<G1tEvent>): Promise<void> {
420 await this.onEvents(batch.messages.map((message) => message.body));
421 }
422}