pr_01m47d24b0e6n91zwymwxg0vpx/services/work/src/lib.rs

2,021 lines78,466 bytesCodeBlame
1//! The work service: issues, pull requests, comments and sessions.
2//!
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
7mod authored;
8mod capture;
9mod checks;
10mod compute;
11mod guardrails;
12mod lifecycle;
13mod memory;
14mod mentions;
15mod mergeability;
16mod plans;
17mod messages;
18mod queue;
19mod retired;
20mod reviews;
21mod rows;
22mod runs;
23mod settings;
24mod statuses;
25
26use g1t_contracts::events::{
27 CommentCreated, Event, IssueEvent, NewEvent, Publish, PullEvent, SessionAppended,
28};
29use g1t_contracts::identity::UsernameArgs;
30use g1t_contracts::repos::{
31 ForkArgs, GetArgs, HeadArgs, LandArgs, Landed, NeedsAgentReason, PullBranchUpdate, Repo, RepoPath,
32 UpdatePullBranchArgs,
33};
34use g1t_contracts::access::{self, Capability, Denied};
35use g1t_contracts::time::rfc3339;
36use g1t_contracts::work::*;
37use futures_util::future::{try_join, try_join3, try_join_all};
38use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
39use g1t_kit::{args, now_ms, reply, rpc_method};
40use serde::Serialize;
41use worker::wasm_bindgen::JsValue;
42use worker::{
43 Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event,
44};
45
46use retired::writable;
47use rows::{CommentRow, IssueRow, MovedRow, NumberRow, PullRow, SessionRow, Snapshot, ValueRow};
48
49const SOURCE: &str = "work";
50const MAX_ENTRY_BATCH: usize = 200;
51const MAX_ENTRY_CHARS: usize = 64_000;
52const MAX_TITLE_CHARS: usize = 200;
53const SESSION_PAGE: u32 = 500;
54const LIST_PAGE: u32 = 100;
55const MAX_ASSIGNEES: usize = 10;
56const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
57
58const ISSUE_COLUMNS: &str = "issues.*,
59 (SELECT count(*) FROM pulls WHERE pulls.issue_id = issues.id) AS pull_count,
60 (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,
64 (SELECT count(*) FROM comments
65 WHERE comments.repo_id = issues.repo_id AND comments.number = issues.number
66 AND comments.kind = 'comment') AS comment_count";
67
68fn no_issue<T>() -> Outcome<T> {
69 Outcome::fail(FailureCode::NotFound, "Issue not found.")
70}
71
72fn no_pull<T>() -> Outcome<T> {
73 Outcome::fail(FailureCode::NotFound, "Pull request not found.")
74}
75
76/// 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
89fn optional(value: &Option<String>) -> JsValue {
90 value.as_deref().map_or(JsValue::NULL, JsValue::from)
91}
92
93fn 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
127struct Work {
128 db: D1Database,
129 identity: Fetcher,
130 repos: Fetcher,
131 events: Fetcher,
132 /// GitHub Actions: runs a merge queue's `merge_group` workflows.
133 actions: Fetcher,
134}
135
136impl Work {
137 async fn publish<T: Serialize>(
138 &self,
139 kind: &'static str,
140 repo_id: &str,
141 actor: &User,
142 data: T,
143 ) -> Result<()> {
144 self.publish_as(kind, repo_id, Some(actor.id.clone()), data)
145 .await
146 }
147
148 /// 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
150 /// would otherwise look missing to them. The membership given reads
151 /// and nothing more: it is for looking, never for acting.
152 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 {
163 author.workspaces.push(g1t_contracts::Membership {
164 base_permission: Some(access::BasePermission::Read),
165 ..g1t_contracts::Membership::member(path.namespace.to_lowercase())
166 });
167 }
168 Ok(Some(author))
169 }
170
171 /// 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<()> {
179 let event = NewEvent {
180 kind,
181 source: SOURCE,
182 repo_id: Some(repo_id.to_owned()),
183 actor,
184 data,
185 };
186 g1t_kit::call(
187 &self.events,
188 "publish",
189 &Publish {
190 events: vec![event],
191 },
192 )
193 .await
194 }
195
196 /// 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>> {
199 g1t_kit::call(
200 &self.repos,
201 "get",
202 &GetArgs {
203 path: path.clone(),
204 viewer: viewer.clone(),
205 },
206 )
207 .await
208 }
209
210 /// 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()))
225 }
226
227 async fn issue(&self, repo_id: &str, number: u32) -> Result<Option<Issue>> {
228 Ok(self
229 .db
230 .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)
235 .await?
236 .map(Issue::from))
237 }
238
239 async fn pull(&self, repo_id: &str, number: u32) -> Result<Option<Pull>> {
240 Ok(self
241 .db
242 .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()
257 .await?
258 .results::<CommentRow>()?;
259 Ok(rows.into_iter().map(Comment::from).collect())
260 }
261
262 /// 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());
271 };
272 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
294 /// 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
362 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()
368 }
369 }
370
371 /// 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
385 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()
392 }
393 }
394
395 // --- Issues ------------------------------------------------------------
396
397 async fn open_issue(&self, a: OpenIssueArgs) -> Result<Outcome<Issue>> {
398 if !a.actor.verified {
399 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
400 }
401 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 {
406 return Ok(Outcome::fail(
407 FailureCode::Invalid,
408 "An issue can have up to 10 labels of up to 40 characters each.",
409 ));
410 };
411 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
412 check!(writable(&repo));
413 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();
421 let id = new_id("iss", now);
422 let number = self.next_number(&repo.id).await?;
423 let timestamp = rfc3339(now);
424 self.db
425 .prepare(
426 "INSERT INTO issues
427 (id, repo_id, number, title, body, labels, checks, author_id, author_name,
428 created_at, updated_at)
429 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
430 )
431 .bind(&[
432 id.as_str().into(),
433 repo.id.as_str().into(),
434 number.into(),
435 title.into(),
436 a.body.trim().into(),
437 serde_json::to_string(&labels)?.into(),
438 serde_json::to_string(&checks)?.into(),
439 a.actor.id.as_str().into(),
440 a.actor.username.as_str().into(),
441 timestamp.as_str().into(),
442 timestamp.as_str().into(),
443 ])?
444 .run()
445 .await?;
446 let Some(issue) = self.issue(&repo.id, number).await? else {
447 return Ok(no_issue());
448 };
449 self.apply_label_rule(&a.actor, &issue, &[]).await?;
450 self.publish(
451 "issue.opened",
452 &repo.id,
453 &a.actor,
454 IssueEvent {
455 title: Some(issue.title.clone()),
456 ..Self::issue_event(&issue)
457 },
458 )
459 .await?;
460 Ok(Outcome::Ok(issue))
461 }
462
463 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());
470 let rows = self
471 .db
472 .prepare(format!(
473 "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 ?"
478 ))
479 .bind(&[
480 repo.id.into(),
481 state.clone(),
482 state,
483 optional(&label),
484 optional(&label),
485 LIST_PAGE.into(),
486 ])?
487 .all()
488 .await?
489 .results::<IssueRow>()?;
490 Ok(Outcome::Ok(rows.into_iter().map(Issue::from).collect()))
491 }
492
493 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
496 .db
497 .prepare("SELECT * FROM pulls WHERE issue_id = ? ORDER BY number")
498 .bind(&[issue.id.as_str().into()])?
499 .all()
500 .await?
501 .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,
506 }))
507 }
508
509 /// The issue, if `actor` wrote it or may triage the repository's issues.
510 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?);
517 check!(writable(&repo));
518 if issue.author.id != actor.id {
519 check!(allowed(Some(actor), &repo, Capability::Triage));
520 }
521 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 };
541 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()?;
546 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),
551 labels = COALESCE(?, labels), assignees = COALESCE(?, assignees),
552 updated_at = ?
553 WHERE id = ?",
554 )
555 .bind(&[
556 optional(&title),
557 optional(&body),
558 optional(&labels),
559 optional(&assigned),
560 rfc3339(now_ms()).into(),
561 issue.id.as_str().into(),
562 ])?
563 .run()
564 .await?;
565 let before = issue.assignees.clone();
566 let labels_before = issue.labels.clone();
567 let Some(issue) = self.issue(&issue.repo_id, issue.number).await? else {
568 return Ok(no_issue());
569 };
570 self.apply_label_rule(&a.actor, &issue, &labels_before).await?;
571 self.publish(
572 "issue.updated",
573 &issue.repo_id,
574 &a.actor,
575 Self::issue_event(&issue),
576 )
577 .await?;
578 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 }
599 Ok(Outcome::Ok(issue))
600 }
601
602 /// 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
656 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 {
659 return Ok(Outcome::fail(
660 FailureCode::Conflict,
661 "This issue is already closed.",
662 ));
663 }
664 let reason = a.reason.unwrap_or(IssueReason::Completed);
665 let now = rfc3339(now_ms());
666 self.db
667 .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 ])?
677 .run()
678 .await?;
679 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)
686 },
687 )
688 .await?;
689 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?;
699 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?;
732 self.note(
733 &issue.repo_id,
734 issue.number,
735 (&a.actor.id, &a.actor.username),
736 "reopened this",
737 )
738 .await?;
739 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))
745 }
746
747 /// 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>> {
792 if !a.actor.verified {
793 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
794 }
795 let body = a.body.trim();
796 // An approval speaks for itself; anything else has to say something.
797 if body.is_empty() && a.verdict != Some(Verdict::Approve) {
798 return Ok(Outcome::fail(
799 FailureCode::Invalid,
800 "A comment cannot be empty.",
801 ));
802 }
803 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());
809 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?);
816 check!(writable(&repo));
817 // The number names an issue or a pull request, never both.
818 let mut pull_id = None;
819 let table = if self.issue(&repo.id, a.number).await?.is_some() {
820 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 }
826 "issues"
827 } 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 }
834 pull_id = Some(pull.id.clone());
835 "pulls"
836 } else {
837 return Ok(Outcome::fail(
838 FailureCode::NotFound,
839 "No issue or pull request has that number.",
840 ));
841 };
842
843 let now = now_ms();
844 let comment = Comment {
845 kind: CommentKind::Comment,
846 id: new_id("cmt", now),
847 author: a.actor.clone(),
848 body: body.to_owned(),
849 path: path.map(str::to_owned),
850 line,
851 verdict: a.verdict,
852 created_at: rfc3339(now),
853 };
854 self.db
855 .batch(vec![
856 self.db
857 .prepare(
858 "INSERT INTO comments
859 (id, repo_id, number, author_id, author_name, body, path, line,
860 verdict, created_at)
861 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
862 )
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(),
870 optional(&comment.path),
871 optional_number(line),
872 a.verdict
873 .map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
874 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?;
887 self.note_mention(&a.actor, &repo, a.number, &comment, pull_id.as_deref()).await?;
888 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,
896 pull_id,
897 verdict: a.verdict,
898 },
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));
909 }
910 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
911 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 }
917 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 };
938 let agent = match a.agent.trim() {
939 "" => "agent",
940 agent => agent,
941 };
942 let runtime = match a.runtime {
943 Runtime::Hosted => "hosted",
944 Runtime::External => "external",
945 };
946
947 let now = now_ms();
948 let id = new_id("pr", now);
949 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());
1014
1015 let number = self.next_number(&repo.id).await?;
1016 let timestamp = rfc3339(now);
1017 self.db
1018 .prepare(
1019 "INSERT INTO pulls
1020 (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 (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
1024 )
1025 .bind(&[
1026 id.as_str().into(),
1027 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(),
1032 optional(&body),
1033 agent.into(),
1034 runtime.into(),
1035 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),
1041 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?;
1048 let Some(pull) = self.pull(&repo.id, number).await? else {
1049 return Ok(no_pull());
1050 };
1051 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 }
1069 self.publish(
1070 "pull.opened",
1071 &repo.id,
1072 &a.actor,
1073 PullEvent {
1074 agent: Some(pull.agent.clone()),
1075 ..Self::pull_event(&pull)
1076 },
1077 )
1078 .await?;
1079 Ok(Outcome::Ok(pull))
1080 }
1081
1082 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 => "",
1088 };
1089 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,
1106 };
1107 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 {
1116 // 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?;
1120 let lifecycle = self
1121 .assess(&pull, &issue, behind)
1122 .await?
1123 .map(|(lifecycle, _)| lifecycle);
1124 Ok::<_, worker::Error>((behind, lifecycle, merge))
1125 };
1126 let (((behind, lifecycle, (mergeable, conflicts)), (landing, stalled), comments), (checks, overlaps, review_pending)) =
1127 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?;
1136 Ok(Outcome::Ok(PullDetail {
1137 comments,
1138 checks,
1139 overlaps,
1140 behind,
1141 review_pending,
1142 lifecycle,
1143 landing,
1144 stalled,
1145 messages: self.messages(&pull.id).await?,
1146 statuses: self.statuses(&repo.id, pull.head_commit.as_deref()).await?,
1147 mergeable,
1148 conflicts,
1149 earlier_checks: self.earlier_checks(&pull.id).await?,
1150 issue,
1151 pull,
1152 }))
1153 }
1154
1155 /// The pull request, if it is still active and `actor` opened it or
1156 /// may triage the repository's pull requests.
1157 async fn manageable_pull(
1158 &self,
1159 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?);
1164 check!(writable(&repo));
1165 if pull.author.id != actor.id {
1166 check!(allowed(Some(actor), &repo, Capability::Triage));
1167 }
1168 if !pull.status.is_active() {
1169 return Ok(Outcome::fail(
1170 FailureCode::Conflict,
1171 format!("This pull request is already {}.", pull.status.as_str()),
1172 ));
1173 }
1174 Ok(Outcome::Ok(pull))
1175 }
1176
1177 /// 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
1180 /// fork; anyone who may push, for a branch of the repository. When it needs a
1181 /// 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?);
1185 check!(writable(&repo));
1186 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 }
1192 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 }
1199 } else {
1200 check!(allowed(Some(&a.actor), &repo, Capability::Push));
1201 }
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
1227 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
1298 /// 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?);
1302 let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
1303 let now = rfc3339(now_ms());
1304 self.db
1305 .prepare(
1306 "UPDATE pulls SET status = 'open', body = COALESCE(?, body), updated_at = ?
1307 WHERE id = ?",
1308 )
1309 .bind(&[
1310 optional(&summary),
1311 now.as_str().into(),
1312 pull.id.as_str().into(),
1313 ])?
1314 .run()
1315 .await?;
1316 if pull.status == PullStatus::Draft {
1317 // 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());
1320 self.publish(
1321 "pull.ready",
1322 &pull.repo_id,
1323 &a.actor,
1324 PullEvent {
1325 commit,
1326 ..Self::pull_event(&pull)
1327 },
1328 )
1329 .await?;
1330 }
1331 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 }
1340 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 )
1360 .await?;
1361 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 }
1383 pull.status = PullStatus::Closed;
1384 pull.updated_at = now;
1385 Ok(Outcome::Ok(pull))
1386 }
1387
1388 /// 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>> {
1393 let viewer = Some(a.actor.clone());
1394 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
1395 check!(writable(&repo));
1396 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 }
1410 }
1411 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 {
1415 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 };
1423 // 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);
1426 if let Some(reason) = waiting {
1427 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 };
1432 return Ok(Outcome::fail(
1433 FailureCode::Conflict,
1434 format!("{reason} {remedy}"),
1435 ));
1436 }
1437 }
1438 if let Some(missing) = self.approvals_gap(&settings, &pull).await? {
1439 return Ok(Outcome::fail(FailureCode::Conflict, missing));
1440 }
1441 // 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 }
1457
1458 // 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 }
1464 check!(allowed(Some(&a.actor), &repo, Capability::Merge));
1465 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 }
1484 check!(allowed(Some(&a.actor), &repo, Capability::Merge));
1485 self.request_landing(&pull, &a.actor, a.keep_issue_open)
1486 .await?;
1487 return Ok(Outcome::Ok(pull));
1488 }
1489
1490 // Whether the actor may write to the repository is decided by repos.
1491 let landed: Outcome<Landed> = g1t_kit::call(
1492 &self.repos,
1493 "land",
1494 &LandArgs {
1495 // 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(),
1498 actor: a.actor.clone(),
1499 },
1500 )
1501 .await?;
1502 let landed = check!(landed);
1503 Ok(Outcome::Ok(
1504 self.record_merge(&repo, pull, &a.actor, a.keep_issue_open, landed)
1505 .await?,
1506 ))
1507 }
1508
1509 /// 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 };
1526 let now = rfc3339(now_ms());
1527 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),
1538 actor.username.as_str().into(),
1539 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(
1546 self.db
1547 .prepare(
1548 "UPDATE issues
1549 SET state = 'closed', reason = 'completed', resolved_by = ?,
1550 closed_at = ?, updated_at = ?
1551 WHERE id = ?",
1552 )
1553 .bind(&[
1554 pull.number.into(),
1555 now.as_str().into(),
1556 now.as_str().into(),
1557 issue.id.as_str().into(),
1558 ])?,
1559 );
1560 statements.push(
1561 self.db
1562 .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,
1579 actor,
1580 PullEvent {
1581 commit: Some(landed.commit.clone()),
1582 ..Self::pull_event(&pull)
1583 },
1584 )
1585 .await?;
1586 if let Some(issue) = &issue {
1587 self.publish(
1588 "issue.closed",
1589 &repo.id,
1590 actor,
1591 IssueEvent {
1592 reason: Some(IssueReason::Completed.as_str()),
1593 resolved_by: Some(pull.number),
1594 ..Self::issue_event(issue)
1595 },
1596 )
1597 .await?;
1598 }
1599
1600 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 }
1611 pull.status = PullStatus::Merged;
1612 pull.head_commit = Some(landed.commit.clone());
1613 pull.merge_base = landed.previous;
1614 pull.merged_by = Some(actor.username.clone());
1615 pull.merged_at = Some(now.clone());
1616 pull.updated_at = now;
1617 Ok(pull)
1618 }
1619
1620 async fn list_active_pulls(&self, a: ViewerArgs) -> Result<Vec<ActivePull>> {
1621 let Some(viewer) = a.viewer else {
1622 return Ok(Vec::new());
1623 };
1624 let found = self
1625 .db
1626 .prepare(
1627 "SELECT * FROM pulls
1628 WHERE author_id = ? AND status IN ('draft', 'open')
1629 ORDER BY updated_at DESC LIMIT 50",
1630 )
1631 .bind(&[viewer.id.into()])?
1632 .all()
1633 .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 {
1639 let issue = match pull.issue {
1640 Some(number) => self.issue(&pull.repo_id, number).await?,
1641 None => None,
1642 };
1643 // 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
1665 }
1666
1667 // --- Sessions ----------------------------------------------------------
1668
1669 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 }
1679 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 }
1687
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(
1699 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
1700 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
1701 FROM session_entries WHERE pull_id = ?",
1702 )
1703 .bind(&[
1704 pull.id.as_str().into(),
1705 kind.as_str().unwrap_or("note").into(),
1706 text.into(),
1707 optional(&entry.tool),
1708 optional(&entry.commit.or_else(|| pull.head_commit.clone())),
1709 now.as_str().into(),
1710 pull.id.as_str().into(),
1711 ])?,
1712 );
1713 }
1714 statements.push(
1715 self.db
1716 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
1717 .bind(&[now.as_str().into(), pull.id.as_str().into()])?,
1718 );
1719 self.db.batch(statements).await?;
1720 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,
1728 count,
1729 },
1730 )
1731 .await?;
1732 Ok(Outcome::Ok(Appended { count }))
1733 }
1734
1735 /// 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
1765 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?);
1767 let rows = self
1768 .db
1769 .prepare(
1770 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
1771 WHERE pull_id = ? AND seq > ? ORDER BY seq LIMIT ?",
1772 )
1773 .bind(&[pull.id.into(), a.after_seq.into(), SESSION_PAGE.into()])?
1774 .all()
1775 .await?
1776 .results::<SessionRow>()?;
1777 Ok(Outcome::Ok(
1778 rows.into_iter().map(SessionEntry::from).collect(),
1779 ))
1780 }
1781
1782 /// 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<()> {
1785 if event.kind != "git.push" {
1786 return Ok(());
1787 }
1788 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 {
1793 return Ok(());
1794 };
1795 let now = rfc3339(now_ms());
1796 // The head moved, so whatever the checks said no longer applies, and
1797 // whatever step g1t was waiting on has been taken.
1798 let moved = "UPDATE pulls
1799 SET head_commit = ?, updated_at = ?, check_status = NULL, check_run_id = NULL,
1800 working_on = NULL, working_until = NULL, stalled = NULL";
1801 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();
1804 // A fork carries its pull request on its default branch.
1805 if event.data["defaultBranch"].as_bool() == Some(true) {
1806 pulls.extend(
1807 self.db
1808 .prepare(format!(
1809 "{moved} WHERE fork_repo_id = ? AND {active} {returning}"
1810 ))
1811 .bind(&[
1812 after.into(),
1813 now.as_str().into(),
1814 repo_id.into(),
1815 after.into(),
1816 ])?
1817 .all()
1818 .await?
1819 .results::<MovedRow>()?,
1820 );
1821 }
1822 if let Some(branch) = git_ref.strip_prefix("refs/heads/") {
1823 pulls.extend(
1824 self.db
1825 .prepare(format!(
1826 "{moved} WHERE repo_id = ? AND source_branch = ? AND {active} {returning}"
1827 ))
1828 .bind(&[
1829 after.into(),
1830 now.as_str().into(),
1831 repo_id.into(),
1832 branch.into(),
1833 after.into(),
1834 ])?
1835 .all()
1836 .await?
1837 .results::<MovedRow>()?,
1838 );
1839 }
1840 // 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 }
1850 // 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;
1855 // 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 }
1875 Ok(())
1876 }
1877}
1878
1879fn service(env: &Env) -> Result<Work> {
1880 Ok(Work {
1881 db: env.d1("DB")?,
1882 identity: env.service("IDENTITY")?,
1883 repos: env.service("REPOS")?,
1884 events: env.service("EVENTS")?,
1885 actions: env.service("ACTIONS")?,
1886 })
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() {
1898 "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?),
1907 "start_checks" => reply(&work.start_checks(args(body)?).await?),
1908 "report_checks" => reply(&work.report_checks(args(body)?).await?),
1909 "set_commit_status" => reply(&work.set_commit_status(args(body)?).await?),
1910 "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?),
1918 "message_agent" => reply(&work.message_agent(args(body)?).await?),
1919 "locate_pull" => reply(&work.locate_pull(args(body)?).await?),
1920 "answer_message" => reply(&work.answer_message(args(body)?).await?),
1921 "take_messages" => reply(&work.take_messages(args(body)?).await?),
1922 "wake_for_messages" => reply(&work.wake_for_messages(args(body)?).await?),
1923 "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?),
1927 "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?),
1930 "update_pull" => reply(&work.update_pull(args(body)?).await?),
1931 "catch_up_pull" => reply(&work.catch_up_pull(args(body)?).await?),
1932 "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?),
1936 "by_author" => reply(&work.by_author(args(body)?).await?),
1937 "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?),
1945 "append_session" => reply(&work.append_session(args(body)?).await?),
1946 "read_session" => reply(&work.read_session(args(body)?).await?),
1947 // 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?),
1961 // 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?),
1965 // 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?),
1976 "start_mergecheck" => reply(&work.start_mergecheck(args(body)?).await?),
1977 "report_mergecheck" => reply(&work.report_mergecheck(args(body)?).await?),
1978 // 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?),
1986 _ => Response::error("Unknown method", 404),
1987 }
1988}
1989
1990/// Events from the bus, delivered on this service's own queue.
1991#[event(queue)]
1992async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
1993 let work = service(&env)?;
1994 for message in batch.messages()? {
1995 // A workspace renamed: its agent runs and memory move to the slug it has now.
1996 if g1t_kit::rename::on_event(&env, &env.d1("DB")?, message.body(), &[memory::RENAMED, guardrails::RENAMED].concat()).await? {
1997 message.ack();
1998 continue;
1999 }
2000 // 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 }
2016 capture::on_event(&work, message.body()).await;
2017 work.on_event(message.body()).await?;
2018 message.ack();
2019 }
2020 Ok(())
2021}