g1t/services/work/src/index.ts

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