pr_01m47d24b0e6n91zwymwxg0vpx/services/work/src/lib.rs

2,021 lines78,466 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.

Issues and pull requests replace intents and attempts1//! The work service: issues, pull requests, comments and sessions.
Work service in Rust, with RFC 3339 timestamps2//!
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
Agents and memory, checks and conflicts, profiles, slug renames, custom domains7mod authored;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API8mod capture;
Acceptance checks in sandboxes, line comments and review verdicts9mod checks;
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look10mod compute;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API11mod guardrails;
Agents as a team: lifecycle, merge queue, billing and a new shell12mod lifecycle;
Agents and memory, checks and conflicts, profiles, slug renames, custom domains13mod memory;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API14mod mentions;
Agents and memory, checks and conflicts, profiles, slug renames, custom domains15mod mergeability;
Agents as a team: lifecycle, merge queue, billing and a new shell16mod plans;
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request17mod messages;
Agents as a team: lifecycle, merge queue, billing and a new shell18mod queue;
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look19mod retired;
Agents as a team: lifecycle, merge queue, billing and a new shell20mod reviews;
Work service in Rust, with RFC 3339 timestamps21mod rows;
Agents and memory, checks and conflicts, profiles, slug renames, custom domains22mod runs;
Agents as a team: lifecycle, merge queue, billing and a new shell23mod settings;
GitHub Actions on g1t, part two: running workflows24mod statuses;
Work service in Rust, with RFC 3339 timestamps25
26use g1t_contracts::events::{
Events service in Rust, with RFC 3339 times and accurate push events27 CommentCreated, Event, IssueEvent, NewEvent, Publish, PullEvent, SessionAppended,
Work service in Rust, with RFC 3339 timestamps28};
Agents as a team: lifecycle, merge queue, billing and a new shell29use g1t_contracts::identity::UsernameArgs;
Catching up with main takes seconds when the two sides touched different files30use g1t_contracts::repos::{
31 ForkArgs, GetArgs, HeadArgs, LandArgs, Landed, NeedsAgentReason, PullBranchUpdate, Repo, RepoPath,
32 UpdatePullBranchArgs,
33};
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look34use g1t_contracts::access::{self, Capability, Denied};
Work service in Rust, with RFC 3339 timestamps35use g1t_contracts::time::rfc3339;
36use g1t_contracts::work::*;
Agents as a team: lifecycle, merge queue, billing and a new shell37use futures_util::future::{try_join, try_join3, try_join_all};
Work service in Rust, with RFC 3339 timestamps38use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
Events service in Rust, with RFC 3339 times and accurate push events39use g1t_kit::{args, now_ms, reply, rpc_method};
Work service in Rust, with RFC 3339 timestamps40use serde::Serialize;
41use worker::wasm_bindgen::JsValue;
42use worker::{
43 Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event,
44};
45
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look46use retired::writable;
Agents as a team: lifecycle, merge queue, billing and a new shell47use rows::{CommentRow, IssueRow, MovedRow, NumberRow, PullRow, SessionRow, Snapshot, ValueRow};
Work service in Rust, with RFC 3339 timestamps48
49const SOURCE: &str = "work";
50const MAX_ENTRY_BATCH: usize = 200;
51const MAX_ENTRY_CHARS: usize = 64_000;
Issues and pull requests replace intents and attempts52const MAX_TITLE_CHARS: usize = 200;
Work service in Rust, with RFC 3339 timestamps53const SESSION_PAGE: u32 = 500;
Issues and pull requests replace intents and attempts54const LIST_PAGE: u32 = 100;
Agents as a team: lifecycle, merge queue, billing and a new shell55const MAX_ASSIGNEES: usize = 10;
Work service in Rust, with RFC 3339 timestamps56const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
57
Issues and pull requests replace intents and attempts58const ISSUE_COLUMNS: &str = "issues.*,
59 (SELECT count(*) FROM pulls WHERE pulls.issue_id = issues.id) AS pull_count,
Agents as a team: lifecycle, merge queue, billing and a new shell60 (SELECT agent FROM pulls
61 WHERE pulls.issue_id = issues.id AND pulls.status IN ('draft', 'open')
62 AND pulls.fork_repo_id IS NOT NULL
63 ORDER BY pulls.number DESC LIMIT 1) AS agent,
Issues and pull requests replace intents and attempts64 (SELECT count(*) FROM comments
Agents as a team: lifecycle, merge queue, billing and a new shell65 WHERE comments.repo_id = issues.repo_id AND comments.number = issues.number
66 AND comments.kind = 'comment') AS comment_count";
Work service in Rust, with RFC 3339 timestamps67
Issues and pull requests replace intents and attempts68fn no_issue<T>() -> Outcome<T> {
69 Outcome::fail(FailureCode::NotFound, "Issue not found.")
Work service in Rust, with RFC 3339 timestamps70}
71
Issues and pull requests replace intents and attempts72fn no_pull<T>() -> Outcome<T> {
73 Outcome::fail(FailureCode::NotFound, "Pull request not found.")
Work service in Rust, with RFC 3339 timestamps74}
75
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look76/// Refuses `actor` unless their role on `repo` has `capability`: not found
77/// when they cannot read it, forbidden with the role it needs otherwise.
78pub(crate) fn allowed(actor: Option<&User>, repo: &Repo, capability: Capability) -> Outcome<()> {
79 match access::check(actor, repo, capability) {
80 Ok(()) => Outcome::Ok(()),
81 Err(Denied::NotFound) => Outcome::fail(FailureCode::NotFound, "Repository not found."),
82 Err(Denied::Forbidden) => Outcome::fail(
83 FailureCode::Forbidden,
84 access::needs(capability, &format!("{}/{}", repo.namespace, repo.name)),
85 ),
86 }
87}
88
Work service in Rust, with RFC 3339 timestamps89fn optional(value: &Option<String>) -> JsValue {
90 value.as_deref().map_or(JsValue::NULL, JsValue::from)
91}
92
Issues and pull requests replace intents and attempts93fn optional_number(value: Option<u32>) -> JsValue {
94 value.map_or(JsValue::NULL, JsValue::from)
95}
96
97/// The lowercase name a `State` is stored and sent as.
98fn state_name(state: Option<State>) -> Option<&'static str> {
99 state.map(|state| match state {
100 State::Open => "open",
101 State::Closed => "closed",
102 })
103}
104
105/// A trimmed title, or why it cannot be used.
106fn valid_title(title: &str) -> std::result::Result<&str, &'static str> {
107 let title = title.trim();
108 if title.is_empty() {
109 Err("A title is required.")
110 } else if title.chars().count() > MAX_TITLE_CHARS {
111 Err("That title is too long.")
112 } else {
113 Ok(title)
114 }
115}
116
117/// Unwraps an `Outcome`, returning its failure from the enclosing method.
118macro_rules! check {
119 ($outcome:expr) => {
120 match $outcome {
121 Outcome::Ok(value) => value,
122 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
123 }
124 };
125}
126
Work service in Rust, with RFC 3339 timestamps127struct Work {
128 db: D1Database,
Agents as a team: lifecycle, merge queue, billing and a new shell129 identity: Fetcher,
Work service in Rust, with RFC 3339 timestamps130 repos: Fetcher,
Events service in Rust, with RFC 3339 times and accurate push events131 events: Fetcher,
Sidebar: the panels really slide132 /// GitHub Actions: runs a merge queue's `merge_group` workflows.
133 actions: Fetcher,
Work service in Rust, with RFC 3339 timestamps134}
135
136impl Work {
Issues and pull requests replace intents and attempts137 async fn publish<T: Serialize>(
138 &self,
139 kind: &'static str,
140 repo_id: &str,
141 actor: &User,
142 data: T,
143 ) -> Result<()> {
Acceptance checks in sandboxes, line comments and review verdicts144 self.publish_as(kind, repo_id, Some(actor.id.clone()), data)
145 .await
146 }
147
Agents move along on private repositories too148 /// A pull request's author as a viewer who can read its repository and
149 /// source. Stored authors carry no memberships, so a private repository
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look150 /// would otherwise look missing to them. The membership given reads
151 /// and nothing more: it is for looking, never for acting.
Agents move along on private repositories too152 pub(crate) async fn author_viewer(&self, pull: &Pull) -> Result<Viewer> {
153 let path: Option<RepoPath> = g1t_kit::call(
154 &self.repos,
155 "path_by_id",
156 &g1t_contracts::repos::PathByIdArgs { id: pull.repo_id.clone() },
157 )
158 .await?;
159 let mut author = pull.author.clone();
160 if let Some(path) = path
161 && !author.is_member(&path.namespace.to_lowercase())
162 {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look163 author.workspaces.push(g1t_contracts::Membership {
164 base_permission: Some(access::BasePermission::Read),
165 ..g1t_contracts::Membership::member(path.namespace.to_lowercase())
166 });
Agents move along on private repositories too167 }
168 Ok(Some(author))
169 }
170
Acceptance checks in sandboxes, line comments and review verdicts171 /// Publishes an event caused by `actor`, or by g1t itself.
172 async fn publish_as<T: Serialize>(
173 &self,
174 kind: &'static str,
175 repo_id: &str,
176 actor: Option<String>,
177 data: T,
178 ) -> Result<()> {
Issues and pull requests replace intents and attempts179 let event = NewEvent {
180 kind,
181 source: SOURCE,
182 repo_id: Some(repo_id.to_owned()),
Acceptance checks in sandboxes, line comments and review verdicts183 actor,
Issues and pull requests replace intents and attempts184 data,
185 };
Events service in Rust, with RFC 3339 times and accurate push events186 g1t_kit::call(
187 &self.events,
188 "publish",
189 &Publish {
190 events: vec![event],
191 },
192 )
193 .await
Work service in Rust, with RFC 3339 timestamps194 }
195
Issues and pull requests replace intents and attempts196 /// The repository, if the viewer may see it. Whether they may is
197 /// decided by the repos service.
198 async fn repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Outcome<Repo>> {
Work service in Rust, with RFC 3339 timestamps199 g1t_kit::call(
200 &self.repos,
201 "get",
202 &GetArgs {
203 path: path.clone(),
204 viewer: viewer.clone(),
205 },
206 )
207 .await
208 }
209
Issues and pull requests replace intents and attempts210 /// The next number in the repository's sequence. Taking it is one
211 /// statement, so concurrent opens cannot be given the same number.
212 async fn next_number(&self, repo_id: &str) -> Result<u32> {
213 let row = self
214 .db
215 .prepare(
216 "INSERT INTO counters (repo_id, last) VALUES (?, 1)
217 ON CONFLICT (repo_id) DO UPDATE SET last = last + 1
218 RETURNING last AS n",
219 )
220 .bind(&[repo_id.into()])?
221 .first::<NumberRow>(None)
222 .await?;
223 row.map(|row| row.n)
224 .ok_or_else(|| worker::Error::RustError("no number was assigned".into()))
Work service in Rust, with RFC 3339 timestamps225 }
226
Issues and pull requests replace intents and attempts227 async fn issue(&self, repo_id: &str, number: u32) -> Result<Option<Issue>> {
Work service in Rust, with RFC 3339 timestamps228 Ok(self
229 .db
Issues and pull requests replace intents and attempts230 .prepare(format!(
231 "SELECT {ISSUE_COLUMNS} FROM issues WHERE repo_id = ? AND number = ?"
232 ))
233 .bind(&[repo_id.into(), number.into()])?
234 .first::<IssueRow>(None)
Work service in Rust, with RFC 3339 timestamps235 .await?
Issues and pull requests replace intents and attempts236 .map(Issue::from))
Work service in Rust, with RFC 3339 timestamps237 }
238
Issues and pull requests replace intents and attempts239 async fn pull(&self, repo_id: &str, number: u32) -> Result<Option<Pull>> {
Work service in Rust, with RFC 3339 timestamps240 Ok(self
241 .db
Issues and pull requests replace intents and attempts242 .prepare("SELECT * FROM pulls WHERE repo_id = ? AND number = ?")
243 .bind(&[repo_id.into(), number.into()])?
244 .first::<PullRow>(None)
245 .await?
246 .map(Pull::from))
247 }
248
249 async fn comments(&self, repo_id: &str, number: u32) -> Result<Vec<Comment>> {
250 let rows = self
251 .db
252 .prepare(
253 "SELECT * FROM comments WHERE repo_id = ? AND number = ? ORDER BY id LIMIT 500",
254 )
255 .bind(&[repo_id.into(), number.into()])?
256 .all()
Work service in Rust, with RFC 3339 timestamps257 .await?
Issues and pull requests replace intents and attempts258 .results::<CommentRow>()?;
259 Ok(rows.into_iter().map(Comment::from).collect())
Work service in Rust, with RFC 3339 timestamps260 }
261
Issues and pull requests replace intents and attempts262 /// The repository and one of its issues, as seen by `viewer`.
263 async fn issue_at(
264 &self,
265 path: &RepoPath,
266 number: u32,
267 viewer: &Viewer,
268 ) -> Result<Outcome<(Repo, Issue)>> {
269 let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
270 return Ok(no_issue());
Work service in Rust, with RFC 3339 timestamps271 };
Issues and pull requests replace intents and attempts272 Ok(match self.issue(&repo.id, number).await? {
273 Some(issue) => Outcome::Ok((repo, issue)),
274 None => no_issue(),
275 })
276 }
277
278 /// The repository and one of its pull requests, as seen by `viewer`.
279 async fn pull_at(
280 &self,
281 path: &RepoPath,
282 number: u32,
283 viewer: &Viewer,
284 ) -> Result<Outcome<(Repo, Pull)>> {
285 let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
286 return Ok(no_pull());
287 };
288 Ok(match self.pull(&repo.id, number).await? {
289 Some(pull) => Outcome::Ok((repo, pull)),
290 None => no_pull(),
291 })
292 }
293
Agents as a team: lifecycle, merge queue, billing and a new shell294 /// Records something that happened to an issue or a pull request, so
295 /// that it shows in the conversation where it happened. `text` is what
296 /// `author` did, as the rest of a sentence starting with their name.
297 pub(crate) async fn note(
298 &self,
299 repo_id: &str,
300 number: u32,
301 author: (&str, &str),
302 text: &str,
303 ) -> Result<()> {
304 let now = now_ms();
305 self.db
306 .prepare(
307 "INSERT INTO comments
308 (id, repo_id, number, author_id, author_name, body, kind, created_at)
309 VALUES (?, ?, ?, ?, ?, ?, 'event', ?)",
310 )
311 .bind(&[
312 new_id("cmt", now).into(),
313 repo_id.into(),
314 number.into(),
315 author.0.into(),
316 author.1.into(),
317 text.into(),
318 rfc3339(now).into(),
319 ])?
320 .run()
321 .await?;
322 Ok(())
323 }
324
325 /// Notes who was added to and removed from a list of people, such as
326 /// "assigned ana" or "requested a review from g1t-agent".
327 async fn note_changes(
328 &self,
329 repo_id: &str,
330 number: u32,
331 actor: &User,
332 before: &[String],
333 after: &[String],
334 (added, removed): (&str, &str),
335 ) -> Result<()> {
336 let joined = |names: Vec<&String>| {
337 names
338 .into_iter()
339 .map(String::as_str)
340 .collect::<Vec<_>>()
341 .join(", ")
342 };
343 let new: Vec<&String> = after.iter().filter(|name| !before.contains(name)).collect();
344 let gone: Vec<&String> = before.iter().filter(|name| !after.contains(name)).collect();
345 let who = (actor.id.as_str(), actor.username.as_str());
346 if !new.is_empty() {
347 // Taking something on oneself reads better said that way.
348 let text = if added == "assigned" && new == [&actor.username] {
349 "self-assigned this".to_owned()
350 } else {
351 format!("{added} {}", joined(new))
352 };
353 self.note(repo_id, number, who, &text).await?;
354 }
355 if !gone.is_empty() {
356 self.note(repo_id, number, who, &format!("{removed} {}", joined(gone)))
357 .await?;
358 }
359 Ok(())
360 }
361
Issues and pull requests replace intents and attempts362 fn issue_event(issue: &Issue) -> IssueEvent {
363 IssueEvent {
364 issue_id: issue.id.clone(),
365 repo_id: issue.repo_id.clone(),
366 number: issue.number,
367 ..IssueEvent::default()
Work service in Rust, with RFC 3339 timestamps368 }
369 }
370
Workflows run when an agent's pull request is marked ready371 /// The commit a pull request's change is at in git right now: its
372 /// fork's default branch, or its branch.
373 async fn live_head(&self, pull: &Pull) -> Result<Option<String>> {
374 g1t_kit::call(
375 &self.repos,
376 "head",
377 &HeadArgs {
378 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| pull.repo_id.clone()),
379 branch: pull.branch.clone().unwrap_or_default(),
380 },
381 )
382 .await
383 }
384
Issues and pull requests replace intents and attempts385 fn pull_event(pull: &Pull) -> PullEvent {
386 PullEvent {
387 pull_id: pull.id.clone(),
388 repo_id: pull.repo_id.clone(),
389 number: pull.number,
390 issue: pull.issue,
391 ..PullEvent::default()
Work service in Rust, with RFC 3339 timestamps392 }
393 }
394
Issues and pull requests replace intents and attempts395 // --- Issues ------------------------------------------------------------
396
397 async fn open_issue(&self, a: OpenIssueArgs) -> Result<Outcome<Issue>> {
Work service in Rust, with RFC 3339 timestamps398 if !a.actor.verified {
399 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
400 }
Issues and pull requests replace intents and attempts401 let title = match valid_title(&a.title) {
402 Ok(title) => title,
403 Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
404 };
405 let Some(labels) = normalize_labels(&a.labels) else {
Work service in Rust, with RFC 3339 timestamps406 return Ok(Outcome::fail(
407 FailureCode::Invalid,
Issues and pull requests replace intents and attempts408 "An issue can have up to 10 labels of up to 40 characters each.",
Work service in Rust, with RFC 3339 timestamps409 ));
410 };
Issues and pull requests replace intents and attempts411 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look412 check!(writable(&repo));
Work service in Rust, with RFC 3339 timestamps413 let checks: Vec<&str> = a
414 .checks
415 .iter()
416 .map(|check| check.trim())
417 .filter(|check| !check.is_empty())
418 .collect();
419
420 let now = now_ms();
Issues and pull requests replace intents and attempts421 let id = new_id("iss", now);
422 let number = self.next_number(&repo.id).await?;
423 let timestamp = rfc3339(now);
Work service in Rust, with RFC 3339 timestamps424 self.db
425 .prepare(
Issues and pull requests replace intents and attempts426 "INSERT INTO issues
427 (id, repo_id, number, title, body, labels, checks, author_id, author_name,
428 created_at, updated_at)
429 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
Work service in Rust, with RFC 3339 timestamps430 )
431 .bind(&[
432 id.as_str().into(),
433 repo.id.as_str().into(),
Issues and pull requests replace intents and attempts434 number.into(),
Work service in Rust, with RFC 3339 timestamps435 title.into(),
Issues and pull requests replace intents and attempts436 a.body.trim().into(),
437 serde_json::to_string(&labels)?.into(),
Work service in Rust, with RFC 3339 timestamps438 serde_json::to_string(&checks)?.into(),
439 a.actor.id.as_str().into(),
440 a.actor.username.as_str().into(),
Issues and pull requests replace intents and attempts441 timestamp.as_str().into(),
442 timestamp.as_str().into(),
Work service in Rust, with RFC 3339 timestamps443 ])?
444 .run()
445 .await?;
Issues and pull requests replace intents and attempts446 let Some(issue) = self.issue(&repo.id, number).await? else {
447 return Ok(no_issue());
Work service in Rust, with RFC 3339 timestamps448 };
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API449 self.apply_label_rule(&a.actor, &issue, &[]).await?;
Issues and pull requests replace intents and attempts450 self.publish(
451 "issue.opened",
452 &repo.id,
453 &a.actor,
454 IssueEvent {
455 title: Some(issue.title.clone()),
456 ..Self::issue_event(&issue)
Work service in Rust, with RFC 3339 timestamps457 },
Issues and pull requests replace intents and attempts458 )
Work service in Rust, with RFC 3339 timestamps459 .await?;
Issues and pull requests replace intents and attempts460 Ok(Outcome::Ok(issue))
Work service in Rust, with RFC 3339 timestamps461 }
462
Issues and pull requests replace intents and attempts463 async fn list_issues(&self, a: ListIssuesArgs) -> Result<Outcome<Vec<Issue>>> {
464 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
465 let state = state_name(a.state).map_or(JsValue::NULL, JsValue::from);
466 let label = a
467 .label
468 .map(|label| label.trim().to_lowercase())
469 .filter(|label| !label.is_empty());
Work service in Rust, with RFC 3339 timestamps470 let rows = self
471 .db
472 .prepare(format!(
Issues and pull requests replace intents and attempts473 "SELECT {ISSUE_COLUMNS} FROM issues
474 WHERE repo_id = ? AND (? IS NULL OR state = ?)
475 AND (? IS NULL OR EXISTS
476 (SELECT 1 FROM json_each(issues.labels) WHERE json_each.value = ?))
477 ORDER BY number DESC LIMIT ?"
Work service in Rust, with RFC 3339 timestamps478 ))
Issues and pull requests replace intents and attempts479 .bind(&[
480 repo.id.into(),
481 state.clone(),
482 state,
483 optional(&label),
484 optional(&label),
485 LIST_PAGE.into(),
486 ])?
Work service in Rust, with RFC 3339 timestamps487 .all()
488 .await?
Issues and pull requests replace intents and attempts489 .results::<IssueRow>()?;
490 Ok(Outcome::Ok(rows.into_iter().map(Issue::from).collect()))
Work service in Rust, with RFC 3339 timestamps491 }
492
Issues and pull requests replace intents and attempts493 async fn get_issue(&self, a: ViewArgs) -> Result<Outcome<IssueDetail>> {
494 let (repo, issue) = check!(self.issue_at(&a.repo, a.number, &a.viewer).await?);
495 let pulls = self
Work service in Rust, with RFC 3339 timestamps496 .db
Issues and pull requests replace intents and attempts497 .prepare("SELECT * FROM pulls WHERE issue_id = ? ORDER BY number")
498 .bind(&[issue.id.as_str().into()])?
Work service in Rust, with RFC 3339 timestamps499 .all()
500 .await?
Issues and pull requests replace intents and attempts501 .results::<PullRow>()?;
502 Ok(Outcome::Ok(IssueDetail {
503 comments: self.comments(&repo.id, issue.number).await?,
504 pulls: pulls.into_iter().map(Pull::from).collect(),
505 issue,
Work service in Rust, with RFC 3339 timestamps506 }))
507 }
508
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look509 /// The issue, if `actor` wrote it or may triage the repository's issues.
Issues and pull requests replace intents and attempts510 async fn manageable_issue(
511 &self,
512 actor: &User,
513 path: &RepoPath,
514 number: u32,
515 ) -> Result<Outcome<Issue>> {
516 let (repo, issue) = check!(self.issue_at(path, number, &Some(actor.clone())).await?);
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look517 check!(writable(&repo));
518 if issue.author.id != actor.id {
519 check!(allowed(Some(actor), &repo, Capability::Triage));
Work service in Rust, with RFC 3339 timestamps520 }
Issues and pull requests replace intents and attempts521 Ok(Outcome::Ok(issue))
522 }
523
524 async fn update_issue(&self, a: UpdateIssueArgs) -> Result<Outcome<Issue>> {
525 let issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
526 let title = match a.title.as_deref().map(valid_title) {
527 Some(Err(message)) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
528 Some(Ok(title)) => Some(title.to_owned()),
529 None => None,
530 };
531 let labels = match a.labels.as_deref().map(normalize_labels) {
532 Some(None) => {
533 return Ok(Outcome::fail(
534 FailureCode::Invalid,
535 "An issue can have up to 10 labels of up to 40 characters each.",
536 ));
537 }
538 Some(Some(labels)) => Some(serde_json::to_string(&labels)?),
539 None => None,
540 };
Agents as a team: lifecycle, merge queue, billing and a new shell541 let assignees = match a.assignees {
542 Some(names) => Some(check!(self.valid_assignees(names).await?)),
543 None => None,
544 };
545 let assigned = assignees.as_ref().map(serde_json::to_string).transpose()?;
Issues and pull requests replace intents and attempts546 let body = a.body.map(|body| body.trim().to_owned());
547 self.db
548 .prepare(
549 "UPDATE issues
550 SET title = COALESCE(?, title), body = COALESCE(?, body),
Agents as a team: lifecycle, merge queue, billing and a new shell551 labels = COALESCE(?, labels), assignees = COALESCE(?, assignees),
552 updated_at = ?
Issues and pull requests replace intents and attempts553 WHERE id = ?",
554 )
555 .bind(&[
556 optional(&title),
557 optional(&body),
558 optional(&labels),
Agents as a team: lifecycle, merge queue, billing and a new shell559 optional(&assigned),
Issues and pull requests replace intents and attempts560 rfc3339(now_ms()).into(),
561 issue.id.as_str().into(),
562 ])?
563 .run()
564 .await?;
Agents as a team: lifecycle, merge queue, billing and a new shell565 let before = issue.assignees.clone();
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API566 let labels_before = issue.labels.clone();
Issues and pull requests replace intents and attempts567 let Some(issue) = self.issue(&issue.repo_id, issue.number).await? else {
568 return Ok(no_issue());
569 };
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API570 self.apply_label_rule(&a.actor, &issue, &labels_before).await?;
Issues and pull requests replace intents and attempts571 self.publish(
572 "issue.updated",
573 &issue.repo_id,
574 &a.actor,
575 Self::issue_event(&issue),
576 )
577 .await?;
Agents as a team: lifecycle, merge queue, billing and a new shell578 if let Some(assignees) = assignees {
579 self.note_changes(
580 &issue.repo_id,
581 issue.number,
582 &a.actor,
583 &before,
584 &assignees,
585 ("assigned", "unassigned"),
586 )
587 .await?;
588 self.publish(
589 "issue.assigned",
590 &issue.repo_id,
591 &a.actor,
592 IssueEvent {
593 assignees: Some(assignees),
594 ..Self::issue_event(&issue)
595 },
596 )
597 .await?;
598 }
Issues and pull requests replace intents and attempts599 Ok(Outcome::Ok(issue))
600 }
601
Agents as a team: lifecycle, merge queue, billing and a new shell602 /// Usernames as given, tidied, if each names an account.
603 async fn valid_assignees(&self, names: Vec<String>) -> Result<Outcome<Vec<String>>> {
604 let mut assignees: Vec<String> = Vec::new();
605 for name in names {
606 let name = name.trim().trim_start_matches('@').to_lowercase();
607 if name.is_empty() || assignees.contains(&name) {
608 continue;
609 }
610 if assignees.len() == MAX_ASSIGNEES {
611 return Ok(Outcome::fail(
612 FailureCode::Invalid,
613 format!("An issue can be assigned to at most {MAX_ASSIGNEES} people."),
614 ));
615 }
616 let account: Viewer = g1t_kit::call(
617 &self.identity,
618 "user_by_username",
619 &UsernameArgs {
620 username: name.clone(),
621 },
622 )
623 .await?;
624 if account.is_none() {
625 return Ok(Outcome::fail(
626 FailureCode::Invalid,
627 format!("There is no account named {name}."),
628 ));
629 }
630 assignees.push(name);
631 }
632 Ok(Outcome::Ok(assignees))
633 }
634
635 /// Open issues assigned to the viewer, in every repository. Callers
636 /// show only those in repositories the viewer can still see.
637 async fn list_assigned_issues(&self, a: ViewerArgs) -> Result<Vec<Issue>> {
638 let Some(viewer) = a.viewer else {
639 return Ok(Vec::new());
640 };
641 let rows = self
642 .db
643 .prepare(format!(
644 "SELECT {ISSUE_COLUMNS} FROM issues
645 WHERE state = 'open' AND EXISTS (
646 SELECT 1 FROM json_each(issues.assignees) WHERE json_each.value = ?)
647 ORDER BY updated_at DESC LIMIT 50"
648 ))
649 .bind(&[viewer.username.into()])?
650 .all()
651 .await?
652 .results::<IssueRow>()?;
653 Ok(rows.into_iter().map(Issue::from).collect())
654 }
655
Issues and pull requests replace intents and attempts656 async fn close_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
657 let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
658 if issue.state == State::Closed {
Work service in Rust, with RFC 3339 timestamps659 return Ok(Outcome::fail(
660 FailureCode::Conflict,
Issues and pull requests replace intents and attempts661 "This issue is already closed.",
Work service in Rust, with RFC 3339 timestamps662 ));
663 }
Issues and pull requests replace intents and attempts664 let reason = a.reason.unwrap_or(IssueReason::Completed);
665 let now = rfc3339(now_ms());
Work service in Rust, with RFC 3339 timestamps666 self.db
Issues and pull requests replace intents and attempts667 .prepare(
668 "UPDATE issues SET state = 'closed', reason = ?, closed_at = ?, updated_at = ?
669 WHERE id = ?",
670 )
671 .bind(&[
672 reason.as_str().into(),
673 now.as_str().into(),
674 now.as_str().into(),
675 issue.id.as_str().into(),
676 ])?
Work service in Rust, with RFC 3339 timestamps677 .run()
678 .await?;
Issues and pull requests replace intents and attempts679 self.publish(
680 "issue.closed",
681 &issue.repo_id,
682 &a.actor,
683 IssueEvent {
684 reason: Some(reason.as_str()),
685 ..Self::issue_event(&issue)
Work service in Rust, with RFC 3339 timestamps686 },
Issues and pull requests replace intents and attempts687 )
Work service in Rust, with RFC 3339 timestamps688 .await?;
Agents as a team: lifecycle, merge queue, billing and a new shell689 self.note(
690 &issue.repo_id,
691 issue.number,
692 (&a.actor.id, &a.actor.username),
693 match reason {
694 IssueReason::Completed => "closed this as completed",
695 IssueReason::NotPlanned => "closed this as not planned",
696 },
697 )
698 .await?;
Issues and pull requests replace intents and attempts699 issue.state = State::Closed;
700 issue.reason = Some(reason);
701 issue.closed_at = Some(now.clone());
702 issue.updated_at = now;
703 Ok(Outcome::Ok(issue))
704 }
705
706 async fn reopen_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
707 let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
708 if issue.state == State::Open {
709 return Ok(Outcome::fail(
710 FailureCode::Conflict,
711 "This issue is already open.",
712 ));
713 }
714 let now = rfc3339(now_ms());
715 self.db
716 .prepare(
717 "UPDATE issues
718 SET state = 'open', reason = NULL, resolved_by = NULL, closed_at = NULL,
719 updated_at = ?
720 WHERE id = ?",
721 )
722 .bind(&[now.as_str().into(), issue.id.as_str().into()])?
723 .run()
724 .await?;
725 self.publish(
726 "issue.reopened",
727 &issue.repo_id,
728 &a.actor,
729 Self::issue_event(&issue),
730 )
731 .await?;
Agents as a team: lifecycle, merge queue, billing and a new shell732 self.note(
733 &issue.repo_id,
734 issue.number,
735 (&a.actor.id, &a.actor.username),
736 "reopened this",
737 )
738 .await?;
Issues and pull requests replace intents and attempts739 issue.state = State::Open;
740 issue.reason = None;
741 issue.resolved_by = None;
742 issue.closed_at = None;
743 issue.updated_at = now;
744 Ok(Outcome::Ok(issue))
Work service in Rust, with RFC 3339 timestamps745 }
746
Issues and pull requests replace intents and attempts747 /// The default labels, then every other label in use on the repository.
748 async fn list_labels(&self, a: ViewArgs) -> Result<Outcome<Vec<String>>> {
749 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
750 let used = self
751 .db
752 .prepare(
753 "SELECT DISTINCT json_each.value AS value
754 FROM issues, json_each(issues.labels)
755 WHERE issues.repo_id = ? ORDER BY 1 LIMIT 200",
756 )
757 .bind(&[repo.id.into()])?
758 .all()
759 .await?
760 .results::<ValueRow>()?;
761 let mut labels: Vec<String> = DEFAULT_LABELS.iter().map(|label| (*label).into()).collect();
762 for row in used {
763 if !labels.contains(&row.value) {
764 labels.push(row.value);
765 }
766 }
767 Ok(Outcome::Ok(labels))
768 }
769
770 async fn counts(&self, a: ViewArgs) -> Result<Outcome<Counts>> {
771 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
772 let counts = self
773 .db
774 .prepare(
775 "SELECT
776 (SELECT count(*) FROM issues WHERE repo_id = ? AND state = 'open') AS issues,
777 (SELECT count(*) FROM pulls
778 WHERE repo_id = ? AND status IN ('draft', 'open')) AS pulls",
779 )
780 .bind(&[repo.id.as_str().into(), repo.id.as_str().into()])?
781 .first::<Counts>(None)
782 .await?;
783 Ok(Outcome::Ok(counts.unwrap_or(Counts {
784 issues: 0,
785 pulls: 0,
786 })))
787 }
788
789 // --- Comments ----------------------------------------------------------
790
791 async fn add_comment(&self, a: AddCommentArgs) -> Result<Outcome<Comment>> {
Work service in Rust, with RFC 3339 timestamps792 if !a.actor.verified {
793 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
794 }
Issues and pull requests replace intents and attempts795 let body = a.body.trim();
Acceptance checks in sandboxes, line comments and review verdicts796 // An approval speaks for itself; anything else has to say something.
797 if body.is_empty() && a.verdict != Some(Verdict::Approve) {
Issues and pull requests replace intents and attempts798 return Ok(Outcome::fail(
799 FailureCode::Invalid,
800 "A comment cannot be empty.",
801 ));
802 }
Acceptance checks in sandboxes, line comments and review verdicts803 let path = a
804 .path
805 .as_deref()
806 .map(str::trim)
807 .filter(|path| !path.is_empty());
808 let line = a.line.filter(|line| *line > 0 && path.is_some());
Issues and pull requests replace intents and attempts809 if body.chars().count() > MAX_ENTRY_CHARS {
810 return Ok(Outcome::fail(
811 FailureCode::Invalid,
812 "That comment is too long.",
813 ));
814 }
815 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look816 check!(writable(&repo));
Issues and pull requests replace intents and attempts817 // The number names an issue or a pull request, never both.
Agents as a team: lifecycle, merge queue, billing and a new shell818 let mut pull_id = None;
Issues and pull requests replace intents and attempts819 let table = if self.issue(&repo.id, a.number).await?.is_some() {
Acceptance checks in sandboxes, line comments and review verdicts820 if path.is_some() || a.verdict.is_some() {
821 return Ok(Outcome::fail(
822 FailureCode::Invalid,
823 "Only a pull request can be reviewed or commented on by line.",
824 ));
825 }
Issues and pull requests replace intents and attempts826 "issues"
Acceptance checks in sandboxes, line comments and review verdicts827 } else if let Some(pull) = self.pull(&repo.id, a.number).await? {
828 if a.verdict.is_some() && pull.author.id == a.actor.id {
829 return Ok(Outcome::fail(
830 FailureCode::Forbidden,
831 "You cannot approve or request changes on your own pull request.",
832 ));
833 }
Agents as a team: lifecycle, merge queue, billing and a new shell834 pull_id = Some(pull.id.clone());
Issues and pull requests replace intents and attempts835 "pulls"
836 } else {
Work service in Rust, with RFC 3339 timestamps837 return Ok(Outcome::fail(
Issues and pull requests replace intents and attempts838 FailureCode::NotFound,
839 "No issue or pull request has that number.",
Work service in Rust, with RFC 3339 timestamps840 ));
Issues and pull requests replace intents and attempts841 };
842
843 let now = now_ms();
844 let comment = Comment {
Agents as a team: lifecycle, merge queue, billing and a new shell845 kind: CommentKind::Comment,
Issues and pull requests replace intents and attempts846 id: new_id("cmt", now),
847 author: a.actor.clone(),
848 body: body.to_owned(),
Acceptance checks in sandboxes, line comments and review verdicts849 path: path.map(str::to_owned),
850 line,
851 verdict: a.verdict,
Issues and pull requests replace intents and attempts852 created_at: rfc3339(now),
853 };
854 self.db
855 .batch(vec![
856 self.db
857 .prepare(
858 "INSERT INTO comments
Acceptance checks in sandboxes, line comments and review verdicts859 (id, repo_id, number, author_id, author_name, body, path, line,
860 verdict, created_at)
861 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
Issues and pull requests replace intents and attempts862 )
863 .bind(&[
864 comment.id.as_str().into(),
865 repo.id.as_str().into(),
866 a.number.into(),
867 a.actor.id.as_str().into(),
868 a.actor.username.as_str().into(),
869 body.into(),
Acceptance checks in sandboxes, line comments and review verdicts870 optional(&comment.path),
871 optional_number(line),
872 a.verdict
873 .map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
Issues and pull requests replace intents and attempts874 comment.created_at.as_str().into(),
875 ])?,
876 self.db
877 .prepare(format!(
878 "UPDATE {table} SET updated_at = ? WHERE repo_id = ? AND number = ?"
879 ))
880 .bind(&[
881 comment.created_at.as_str().into(),
882 repo.id.as_str().into(),
883 a.number.into(),
884 ])?,
885 ])
886 .await?;
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API887 self.note_mention(&a.actor, &repo, a.number, &comment, pull_id.as_deref()).await?;
Issues and pull requests replace intents and attempts888 self.publish(
889 "comment.created",
890 &repo.id,
891 &a.actor,
892 CommentCreated {
893 comment_id: comment.id.clone(),
894 repo_id: repo.id.clone(),
895 number: a.number,
Agents as a team: lifecycle, merge queue, billing and a new shell896 pull_id,
897 verdict: a.verdict,
Issues and pull requests replace intents and attempts898 },
899 )
900 .await?;
901 Ok(Outcome::Ok(comment))
902 }
903
904 // --- Pull requests -----------------------------------------------------
905
906 async fn open_pull(&self, a: OpenPullArgs) -> Result<Outcome<Pull>> {
907 if !a.actor.verified {
908 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
Work service in Rust, with RFC 3339 timestamps909 }
Issues and pull requests replace intents and attempts910 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look911 check!(writable(&repo));
912 // g1t's own agent at work spends the workspace's compute; a pull
913 // request anyone else's agent makes is like any other.
914 if matches!(a.runtime, Runtime::Hosted) {
915 check!(allowed(Some(&a.actor), &repo, Capability::Run));
916 }
Issues and pull requests replace intents and attempts917 let issue = match a.issue {
918 Some(number) => match self.issue(&repo.id, number).await? {
919 Some(issue) if issue.state == State::Open => Some(issue),
920 Some(_) => {
921 return Ok(Outcome::fail(
922 FailureCode::Conflict,
923 "This issue is closed.",
924 ));
925 }
926 None => return Ok(no_issue()),
927 },
928 None => None,
929 };
930 // A pull request for an issue takes the issue's title unless given one.
931 let title = match (a.title.trim(), &issue) {
932 ("", Some(issue)) => issue.title.clone(),
933 (title, _) => match valid_title(title) {
934 Ok(title) => title.to_owned(),
935 Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
936 },
937 };
Work service in Rust, with RFC 3339 timestamps938 let agent = match a.agent.trim() {
939 "" => "agent",
940 agent => agent,
941 };
942 let runtime = match a.runtime {
Issues and pull requests replace intents and attempts943 Runtime::Hosted => "hosted",
944 Runtime::External => "external",
Work service in Rust, with RFC 3339 timestamps945 };
946
947 let now = now_ms();
Issues and pull requests replace intents and attempts948 let id = new_id("pr", now);
Pull requests from branches949 let branch = a
950 .branch
951 .as_deref()
952 .map(str::trim)
953 .filter(|branch| !branch.is_empty());
954 // The change is on a branch already pushed to the repository, or
955 // will be made in a fork created for this pull request.
956 let (fork, head) = match branch {
957 Some(branch) => {
958 if branch == repo.default_branch {
959 return Ok(Outcome::fail(
960 FailureCode::Invalid,
961 format!("Choose a branch other than {branch}."),
962 ));
963 }
964 let head: Option<String> = g1t_kit::call(
965 &self.repos,
966 "head",
967 &HeadArgs {
968 repo_id: repo.id.clone(),
969 branch: branch.to_owned(),
970 },
971 )
972 .await?;
973 let Some(head) = head else {
974 return Ok(Outcome::fail(
975 FailureCode::NotFound,
976 format!("There is no branch named {branch}. Push it first."),
977 ));
978 };
979 let existing = self
980 .db
981 .prepare(
982 "SELECT number AS n FROM pulls
983 WHERE repo_id = ? AND source_branch = ? AND status IN ('draft', 'open')",
984 )
985 .bind(&[repo.id.as_str().into(), branch.into()])?
986 .first::<NumberRow>(None)
987 .await?;
988 if let Some(existing) = existing {
989 return Ok(Outcome::fail(
990 FailureCode::Conflict,
991 format!("Pull request #{} is already open for {branch}.", existing.n),
992 ));
993 }
994 (None, Some(head))
995 }
996 None => {
997 let fork: Outcome<Repo> = g1t_kit::call(
998 &self.repos,
999 "fork_for_pull",
1000 &ForkArgs {
1001 source_id: repo.id.clone(),
1002 pull_id: id.clone(),
1003 actor: a.actor.clone(),
1004 },
1005 )
1006 .await?;
1007 (Some(check!(fork)), None)
1008 }
1009 };
1010 // A branch already holds the work, so its pull request is ready for
1011 // review from the start; one with a fork starts as a draft.
1012 let status = if branch.is_some() { "open" } else { "draft" };
1013 let body = Some(a.body.trim().to_owned()).filter(|body| !body.is_empty());
Work service in Rust, with RFC 3339 timestamps1014
Issues and pull requests replace intents and attempts1015 let number = self.next_number(&repo.id).await?;
Work service in Rust, with RFC 3339 timestamps1016 let timestamp = rfc3339(now);
1017 self.db
1018 .prepare(
Issues and pull requests replace intents and attempts1019 "INSERT INTO pulls
Pull requests from branches1020 (id, repo_id, number, issue_id, issue_number, title, body, agent, runtime,
1021 status, fork_repo_id, fork_namespace, fork_name, source_branch, head_commit,
1022 author_id, author_name, created_at, updated_at)
1023 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
Work service in Rust, with RFC 3339 timestamps1024 )
1025 .bind(&[
1026 id.as_str().into(),
Issues and pull requests replace intents and attempts1027 repo.id.as_str().into(),
1028 number.into(),
1029 optional(&issue.as_ref().map(|issue| issue.id.clone())),
1030 optional_number(issue.as_ref().map(|issue| issue.number)),
1031 title.into(),
Pull requests from branches1032 optional(&body),
Work service in Rust, with RFC 3339 timestamps1033 agent.into(),
1034 runtime.into(),
Pull requests from branches1035 status.into(),
1036 optional(&fork.as_ref().map(|fork| fork.id.clone())),
1037 optional(&fork.as_ref().map(|fork| fork.namespace.clone())),
1038 optional(&fork.as_ref().map(|fork| fork.name.clone())),
1039 optional(&branch.map(str::to_owned)),
1040 optional(&head),
Work service in Rust, with RFC 3339 timestamps1041 a.actor.id.as_str().into(),
1042 a.actor.username.as_str().into(),
1043 timestamp.as_str().into(),
1044 timestamp.as_str().into(),
1045 ])?
1046 .run()
1047 .await?;
Issues and pull requests replace intents and attempts1048 let Some(pull) = self.pull(&repo.id, number).await? else {
1049 return Ok(no_pull());
Work service in Rust, with RFC 3339 timestamps1050 };
Agents as a team: lifecycle, merge queue, billing and a new shell1051 self.manage(&pull).await?;
1052 // Someone is on it now, so it is no longer waiting for an agent.
1053 if let Some(issue) = pull.issue {
1054 self.db
1055 .prepare("UPDATE issues SET queued_by = NULL WHERE repo_id = ? AND number = ?")
1056 .bind(&[repo.id.as_str().into(), issue.into()])?
1057 .run()
1058 .await?;
1059 }
1060 if let Some(issue) = pull.issue {
1061 let text = if lifecycle::made_by_g1t(&pull) {
1062 format!("assigned this to g1t-agent, which opened #{}", pull.number)
1063 } else {
1064 format!("opened #{} for this", pull.number)
1065 };
1066 self.note(&repo.id, issue, (&a.actor.id, &a.actor.username), &text)
1067 .await?;
1068 }
Issues and pull requests replace intents and attempts1069 self.publish(
1070 "pull.opened",
1071 &repo.id,
1072 &a.actor,
1073 PullEvent {
1074 agent: Some(pull.agent.clone()),
1075 ..Self::pull_event(&pull)
Work service in Rust, with RFC 3339 timestamps1076 },
Issues and pull requests replace intents and attempts1077 )
Work service in Rust, with RFC 3339 timestamps1078 .await?;
Issues and pull requests replace intents and attempts1079 Ok(Outcome::Ok(pull))
Work service in Rust, with RFC 3339 timestamps1080 }
1081
Issues and pull requests replace intents and attempts1082 async fn list_pulls(&self, a: ListPullsArgs) -> Result<Outcome<Vec<Pull>>> {
1083 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
1084 let filter = match a.state {
1085 Some(State::Open) => "AND status IN ('draft', 'open')",
1086 Some(State::Closed) => "AND status IN ('merged', 'closed')",
1087 None => "",
Work service in Rust, with RFC 3339 timestamps1088 };
Issues and pull requests replace intents and attempts1089 let rows = self
1090 .db
1091 .prepare(format!(
1092 "SELECT * FROM pulls WHERE repo_id = ? {filter} ORDER BY number DESC LIMIT ?"
1093 ))
1094 .bind(&[repo.id.into(), LIST_PAGE.into()])?
1095 .all()
1096 .await?
1097 .results::<PullRow>()?;
1098 Ok(Outcome::Ok(rows.into_iter().map(Pull::from).collect()))
1099 }
1100
1101 async fn get_pull(&self, a: ViewArgs) -> Result<Outcome<PullDetail>> {
1102 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
1103 let issue = match pull.issue {
1104 Some(number) => self.issue(&repo.id, number).await?,
1105 None => None,
Work service in Rust, with RFC 3339 timestamps1106 };
Agents as a team: lifecycle, merge queue, billing and a new shell1107 let mut pull = pull;
1108 // Worked out on each push; this covers a pull request from before
1109 // that was recorded.
1110 if pull.files.is_empty() && pull.head_commit.is_some() {
1111 pull.files = self.refresh_files(&pull).await?;
1112 }
1113 // Everything else at once: none of it depends on the rest, and each
1114 // is a round trip of its own.
1115 let standing = async {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains1116 // Mergeability first: where g1t sees a pull request through, a
1117 // conflict decides its next step.
1118 let (merge, behind) =
1119 try_join(self.mergeability(&pull), self.is_behind(&repo.id, &pull)).await?;
Agents as a team: lifecycle, merge queue, billing and a new shell1120 let lifecycle = self
1121 .assess(&pull, &issue, behind)
1122 .await?
1123 .map(|(lifecycle, _)| lifecycle);
Agents and memory, checks and conflicts, profiles, slug renames, custom domains1124 Ok::<_, worker::Error>((behind, lifecycle, merge))
Agents as a team: lifecycle, merge queue, billing and a new shell1125 };
Agents and memory, checks and conflicts, profiles, slug renames, custom domains1126 let (((behind, lifecycle, (mergeable, conflicts)), (landing, stalled), comments), (checks, overlaps, review_pending)) =
Agents as a team: lifecycle, merge queue, billing and a new shell1127 try_join(
1128 try_join3(standing, self.landing_state(&pull.id), self.comments(&repo.id, pull.number)),
1129 try_join3(
1130 self.latest_checks(&pull.id),
1131 self.overlaps(&pull),
1132 self.review_pending(&pull.id),
1133 ),
1134 )
1135 .await?;
Issues and pull requests replace intents and attempts1136 Ok(Outcome::Ok(PullDetail {
Agents as a team: lifecycle, merge queue, billing and a new shell1137 comments,
1138 checks,
1139 overlaps,
1140 behind,
1141 review_pending,
1142 lifecycle,
1143 landing,
1144 stalled,
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request1145 messages: self.messages(&pull.id).await?,
GitHub Actions on g1t, part two: running workflows1146 statuses: self.statuses(&repo.id, pull.head_commit.as_deref()).await?,
Agents and memory, checks and conflicts, profiles, slug renames, custom domains1147 mergeable,
1148 conflicts,
1149 earlier_checks: self.earlier_checks(&pull.id).await?,
Issues and pull requests replace intents and attempts1150 issue,
1151 pull,
1152 }))
Work service in Rust, with RFC 3339 timestamps1153 }
1154
Issues and pull requests replace intents and attempts1155 /// The pull request, if it is still active and `actor` opened it or
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1156 /// may triage the repository's pull requests.
Issues and pull requests replace intents and attempts1157 async fn manageable_pull(
Work service in Rust, with RFC 3339 timestamps1158 &self,
Issues and pull requests replace intents and attempts1159 actor: &User,
1160 path: &RepoPath,
1161 number: u32,
1162 ) -> Result<Outcome<Pull>> {
1163 let (repo, pull) = check!(self.pull_at(path, number, &Some(actor.clone())).await?);
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1164 check!(writable(&repo));
1165 if pull.author.id != actor.id {
1166 check!(allowed(Some(actor), &repo, Capability::Triage));
Issues and pull requests replace intents and attempts1167 }
1168 if !pull.status.is_active() {
1169 return Ok(Outcome::fail(
Work service in Rust, with RFC 3339 timestamps1170 FailureCode::Conflict,
Issues and pull requests replace intents and attempts1171 format!("This pull request is already {}.", pull.status.as_str()),
Work service in Rust, with RFC 3339 timestamps1172 ));
1173 }
Issues and pull requests replace intents and attempts1174 Ok(Outcome::Ok(pull))
1175 }
1176
Catching up with main takes seconds when the two sides touched different files1177 /// Brings a pull request up to date with the default branch without a
1178 /// sandbox, where the repos service can do that safely. Whoever could
1179 /// have pushed the merge themselves may ask: whoever opened it, for a
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1180 /// fork; anyone who may push, for a branch of the repository. When it needs a
Catching up with main takes seconds when the two sides touched different files1181 /// real merge, says so, naming the conflicting files if a probe found
1182 /// them, and pushes nothing.
1183 async fn catch_up_pull(&self, a: PullActionArgs) -> Result<Outcome<PullBranchUpdate>> {
1184 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &Some(a.actor.clone())).await?);
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1185 check!(writable(&repo));
Catching up with main takes seconds when the two sides touched different files1186 if !pull.status.is_active() {
1187 return Ok(Outcome::fail(
1188 FailureCode::Conflict,
1189 format!("This pull request is already {}.", pull.status.as_str()),
1190 ));
1191 }
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1192 if pull.fork_repo_id.is_some() {
1193 if pull.author.id != a.actor.id {
1194 return Ok(Outcome::fail(
1195 FailureCode::Forbidden,
1196 "Only whoever opened this pull request can update it.",
1197 ));
1198 }
Catching up with main takes seconds when the two sides touched different files1199 } else {
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1200 check!(allowed(Some(&a.actor), &repo, Capability::Push));
Catching up with main takes seconds when the two sides touched different files1201 }
1202 let updated: Outcome<PullBranchUpdate> = g1t_kit::call(
1203 &self.repos,
1204 "update_pull_branch",
1205 &UpdatePullBranchArgs {
1206 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
1207 branch: pull.branch.clone(),
1208 number: pull.number,
1209 actor: a.actor,
1210 },
1211 )
1212 .await?;
1213 // A probe that found conflicts says more than "both changed it".
1214 if let Outcome::Ok(PullBranchUpdate::NeedsAgent { .. }) = &updated
1215 && let Some(files) = self.conflicting_files(&pull).await?
1216 && !files.is_empty()
1217 {
1218 return Ok(Outcome::Ok(PullBranchUpdate::NeedsAgent {
1219 reason: NeedsAgentReason::Conflicting,
1220 detail: "Merging it conflicts.".to_owned(),
1221 paths: files,
1222 }));
1223 }
1224 Ok(updated)
1225 }
1226
Agents as a team: lifecycle, merge queue, billing and a new shell1227 async fn update_pull(&self, a: UpdatePullArgs) -> Result<Outcome<Pull>> {
1228 let pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1229 let assignees = match a.assignees {
1230 Some(names) => Some(check!(self.valid_assignees(names).await?)),
1231 None => None,
1232 };
1233 let reviewers = match a.reviewers {
1234 Some(names) => {
1235 // A g1t agent is not an account; everyone else has to be.
1236 let agent = names
1237 .iter()
1238 .any(|name| name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME));
1239 let people = names
1240 .into_iter()
1241 .filter(|name| !name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME))
1242 .collect();
1243 let mut reviewers = check!(self.valid_assignees(people).await?);
1244 reviewers.retain(|name| *name != pull.author.username);
1245 if agent {
1246 reviewers.insert(0, reviews::AGENT_NAME.to_owned());
1247 }
1248 Some(reviewers)
1249 }
1250 None => None,
1251 };
1252 self.db
1253 .prepare(
1254 "UPDATE pulls
1255 SET assignees = COALESCE(?, assignees), reviewers = COALESCE(?, reviewers),
1256 updated_at = ?
1257 WHERE id = ?",
1258 )
1259 .bind(&[
1260 optional(&assignees.as_ref().map(serde_json::to_string).transpose()?),
1261 optional(&reviewers.as_ref().map(serde_json::to_string).transpose()?),
1262 rfc3339(now_ms()).into(),
1263 pull.id.as_str().into(),
1264 ])?
1265 .run()
1266 .await?;
1267 if let Some(assignees) = &assignees {
1268 self.note_changes(
1269 &pull.repo_id,
1270 pull.number,
1271 &a.actor,
1272 &pull.assignees,
1273 assignees,
1274 ("assigned", "unassigned"),
1275 )
1276 .await?;
1277 }
1278 if let Some(reviewers) = &reviewers {
1279 self.note_changes(
1280 &pull.repo_id,
1281 pull.number,
1282 &a.actor,
1283 &pull.reviewers,
1284 reviewers,
1285 (
1286 "requested a review from",
1287 "withdrew the request for a review from",
1288 ),
1289 )
1290 .await?;
1291 }
1292 Ok(match self.pull(&pull.repo_id, pull.number).await? {
1293 Some(pull) => Outcome::Ok(pull),
1294 None => no_pull(),
1295 })
1296 }
1297
Issues and pull requests replace intents and attempts1298 /// Marks a draft ready for review, or updates the description of one
1299 /// that already is.
1300 async fn ready_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1301 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
Work service in Rust, with RFC 3339 timestamps1302 let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
1303 let now = rfc3339(now_ms());
1304 self.db
1305 .prepare(
Issues and pull requests replace intents and attempts1306 "UPDATE pulls SET status = 'open', body = COALESCE(?, body), updated_at = ?
Work service in Rust, with RFC 3339 timestamps1307 WHERE id = ?",
1308 )
1309 .bind(&[
1310 optional(&summary),
1311 now.as_str().into(),
Issues and pull requests replace intents and attempts1312 pull.id.as_str().into(),
Work service in Rust, with RFC 3339 timestamps1313 ])?
1314 .run()
1315 .await?;
Issues and pull requests replace intents and attempts1316 if pull.status == PullStatus::Draft {
Workflows run when an agent's pull request is marked ready1317 // The head as it is now: the push that came just before may not
1318 // have reached `head_commit` yet, and workflows run on it.
1319 let commit = self.live_head(&pull).await?.or_else(|| pull.head_commit.clone());
Issues and pull requests replace intents and attempts1320 self.publish(
1321 "pull.ready",
1322 &pull.repo_id,
1323 &a.actor,
Workflows run when an agent's pull request is marked ready1324 PullEvent {
1325 commit,
1326 ..Self::pull_event(&pull)
1327 },
Issues and pull requests replace intents and attempts1328 )
1329 .await?;
1330 }
Agents as a team: lifecycle, merge queue, billing and a new shell1331 if pull.status == PullStatus::Draft {
1332 self.note(
1333 &pull.repo_id,
1334 pull.number,
1335 (&a.actor.id, &a.actor.username),
1336 "marked this ready for review",
1337 )
1338 .await?;
1339 }
Issues and pull requests replace intents and attempts1340 pull.status = PullStatus::Open;
1341 pull.body = summary.or(pull.body);
1342 pull.updated_at = now;
1343 Ok(Outcome::Ok(pull))
1344 }
1345
1346 async fn close_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1347 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1348 let now = rfc3339(now_ms());
1349 self.db
1350 .prepare("UPDATE pulls SET status = 'closed', updated_at = ? WHERE id = ?")
1351 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
1352 .run()
1353 .await?;
1354 self.publish(
1355 "pull.closed",
1356 &pull.repo_id,
1357 &a.actor,
1358 Self::pull_event(&pull),
1359 )
Work service in Rust, with RFC 3339 timestamps1360 .await?;
Agents as a team: lifecycle, merge queue, billing and a new shell1361 self.note(
1362 &pull.repo_id,
1363 pull.number,
1364 (&a.actor.id, &a.actor.username),
1365 "closed this",
1366 )
1367 .await?;
1368 // A closed pull request leaves the merge queue.
1369 if self
1370 .leave(&pull.repo_id, &pull, QueueState::Removed, Some("It was closed."))
1371 .await?
1372 {
1373 self.publish_as(
1374 "queue.changed",
1375 &pull.repo_id,
1376 None,
1377 g1t_contracts::events::QueueChanged {
1378 repo_id: pull.repo_id.clone(),
1379 },
1380 )
1381 .await?;
1382 }
Issues and pull requests replace intents and attempts1383 pull.status = PullStatus::Closed;
1384 pull.updated_at = now;
1385 Ok(Outcome::Ok(pull))
Work service in Rust, with RFC 3339 timestamps1386 }
1387
Issues and pull requests replace intents and attempts1388 /// Lands the pull request on the repository's default branch. Unless
1389 /// told to keep it open, that resolves the issue it was for: the issue
1390 /// closes naming this pull request, and the others still in progress
1391 /// for it close as superseded.
1392 async fn merge_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
Work service in Rust, with RFC 3339 timestamps1393 let viewer = Some(a.actor.clone());
Agents as a team: lifecycle, merge queue, billing and a new shell1394 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1395 check!(writable(&repo));
Issues and pull requests replace intents and attempts1396 match pull.status {
1397 PullStatus::Open => {}
1398 PullStatus::Draft => {
1399 return Ok(Outcome::fail(
1400 FailureCode::Conflict,
1401 "This pull request is still a draft. Mark it ready for review first.",
1402 ));
1403 }
1404 status => {
1405 return Ok(Outcome::fail(
1406 FailureCode::Conflict,
1407 format!("This pull request is already {}.", status.as_str()),
1408 ));
1409 }
Work service in Rust, with RFC 3339 timestamps1410 }
Agents as a team: lifecycle, merge queue, billing and a new shell1411 let settings = self.settings(&repo.id).await?;
1412 // Where the repository does not allow it, asking to ignore the
1413 // checks changes nothing.
1414 if !a.ignore_checks || !settings.allow_ignoring_checks {
Acceptance checks in sandboxes, line comments and review verdicts1415 let waiting = match pull.check_status {
1416 Some(CheckStatus::Queued | CheckStatus::Running) => {
1417 Some("The acceptance checks are still running.")
1418 }
1419 Some(CheckStatus::Failed) => Some("The acceptance checks did not pass."),
1420 Some(CheckStatus::Errored) => Some("The acceptance checks could not be run."),
1421 Some(CheckStatus::Passed) | None => None,
1422 };
GitHub Actions on g1t, part two: running workflows1423 // Workflows run on its head count as checks too.
1424 let workflows = statuses::WorkflowFacts::of(&self.statuses(&repo.id, pull.head_commit.as_deref()).await?).refusal();
1425 let waiting = waiting.map(str::to_owned).or(workflows);
Acceptance checks in sandboxes, line comments and review verdicts1426 if let Some(reason) = waiting {
Agents as a team: lifecycle, merge queue, billing and a new shell1427 let remedy = if settings.allow_ignoring_checks {
1428 "Wait or fix them, or merge anyway by ignoring the checks."
1429 } else {
1430 "This repository only merges pull requests whose checks pass."
1431 };
Acceptance checks in sandboxes, line comments and review verdicts1432 return Ok(Outcome::fail(
1433 FailureCode::Conflict,
Agents as a team: lifecycle, merge queue, billing and a new shell1434 format!("{reason} {remedy}"),
Acceptance checks in sandboxes, line comments and review verdicts1435 ));
1436 }
1437 }
Agents as a team: lifecycle, merge queue, billing and a new shell1438 if let Some(missing) = self.approvals_gap(&settings, &pull).await? {
1439 return Ok(Outcome::fail(FailureCode::Conflict, missing));
1440 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains1441 // Known ahead of time to conflict: neither a merge nor the queue
1442 // would get through, so say what has to be resolved now.
1443 if let Some(files) = self.conflicting_files(&pull).await? {
1444 let named = if files.is_empty() {
1445 String::new()
1446 } else {
1447 format!(" in {}", files.join(", "))
1448 };
1449 return Ok(Outcome::fail(
1450 FailureCode::Conflict,
1451 format!(
1452 "This branch has conflicts with {}{named} that must be resolved first. Have the g1t agent resolve them, or merge {0} into it, fix them and push.",
1453 repo.default_branch
1454 ),
1455 ));
1456 }
Work service in Rust, with RFC 3339 timestamps1457
Agents as a team: lifecycle, merge queue, billing and a new shell1458 // A repository that merges through a queue: it joins the queue, and
1459 // lands once its state together with everything ahead has passed.
1460 if settings.merge_queue {
1461 if !a.actor.verified {
1462 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1463 }
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1464 check!(allowed(Some(&a.actor), &repo, Capability::Merge));
Agents as a team: lifecycle, merge queue, billing and a new shell1465 return self.enqueue(&repo, &pull, &a.actor, a.keep_issue_open).await;
1466 }
1467
1468 // The default branch has moved under it. Unless the repository
1469 // insists on that being dealt with first, bring it up to date and
1470 // land it when that is done.
1471 if self.is_behind(&repo.id, &pull).await? {
1472 if settings.require_up_to_date {
1473 return Ok(Outcome::fail(
1474 FailureCode::Conflict,
1475 format!(
1476 "{} has moved since this pull request was made, and this repository requires pull requests to be up to date before they merge. Catch up with {0} first.",
1477 repo.default_branch
1478 ),
1479 ));
1480 }
1481 if !a.actor.verified {
1482 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1483 }
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1484 check!(allowed(Some(&a.actor), &repo, Capability::Merge));
Agents as a team: lifecycle, merge queue, billing and a new shell1485 self.request_landing(&pull, &a.actor, a.keep_issue_open)
1486 .await?;
1487 return Ok(Outcome::Ok(pull));
1488 }
1489
Issues and pull requests replace intents and attempts1490 // Whether the actor may write to the repository is decided by repos.
Work service in Rust, with RFC 3339 timestamps1491 let landed: Outcome<Landed> = g1t_kit::call(
1492 &self.repos,
1493 "land",
1494 &LandArgs {
Pull requests from branches1495 // A pull request from a branch lands from the repository itself.
1496 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
1497 branch: pull.branch.clone(),
Work service in Rust, with RFC 3339 timestamps1498 actor: a.actor.clone(),
1499 },
1500 )
1501 .await?;
Issues and pull requests replace intents and attempts1502 let landed = check!(landed);
Agents as a team: lifecycle, merge queue, billing and a new shell1503 Ok(Outcome::Ok(
1504 self.record_merge(&repo, pull, &a.actor, a.keep_issue_open, landed)
1505 .await?,
1506 ))
1507 }
Work service in Rust, with RFC 3339 timestamps1508
Agents as a team: lifecycle, merge queue, billing and a new shell1509 /// Records a pull request as merged once the default branch holds it:
1510 /// closes its issue, supersedes the others for it, and says so.
1511 pub(crate) async fn record_merge(
1512 &self,
1513 repo: &Repo,
1514 mut pull: Pull,
1515 actor: &User,
1516 keep_issue_open: bool,
1517 landed: Landed,
1518 ) -> Result<Pull> {
1519 let issue = match pull.issue {
1520 Some(number) if !keep_issue_open => self
1521 .issue(&repo.id, number)
1522 .await?
1523 .filter(|issue| issue.state == State::Open),
1524 _ => None,
1525 };
Work service in Rust, with RFC 3339 timestamps1526 let now = rfc3339(now_ms());
Issues and pull requests replace intents and attempts1527 let mut statements = vec![
1528 self.db
1529 .prepare(
1530 "UPDATE pulls
1531 SET status = 'merged', head_commit = ?, merge_base = ?, merged_by = ?,
1532 merged_at = ?, updated_at = ?
1533 WHERE id = ?",
1534 )
1535 .bind(&[
1536 landed.commit.as_str().into(),
1537 optional(&landed.previous),
Agents as a team: lifecycle, merge queue, billing and a new shell1538 actor.username.as_str().into(),
Issues and pull requests replace intents and attempts1539 now.as_str().into(),
1540 now.as_str().into(),
1541 pull.id.as_str().into(),
1542 ])?,
1543 ];
1544 if let Some(issue) = &issue {
1545 statements.push(
Work service in Rust, with RFC 3339 timestamps1546 self.db
1547 .prepare(
Issues and pull requests replace intents and attempts1548 "UPDATE issues
1549 SET state = 'closed', reason = 'completed', resolved_by = ?,
1550 closed_at = ?, updated_at = ?
Work service in Rust, with RFC 3339 timestamps1551 WHERE id = ?",
1552 )
1553 .bind(&[
Issues and pull requests replace intents and attempts1554 pull.number.into(),
1555 now.as_str().into(),
Work service in Rust, with RFC 3339 timestamps1556 now.as_str().into(),
Issues and pull requests replace intents and attempts1557 issue.id.as_str().into(),
Work service in Rust, with RFC 3339 timestamps1558 ])?,
Issues and pull requests replace intents and attempts1559 );
1560 statements.push(
Work service in Rust, with RFC 3339 timestamps1561 self.db
Issues and pull requests replace intents and attempts1562 .prepare(
1563 "UPDATE pulls SET status = 'closed', superseded_by = ?, updated_at = ?
1564 WHERE issue_id = ? AND id != ? AND status IN ('draft', 'open')",
1565 )
1566 .bind(&[
1567 pull.number.into(),
1568 now.as_str().into(),
1569 issue.id.as_str().into(),
1570 pull.id.as_str().into(),
1571 ])?,
1572 );
1573 }
1574 self.db.batch(statements).await?;
1575
1576 self.publish(
1577 "pull.merged",
1578 &repo.id,
Agents as a team: lifecycle, merge queue, billing and a new shell1579 actor,
Issues and pull requests replace intents and attempts1580 PullEvent {
Work service in Rust, with RFC 3339 timestamps1581 commit: Some(landed.commit.clone()),
Issues and pull requests replace intents and attempts1582 ..Self::pull_event(&pull)
Work service in Rust, with RFC 3339 timestamps1583 },
Issues and pull requests replace intents and attempts1584 )
Work service in Rust, with RFC 3339 timestamps1585 .await?;
Issues and pull requests replace intents and attempts1586 if let Some(issue) = &issue {
1587 self.publish(
1588 "issue.closed",
1589 &repo.id,
Agents as a team: lifecycle, merge queue, billing and a new shell1590 actor,
Issues and pull requests replace intents and attempts1591 IssueEvent {
1592 reason: Some(IssueReason::Completed.as_str()),
1593 resolved_by: Some(pull.number),
1594 ..Self::issue_event(issue)
1595 },
1596 )
1597 .await?;
1598 }
Work service in Rust, with RFC 3339 timestamps1599
Agents as a team: lifecycle, merge queue, billing and a new shell1600 let who = (actor.id.as_str(), actor.username.as_str());
1601 self.note(&repo.id, pull.number, who, "merged this").await?;
1602 if let Some(issue) = &issue {
1603 self.note(
1604 &repo.id,
1605 issue.number,
1606 who,
1607 &format!("closed this by merging #{}", pull.number),
1608 )
1609 .await?;
1610 }
Issues and pull requests replace intents and attempts1611 pull.status = PullStatus::Merged;
Agents as a team: lifecycle, merge queue, billing and a new shell1612 pull.head_commit = Some(landed.commit.clone());
Issues and pull requests replace intents and attempts1613 pull.merge_base = landed.previous;
Agents as a team: lifecycle, merge queue, billing and a new shell1614 pull.merged_by = Some(actor.username.clone());
Issues and pull requests replace intents and attempts1615 pull.merged_at = Some(now.clone());
1616 pull.updated_at = now;
Agents as a team: lifecycle, merge queue, billing and a new shell1617 Ok(pull)
Work service in Rust, with RFC 3339 timestamps1618 }
1619
Issues and pull requests replace intents and attempts1620 async fn list_active_pulls(&self, a: ViewerArgs) -> Result<Vec<ActivePull>> {
Work service in Rust, with RFC 3339 timestamps1621 let Some(viewer) = a.viewer else {
1622 return Ok(Vec::new());
1623 };
Agents as a team: lifecycle, merge queue, billing and a new shell1624 let found = self
Work service in Rust, with RFC 3339 timestamps1625 .db
1626 .prepare(
Issues and pull requests replace intents and attempts1627 "SELECT * FROM pulls
1628 WHERE author_id = ? AND status IN ('draft', 'open')
Work service in Rust, with RFC 3339 timestamps1629 ORDER BY updated_at DESC LIMIT 50",
1630 )
1631 .bind(&[viewer.id.into()])?
1632 .all()
Agents as a team: lifecycle, merge queue, billing and a new shell1633 .await?;
1634 let snapshots = found.results::<Snapshot>()?;
1635 let pulls: Vec<Pull> = found.results::<PullRow>()?.into_iter().map(Pull::from).collect();
1636 // Each one at once: its issue, and where it stands. That is the
1637 // remembered assessment when there is one, and worked out otherwise.
1638 try_join_all(pulls.into_iter().zip(snapshots).map(|(pull, snapshot)| async move {
Issues and pull requests replace intents and attempts1639 let issue = match pull.issue {
1640 Some(number) => self.issue(&pull.repo_id, number).await?,
1641 None => None,
1642 };
Agents as a team: lifecycle, merge queue, billing and a new shell1643 // Only a pull request g1t is seeing through has a lifecycle.
1644 let lifecycle = if !lifecycle::made_by_g1t(&pull) || snapshot.managed == 0 {
1645 None
1646 } else if let (Some(stage), Some(detail)) = (snapshot.stage, snapshot.stage_detail) {
1647 Some(Lifecycle {
1648 stage,
1649 detail,
1650 revisions: snapshot.revisions,
1651 })
1652 } else {
1653 let behind = self.is_behind(&pull.repo_id, &pull).await?;
1654 self.assess(&pull, &issue, behind)
1655 .await?
1656 .map(|(lifecycle, _)| lifecycle)
1657 };
1658 Ok::<_, worker::Error>(ActivePull {
1659 pull,
1660 issue,
1661 lifecycle,
1662 })
1663 }))
1664 .await
Work service in Rust, with RFC 3339 timestamps1665 }
1666
Issues and pull requests replace intents and attempts1667 // --- Sessions ----------------------------------------------------------
1668
Work service in Rust, with RFC 3339 timestamps1669 async fn append_session(&self, a: AppendSessionArgs) -> Result<Outcome<Appended>> {
1670 if a.entries.is_empty() {
1671 return Ok(Outcome::Ok(Appended { count: 0 }));
1672 }
1673 if a.entries.len() > MAX_ENTRY_BATCH {
1674 return Ok(Outcome::fail(
1675 FailureCode::Invalid,
1676 format!("Send at most {MAX_ENTRY_BATCH} entries at a time."),
1677 ));
1678 }
Issues and pull requests replace intents and attempts1679 let viewer = Some(a.actor.clone());
1680 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
1681 if pull.author.id != a.actor.id {
1682 return Ok(Outcome::fail(
1683 FailureCode::Forbidden,
1684 "Only whoever opened a pull request can record its session.",
1685 ));
1686 }
Work service in Rust, with RFC 3339 timestamps1687
1688 let now = rfc3339(now_ms());
1689 let count = a.entries.len() as u32;
1690 let mut statements = Vec::with_capacity(a.entries.len() + 1);
1691 for entry in a.entries {
1692 let kind = serde_json::to_value(entry.kind)?;
1693 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
1694 // Each insert takes the next sequence number itself, so two
1695 // writers appending at once cannot collide.
1696 statements.push(
1697 self.db
1698 .prepare(
Issues and pull requests replace intents and attempts1699 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
Work service in Rust, with RFC 3339 timestamps1700 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
Issues and pull requests replace intents and attempts1701 FROM session_entries WHERE pull_id = ?",
Work service in Rust, with RFC 3339 timestamps1702 )
1703 .bind(&[
Issues and pull requests replace intents and attempts1704 pull.id.as_str().into(),
Work service in Rust, with RFC 3339 timestamps1705 kind.as_str().unwrap_or("note").into(),
1706 text.into(),
1707 optional(&entry.tool),
Issues and pull requests replace intents and attempts1708 optional(&entry.commit.or_else(|| pull.head_commit.clone())),
Work service in Rust, with RFC 3339 timestamps1709 now.as_str().into(),
Issues and pull requests replace intents and attempts1710 pull.id.as_str().into(),
Work service in Rust, with RFC 3339 timestamps1711 ])?,
1712 );
1713 }
1714 statements.push(
1715 self.db
Issues and pull requests replace intents and attempts1716 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
1717 .bind(&[now.as_str().into(), pull.id.as_str().into()])?,
Work service in Rust, with RFC 3339 timestamps1718 );
1719 self.db.batch(statements).await?;
Issues and pull requests replace intents and attempts1720 self.publish(
1721 "session.appended",
1722 &pull.repo_id,
1723 &a.actor,
1724 SessionAppended {
1725 pull_id: pull.id.clone(),
1726 repo_id: pull.repo_id.clone(),
1727 number: pull.number,
Work service in Rust, with RFC 3339 timestamps1728 count,
1729 },
Issues and pull requests replace intents and attempts1730 )
Work service in Rust, with RFC 3339 timestamps1731 .await?;
1732 Ok(Outcome::Ok(Appended { count }))
1733 }
1734
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request1735 /// Adds entries to a pull request's session, each taking the next
1736 /// sequence number, without announcing it.
1737 pub(crate) async fn append_entries(&self, pull: &Pull, entries: &[NewSessionEntry]) -> Result<()> {
1738 let now = rfc3339(now_ms());
1739 let mut statements = Vec::with_capacity(entries.len());
1740 for entry in entries {
1741 let kind = serde_json::to_value(entry.kind)?;
1742 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
1743 statements.push(
1744 self.db
1745 .prepare(
1746 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
1747 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
1748 FROM session_entries WHERE pull_id = ?",
1749 )
1750 .bind(&[
1751 pull.id.as_str().into(),
1752 kind.as_str().unwrap_or("note").into(),
1753 text.into(),
1754 optional(&entry.tool),
1755 optional(&entry.commit.clone().or_else(|| pull.head_commit.clone())),
1756 now.as_str().into(),
1757 pull.id.as_str().into(),
1758 ])?,
1759 );
1760 }
1761 self.db.batch(statements).await?;
1762 Ok(())
1763 }
1764
Issues and pull requests replace intents and attempts1765 async fn read_session(&self, a: ViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
1766 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
Work service in Rust, with RFC 3339 timestamps1767 let rows = self
1768 .db
1769 .prepare(
1770 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
Issues and pull requests replace intents and attempts1771 WHERE pull_id = ? AND seq > ? ORDER BY seq LIMIT ?",
Work service in Rust, with RFC 3339 timestamps1772 )
Issues and pull requests replace intents and attempts1773 .bind(&[pull.id.into(), a.after_seq.into(), SESSION_PAGE.into()])?
Work service in Rust, with RFC 3339 timestamps1774 .all()
1775 .await?
1776 .results::<SessionRow>()?;
1777 Ok(Outcome::Ok(
1778 rows.into_iter().map(SessionEntry::from).collect(),
1779 ))
1780 }
1781
Events service in Rust, with RFC 3339 times and accurate push events1782 /// A push moves the head of the pull request it concerns: the one whose
1783 /// fork was pushed to, or the one opened from the branch that moved.
1784 async fn on_event(&self, event: &Event) -> Result<()> {
Work service in Rust, with RFC 3339 timestamps1785 if event.kind != "git.push" {
1786 return Ok(());
1787 }
Events service in Rust, with RFC 3339 times and accurate push events1788 let (Some(repo_id), Some(after), Some(git_ref)) = (
1789 event.repo_id.as_deref(),
1790 event.data["after"].as_str(),
1791 event.data["ref"].as_str(),
1792 ) else {
Work service in Rust, with RFC 3339 timestamps1793 return Ok(());
1794 };
Pull requests from branches1795 let now = rfc3339(now_ms());
Agents as a team: lifecycle, merge queue, billing and a new shell1796 // The head moved, so whatever the checks said no longer applies, and
1797 // whatever step g1t was waiting on has been taken.
Acceptance checks in sandboxes, line comments and review verdicts1798 let moved = "UPDATE pulls
Agents as a team: lifecycle, merge queue, billing and a new shell1799 SET head_commit = ?, updated_at = ?, check_status = NULL, check_run_id = NULL,
1800 working_on = NULL, working_until = NULL, stalled = NULL";
Acceptance checks in sandboxes, line comments and review verdicts1801 let active = "status IN ('draft', 'open') AND head_commit IS NOT ?";
1802 let returning = "RETURNING id, repo_id, number, issue_number, status";
1803 let mut pulls: Vec<MovedRow> = Vec::new();
Events service in Rust, with RFC 3339 times and accurate push events1804 // A fork carries its pull request on its default branch.
1805 if event.data["defaultBranch"].as_bool() == Some(true) {
Acceptance checks in sandboxes, line comments and review verdicts1806 pulls.extend(
Events service in Rust, with RFC 3339 times and accurate push events1807 self.db
1808 .prepare(format!(
Acceptance checks in sandboxes, line comments and review verdicts1809 "{moved} WHERE fork_repo_id = ? AND {active} {returning}"
Events service in Rust, with RFC 3339 times and accurate push events1810 ))
Acceptance checks in sandboxes, line comments and review verdicts1811 .bind(&[
1812 after.into(),
1813 now.as_str().into(),
1814 repo_id.into(),
1815 after.into(),
1816 ])?
1817 .all()
1818 .await?
1819 .results::<MovedRow>()?,
Events service in Rust, with RFC 3339 times and accurate push events1820 );
1821 }
1822 if let Some(branch) = git_ref.strip_prefix("refs/heads/") {
Acceptance checks in sandboxes, line comments and review verdicts1823 pulls.extend(
Pull requests from branches1824 self.db
Events service in Rust, with RFC 3339 times and accurate push events1825 .prepare(format!(
Acceptance checks in sandboxes, line comments and review verdicts1826 "{moved} WHERE repo_id = ? AND source_branch = ? AND {active} {returning}"
Events service in Rust, with RFC 3339 times and accurate push events1827 ))
1828 .bind(&[
1829 after.into(),
1830 now.as_str().into(),
1831 repo_id.into(),
1832 branch.into(),
Acceptance checks in sandboxes, line comments and review verdicts1833 after.into(),
1834 ])?
1835 .all()
1836 .await?
1837 .results::<MovedRow>()?,
Events service in Rust, with RFC 3339 times and accurate push events1838 );
Pull requests from branches1839 }
Agents as a team: lifecycle, merge queue, billing and a new shell1840 // What each now changes, so overlaps show while the work is under way.
1841 for moved in &pulls {
1842 if let Some(pull) = self.pull_by_id(&moved.id).await? {
1843 self.refresh_files(&pull).await?;
1844 }
1845 }
1846 // A merge that was waiting for this push to bring it up to date.
1847 for moved in &pulls {
1848 self.land_if_requested(&moved.id).await?;
1849 }
Agents and memory, checks and conflicts, profiles, slug renames, custom domains1850 // Whether each still merges cleanly, and, when a default branch
1851 // moved, every open pull request into it.
1852 let moved_ids: Vec<String> = pulls.iter().map(|pull| pull.id.clone()).collect();
1853 self.after_push(repo_id, event.data["defaultBranch"].as_bool() == Some(true), &moved_ids)
1854 .await;
Acceptance checks in sandboxes, line comments and review verdicts1855 // A draft is announced when it is marked ready instead.
1856 for pull in pulls
1857 .into_iter()
1858 .filter(|pull| pull.status == PullStatus::Open)
1859 {
1860 self.publish_as(
1861 "pull.updated",
1862 &pull.repo_id,
1863 event.actor.clone(),
1864 PullEvent {
1865 pull_id: pull.id,
1866 repo_id: pull.repo_id.clone(),
1867 number: pull.number,
1868 issue: pull.issue_number,
1869 commit: Some(after.to_owned()),
1870 ..PullEvent::default()
1871 },
1872 )
1873 .await?;
1874 }
Work service in Rust, with RFC 3339 timestamps1875 Ok(())
1876 }
1877}
1878
1879fn service(env: &Env) -> Result<Work> {
1880 Ok(Work {
1881 db: env.d1("DB")?,
Agents as a team: lifecycle, merge queue, billing and a new shell1882 identity: env.service("IDENTITY")?,
Work service in Rust, with RFC 3339 timestamps1883 repos: env.service("REPOS")?,
Events service in Rust, with RFC 3339 times and accurate push events1884 events: env.service("EVENTS")?,
Sidebar: the panels really slide1885 actions: env.service("ACTIONS")?,
Work service in Rust, with RFC 3339 timestamps1886 })
1887}
1888
1889#[event(fetch)]
1890async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
1891 let Some(method) = rpc_method(&request) else {
1892 return Response::error("Not found", 404);
1893 };
1894 let body: serde_json::Value = request.json().await?;
1895 let work = service(&env)?;
1896
1897 match method.as_str() {
Issues and pull requests replace intents and attempts1898 "open_issue" => reply(&work.open_issue(args(body)?).await?),
1899 "list_issues" => reply(&work.list_issues(args(body)?).await?),
1900 "get_issue" => reply(&work.get_issue(args(body)?).await?),
1901 "update_issue" => reply(&work.update_issue(args(body)?).await?),
1902 "close_issue" => reply(&work.close_issue(args(body)?).await?),
1903 "reopen_issue" => reply(&work.reopen_issue(args(body)?).await?),
1904 "list_labels" => reply(&work.list_labels(args(body)?).await?),
1905 "counts" => reply(&work.counts(args(body)?).await?),
1906 "add_comment" => reply(&work.add_comment(args(body)?).await?),
Acceptance checks in sandboxes, line comments and review verdicts1907 "start_checks" => reply(&work.start_checks(args(body)?).await?),
1908 "report_checks" => reply(&work.report_checks(args(body)?).await?),
GitHub Actions on g1t, part two: running workflows1909 "set_commit_status" => reply(&work.set_commit_status(args(body)?).await?),
Agents as a team: lifecycle, merge queue, billing and a new shell1910 "start_review" => reply(&work.start_review(args(body)?).await?),
1911 "advance" => reply(&work.advance(args(body)?).await?),
1912 "stall" => reply(&work.stall(args(body)?).await?),
1913 "managed_pulls" => reply(&work.managed_pulls(args(body)?).await?),
1914 "queue" => reply(&work.queue(args(body)?).await?),
1915 "queue_build" => reply(&work.queue_build(args(body)?).await?),
1916 "report_queue" => reply(&work.report_queue(args(body)?).await?),
1917 "remove_from_queue" => reply(&work.remove_from_queue(args(body)?).await?),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request1918 "message_agent" => reply(&work.message_agent(args(body)?).await?),
Record your own agent's sessions automatically1919 "locate_pull" => reply(&work.locate_pull(args(body)?).await?),
Agents ask each other, hand each other work, and answer1920 "answer_message" => reply(&work.answer_message(args(body)?).await?),
Usage, like a hosting provider's: what agents cost, per day, task, repository and pull request1921 "take_messages" => reply(&work.take_messages(args(body)?).await?),
Agents asked while not at work are woken to answer1922 "wake_for_messages" => reply(&work.wake_for_messages(args(body)?).await?),
Agents as a team: lifecycle, merge queue, billing and a new shell1923 "catch_up_job" => reply(&work.catch_up_job(args(body)?).await?),
1924 "get_settings" => reply(&work.get_settings(args(body)?).await?),
1925 "update_settings" => reply(&work.update_settings(args(body)?).await?),
1926 "report_review" => reply(&work.report_review(args(body)?).await?),
Issues and pull requests replace intents and attempts1927 "open_pull" => reply(&work.open_pull(args(body)?).await?),
1928 "list_pulls" => reply(&work.list_pulls(args(body)?).await?),
1929 "get_pull" => reply(&work.get_pull(args(body)?).await?),
Agents as a team: lifecycle, merge queue, billing and a new shell1930 "update_pull" => reply(&work.update_pull(args(body)?).await?),
Catching up with main takes seconds when the two sides touched different files1931 "catch_up_pull" => reply(&work.catch_up_pull(args(body)?).await?),
Issues and pull requests replace intents and attempts1932 "ready_pull" => reply(&work.ready_pull(args(body)?).await?),
1933 "close_pull" => reply(&work.close_pull(args(body)?).await?),
1934 "merge_pull" => reply(&work.merge_pull(args(body)?).await?),
1935 "list_active_pulls" => reply(&work.list_active_pulls(args(body)?).await?),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains1936 "by_author" => reply(&work.by_author(args(body)?).await?),
Agents as a team: lifecycle, merge queue, billing and a new shell1937 "start_plan" => reply(&work.start_plan(args(body)?).await?),
1938 "report_plan" => reply(&work.report_plan(args(body)?).await?),
1939 "get_plan" => reply(&work.get_plan(args(body)?).await?),
1940 "list_plans" => reply(&work.list_plans(args(body)?).await?),
1941 "apply_plan" => reply(&work.apply_plan(args(body)?).await?),
1942 "queue_issue" => reply(&work.queue_issue(args(body)?).await?),
1943 "ready_issues" => reply(&work.ready_issues(args(body)?).await?),
1944 "list_assigned_issues" => reply(&work.list_assigned_issues(args(body)?).await?),
Work service in Rust, with RFC 3339 timestamps1945 "append_session" => reply(&work.append_session(args(body)?).await?),
1946 "read_session" => reply(&work.read_session(args(body)?).await?),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains1947 // Agents at work, their sessions, and memory (runs.rs, memory.rs).
1948 "open_run" => reply(&work.open_run(args(body)?).await?),
1949 "report_run" => reply(&work.report_run(args(body)?).await?),
1950 "stop_run" => reply(&work.stop_run(args(body)?).await?),
1951 "list_runs" => reply(&work.list_runs(args(body)?).await?),
1952 "get_run" => reply(&work.get_run(args(body)?).await?),
1953 "list_sessions" => reply(&work.list_sessions(args(body)?).await?),
1954 "get_session" => reply(&work.get_session(args(body)?).await?),
1955 "list_memories" => reply(&work.list_memories(args(body)?).await?),
1956 "add_memory" => reply(&work.add_memory(args(body)?).await?),
1957 "update_memory" => reply(&work.update_memory(args(body)?).await?),
1958 "delete_memory" => reply(&work.delete_memory(args(body)?).await?),
1959 "recall" => reply(&work.recall(args(body)?).await?),
1960 "memory_context" => reply(&work.memory_context(args(body)?).await?),
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API1961 // What agents may do in a sandbox (guardrails.rs).
1962 "get_guardrails" => reply(&work.get_guardrails(args(body)?).await?),
1963 "update_guardrails" => reply(&work.update_guardrails(args(body)?).await?),
1964 "run_guardrails" => reply(&work.run_guardrails(args(body)?).await?),
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look1965 // Plan caps the runner applies (compute.rs).
1966 "active_agents" => reply(&work.active_agents(args(body)?).await?),
1967 "issue_spend" => reply(&work.issue_spend(args(body)?).await?),
1968 "wait_for_slot" => reply(&work.wait_for_slot(args(body)?).await?),
1969 "agent_comment" => reply(&work.agent_comment(args(body)?).await?),
1970 "add_wait" => reply(&work.add_wait(args(body)?).await?),
1971 "waiting_workspaces" => reply(&work.waiting_workspaces(args(body)?).await?),
1972 "take_wait" => reply(&work.take_wait(args(body)?).await?),
1973 // The runs whose sandboxes stop with their repository (retired.rs).
1974 "runs_in_repo" => reply(&work.runs_in_repo(args(body)?).await?),
1975 "run_cost" => reply(&work.run_cost(args(body)?).await?),
Agents and memory, checks and conflicts, profiles, slug renames, custom domains1976 "start_mergecheck" => reply(&work.start_mergecheck(args(body)?).await?),
1977 "report_mergecheck" => reply(&work.report_mergecheck(args(body)?).await?),
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API1978 // Memory that fills itself, and its review queue (capture.rs).
1979 method if capture::METHODS.contains(&method) => capture::dispatch(&work, method, body).await,
1980 // @g1t-agent in comments, and the label rule (mentions.rs).
1981 "take_mention" => reply(&work.take_mention(args(body)?).await?),
1982 "mention_revision" => reply(&work.mention_revision(args(body)?).await?),
1983 "reply_mention" => reply(&work.reply_mention(args(body)?).await?),
1984 "get_agent_rules" => reply(&work.get_agent_rules(args(body)?).await?),
1985 "set_agent_rules" => reply(&work.set_agent_rules(args(body)?).await?),
Work service in Rust, with RFC 3339 timestamps1986 _ => Response::error("Unknown method", 404),
1987 }
1988}
1989
1990/// Events from the bus, delivered on this service's own queue.
1991#[event(queue)]
Events service in Rust, with RFC 3339 times and accurate push events1992async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
Work service in Rust, with RFC 3339 timestamps1993 let work = service(&env)?;
1994 for message in batch.messages()? {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains1995 // A workspace renamed: its agent runs and memory move to the slug it has now.
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API1996 if g1t_kit::rename::on_event(&env, &env.d1("DB")?, message.body(), &[memory::RENAMED, guardrails::RENAMED].concat()).await? {
Agents and memory, checks and conflicts, profiles, slug renames, custom domains1997 message.ack();
1998 continue;
1999 }
Invite-only launch: sign in with GitHub, repository access and lifecycle, many emails, a new look2000 // A repository renamed or transferred: its runs, memory, guardrails
2001 // and runs waiting for a slot follow.
2002 if g1t_kit::transfer::on_event(&env, &env.d1("DB")?, message.body(), &[memory::TRANSFERRED, guardrails::TRANSFERRED, retired::WAITS_MOVED].concat()).await? {
2003 message.ack();
2004 continue;
2005 }
2006 // A workspace deleted: what it kept for itself goes.
2007 if g1t_kit::deleted::on_event(&env.d1("DB")?, message.body(), memory::DELETED).await? {
2008 message.ack();
2009 continue;
2010 }
2011 // A repository deleted, archived or purged, or a branch renamed (retired.rs).
2012 if work.on_retired(message.body()).await? {
2013 message.ack();
2014 continue;
2015 }
Agents get guardrails, run credentials, an audit log, a context hub, repository instructions and mentions; security upkeep; snake_case API2016 capture::on_event(&work, message.body()).await;
Work service in Rust, with RFC 3339 timestamps2017 work.on_event(message.body()).await?;
2018 message.ack();
2019 }
2020 Ok(())
2021}