pr_01m47d15m3e54sn21z27rpy5n9/services/work/src/index.ts
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 attempts | 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, | |
| Rust repos service with shipping; pull requests kept in the model | 15 | type ServiceBinding, |
| Initial g1t: services, event bus, intents and attempts | 16 | 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 positioning | 22 | UNVERIFIED, |
| Initial g1t: services, event bus, intents and attempts | 23 | fail, |
| 24 | newId, | |
| 25 | ok, | |
| Rust repos service with shipping; pull requests kept in the model | 26 | reposClient, |
| Initial g1t: services, event bus, intents and attempts | 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; | |
| Rust repos service with shipping; pull requests kept in the model | 40 | REPOS: ServiceBinding; |
| Initial g1t: services, event bus, intents and attempts | 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 | ||
| Rust repos service with shipping; pull requests kept in the model | 64 | private get repos(): ReposApi { |
| 65 | return reposClient(this.env.REPOS); | |
| 66 | } | |
| 67 | ||
| Initial g1t: services, event bus, intents and attempts | 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. | |
| Rust repos service with shipping; pull requests kept in the model | 89 | const repo = await this.repos.getById(attempt.repoId, actor); |
| Initial g1t: services, event bus, intents and attempts | 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>> { | |
| Email verification, password reset, and Git for AI scale positioning | 106 | if (!actor.verified) return UNVERIFIED; |
| Initial g1t: services, event bus, intents and attempts | 107 | 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 model | 110 | const repo = await this.repos.get(repoPath, actor); |
| Initial g1t: services, event bus, intents and attempts | 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[]>> { | |
| Rust repos service with shipping; pull requests kept in the model | 157 | const repo = await this.repos.get(repoPath, viewer); |
| Initial g1t: services, event bus, intents and attempts | 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>> { | |
| Rust repos service with shipping; pull requests kept in the model | 175 | const repo = await this.repos.get(repoPath, viewer); |
| Initial g1t: services, event bus, intents and attempts | 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; | |
| Rust repos service with shipping; pull requests kept in the model | 194 | const repo = await this.repos.getById(intent.repoId, actor); |
| Initial g1t: services, event bus, intents and attempts | 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>> { | |
| Email verification, password reset, and Git for AI scale positioning | 221 | if (!actor.verified) return UNVERIFIED; |
| Initial g1t: services, event bus, intents and attempts | 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"); | |
| Rust repos service with shipping; pull requests kept in the model | 230 | const fork = await this.repos.forkForAttempt(intent.repoId, id, actor); |
| Initial g1t: services, event bus, intents and attempts | 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; | |
| Rust repos service with shipping; pull requests kept in the model | 281 | const repo = await this.repos.getById(attempt.repoId, viewer); |
| Initial g1t: services, event bus, intents and attempts | 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 | ||
| Rust repos service with shipping; pull requests kept in the model | 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 = ?, 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 attempts | 392 | 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 | } |