pr_01m47d15m3e54sn21z27rpy5n9/services/work/src/index.ts

492 lines14,919 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Initial g1t: services, event bus, intents and attempts1import { 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,
Rust repos service with shipping; pull requests kept in the model15 type ServiceBinding,
Initial g1t: services, event bus, intents and attempts16 type Result,
17 type SessionEntry,
18 type StartAttemptInput,
19 type User,
20 type Viewer,
21 type WorkApi,
Email verification, password reset, and Git for AI scale positioning22 UNVERIFIED,
Initial g1t: services, event bus, intents and attempts23 fail,
24 newId,
25 ok,
Rust repos service with shipping; pull requests kept in the model26 reposClient,
Initial g1t: services, event bus, intents and attempts27} from "@g1t/contracts";
28
29import {
30 type AttemptRow,
31 type IntentRow,
32 type SessionRow,
33 toAttempt,
34 toIntent,
35 toSessionEntry,
36} from "./rows";
37
38export interface WorkEnv {
39 DB: D1Database;
Rust repos service with shipping; pull requests kept in the model40 REPOS: ServiceBinding;
Initial g1t: services, event bus, intents and attempts41 EVENTS: EventsApi;
42}
43
44const SOURCE = "work";
45const MAX_ENTRY_BATCH = 200;
46const MAX_ENTRY_CHARS = 64_000;
47const SESSION_PAGE = 500;
48
49const INTENT_COLUMNS = `intents.*,
50 (SELECT count(*) FROM attempts WHERE attempts.intent_id = intents.id) AS attempt_count`;
51
52const NO_INTENT = fail("not_found", "Intent not found.");
53const NO_ATTEMPT = fail("not_found", "Attempt not found.");
54const SIGN_IN = fail("unauthenticated", "Sign in to do that.");
55
56export default class WorkService
57 extends WorkerEntrypoint<WorkEnv>
58 implements WorkApi
59{
60 private get db(): D1Database {
61 return this.env.DB;
62 }
63
Rust repos service with shipping; pull requests kept in the model64 private get repos(): ReposApi {
65 return reposClient(this.env.REPOS);
66 }
67
Initial g1t: services, event bus, intents and attempts68 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.
Rust repos service with shipping; pull requests kept in the model89 const repo = await this.repos.getById(attempt.repoId, actor);
Initial g1t: services, event bus, intents and attempts90 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>> {
Email verification, password reset, and Git for AI scale positioning106 if (!actor.verified) return UNVERIFIED;
Initial g1t: services, event bus, intents and attempts107 const title = input.title.trim();
108 const brief = input.brief.trim();
109 if (!title) return fail("invalid", "An intent needs a title.");
Rust repos service with shipping; pull requests kept in the model110 const repo = await this.repos.get(repoPath, actor);
Initial g1t: services, event bus, intents and attempts111 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[]>> {
Rust repos service with shipping; pull requests kept in the model157 const repo = await this.repos.get(repoPath, viewer);
Initial g1t: services, event bus, intents and attempts158 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>> {
Rust repos service with shipping; pull requests kept in the model175 const repo = await this.repos.get(repoPath, viewer);
Initial g1t: services, event bus, intents and attempts176 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;
Rust repos service with shipping; pull requests kept in the model194 const repo = await this.repos.getById(intent.repoId, actor);
Initial g1t: services, event bus, intents and attempts195 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>> {
Email verification, password reset, and Git for AI scale positioning221 if (!actor.verified) return UNVERIFIED;
Initial g1t: services, event bus, intents and attempts222 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");
Rust repos service with shipping; pull requests kept in the model230 const fork = await this.repos.forkForAttempt(intent.repoId, id, actor);
Initial g1t: services, event bus, intents and attempts231 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;
Rust repos service with shipping; pull requests kept in the model281 const repo = await this.repos.getById(attempt.repoId, viewer);
Initial g1t: services, event bus, intents and attempts282 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
Rust repos service with shipping; pull requests kept in the model331 async shipAttempt(actor: User, attemptId: string): Promise<Result<Attempt>> {
332 const row = await this.db
333 .prepare("SELECT * FROM attempts WHERE id = ?")
334 .bind(attemptId)
335 .first<AttemptRow>();
336 if (!row) return NO_ATTEMPT;
337 const attempt = toAttempt(row);
338 // Whether the actor may see and write the repo is decided by repos.
339 const repo = await this.repos.getById(attempt.repoId, actor);
340 if (!repo.ok) return NO_ATTEMPT;
341 if (attempt.status !== "working" && attempt.status !== "submitted") {
342 return fail("conflict", `This attempt is already ${attempt.status}.`);
343 }
344 const intent = (await this.intentById(attempt.intentId))!;
345 if (intent.status !== "open") {
346 return fail("conflict", `This intent is already ${intent.status}.`);
347 }
348
349 const landed = await this.repos.land(row.fork_repo_id, actor);
350 if (!landed.ok) return landed;
351
352 const now = Date.now();
353 await this.db.batch([
354 this.db
355 .prepare(
356 "UPDATE attempts SET status = 'shipped', head_commit = ?, updated_at = ? WHERE id = ?",
357 )
358 .bind(landed.value.commit, now, attempt.id),
359 this.db
360 .prepare("UPDATE intents SET status = 'shipped' WHERE id = ?")
361 .bind(intent.id),
362 ]);
363 const ids = {
364 attemptId: attempt.id,
365 intentId: intent.id,
366 repoId: attempt.repoId,
367 };
368 await this.publish(
369 {
370 type: "attempt.shipped",
371 source: SOURCE,
372 repoId: attempt.repoId,
373 actor: actor.id,
374 data: { ...ids, commit: landed.value.commit },
375 },
376 {
377 type: "intent.closed",
378 source: SOURCE,
379 repoId: attempt.repoId,
380 actor: actor.id,
381 data: { intentId: intent.id, repoId: attempt.repoId, reason: "shipped" },
382 },
383 );
384 return ok({
385 ...attempt,
386 status: "shipped",
387 headCommit: landed.value.commit,
388 updatedAt: now,
389 });
390 }
391
Initial g1t: services, event bus, intents and attempts392 async listActiveAttempts(
393 viewer: Viewer,
394 ): Promise<{ attempt: Attempt; intent: Intent }[]> {
395 if (!viewer) return [];
396 const { results } = await this.db
397 .prepare(
398 `SELECT * FROM attempts
399 WHERE started_by_id = ? AND status IN ('working', 'submitted')
400 ORDER BY updated_at DESC LIMIT 50`,
401 )
402 .bind(viewer.id)
403 .all<AttemptRow>();
404 return Promise.all(
405 results.map(async (row) => ({
406 attempt: toAttempt(row),
407 intent: (await this.intentById(row.intent_id))!,
408 })),
409 );
410 }
411
412 async appendSession(
413 actor: User,
414 attemptId: string,
415 entries: NewSessionEntry[],
416 ): Promise<Result<{ count: number }>> {
417 if (!actor) return SIGN_IN;
418 if (entries.length === 0) return ok({ count: 0 });
419 if (entries.length > MAX_ENTRY_BATCH) {
420 return fail("invalid", `Send at most ${MAX_ENTRY_BATCH} entries at a time.`);
421 }
422 const own = await this.ownAttempt(actor, attemptId);
423 if (!own.ok) return own;
424 const attempt = own.value;
425
426 const now = Date.now();
427 // Each insert takes the next sequence number itself, so two writers
428 // appending at once cannot collide.
429 const insert = this.db.prepare(
430 `INSERT INTO session_entries (attempt_id, seq, kind, text, tool, "commit", at)
431 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
432 FROM session_entries WHERE attempt_id = ?`,
433 );
434 await this.db.batch([
435 ...entries.map((entry) =>
436 insert.bind(
437 attempt.id,
438 entry.kind,
439 entry.text.slice(0, MAX_ENTRY_CHARS),
440 entry.tool ?? null,
441 entry.commit ?? attempt.headCommit,
442 entry.at ?? now,
443 attempt.id,
444 ),
445 ),
446 this.db
447 .prepare("UPDATE attempts SET updated_at = ? WHERE id = ?")
448 .bind(now, attempt.id),
449 ]);
450 await this.publish({
451 type: "session.appended",
452 source: SOURCE,
453 repoId: attempt.repoId,
454 actor: actor.id,
455 data: { attemptId: attempt.id, sessionId: attempt.id, count: entries.length },
456 });
457 return ok({ count: entries.length });
458 }
459
460 async readSession(
461 attemptId: string,
462 viewer: Viewer,
463 afterSeq = 0,
464 ): Promise<Result<SessionEntry[]>> {
465 const found = await this.getAttempt(attemptId, viewer);
466 if (!found.ok) return found;
467 const { results } = await this.db
468 .prepare(
469 "SELECT * FROM session_entries WHERE attempt_id = ? AND seq > ? ORDER BY seq LIMIT ?",
470 )
471 .bind(attemptId, afterSeq, SESSION_PAGE)
472 .all<SessionRow>();
473 return ok(results.map(toSessionEntry));
474 }
475
476 /** A push to an attempt's fork moves that attempt's head. */
477 async onEvents(events: G1tEvent[]): Promise<void> {
478 for (const event of events) {
479 if (event.type !== "git.push") continue;
480 await this.db
481 .prepare(
482 "UPDATE attempts SET head_commit = ?, updated_at = ? WHERE fork_repo_id = ?",
483 )
484 .bind(event.data.after, event.time, event.data.repoId)
485 .run();
486 }
487 }
488
489 async queue(batch: MessageBatch<G1tEvent>): Promise<void> {
490 await this.onEvents(batch.messages.map((message) => message.body));
491 }
492}