pr_01m47d15m3e54sn21z27rpy5n9/services/work/src/lib.rs

1,766 lines65,673 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
Acceptance checks in sandboxes, line comments and review verdicts7mod checks;
Agents as a team: lifecycle, merge queue, billing and a new shell8mod lifecycle;
9mod plans;
10mod queue;
11mod reviews;
Work service in Rust, with RFC 3339 timestamps12mod rows;
Agents as a team: lifecycle, merge queue, billing and a new shell13mod settings;
Work service in Rust, with RFC 3339 timestamps14
15use g1t_contracts::events::{
Events service in Rust, with RFC 3339 times and accurate push events16 CommentCreated, Event, IssueEvent, NewEvent, Publish, PullEvent, SessionAppended,
Work service in Rust, with RFC 3339 timestamps17};
Agents as a team: lifecycle, merge queue, billing and a new shell18use g1t_contracts::identity::UsernameArgs;
Pull requests from branches19use g1t_contracts::repos::{ForkArgs, GetArgs, HeadArgs, LandArgs, Landed, Repo, RepoPath};
Work service in Rust, with RFC 3339 timestamps20use g1t_contracts::time::rfc3339;
21use g1t_contracts::work::*;
Agents as a team: lifecycle, merge queue, billing and a new shell22use futures_util::future::{try_join, try_join3, try_join_all};
Work service in Rust, with RFC 3339 timestamps23use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
Events service in Rust, with RFC 3339 times and accurate push events24use g1t_kit::{args, now_ms, reply, rpc_method};
Work service in Rust, with RFC 3339 timestamps25use serde::Serialize;
26use worker::wasm_bindgen::JsValue;
27use worker::{
28 Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event,
29};
30
Agents as a team: lifecycle, merge queue, billing and a new shell31use rows::{CommentRow, IssueRow, MovedRow, NumberRow, PullRow, SessionRow, Snapshot, ValueRow};
Work service in Rust, with RFC 3339 timestamps32
33const SOURCE: &str = "work";
34const MAX_ENTRY_BATCH: usize = 200;
35const MAX_ENTRY_CHARS: usize = 64_000;
Issues and pull requests replace intents and attempts36const MAX_TITLE_CHARS: usize = 200;
Work service in Rust, with RFC 3339 timestamps37const SESSION_PAGE: u32 = 500;
Issues and pull requests replace intents and attempts38const LIST_PAGE: u32 = 100;
Agents as a team: lifecycle, merge queue, billing and a new shell39const MAX_ASSIGNEES: usize = 10;
Work service in Rust, with RFC 3339 timestamps40const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
41
Issues and pull requests replace intents and attempts42const ISSUE_COLUMNS: &str = "issues.*,
43 (SELECT count(*) FROM pulls WHERE pulls.issue_id = issues.id) AS pull_count,
Agents as a team: lifecycle, merge queue, billing and a new shell44 (SELECT agent FROM pulls
45 WHERE pulls.issue_id = issues.id AND pulls.status IN ('draft', 'open')
46 AND pulls.fork_repo_id IS NOT NULL
47 ORDER BY pulls.number DESC LIMIT 1) AS agent,
Issues and pull requests replace intents and attempts48 (SELECT count(*) FROM comments
Agents as a team: lifecycle, merge queue, billing and a new shell49 WHERE comments.repo_id = issues.repo_id AND comments.number = issues.number
50 AND comments.kind = 'comment') AS comment_count";
Work service in Rust, with RFC 3339 timestamps51
Issues and pull requests replace intents and attempts52fn no_issue<T>() -> Outcome<T> {
53 Outcome::fail(FailureCode::NotFound, "Issue not found.")
Work service in Rust, with RFC 3339 timestamps54}
55
Issues and pull requests replace intents and attempts56fn no_pull<T>() -> Outcome<T> {
57 Outcome::fail(FailureCode::NotFound, "Pull request not found.")
Work service in Rust, with RFC 3339 timestamps58}
59
60fn optional(value: &Option<String>) -> JsValue {
61 value.as_deref().map_or(JsValue::NULL, JsValue::from)
62}
63
Issues and pull requests replace intents and attempts64fn optional_number(value: Option<u32>) -> JsValue {
65 value.map_or(JsValue::NULL, JsValue::from)
66}
67
68/// The lowercase name a `State` is stored and sent as.
69fn state_name(state: Option<State>) -> Option<&'static str> {
70 state.map(|state| match state {
71 State::Open => "open",
72 State::Closed => "closed",
73 })
74}
75
76/// A trimmed title, or why it cannot be used.
77fn valid_title(title: &str) -> std::result::Result<&str, &'static str> {
78 let title = title.trim();
79 if title.is_empty() {
80 Err("A title is required.")
81 } else if title.chars().count() > MAX_TITLE_CHARS {
82 Err("That title is too long.")
83 } else {
84 Ok(title)
85 }
86}
87
88/// Unwraps an `Outcome`, returning its failure from the enclosing method.
89macro_rules! check {
90 ($outcome:expr) => {
91 match $outcome {
92 Outcome::Ok(value) => value,
93 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
94 }
95 };
96}
97
Work service in Rust, with RFC 3339 timestamps98struct Work {
99 db: D1Database,
Agents as a team: lifecycle, merge queue, billing and a new shell100 identity: Fetcher,
Work service in Rust, with RFC 3339 timestamps101 repos: Fetcher,
Events service in Rust, with RFC 3339 times and accurate push events102 events: Fetcher,
Work service in Rust, with RFC 3339 timestamps103}
104
105impl Work {
Issues and pull requests replace intents and attempts106 async fn publish<T: Serialize>(
107 &self,
108 kind: &'static str,
109 repo_id: &str,
110 actor: &User,
111 data: T,
112 ) -> Result<()> {
Acceptance checks in sandboxes, line comments and review verdicts113 self.publish_as(kind, repo_id, Some(actor.id.clone()), data)
114 .await
115 }
116
117 /// Publishes an event caused by `actor`, or by g1t itself.
118 async fn publish_as<T: Serialize>(
119 &self,
120 kind: &'static str,
121 repo_id: &str,
122 actor: Option<String>,
123 data: T,
124 ) -> Result<()> {
Issues and pull requests replace intents and attempts125 let event = NewEvent {
126 kind,
127 source: SOURCE,
128 repo_id: Some(repo_id.to_owned()),
Acceptance checks in sandboxes, line comments and review verdicts129 actor,
Issues and pull requests replace intents and attempts130 data,
131 };
Events service in Rust, with RFC 3339 times and accurate push events132 g1t_kit::call(
133 &self.events,
134 "publish",
135 &Publish {
136 events: vec![event],
137 },
138 )
139 .await
Work service in Rust, with RFC 3339 timestamps140 }
141
Issues and pull requests replace intents and attempts142 /// The repository, if the viewer may see it. Whether they may is
143 /// decided by the repos service.
144 async fn repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Outcome<Repo>> {
Work service in Rust, with RFC 3339 timestamps145 g1t_kit::call(
146 &self.repos,
147 "get",
148 &GetArgs {
149 path: path.clone(),
150 viewer: viewer.clone(),
151 },
152 )
153 .await
154 }
155
Issues and pull requests replace intents and attempts156 /// The next number in the repository's sequence. Taking it is one
157 /// statement, so concurrent opens cannot be given the same number.
158 async fn next_number(&self, repo_id: &str) -> Result<u32> {
159 let row = self
160 .db
161 .prepare(
162 "INSERT INTO counters (repo_id, last) VALUES (?, 1)
163 ON CONFLICT (repo_id) DO UPDATE SET last = last + 1
164 RETURNING last AS n",
165 )
166 .bind(&[repo_id.into()])?
167 .first::<NumberRow>(None)
168 .await?;
169 row.map(|row| row.n)
170 .ok_or_else(|| worker::Error::RustError("no number was assigned".into()))
Work service in Rust, with RFC 3339 timestamps171 }
172
Issues and pull requests replace intents and attempts173 async fn issue(&self, repo_id: &str, number: u32) -> Result<Option<Issue>> {
Work service in Rust, with RFC 3339 timestamps174 Ok(self
175 .db
Issues and pull requests replace intents and attempts176 .prepare(format!(
177 "SELECT {ISSUE_COLUMNS} FROM issues WHERE repo_id = ? AND number = ?"
178 ))
179 .bind(&[repo_id.into(), number.into()])?
180 .first::<IssueRow>(None)
Work service in Rust, with RFC 3339 timestamps181 .await?
Issues and pull requests replace intents and attempts182 .map(Issue::from))
Work service in Rust, with RFC 3339 timestamps183 }
184
Issues and pull requests replace intents and attempts185 async fn pull(&self, repo_id: &str, number: u32) -> Result<Option<Pull>> {
Work service in Rust, with RFC 3339 timestamps186 Ok(self
187 .db
Issues and pull requests replace intents and attempts188 .prepare("SELECT * FROM pulls WHERE repo_id = ? AND number = ?")
189 .bind(&[repo_id.into(), number.into()])?
190 .first::<PullRow>(None)
191 .await?
192 .map(Pull::from))
193 }
194
195 async fn comments(&self, repo_id: &str, number: u32) -> Result<Vec<Comment>> {
196 let rows = self
197 .db
198 .prepare(
199 "SELECT * FROM comments WHERE repo_id = ? AND number = ? ORDER BY id LIMIT 500",
200 )
201 .bind(&[repo_id.into(), number.into()])?
202 .all()
Work service in Rust, with RFC 3339 timestamps203 .await?
Issues and pull requests replace intents and attempts204 .results::<CommentRow>()?;
205 Ok(rows.into_iter().map(Comment::from).collect())
Work service in Rust, with RFC 3339 timestamps206 }
207
Issues and pull requests replace intents and attempts208 /// The repository and one of its issues, as seen by `viewer`.
209 async fn issue_at(
210 &self,
211 path: &RepoPath,
212 number: u32,
213 viewer: &Viewer,
214 ) -> Result<Outcome<(Repo, Issue)>> {
215 let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
216 return Ok(no_issue());
Work service in Rust, with RFC 3339 timestamps217 };
Issues and pull requests replace intents and attempts218 Ok(match self.issue(&repo.id, number).await? {
219 Some(issue) => Outcome::Ok((repo, issue)),
220 None => no_issue(),
221 })
222 }
223
224 /// The repository and one of its pull requests, as seen by `viewer`.
225 async fn pull_at(
226 &self,
227 path: &RepoPath,
228 number: u32,
229 viewer: &Viewer,
230 ) -> Result<Outcome<(Repo, Pull)>> {
231 let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
232 return Ok(no_pull());
233 };
234 Ok(match self.pull(&repo.id, number).await? {
235 Some(pull) => Outcome::Ok((repo, pull)),
236 None => no_pull(),
237 })
238 }
239
Agents as a team: lifecycle, merge queue, billing and a new shell240 /// Records something that happened to an issue or a pull request, so
241 /// that it shows in the conversation where it happened. `text` is what
242 /// `author` did, as the rest of a sentence starting with their name.
243 pub(crate) async fn note(
244 &self,
245 repo_id: &str,
246 number: u32,
247 author: (&str, &str),
248 text: &str,
249 ) -> Result<()> {
250 let now = now_ms();
251 self.db
252 .prepare(
253 "INSERT INTO comments
254 (id, repo_id, number, author_id, author_name, body, kind, created_at)
255 VALUES (?, ?, ?, ?, ?, ?, 'event', ?)",
256 )
257 .bind(&[
258 new_id("cmt", now).into(),
259 repo_id.into(),
260 number.into(),
261 author.0.into(),
262 author.1.into(),
263 text.into(),
264 rfc3339(now).into(),
265 ])?
266 .run()
267 .await?;
268 Ok(())
269 }
270
271 /// Notes who was added to and removed from a list of people, such as
272 /// "assigned ana" or "requested a review from g1t-agent".
273 async fn note_changes(
274 &self,
275 repo_id: &str,
276 number: u32,
277 actor: &User,
278 before: &[String],
279 after: &[String],
280 (added, removed): (&str, &str),
281 ) -> Result<()> {
282 let joined = |names: Vec<&String>| {
283 names
284 .into_iter()
285 .map(String::as_str)
286 .collect::<Vec<_>>()
287 .join(", ")
288 };
289 let new: Vec<&String> = after.iter().filter(|name| !before.contains(name)).collect();
290 let gone: Vec<&String> = before.iter().filter(|name| !after.contains(name)).collect();
291 let who = (actor.id.as_str(), actor.username.as_str());
292 if !new.is_empty() {
293 // Taking something on oneself reads better said that way.
294 let text = if added == "assigned" && new == [&actor.username] {
295 "self-assigned this".to_owned()
296 } else {
297 format!("{added} {}", joined(new))
298 };
299 self.note(repo_id, number, who, &text).await?;
300 }
301 if !gone.is_empty() {
302 self.note(repo_id, number, who, &format!("{removed} {}", joined(gone)))
303 .await?;
304 }
305 Ok(())
306 }
307
Issues and pull requests replace intents and attempts308 fn issue_event(issue: &Issue) -> IssueEvent {
309 IssueEvent {
310 issue_id: issue.id.clone(),
311 repo_id: issue.repo_id.clone(),
312 number: issue.number,
313 ..IssueEvent::default()
Work service in Rust, with RFC 3339 timestamps314 }
315 }
316
Issues and pull requests replace intents and attempts317 fn pull_event(pull: &Pull) -> PullEvent {
318 PullEvent {
319 pull_id: pull.id.clone(),
320 repo_id: pull.repo_id.clone(),
321 number: pull.number,
322 issue: pull.issue,
323 ..PullEvent::default()
Work service in Rust, with RFC 3339 timestamps324 }
325 }
326
Issues and pull requests replace intents and attempts327 // --- Issues ------------------------------------------------------------
328
329 async fn open_issue(&self, a: OpenIssueArgs) -> Result<Outcome<Issue>> {
Work service in Rust, with RFC 3339 timestamps330 if !a.actor.verified {
331 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
332 }
Issues and pull requests replace intents and attempts333 let title = match valid_title(&a.title) {
334 Ok(title) => title,
335 Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
336 };
337 let Some(labels) = normalize_labels(&a.labels) else {
Work service in Rust, with RFC 3339 timestamps338 return Ok(Outcome::fail(
339 FailureCode::Invalid,
Issues and pull requests replace intents and attempts340 "An issue can have up to 10 labels of up to 40 characters each.",
Work service in Rust, with RFC 3339 timestamps341 ));
342 };
Issues and pull requests replace intents and attempts343 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
Work service in Rust, with RFC 3339 timestamps344 let checks: Vec<&str> = a
345 .checks
346 .iter()
347 .map(|check| check.trim())
348 .filter(|check| !check.is_empty())
349 .collect();
350
351 let now = now_ms();
Issues and pull requests replace intents and attempts352 let id = new_id("iss", now);
353 let number = self.next_number(&repo.id).await?;
354 let timestamp = rfc3339(now);
Work service in Rust, with RFC 3339 timestamps355 self.db
356 .prepare(
Issues and pull requests replace intents and attempts357 "INSERT INTO issues
358 (id, repo_id, number, title, body, labels, checks, author_id, author_name,
359 created_at, updated_at)
360 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
Work service in Rust, with RFC 3339 timestamps361 )
362 .bind(&[
363 id.as_str().into(),
364 repo.id.as_str().into(),
Issues and pull requests replace intents and attempts365 number.into(),
Work service in Rust, with RFC 3339 timestamps366 title.into(),
Issues and pull requests replace intents and attempts367 a.body.trim().into(),
368 serde_json::to_string(&labels)?.into(),
Work service in Rust, with RFC 3339 timestamps369 serde_json::to_string(&checks)?.into(),
370 a.actor.id.as_str().into(),
371 a.actor.username.as_str().into(),
Issues and pull requests replace intents and attempts372 timestamp.as_str().into(),
373 timestamp.as_str().into(),
Work service in Rust, with RFC 3339 timestamps374 ])?
375 .run()
376 .await?;
Issues and pull requests replace intents and attempts377 let Some(issue) = self.issue(&repo.id, number).await? else {
378 return Ok(no_issue());
Work service in Rust, with RFC 3339 timestamps379 };
Issues and pull requests replace intents and attempts380 self.publish(
381 "issue.opened",
382 &repo.id,
383 &a.actor,
384 IssueEvent {
385 title: Some(issue.title.clone()),
386 ..Self::issue_event(&issue)
Work service in Rust, with RFC 3339 timestamps387 },
Issues and pull requests replace intents and attempts388 )
Work service in Rust, with RFC 3339 timestamps389 .await?;
Issues and pull requests replace intents and attempts390 Ok(Outcome::Ok(issue))
Work service in Rust, with RFC 3339 timestamps391 }
392
Issues and pull requests replace intents and attempts393 async fn list_issues(&self, a: ListIssuesArgs) -> Result<Outcome<Vec<Issue>>> {
394 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
395 let state = state_name(a.state).map_or(JsValue::NULL, JsValue::from);
396 let label = a
397 .label
398 .map(|label| label.trim().to_lowercase())
399 .filter(|label| !label.is_empty());
Work service in Rust, with RFC 3339 timestamps400 let rows = self
401 .db
402 .prepare(format!(
Issues and pull requests replace intents and attempts403 "SELECT {ISSUE_COLUMNS} FROM issues
404 WHERE repo_id = ? AND (? IS NULL OR state = ?)
405 AND (? IS NULL OR EXISTS
406 (SELECT 1 FROM json_each(issues.labels) WHERE json_each.value = ?))
407 ORDER BY number DESC LIMIT ?"
Work service in Rust, with RFC 3339 timestamps408 ))
Issues and pull requests replace intents and attempts409 .bind(&[
410 repo.id.into(),
411 state.clone(),
412 state,
413 optional(&label),
414 optional(&label),
415 LIST_PAGE.into(),
416 ])?
Work service in Rust, with RFC 3339 timestamps417 .all()
418 .await?
Issues and pull requests replace intents and attempts419 .results::<IssueRow>()?;
420 Ok(Outcome::Ok(rows.into_iter().map(Issue::from).collect()))
Work service in Rust, with RFC 3339 timestamps421 }
422
Issues and pull requests replace intents and attempts423 async fn get_issue(&self, a: ViewArgs) -> Result<Outcome<IssueDetail>> {
424 let (repo, issue) = check!(self.issue_at(&a.repo, a.number, &a.viewer).await?);
425 let pulls = self
Work service in Rust, with RFC 3339 timestamps426 .db
Issues and pull requests replace intents and attempts427 .prepare("SELECT * FROM pulls WHERE issue_id = ? ORDER BY number")
428 .bind(&[issue.id.as_str().into()])?
Work service in Rust, with RFC 3339 timestamps429 .all()
430 .await?
Issues and pull requests replace intents and attempts431 .results::<PullRow>()?;
432 Ok(Outcome::Ok(IssueDetail {
433 comments: self.comments(&repo.id, issue.number).await?,
434 pulls: pulls.into_iter().map(Pull::from).collect(),
435 issue,
Work service in Rust, with RFC 3339 timestamps436 }))
437 }
438
Issues and pull requests replace intents and attempts439 /// The issue, if `actor` wrote it or belongs to the repository's workspace.
440 async fn manageable_issue(
441 &self,
442 actor: &User,
443 path: &RepoPath,
444 number: u32,
445 ) -> Result<Outcome<Issue>> {
446 let (repo, issue) = check!(self.issue_at(path, number, &Some(actor.clone())).await?);
447 if issue.author.id != actor.id && !actor.is_member(&repo.namespace) {
Work service in Rust, with RFC 3339 timestamps448 return Ok(Outcome::fail(
449 FailureCode::Forbidden,
Issues and pull requests replace intents and attempts450 "Only the author or a member of the workspace can change an issue.",
Work service in Rust, with RFC 3339 timestamps451 ));
452 }
Issues and pull requests replace intents and attempts453 Ok(Outcome::Ok(issue))
454 }
455
456 async fn update_issue(&self, a: UpdateIssueArgs) -> Result<Outcome<Issue>> {
457 let issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
458 let title = match a.title.as_deref().map(valid_title) {
459 Some(Err(message)) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
460 Some(Ok(title)) => Some(title.to_owned()),
461 None => None,
462 };
463 let labels = match a.labels.as_deref().map(normalize_labels) {
464 Some(None) => {
465 return Ok(Outcome::fail(
466 FailureCode::Invalid,
467 "An issue can have up to 10 labels of up to 40 characters each.",
468 ));
469 }
470 Some(Some(labels)) => Some(serde_json::to_string(&labels)?),
471 None => None,
472 };
Agents as a team: lifecycle, merge queue, billing and a new shell473 let assignees = match a.assignees {
474 Some(names) => Some(check!(self.valid_assignees(names).await?)),
475 None => None,
476 };
477 let assigned = assignees.as_ref().map(serde_json::to_string).transpose()?;
Issues and pull requests replace intents and attempts478 let body = a.body.map(|body| body.trim().to_owned());
479 self.db
480 .prepare(
481 "UPDATE issues
482 SET title = COALESCE(?, title), body = COALESCE(?, body),
Agents as a team: lifecycle, merge queue, billing and a new shell483 labels = COALESCE(?, labels), assignees = COALESCE(?, assignees),
484 updated_at = ?
Issues and pull requests replace intents and attempts485 WHERE id = ?",
486 )
487 .bind(&[
488 optional(&title),
489 optional(&body),
490 optional(&labels),
Agents as a team: lifecycle, merge queue, billing and a new shell491 optional(&assigned),
Issues and pull requests replace intents and attempts492 rfc3339(now_ms()).into(),
493 issue.id.as_str().into(),
494 ])?
495 .run()
496 .await?;
Agents as a team: lifecycle, merge queue, billing and a new shell497 let before = issue.assignees.clone();
Issues and pull requests replace intents and attempts498 let Some(issue) = self.issue(&issue.repo_id, issue.number).await? else {
499 return Ok(no_issue());
500 };
501 self.publish(
502 "issue.updated",
503 &issue.repo_id,
504 &a.actor,
505 Self::issue_event(&issue),
506 )
507 .await?;
Agents as a team: lifecycle, merge queue, billing and a new shell508 if let Some(assignees) = assignees {
509 self.note_changes(
510 &issue.repo_id,
511 issue.number,
512 &a.actor,
513 &before,
514 &assignees,
515 ("assigned", "unassigned"),
516 )
517 .await?;
518 self.publish(
519 "issue.assigned",
520 &issue.repo_id,
521 &a.actor,
522 IssueEvent {
523 assignees: Some(assignees),
524 ..Self::issue_event(&issue)
525 },
526 )
527 .await?;
528 }
Issues and pull requests replace intents and attempts529 Ok(Outcome::Ok(issue))
530 }
531
Agents as a team: lifecycle, merge queue, billing and a new shell532 /// Usernames as given, tidied, if each names an account.
533 async fn valid_assignees(&self, names: Vec<String>) -> Result<Outcome<Vec<String>>> {
534 let mut assignees: Vec<String> = Vec::new();
535 for name in names {
536 let name = name.trim().trim_start_matches('@').to_lowercase();
537 if name.is_empty() || assignees.contains(&name) {
538 continue;
539 }
540 if assignees.len() == MAX_ASSIGNEES {
541 return Ok(Outcome::fail(
542 FailureCode::Invalid,
543 format!("An issue can be assigned to at most {MAX_ASSIGNEES} people."),
544 ));
545 }
546 let account: Viewer = g1t_kit::call(
547 &self.identity,
548 "user_by_username",
549 &UsernameArgs {
550 username: name.clone(),
551 },
552 )
553 .await?;
554 if account.is_none() {
555 return Ok(Outcome::fail(
556 FailureCode::Invalid,
557 format!("There is no account named {name}."),
558 ));
559 }
560 assignees.push(name);
561 }
562 Ok(Outcome::Ok(assignees))
563 }
564
565 /// Open issues assigned to the viewer, in every repository. Callers
566 /// show only those in repositories the viewer can still see.
567 async fn list_assigned_issues(&self, a: ViewerArgs) -> Result<Vec<Issue>> {
568 let Some(viewer) = a.viewer else {
569 return Ok(Vec::new());
570 };
571 let rows = self
572 .db
573 .prepare(format!(
574 "SELECT {ISSUE_COLUMNS} FROM issues
575 WHERE state = 'open' AND EXISTS (
576 SELECT 1 FROM json_each(issues.assignees) WHERE json_each.value = ?)
577 ORDER BY updated_at DESC LIMIT 50"
578 ))
579 .bind(&[viewer.username.into()])?
580 .all()
581 .await?
582 .results::<IssueRow>()?;
583 Ok(rows.into_iter().map(Issue::from).collect())
584 }
585
Issues and pull requests replace intents and attempts586 async fn close_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
587 let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
588 if issue.state == State::Closed {
Work service in Rust, with RFC 3339 timestamps589 return Ok(Outcome::fail(
590 FailureCode::Conflict,
Issues and pull requests replace intents and attempts591 "This issue is already closed.",
Work service in Rust, with RFC 3339 timestamps592 ));
593 }
Issues and pull requests replace intents and attempts594 let reason = a.reason.unwrap_or(IssueReason::Completed);
595 let now = rfc3339(now_ms());
Work service in Rust, with RFC 3339 timestamps596 self.db
Issues and pull requests replace intents and attempts597 .prepare(
598 "UPDATE issues SET state = 'closed', reason = ?, closed_at = ?, updated_at = ?
599 WHERE id = ?",
600 )
601 .bind(&[
602 reason.as_str().into(),
603 now.as_str().into(),
604 now.as_str().into(),
605 issue.id.as_str().into(),
606 ])?
Work service in Rust, with RFC 3339 timestamps607 .run()
608 .await?;
Issues and pull requests replace intents and attempts609 self.publish(
610 "issue.closed",
611 &issue.repo_id,
612 &a.actor,
613 IssueEvent {
614 reason: Some(reason.as_str()),
615 ..Self::issue_event(&issue)
Work service in Rust, with RFC 3339 timestamps616 },
Issues and pull requests replace intents and attempts617 )
Work service in Rust, with RFC 3339 timestamps618 .await?;
Agents as a team: lifecycle, merge queue, billing and a new shell619 self.note(
620 &issue.repo_id,
621 issue.number,
622 (&a.actor.id, &a.actor.username),
623 match reason {
624 IssueReason::Completed => "closed this as completed",
625 IssueReason::NotPlanned => "closed this as not planned",
626 },
627 )
628 .await?;
Issues and pull requests replace intents and attempts629 issue.state = State::Closed;
630 issue.reason = Some(reason);
631 issue.closed_at = Some(now.clone());
632 issue.updated_at = now;
633 Ok(Outcome::Ok(issue))
634 }
635
636 async fn reopen_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
637 let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
638 if issue.state == State::Open {
639 return Ok(Outcome::fail(
640 FailureCode::Conflict,
641 "This issue is already open.",
642 ));
643 }
644 let now = rfc3339(now_ms());
645 self.db
646 .prepare(
647 "UPDATE issues
648 SET state = 'open', reason = NULL, resolved_by = NULL, closed_at = NULL,
649 updated_at = ?
650 WHERE id = ?",
651 )
652 .bind(&[now.as_str().into(), issue.id.as_str().into()])?
653 .run()
654 .await?;
655 self.publish(
656 "issue.reopened",
657 &issue.repo_id,
658 &a.actor,
659 Self::issue_event(&issue),
660 )
661 .await?;
Agents as a team: lifecycle, merge queue, billing and a new shell662 self.note(
663 &issue.repo_id,
664 issue.number,
665 (&a.actor.id, &a.actor.username),
666 "reopened this",
667 )
668 .await?;
Issues and pull requests replace intents and attempts669 issue.state = State::Open;
670 issue.reason = None;
671 issue.resolved_by = None;
672 issue.closed_at = None;
673 issue.updated_at = now;
674 Ok(Outcome::Ok(issue))
Work service in Rust, with RFC 3339 timestamps675 }
676
Issues and pull requests replace intents and attempts677 /// The default labels, then every other label in use on the repository.
678 async fn list_labels(&self, a: ViewArgs) -> Result<Outcome<Vec<String>>> {
679 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
680 let used = self
681 .db
682 .prepare(
683 "SELECT DISTINCT json_each.value AS value
684 FROM issues, json_each(issues.labels)
685 WHERE issues.repo_id = ? ORDER BY 1 LIMIT 200",
686 )
687 .bind(&[repo.id.into()])?
688 .all()
689 .await?
690 .results::<ValueRow>()?;
691 let mut labels: Vec<String> = DEFAULT_LABELS.iter().map(|label| (*label).into()).collect();
692 for row in used {
693 if !labels.contains(&row.value) {
694 labels.push(row.value);
695 }
696 }
697 Ok(Outcome::Ok(labels))
698 }
699
700 async fn counts(&self, a: ViewArgs) -> Result<Outcome<Counts>> {
701 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
702 let counts = self
703 .db
704 .prepare(
705 "SELECT
706 (SELECT count(*) FROM issues WHERE repo_id = ? AND state = 'open') AS issues,
707 (SELECT count(*) FROM pulls
708 WHERE repo_id = ? AND status IN ('draft', 'open')) AS pulls",
709 )
710 .bind(&[repo.id.as_str().into(), repo.id.as_str().into()])?
711 .first::<Counts>(None)
712 .await?;
713 Ok(Outcome::Ok(counts.unwrap_or(Counts {
714 issues: 0,
715 pulls: 0,
716 })))
717 }
718
719 // --- Comments ----------------------------------------------------------
720
721 async fn add_comment(&self, a: AddCommentArgs) -> Result<Outcome<Comment>> {
Work service in Rust, with RFC 3339 timestamps722 if !a.actor.verified {
723 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
724 }
Issues and pull requests replace intents and attempts725 let body = a.body.trim();
Acceptance checks in sandboxes, line comments and review verdicts726 // An approval speaks for itself; anything else has to say something.
727 if body.is_empty() && a.verdict != Some(Verdict::Approve) {
Issues and pull requests replace intents and attempts728 return Ok(Outcome::fail(
729 FailureCode::Invalid,
730 "A comment cannot be empty.",
731 ));
732 }
Acceptance checks in sandboxes, line comments and review verdicts733 let path = a
734 .path
735 .as_deref()
736 .map(str::trim)
737 .filter(|path| !path.is_empty());
738 let line = a.line.filter(|line| *line > 0 && path.is_some());
Issues and pull requests replace intents and attempts739 if body.chars().count() > MAX_ENTRY_CHARS {
740 return Ok(Outcome::fail(
741 FailureCode::Invalid,
742 "That comment is too long.",
743 ));
744 }
745 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
746 // The number names an issue or a pull request, never both.
Agents as a team: lifecycle, merge queue, billing and a new shell747 let mut pull_id = None;
Issues and pull requests replace intents and attempts748 let table = if self.issue(&repo.id, a.number).await?.is_some() {
Acceptance checks in sandboxes, line comments and review verdicts749 if path.is_some() || a.verdict.is_some() {
750 return Ok(Outcome::fail(
751 FailureCode::Invalid,
752 "Only a pull request can be reviewed or commented on by line.",
753 ));
754 }
Issues and pull requests replace intents and attempts755 "issues"
Acceptance checks in sandboxes, line comments and review verdicts756 } else if let Some(pull) = self.pull(&repo.id, a.number).await? {
757 if a.verdict.is_some() && pull.author.id == a.actor.id {
758 return Ok(Outcome::fail(
759 FailureCode::Forbidden,
760 "You cannot approve or request changes on your own pull request.",
761 ));
762 }
Agents as a team: lifecycle, merge queue, billing and a new shell763 pull_id = Some(pull.id.clone());
Issues and pull requests replace intents and attempts764 "pulls"
765 } else {
Work service in Rust, with RFC 3339 timestamps766 return Ok(Outcome::fail(
Issues and pull requests replace intents and attempts767 FailureCode::NotFound,
768 "No issue or pull request has that number.",
Work service in Rust, with RFC 3339 timestamps769 ));
Issues and pull requests replace intents and attempts770 };
771
772 let now = now_ms();
773 let comment = Comment {
Agents as a team: lifecycle, merge queue, billing and a new shell774 kind: CommentKind::Comment,
Issues and pull requests replace intents and attempts775 id: new_id("cmt", now),
776 author: a.actor.clone(),
777 body: body.to_owned(),
Acceptance checks in sandboxes, line comments and review verdicts778 path: path.map(str::to_owned),
779 line,
780 verdict: a.verdict,
Issues and pull requests replace intents and attempts781 created_at: rfc3339(now),
782 };
783 self.db
784 .batch(vec![
785 self.db
786 .prepare(
787 "INSERT INTO comments
Acceptance checks in sandboxes, line comments and review verdicts788 (id, repo_id, number, author_id, author_name, body, path, line,
789 verdict, created_at)
790 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
Issues and pull requests replace intents and attempts791 )
792 .bind(&[
793 comment.id.as_str().into(),
794 repo.id.as_str().into(),
795 a.number.into(),
796 a.actor.id.as_str().into(),
797 a.actor.username.as_str().into(),
798 body.into(),
Acceptance checks in sandboxes, line comments and review verdicts799 optional(&comment.path),
800 optional_number(line),
801 a.verdict
802 .map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
Issues and pull requests replace intents and attempts803 comment.created_at.as_str().into(),
804 ])?,
805 self.db
806 .prepare(format!(
807 "UPDATE {table} SET updated_at = ? WHERE repo_id = ? AND number = ?"
808 ))
809 .bind(&[
810 comment.created_at.as_str().into(),
811 repo.id.as_str().into(),
812 a.number.into(),
813 ])?,
814 ])
815 .await?;
816 self.publish(
817 "comment.created",
818 &repo.id,
819 &a.actor,
820 CommentCreated {
821 comment_id: comment.id.clone(),
822 repo_id: repo.id.clone(),
823 number: a.number,
Agents as a team: lifecycle, merge queue, billing and a new shell824 pull_id,
825 verdict: a.verdict,
Issues and pull requests replace intents and attempts826 },
827 )
828 .await?;
829 Ok(Outcome::Ok(comment))
830 }
831
832 // --- Pull requests -----------------------------------------------------
833
834 async fn open_pull(&self, a: OpenPullArgs) -> Result<Outcome<Pull>> {
835 if !a.actor.verified {
836 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
Work service in Rust, with RFC 3339 timestamps837 }
Issues and pull requests replace intents and attempts838 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
839 let issue = match a.issue {
840 Some(number) => match self.issue(&repo.id, number).await? {
841 Some(issue) if issue.state == State::Open => Some(issue),
842 Some(_) => {
843 return Ok(Outcome::fail(
844 FailureCode::Conflict,
845 "This issue is closed.",
846 ));
847 }
848 None => return Ok(no_issue()),
849 },
850 None => None,
851 };
852 // A pull request for an issue takes the issue's title unless given one.
853 let title = match (a.title.trim(), &issue) {
854 ("", Some(issue)) => issue.title.clone(),
855 (title, _) => match valid_title(title) {
856 Ok(title) => title.to_owned(),
857 Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
858 },
859 };
Work service in Rust, with RFC 3339 timestamps860 let agent = match a.agent.trim() {
861 "" => "agent",
862 agent => agent,
863 };
864 let runtime = match a.runtime {
Issues and pull requests replace intents and attempts865 Runtime::Hosted => "hosted",
866 Runtime::External => "external",
Work service in Rust, with RFC 3339 timestamps867 };
868
869 let now = now_ms();
Issues and pull requests replace intents and attempts870 let id = new_id("pr", now);
Pull requests from branches871 let branch = a
872 .branch
873 .as_deref()
874 .map(str::trim)
875 .filter(|branch| !branch.is_empty());
876 // The change is on a branch already pushed to the repository, or
877 // will be made in a fork created for this pull request.
878 let (fork, head) = match branch {
879 Some(branch) => {
880 if branch == repo.default_branch {
881 return Ok(Outcome::fail(
882 FailureCode::Invalid,
883 format!("Choose a branch other than {branch}."),
884 ));
885 }
886 let head: Option<String> = g1t_kit::call(
887 &self.repos,
888 "head",
889 &HeadArgs {
890 repo_id: repo.id.clone(),
891 branch: branch.to_owned(),
892 },
893 )
894 .await?;
895 let Some(head) = head else {
896 return Ok(Outcome::fail(
897 FailureCode::NotFound,
898 format!("There is no branch named {branch}. Push it first."),
899 ));
900 };
901 let existing = self
902 .db
903 .prepare(
904 "SELECT number AS n FROM pulls
905 WHERE repo_id = ? AND source_branch = ? AND status IN ('draft', 'open')",
906 )
907 .bind(&[repo.id.as_str().into(), branch.into()])?
908 .first::<NumberRow>(None)
909 .await?;
910 if let Some(existing) = existing {
911 return Ok(Outcome::fail(
912 FailureCode::Conflict,
913 format!("Pull request #{} is already open for {branch}.", existing.n),
914 ));
915 }
916 (None, Some(head))
917 }
918 None => {
919 let fork: Outcome<Repo> = g1t_kit::call(
920 &self.repos,
921 "fork_for_pull",
922 &ForkArgs {
923 source_id: repo.id.clone(),
924 pull_id: id.clone(),
925 actor: a.actor.clone(),
926 },
927 )
928 .await?;
929 (Some(check!(fork)), None)
930 }
931 };
932 // A branch already holds the work, so its pull request is ready for
933 // review from the start; one with a fork starts as a draft.
934 let status = if branch.is_some() { "open" } else { "draft" };
935 let body = Some(a.body.trim().to_owned()).filter(|body| !body.is_empty());
Work service in Rust, with RFC 3339 timestamps936
Issues and pull requests replace intents and attempts937 let number = self.next_number(&repo.id).await?;
Work service in Rust, with RFC 3339 timestamps938 let timestamp = rfc3339(now);
939 self.db
940 .prepare(
Issues and pull requests replace intents and attempts941 "INSERT INTO pulls
Pull requests from branches942 (id, repo_id, number, issue_id, issue_number, title, body, agent, runtime,
943 status, fork_repo_id, fork_namespace, fork_name, source_branch, head_commit,
944 author_id, author_name, created_at, updated_at)
945 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
Work service in Rust, with RFC 3339 timestamps946 )
947 .bind(&[
948 id.as_str().into(),
Issues and pull requests replace intents and attempts949 repo.id.as_str().into(),
950 number.into(),
951 optional(&issue.as_ref().map(|issue| issue.id.clone())),
952 optional_number(issue.as_ref().map(|issue| issue.number)),
953 title.into(),
Pull requests from branches954 optional(&body),
Work service in Rust, with RFC 3339 timestamps955 agent.into(),
956 runtime.into(),
Pull requests from branches957 status.into(),
958 optional(&fork.as_ref().map(|fork| fork.id.clone())),
959 optional(&fork.as_ref().map(|fork| fork.namespace.clone())),
960 optional(&fork.as_ref().map(|fork| fork.name.clone())),
961 optional(&branch.map(str::to_owned)),
962 optional(&head),
Work service in Rust, with RFC 3339 timestamps963 a.actor.id.as_str().into(),
964 a.actor.username.as_str().into(),
965 timestamp.as_str().into(),
966 timestamp.as_str().into(),
967 ])?
968 .run()
969 .await?;
Issues and pull requests replace intents and attempts970 let Some(pull) = self.pull(&repo.id, number).await? else {
971 return Ok(no_pull());
Work service in Rust, with RFC 3339 timestamps972 };
Agents as a team: lifecycle, merge queue, billing and a new shell973 self.manage(&pull).await?;
974 // Someone is on it now, so it is no longer waiting for an agent.
975 if let Some(issue) = pull.issue {
976 self.db
977 .prepare("UPDATE issues SET queued_by = NULL WHERE repo_id = ? AND number = ?")
978 .bind(&[repo.id.as_str().into(), issue.into()])?
979 .run()
980 .await?;
981 }
982 if let Some(issue) = pull.issue {
983 let text = if lifecycle::made_by_g1t(&pull) {
984 format!("assigned this to g1t-agent, which opened #{}", pull.number)
985 } else {
986 format!("opened #{} for this", pull.number)
987 };
988 self.note(&repo.id, issue, (&a.actor.id, &a.actor.username), &text)
989 .await?;
990 }
Issues and pull requests replace intents and attempts991 self.publish(
992 "pull.opened",
993 &repo.id,
994 &a.actor,
995 PullEvent {
996 agent: Some(pull.agent.clone()),
997 ..Self::pull_event(&pull)
Work service in Rust, with RFC 3339 timestamps998 },
Issues and pull requests replace intents and attempts999 )
Work service in Rust, with RFC 3339 timestamps1000 .await?;
Issues and pull requests replace intents and attempts1001 Ok(Outcome::Ok(pull))
Work service in Rust, with RFC 3339 timestamps1002 }
1003
Issues and pull requests replace intents and attempts1004 async fn list_pulls(&self, a: ListPullsArgs) -> Result<Outcome<Vec<Pull>>> {
1005 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
1006 let filter = match a.state {
1007 Some(State::Open) => "AND status IN ('draft', 'open')",
1008 Some(State::Closed) => "AND status IN ('merged', 'closed')",
1009 None => "",
Work service in Rust, with RFC 3339 timestamps1010 };
Issues and pull requests replace intents and attempts1011 let rows = self
1012 .db
1013 .prepare(format!(
1014 "SELECT * FROM pulls WHERE repo_id = ? {filter} ORDER BY number DESC LIMIT ?"
1015 ))
1016 .bind(&[repo.id.into(), LIST_PAGE.into()])?
1017 .all()
1018 .await?
1019 .results::<PullRow>()?;
1020 Ok(Outcome::Ok(rows.into_iter().map(Pull::from).collect()))
1021 }
1022
1023 async fn get_pull(&self, a: ViewArgs) -> Result<Outcome<PullDetail>> {
1024 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
1025 let issue = match pull.issue {
1026 Some(number) => self.issue(&repo.id, number).await?,
1027 None => None,
Work service in Rust, with RFC 3339 timestamps1028 };
Agents as a team: lifecycle, merge queue, billing and a new shell1029 let mut pull = pull;
1030 // Worked out on each push; this covers a pull request from before
1031 // that was recorded.
1032 if pull.files.is_empty() && pull.head_commit.is_some() {
1033 pull.files = self.refresh_files(&pull).await?;
1034 }
1035 // Everything else at once: none of it depends on the rest, and each
1036 // is a round trip of its own.
1037 let standing = async {
1038 let behind = self.is_behind(&repo.id, &pull).await?;
1039 let lifecycle = self
1040 .assess(&pull, &issue, behind)
1041 .await?
1042 .map(|(lifecycle, _)| lifecycle);
1043 Ok::<_, worker::Error>((behind, lifecycle))
1044 };
1045 let (((behind, lifecycle), (landing, stalled), comments), (checks, overlaps, review_pending)) =
1046 try_join(
1047 try_join3(standing, self.landing_state(&pull.id), self.comments(&repo.id, pull.number)),
1048 try_join3(
1049 self.latest_checks(&pull.id),
1050 self.overlaps(&pull),
1051 self.review_pending(&pull.id),
1052 ),
1053 )
1054 .await?;
Issues and pull requests replace intents and attempts1055 Ok(Outcome::Ok(PullDetail {
Agents as a team: lifecycle, merge queue, billing and a new shell1056 comments,
1057 checks,
1058 overlaps,
1059 behind,
1060 review_pending,
1061 lifecycle,
1062 landing,
1063 stalled,
Issues and pull requests replace intents and attempts1064 issue,
1065 pull,
1066 }))
Work service in Rust, with RFC 3339 timestamps1067 }
1068
Issues and pull requests replace intents and attempts1069 /// The pull request, if it is still active and `actor` opened it or
1070 /// belongs to the repository's workspace.
1071 async fn manageable_pull(
Work service in Rust, with RFC 3339 timestamps1072 &self,
Issues and pull requests replace intents and attempts1073 actor: &User,
1074 path: &RepoPath,
1075 number: u32,
1076 ) -> Result<Outcome<Pull>> {
1077 let (repo, pull) = check!(self.pull_at(path, number, &Some(actor.clone())).await?);
1078 if pull.author.id != actor.id && !actor.is_member(&repo.namespace) {
Work service in Rust, with RFC 3339 timestamps1079 return Ok(Outcome::fail(
Issues and pull requests replace intents and attempts1080 FailureCode::Forbidden,
1081 "Only whoever opened a pull request, or a member of the workspace, can change it.",
1082 ));
1083 }
1084 if !pull.status.is_active() {
1085 return Ok(Outcome::fail(
Work service in Rust, with RFC 3339 timestamps1086 FailureCode::Conflict,
Issues and pull requests replace intents and attempts1087 format!("This pull request is already {}.", pull.status.as_str()),
Work service in Rust, with RFC 3339 timestamps1088 ));
1089 }
Issues and pull requests replace intents and attempts1090 Ok(Outcome::Ok(pull))
1091 }
1092
Agents as a team: lifecycle, merge queue, billing and a new shell1093 async fn update_pull(&self, a: UpdatePullArgs) -> Result<Outcome<Pull>> {
1094 let pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1095 let assignees = match a.assignees {
1096 Some(names) => Some(check!(self.valid_assignees(names).await?)),
1097 None => None,
1098 };
1099 let reviewers = match a.reviewers {
1100 Some(names) => {
1101 // A g1t agent is not an account; everyone else has to be.
1102 let agent = names
1103 .iter()
1104 .any(|name| name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME));
1105 let people = names
1106 .into_iter()
1107 .filter(|name| !name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME))
1108 .collect();
1109 let mut reviewers = check!(self.valid_assignees(people).await?);
1110 reviewers.retain(|name| *name != pull.author.username);
1111 if agent {
1112 reviewers.insert(0, reviews::AGENT_NAME.to_owned());
1113 }
1114 Some(reviewers)
1115 }
1116 None => None,
1117 };
1118 self.db
1119 .prepare(
1120 "UPDATE pulls
1121 SET assignees = COALESCE(?, assignees), reviewers = COALESCE(?, reviewers),
1122 updated_at = ?
1123 WHERE id = ?",
1124 )
1125 .bind(&[
1126 optional(&assignees.as_ref().map(serde_json::to_string).transpose()?),
1127 optional(&reviewers.as_ref().map(serde_json::to_string).transpose()?),
1128 rfc3339(now_ms()).into(),
1129 pull.id.as_str().into(),
1130 ])?
1131 .run()
1132 .await?;
1133 if let Some(assignees) = &assignees {
1134 self.note_changes(
1135 &pull.repo_id,
1136 pull.number,
1137 &a.actor,
1138 &pull.assignees,
1139 assignees,
1140 ("assigned", "unassigned"),
1141 )
1142 .await?;
1143 }
1144 if let Some(reviewers) = &reviewers {
1145 self.note_changes(
1146 &pull.repo_id,
1147 pull.number,
1148 &a.actor,
1149 &pull.reviewers,
1150 reviewers,
1151 (
1152 "requested a review from",
1153 "withdrew the request for a review from",
1154 ),
1155 )
1156 .await?;
1157 }
1158 Ok(match self.pull(&pull.repo_id, pull.number).await? {
1159 Some(pull) => Outcome::Ok(pull),
1160 None => no_pull(),
1161 })
1162 }
1163
Issues and pull requests replace intents and attempts1164 /// Marks a draft ready for review, or updates the description of one
1165 /// that already is.
1166 async fn ready_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1167 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
Work service in Rust, with RFC 3339 timestamps1168 let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
1169 let now = rfc3339(now_ms());
1170 self.db
1171 .prepare(
Issues and pull requests replace intents and attempts1172 "UPDATE pulls SET status = 'open', body = COALESCE(?, body), updated_at = ?
Work service in Rust, with RFC 3339 timestamps1173 WHERE id = ?",
1174 )
1175 .bind(&[
1176 optional(&summary),
1177 now.as_str().into(),
Issues and pull requests replace intents and attempts1178 pull.id.as_str().into(),
Work service in Rust, with RFC 3339 timestamps1179 ])?
1180 .run()
1181 .await?;
Issues and pull requests replace intents and attempts1182 if pull.status == PullStatus::Draft {
1183 self.publish(
1184 "pull.ready",
1185 &pull.repo_id,
1186 &a.actor,
1187 Self::pull_event(&pull),
1188 )
1189 .await?;
1190 }
Agents as a team: lifecycle, merge queue, billing and a new shell1191 if pull.status == PullStatus::Draft {
1192 self.note(
1193 &pull.repo_id,
1194 pull.number,
1195 (&a.actor.id, &a.actor.username),
1196 "marked this ready for review",
1197 )
1198 .await?;
1199 }
Issues and pull requests replace intents and attempts1200 pull.status = PullStatus::Open;
1201 pull.body = summary.or(pull.body);
1202 pull.updated_at = now;
1203 Ok(Outcome::Ok(pull))
1204 }
1205
1206 async fn close_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1207 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1208 let now = rfc3339(now_ms());
1209 self.db
1210 .prepare("UPDATE pulls SET status = 'closed', updated_at = ? WHERE id = ?")
1211 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
1212 .run()
1213 .await?;
1214 self.publish(
1215 "pull.closed",
1216 &pull.repo_id,
1217 &a.actor,
1218 Self::pull_event(&pull),
1219 )
Work service in Rust, with RFC 3339 timestamps1220 .await?;
Agents as a team: lifecycle, merge queue, billing and a new shell1221 self.note(
1222 &pull.repo_id,
1223 pull.number,
1224 (&a.actor.id, &a.actor.username),
1225 "closed this",
1226 )
1227 .await?;
1228 // A closed pull request leaves the merge queue.
1229 if self
1230 .leave(&pull.repo_id, &pull, QueueState::Removed, Some("It was closed."))
1231 .await?
1232 {
1233 self.publish_as(
1234 "queue.changed",
1235 &pull.repo_id,
1236 None,
1237 g1t_contracts::events::QueueChanged {
1238 repo_id: pull.repo_id.clone(),
1239 },
1240 )
1241 .await?;
1242 }
Issues and pull requests replace intents and attempts1243 pull.status = PullStatus::Closed;
1244 pull.updated_at = now;
1245 Ok(Outcome::Ok(pull))
Work service in Rust, with RFC 3339 timestamps1246 }
1247
Issues and pull requests replace intents and attempts1248 /// Lands the pull request on the repository's default branch. Unless
1249 /// told to keep it open, that resolves the issue it was for: the issue
1250 /// closes naming this pull request, and the others still in progress
1251 /// for it close as superseded.
1252 async fn merge_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
Work service in Rust, with RFC 3339 timestamps1253 let viewer = Some(a.actor.clone());
Agents as a team: lifecycle, merge queue, billing and a new shell1254 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
Issues and pull requests replace intents and attempts1255 match pull.status {
1256 PullStatus::Open => {}
1257 PullStatus::Draft => {
1258 return Ok(Outcome::fail(
1259 FailureCode::Conflict,
1260 "This pull request is still a draft. Mark it ready for review first.",
1261 ));
1262 }
1263 status => {
1264 return Ok(Outcome::fail(
1265 FailureCode::Conflict,
1266 format!("This pull request is already {}.", status.as_str()),
1267 ));
1268 }
Work service in Rust, with RFC 3339 timestamps1269 }
Agents as a team: lifecycle, merge queue, billing and a new shell1270 let settings = self.settings(&repo.id).await?;
1271 // Where the repository does not allow it, asking to ignore the
1272 // checks changes nothing.
1273 if !a.ignore_checks || !settings.allow_ignoring_checks {
Acceptance checks in sandboxes, line comments and review verdicts1274 let waiting = match pull.check_status {
1275 Some(CheckStatus::Queued | CheckStatus::Running) => {
1276 Some("The acceptance checks are still running.")
1277 }
1278 Some(CheckStatus::Failed) => Some("The acceptance checks did not pass."),
1279 Some(CheckStatus::Errored) => Some("The acceptance checks could not be run."),
1280 Some(CheckStatus::Passed) | None => None,
1281 };
1282 if let Some(reason) = waiting {
Agents as a team: lifecycle, merge queue, billing and a new shell1283 let remedy = if settings.allow_ignoring_checks {
1284 "Wait or fix them, or merge anyway by ignoring the checks."
1285 } else {
1286 "This repository only merges pull requests whose checks pass."
1287 };
Acceptance checks in sandboxes, line comments and review verdicts1288 return Ok(Outcome::fail(
1289 FailureCode::Conflict,
Agents as a team: lifecycle, merge queue, billing and a new shell1290 format!("{reason} {remedy}"),
Acceptance checks in sandboxes, line comments and review verdicts1291 ));
1292 }
1293 }
Agents as a team: lifecycle, merge queue, billing and a new shell1294 if let Some(missing) = self.approvals_gap(&settings, &pull).await? {
1295 return Ok(Outcome::fail(FailureCode::Conflict, missing));
1296 }
Work service in Rust, with RFC 3339 timestamps1297
Agents as a team: lifecycle, merge queue, billing and a new shell1298 // A repository that merges through a queue: it joins the queue, and
1299 // lands once its state together with everything ahead has passed.
1300 if settings.merge_queue {
1301 if !a.actor.verified {
1302 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1303 }
1304 if !a.actor.is_member(&repo.namespace) {
1305 return Ok(Outcome::fail(
1306 FailureCode::Forbidden,
1307 "Only members of the repository's workspace can merge a pull request.",
1308 ));
1309 }
1310 return self.enqueue(&repo, &pull, &a.actor, a.keep_issue_open).await;
1311 }
1312
1313 // The default branch has moved under it. Unless the repository
1314 // insists on that being dealt with first, bring it up to date and
1315 // land it when that is done.
1316 if self.is_behind(&repo.id, &pull).await? {
1317 if settings.require_up_to_date {
1318 return Ok(Outcome::fail(
1319 FailureCode::Conflict,
1320 format!(
1321 "{} 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.",
1322 repo.default_branch
1323 ),
1324 ));
1325 }
1326 if !a.actor.verified {
1327 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1328 }
1329 if !a.actor.is_member(&repo.namespace) {
1330 return Ok(Outcome::fail(
1331 FailureCode::Forbidden,
1332 "Only members of the repository's workspace can merge a pull request.",
1333 ));
1334 }
1335 self.request_landing(&pull, &a.actor, a.keep_issue_open)
1336 .await?;
1337 return Ok(Outcome::Ok(pull));
1338 }
1339
Issues and pull requests replace intents and attempts1340 // Whether the actor may write to the repository is decided by repos.
Work service in Rust, with RFC 3339 timestamps1341 let landed: Outcome<Landed> = g1t_kit::call(
1342 &self.repos,
1343 "land",
1344 &LandArgs {
Pull requests from branches1345 // A pull request from a branch lands from the repository itself.
1346 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
1347 branch: pull.branch.clone(),
Work service in Rust, with RFC 3339 timestamps1348 actor: a.actor.clone(),
1349 },
1350 )
1351 .await?;
Issues and pull requests replace intents and attempts1352 let landed = check!(landed);
Agents as a team: lifecycle, merge queue, billing and a new shell1353 Ok(Outcome::Ok(
1354 self.record_merge(&repo, pull, &a.actor, a.keep_issue_open, landed)
1355 .await?,
1356 ))
1357 }
Work service in Rust, with RFC 3339 timestamps1358
Agents as a team: lifecycle, merge queue, billing and a new shell1359 /// Records a pull request as merged once the default branch holds it:
1360 /// closes its issue, supersedes the others for it, and says so.
1361 pub(crate) async fn record_merge(
1362 &self,
1363 repo: &Repo,
1364 mut pull: Pull,
1365 actor: &User,
1366 keep_issue_open: bool,
1367 landed: Landed,
1368 ) -> Result<Pull> {
1369 let issue = match pull.issue {
1370 Some(number) if !keep_issue_open => self
1371 .issue(&repo.id, number)
1372 .await?
1373 .filter(|issue| issue.state == State::Open),
1374 _ => None,
1375 };
Work service in Rust, with RFC 3339 timestamps1376 let now = rfc3339(now_ms());
Issues and pull requests replace intents and attempts1377 let mut statements = vec![
1378 self.db
1379 .prepare(
1380 "UPDATE pulls
1381 SET status = 'merged', head_commit = ?, merge_base = ?, merged_by = ?,
1382 merged_at = ?, updated_at = ?
1383 WHERE id = ?",
1384 )
1385 .bind(&[
1386 landed.commit.as_str().into(),
1387 optional(&landed.previous),
Agents as a team: lifecycle, merge queue, billing and a new shell1388 actor.username.as_str().into(),
Issues and pull requests replace intents and attempts1389 now.as_str().into(),
1390 now.as_str().into(),
1391 pull.id.as_str().into(),
1392 ])?,
1393 ];
1394 if let Some(issue) = &issue {
1395 statements.push(
Work service in Rust, with RFC 3339 timestamps1396 self.db
1397 .prepare(
Issues and pull requests replace intents and attempts1398 "UPDATE issues
1399 SET state = 'closed', reason = 'completed', resolved_by = ?,
1400 closed_at = ?, updated_at = ?
Work service in Rust, with RFC 3339 timestamps1401 WHERE id = ?",
1402 )
1403 .bind(&[
Issues and pull requests replace intents and attempts1404 pull.number.into(),
1405 now.as_str().into(),
Work service in Rust, with RFC 3339 timestamps1406 now.as_str().into(),
Issues and pull requests replace intents and attempts1407 issue.id.as_str().into(),
Work service in Rust, with RFC 3339 timestamps1408 ])?,
Issues and pull requests replace intents and attempts1409 );
1410 statements.push(
Work service in Rust, with RFC 3339 timestamps1411 self.db
Issues and pull requests replace intents and attempts1412 .prepare(
1413 "UPDATE pulls SET status = 'closed', superseded_by = ?, updated_at = ?
1414 WHERE issue_id = ? AND id != ? AND status IN ('draft', 'open')",
1415 )
1416 .bind(&[
1417 pull.number.into(),
1418 now.as_str().into(),
1419 issue.id.as_str().into(),
1420 pull.id.as_str().into(),
1421 ])?,
1422 );
1423 }
1424 self.db.batch(statements).await?;
1425
1426 self.publish(
1427 "pull.merged",
1428 &repo.id,
Agents as a team: lifecycle, merge queue, billing and a new shell1429 actor,
Issues and pull requests replace intents and attempts1430 PullEvent {
Work service in Rust, with RFC 3339 timestamps1431 commit: Some(landed.commit.clone()),
Issues and pull requests replace intents and attempts1432 ..Self::pull_event(&pull)
Work service in Rust, with RFC 3339 timestamps1433 },
Issues and pull requests replace intents and attempts1434 )
Work service in Rust, with RFC 3339 timestamps1435 .await?;
Issues and pull requests replace intents and attempts1436 if let Some(issue) = &issue {
1437 self.publish(
1438 "issue.closed",
1439 &repo.id,
Agents as a team: lifecycle, merge queue, billing and a new shell1440 actor,
Issues and pull requests replace intents and attempts1441 IssueEvent {
1442 reason: Some(IssueReason::Completed.as_str()),
1443 resolved_by: Some(pull.number),
1444 ..Self::issue_event(issue)
1445 },
1446 )
1447 .await?;
1448 }
Work service in Rust, with RFC 3339 timestamps1449
Agents as a team: lifecycle, merge queue, billing and a new shell1450 let who = (actor.id.as_str(), actor.username.as_str());
1451 self.note(&repo.id, pull.number, who, "merged this").await?;
1452 if let Some(issue) = &issue {
1453 self.note(
1454 &repo.id,
1455 issue.number,
1456 who,
1457 &format!("closed this by merging #{}", pull.number),
1458 )
1459 .await?;
1460 }
Issues and pull requests replace intents and attempts1461 pull.status = PullStatus::Merged;
Agents as a team: lifecycle, merge queue, billing and a new shell1462 pull.head_commit = Some(landed.commit.clone());
Issues and pull requests replace intents and attempts1463 pull.merge_base = landed.previous;
Agents as a team: lifecycle, merge queue, billing and a new shell1464 pull.merged_by = Some(actor.username.clone());
Issues and pull requests replace intents and attempts1465 pull.merged_at = Some(now.clone());
1466 pull.updated_at = now;
Agents as a team: lifecycle, merge queue, billing and a new shell1467 Ok(pull)
Work service in Rust, with RFC 3339 timestamps1468 }
1469
Issues and pull requests replace intents and attempts1470 async fn list_active_pulls(&self, a: ViewerArgs) -> Result<Vec<ActivePull>> {
Work service in Rust, with RFC 3339 timestamps1471 let Some(viewer) = a.viewer else {
1472 return Ok(Vec::new());
1473 };
Agents as a team: lifecycle, merge queue, billing and a new shell1474 let found = self
Work service in Rust, with RFC 3339 timestamps1475 .db
1476 .prepare(
Issues and pull requests replace intents and attempts1477 "SELECT * FROM pulls
1478 WHERE author_id = ? AND status IN ('draft', 'open')
Work service in Rust, with RFC 3339 timestamps1479 ORDER BY updated_at DESC LIMIT 50",
1480 )
1481 .bind(&[viewer.id.into()])?
1482 .all()
Agents as a team: lifecycle, merge queue, billing and a new shell1483 .await?;
1484 let snapshots = found.results::<Snapshot>()?;
1485 let pulls: Vec<Pull> = found.results::<PullRow>()?.into_iter().map(Pull::from).collect();
1486 // Each one at once: its issue, and where it stands. That is the
1487 // remembered assessment when there is one, and worked out otherwise.
1488 try_join_all(pulls.into_iter().zip(snapshots).map(|(pull, snapshot)| async move {
Issues and pull requests replace intents and attempts1489 let issue = match pull.issue {
1490 Some(number) => self.issue(&pull.repo_id, number).await?,
1491 None => None,
1492 };
Agents as a team: lifecycle, merge queue, billing and a new shell1493 // Only a pull request g1t is seeing through has a lifecycle.
1494 let lifecycle = if !lifecycle::made_by_g1t(&pull) || snapshot.managed == 0 {
1495 None
1496 } else if let (Some(stage), Some(detail)) = (snapshot.stage, snapshot.stage_detail) {
1497 Some(Lifecycle {
1498 stage,
1499 detail,
1500 revisions: snapshot.revisions,
1501 })
1502 } else {
1503 let behind = self.is_behind(&pull.repo_id, &pull).await?;
1504 self.assess(&pull, &issue, behind)
1505 .await?
1506 .map(|(lifecycle, _)| lifecycle)
1507 };
1508 Ok::<_, worker::Error>(ActivePull {
1509 pull,
1510 issue,
1511 lifecycle,
1512 })
1513 }))
1514 .await
Work service in Rust, with RFC 3339 timestamps1515 }
1516
Issues and pull requests replace intents and attempts1517 // --- Sessions ----------------------------------------------------------
1518
Work service in Rust, with RFC 3339 timestamps1519 async fn append_session(&self, a: AppendSessionArgs) -> Result<Outcome<Appended>> {
1520 if a.entries.is_empty() {
1521 return Ok(Outcome::Ok(Appended { count: 0 }));
1522 }
1523 if a.entries.len() > MAX_ENTRY_BATCH {
1524 return Ok(Outcome::fail(
1525 FailureCode::Invalid,
1526 format!("Send at most {MAX_ENTRY_BATCH} entries at a time."),
1527 ));
1528 }
Issues and pull requests replace intents and attempts1529 let viewer = Some(a.actor.clone());
1530 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
1531 if pull.author.id != a.actor.id {
1532 return Ok(Outcome::fail(
1533 FailureCode::Forbidden,
1534 "Only whoever opened a pull request can record its session.",
1535 ));
1536 }
Work service in Rust, with RFC 3339 timestamps1537
1538 let now = rfc3339(now_ms());
1539 let count = a.entries.len() as u32;
1540 let mut statements = Vec::with_capacity(a.entries.len() + 1);
1541 for entry in a.entries {
1542 let kind = serde_json::to_value(entry.kind)?;
1543 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
1544 // Each insert takes the next sequence number itself, so two
1545 // writers appending at once cannot collide.
1546 statements.push(
1547 self.db
1548 .prepare(
Issues and pull requests replace intents and attempts1549 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
Work service in Rust, with RFC 3339 timestamps1550 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
Issues and pull requests replace intents and attempts1551 FROM session_entries WHERE pull_id = ?",
Work service in Rust, with RFC 3339 timestamps1552 )
1553 .bind(&[
Issues and pull requests replace intents and attempts1554 pull.id.as_str().into(),
Work service in Rust, with RFC 3339 timestamps1555 kind.as_str().unwrap_or("note").into(),
1556 text.into(),
1557 optional(&entry.tool),
Issues and pull requests replace intents and attempts1558 optional(&entry.commit.or_else(|| pull.head_commit.clone())),
Work service in Rust, with RFC 3339 timestamps1559 now.as_str().into(),
Issues and pull requests replace intents and attempts1560 pull.id.as_str().into(),
Work service in Rust, with RFC 3339 timestamps1561 ])?,
1562 );
1563 }
1564 statements.push(
1565 self.db
Issues and pull requests replace intents and attempts1566 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
1567 .bind(&[now.as_str().into(), pull.id.as_str().into()])?,
Work service in Rust, with RFC 3339 timestamps1568 );
1569 self.db.batch(statements).await?;
Issues and pull requests replace intents and attempts1570 self.publish(
1571 "session.appended",
1572 &pull.repo_id,
1573 &a.actor,
1574 SessionAppended {
1575 pull_id: pull.id.clone(),
1576 repo_id: pull.repo_id.clone(),
1577 number: pull.number,
Work service in Rust, with RFC 3339 timestamps1578 count,
1579 },
Issues and pull requests replace intents and attempts1580 )
Work service in Rust, with RFC 3339 timestamps1581 .await?;
1582 Ok(Outcome::Ok(Appended { count }))
1583 }
1584
Issues and pull requests replace intents and attempts1585 async fn read_session(&self, a: ViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
1586 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
Work service in Rust, with RFC 3339 timestamps1587 let rows = self
1588 .db
1589 .prepare(
1590 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
Issues and pull requests replace intents and attempts1591 WHERE pull_id = ? AND seq > ? ORDER BY seq LIMIT ?",
Work service in Rust, with RFC 3339 timestamps1592 )
Issues and pull requests replace intents and attempts1593 .bind(&[pull.id.into(), a.after_seq.into(), SESSION_PAGE.into()])?
Work service in Rust, with RFC 3339 timestamps1594 .all()
1595 .await?
1596 .results::<SessionRow>()?;
1597 Ok(Outcome::Ok(
1598 rows.into_iter().map(SessionEntry::from).collect(),
1599 ))
1600 }
1601
Events service in Rust, with RFC 3339 times and accurate push events1602 /// A push moves the head of the pull request it concerns: the one whose
1603 /// fork was pushed to, or the one opened from the branch that moved.
1604 async fn on_event(&self, event: &Event) -> Result<()> {
Work service in Rust, with RFC 3339 timestamps1605 if event.kind != "git.push" {
1606 return Ok(());
1607 }
Events service in Rust, with RFC 3339 times and accurate push events1608 let (Some(repo_id), Some(after), Some(git_ref)) = (
1609 event.repo_id.as_deref(),
1610 event.data["after"].as_str(),
1611 event.data["ref"].as_str(),
1612 ) else {
Work service in Rust, with RFC 3339 timestamps1613 return Ok(());
1614 };
Pull requests from branches1615 let now = rfc3339(now_ms());
Agents as a team: lifecycle, merge queue, billing and a new shell1616 // The head moved, so whatever the checks said no longer applies, and
1617 // whatever step g1t was waiting on has been taken.
Acceptance checks in sandboxes, line comments and review verdicts1618 let moved = "UPDATE pulls
Agents as a team: lifecycle, merge queue, billing and a new shell1619 SET head_commit = ?, updated_at = ?, check_status = NULL, check_run_id = NULL,
1620 working_on = NULL, working_until = NULL, stalled = NULL";
Acceptance checks in sandboxes, line comments and review verdicts1621 let active = "status IN ('draft', 'open') AND head_commit IS NOT ?";
1622 let returning = "RETURNING id, repo_id, number, issue_number, status";
1623 let mut pulls: Vec<MovedRow> = Vec::new();
Events service in Rust, with RFC 3339 times and accurate push events1624 // A fork carries its pull request on its default branch.
1625 if event.data["defaultBranch"].as_bool() == Some(true) {
Acceptance checks in sandboxes, line comments and review verdicts1626 pulls.extend(
Events service in Rust, with RFC 3339 times and accurate push events1627 self.db
1628 .prepare(format!(
Acceptance checks in sandboxes, line comments and review verdicts1629 "{moved} WHERE fork_repo_id = ? AND {active} {returning}"
Events service in Rust, with RFC 3339 times and accurate push events1630 ))
Acceptance checks in sandboxes, line comments and review verdicts1631 .bind(&[
1632 after.into(),
1633 now.as_str().into(),
1634 repo_id.into(),
1635 after.into(),
1636 ])?
1637 .all()
1638 .await?
1639 .results::<MovedRow>()?,
Events service in Rust, with RFC 3339 times and accurate push events1640 );
1641 }
1642 if let Some(branch) = git_ref.strip_prefix("refs/heads/") {
Acceptance checks in sandboxes, line comments and review verdicts1643 pulls.extend(
Pull requests from branches1644 self.db
Events service in Rust, with RFC 3339 times and accurate push events1645 .prepare(format!(
Acceptance checks in sandboxes, line comments and review verdicts1646 "{moved} WHERE repo_id = ? AND source_branch = ? AND {active} {returning}"
Events service in Rust, with RFC 3339 times and accurate push events1647 ))
1648 .bind(&[
1649 after.into(),
1650 now.as_str().into(),
1651 repo_id.into(),
1652 branch.into(),
Acceptance checks in sandboxes, line comments and review verdicts1653 after.into(),
1654 ])?
1655 .all()
1656 .await?
1657 .results::<MovedRow>()?,
Events service in Rust, with RFC 3339 times and accurate push events1658 );
Pull requests from branches1659 }
Agents as a team: lifecycle, merge queue, billing and a new shell1660 // What each now changes, so overlaps show while the work is under way.
1661 for moved in &pulls {
1662 if let Some(pull) = self.pull_by_id(&moved.id).await? {
1663 self.refresh_files(&pull).await?;
1664 }
1665 }
1666 // A merge that was waiting for this push to bring it up to date.
1667 for moved in &pulls {
1668 self.land_if_requested(&moved.id).await?;
1669 }
Acceptance checks in sandboxes, line comments and review verdicts1670 // A draft is announced when it is marked ready instead.
1671 for pull in pulls
1672 .into_iter()
1673 .filter(|pull| pull.status == PullStatus::Open)
1674 {
1675 self.publish_as(
1676 "pull.updated",
1677 &pull.repo_id,
1678 event.actor.clone(),
1679 PullEvent {
1680 pull_id: pull.id,
1681 repo_id: pull.repo_id.clone(),
1682 number: pull.number,
1683 issue: pull.issue_number,
1684 commit: Some(after.to_owned()),
1685 ..PullEvent::default()
1686 },
1687 )
1688 .await?;
1689 }
Work service in Rust, with RFC 3339 timestamps1690 Ok(())
1691 }
1692}
1693
1694fn service(env: &Env) -> Result<Work> {
1695 Ok(Work {
1696 db: env.d1("DB")?,
Agents as a team: lifecycle, merge queue, billing and a new shell1697 identity: env.service("IDENTITY")?,
Work service in Rust, with RFC 3339 timestamps1698 repos: env.service("REPOS")?,
Events service in Rust, with RFC 3339 times and accurate push events1699 events: env.service("EVENTS")?,
Work service in Rust, with RFC 3339 timestamps1700 })
1701}
1702
1703#[event(fetch)]
1704async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
1705 let Some(method) = rpc_method(&request) else {
1706 return Response::error("Not found", 404);
1707 };
1708 let body: serde_json::Value = request.json().await?;
1709 let work = service(&env)?;
1710
1711 match method.as_str() {
Issues and pull requests replace intents and attempts1712 "open_issue" => reply(&work.open_issue(args(body)?).await?),
1713 "list_issues" => reply(&work.list_issues(args(body)?).await?),
1714 "get_issue" => reply(&work.get_issue(args(body)?).await?),
1715 "update_issue" => reply(&work.update_issue(args(body)?).await?),
1716 "close_issue" => reply(&work.close_issue(args(body)?).await?),
1717 "reopen_issue" => reply(&work.reopen_issue(args(body)?).await?),
1718 "list_labels" => reply(&work.list_labels(args(body)?).await?),
1719 "counts" => reply(&work.counts(args(body)?).await?),
1720 "add_comment" => reply(&work.add_comment(args(body)?).await?),
Acceptance checks in sandboxes, line comments and review verdicts1721 "start_checks" => reply(&work.start_checks(args(body)?).await?),
1722 "report_checks" => reply(&work.report_checks(args(body)?).await?),
Agents as a team: lifecycle, merge queue, billing and a new shell1723 "start_review" => reply(&work.start_review(args(body)?).await?),
1724 "advance" => reply(&work.advance(args(body)?).await?),
1725 "stall" => reply(&work.stall(args(body)?).await?),
1726 "managed_pulls" => reply(&work.managed_pulls(args(body)?).await?),
1727 "queue" => reply(&work.queue(args(body)?).await?),
1728 "queue_build" => reply(&work.queue_build(args(body)?).await?),
1729 "report_queue" => reply(&work.report_queue(args(body)?).await?),
1730 "remove_from_queue" => reply(&work.remove_from_queue(args(body)?).await?),
1731 "catch_up_job" => reply(&work.catch_up_job(args(body)?).await?),
1732 "get_settings" => reply(&work.get_settings(args(body)?).await?),
1733 "update_settings" => reply(&work.update_settings(args(body)?).await?),
1734 "report_review" => reply(&work.report_review(args(body)?).await?),
Issues and pull requests replace intents and attempts1735 "open_pull" => reply(&work.open_pull(args(body)?).await?),
1736 "list_pulls" => reply(&work.list_pulls(args(body)?).await?),
1737 "get_pull" => reply(&work.get_pull(args(body)?).await?),
Agents as a team: lifecycle, merge queue, billing and a new shell1738 "update_pull" => reply(&work.update_pull(args(body)?).await?),
Issues and pull requests replace intents and attempts1739 "ready_pull" => reply(&work.ready_pull(args(body)?).await?),
1740 "close_pull" => reply(&work.close_pull(args(body)?).await?),
1741 "merge_pull" => reply(&work.merge_pull(args(body)?).await?),
1742 "list_active_pulls" => reply(&work.list_active_pulls(args(body)?).await?),
Agents as a team: lifecycle, merge queue, billing and a new shell1743 "start_plan" => reply(&work.start_plan(args(body)?).await?),
1744 "report_plan" => reply(&work.report_plan(args(body)?).await?),
1745 "get_plan" => reply(&work.get_plan(args(body)?).await?),
1746 "list_plans" => reply(&work.list_plans(args(body)?).await?),
1747 "apply_plan" => reply(&work.apply_plan(args(body)?).await?),
1748 "queue_issue" => reply(&work.queue_issue(args(body)?).await?),
1749 "ready_issues" => reply(&work.ready_issues(args(body)?).await?),
1750 "list_assigned_issues" => reply(&work.list_assigned_issues(args(body)?).await?),
Work service in Rust, with RFC 3339 timestamps1751 "append_session" => reply(&work.append_session(args(body)?).await?),
1752 "read_session" => reply(&work.read_session(args(body)?).await?),
1753 _ => Response::error("Unknown method", 404),
1754 }
1755}
1756
1757/// Events from the bus, delivered on this service's own queue.
1758#[event(queue)]
Events service in Rust, with RFC 3339 times and accurate push events1759async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
Work service in Rust, with RFC 3339 timestamps1760 let work = service(&env)?;
1761 for message in batch.messages()? {
1762 work.on_event(message.body()).await?;
1763 message.ack();
1764 }
1765 Ok(())
1766}