flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/work/src/lib.rs

1,982 lines76,329 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 guardrails;
11mod lifecycle;
12mod memory;
13mod mentions;
14mod mergeability;
15mod plans;
16mod messages;
17mod queue;
18mod reviews;
19mod rows;
20mod runs;
21mod settings;
22mod statuses;
23
24use g1t_contracts::events::{
25 CommentCreated, Event, IssueEvent, NewEvent, Publish, PullEvent, SessionAppended,
26};
27use g1t_contracts::identity::UsernameArgs;
28use g1t_contracts::repos::{
29 ForkArgs, GetArgs, HeadArgs, LandArgs, Landed, NeedsAgentReason, PullBranchUpdate, Repo, RepoPath,
30 UpdatePullBranchArgs,
31};
32use g1t_contracts::time::rfc3339;
33use g1t_contracts::work::*;
34use futures_util::future::{try_join, try_join3, try_join_all};
35use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
36use g1t_kit::{args, now_ms, reply, rpc_method};
37use serde::Serialize;
38use worker::wasm_bindgen::JsValue;
39use worker::{
40 Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event,
41};
42
43use rows::{CommentRow, IssueRow, MovedRow, NumberRow, PullRow, SessionRow, Snapshot, ValueRow};
44
45const SOURCE: &str = "work";
46const MAX_ENTRY_BATCH: usize = 200;
47const MAX_ENTRY_CHARS: usize = 64_000;
48const MAX_TITLE_CHARS: usize = 200;
49const SESSION_PAGE: u32 = 500;
50const LIST_PAGE: u32 = 100;
51const MAX_ASSIGNEES: usize = 10;
52const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
53
54const ISSUE_COLUMNS: &str = "issues.*,
55 (SELECT count(*) FROM pulls WHERE pulls.issue_id = issues.id) AS pull_count,
56 (SELECT agent FROM pulls
57 WHERE pulls.issue_id = issues.id AND pulls.status IN ('draft', 'open')
58 AND pulls.fork_repo_id IS NOT NULL
59 ORDER BY pulls.number DESC LIMIT 1) AS agent,
60 (SELECT count(*) FROM comments
61 WHERE comments.repo_id = issues.repo_id AND comments.number = issues.number
62 AND comments.kind = 'comment') AS comment_count";
63
64fn no_issue<T>() -> Outcome<T> {
65 Outcome::fail(FailureCode::NotFound, "Issue not found.")
66}
67
68fn no_pull<T>() -> Outcome<T> {
69 Outcome::fail(FailureCode::NotFound, "Pull request not found.")
70}
71
72fn optional(value: &Option<String>) -> JsValue {
73 value.as_deref().map_or(JsValue::NULL, JsValue::from)
74}
75
76fn optional_number(value: Option<u32>) -> JsValue {
77 value.map_or(JsValue::NULL, JsValue::from)
78}
79
80/// The lowercase name a `State` is stored and sent as.
81fn state_name(state: Option<State>) -> Option<&'static str> {
82 state.map(|state| match state {
83 State::Open => "open",
84 State::Closed => "closed",
85 })
86}
87
88/// A trimmed title, or why it cannot be used.
89fn valid_title(title: &str) -> std::result::Result<&str, &'static str> {
90 let title = title.trim();
91 if title.is_empty() {
92 Err("A title is required.")
93 } else if title.chars().count() > MAX_TITLE_CHARS {
94 Err("That title is too long.")
95 } else {
96 Ok(title)
97 }
98}
99
100/// Unwraps an `Outcome`, returning its failure from the enclosing method.
101macro_rules! check {
102 ($outcome:expr) => {
103 match $outcome {
104 Outcome::Ok(value) => value,
105 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
106 }
107 };
108}
109
110struct Work {
111 db: D1Database,
112 identity: Fetcher,
113 repos: Fetcher,
114 events: Fetcher,
115 /// GitHub Actions: runs a merge queue's `merge_group` workflows.
116 actions: Fetcher,
117}
118
119impl Work {
120 async fn publish<T: Serialize>(
121 &self,
122 kind: &'static str,
123 repo_id: &str,
124 actor: &User,
125 data: T,
126 ) -> Result<()> {
127 self.publish_as(kind, repo_id, Some(actor.id.clone()), data)
128 .await
129 }
130
131 /// A pull request's author as a viewer who can read its repository and
132 /// source. Stored authors carry no memberships, so a private repository
133 /// would otherwise look missing to them.
134 pub(crate) async fn author_viewer(&self, pull: &Pull) -> Result<Viewer> {
135 let path: Option<RepoPath> = g1t_kit::call(
136 &self.repos,
137 "path_by_id",
138 &g1t_contracts::repos::PathByIdArgs { id: pull.repo_id.clone() },
139 )
140 .await?;
141 let mut author = pull.author.clone();
142 if let Some(path) = path
143 && !author.is_member(&path.namespace.to_lowercase())
144 {
145 author.workspaces.push(g1t_contracts::Membership::member(path.namespace.to_lowercase()));
146 }
147 Ok(Some(author))
148 }
149
150 /// Publishes an event caused by `actor`, or by g1t itself.
151 async fn publish_as<T: Serialize>(
152 &self,
153 kind: &'static str,
154 repo_id: &str,
155 actor: Option<String>,
156 data: T,
157 ) -> Result<()> {
158 let event = NewEvent {
159 kind,
160 source: SOURCE,
161 repo_id: Some(repo_id.to_owned()),
162 actor,
163 data,
164 };
165 g1t_kit::call(
166 &self.events,
167 "publish",
168 &Publish {
169 events: vec![event],
170 },
171 )
172 .await
173 }
174
175 /// The repository, if the viewer may see it. Whether they may is
176 /// decided by the repos service.
177 async fn repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Outcome<Repo>> {
178 g1t_kit::call(
179 &self.repos,
180 "get",
181 &GetArgs {
182 path: path.clone(),
183 viewer: viewer.clone(),
184 },
185 )
186 .await
187 }
188
189 /// The next number in the repository's sequence. Taking it is one
190 /// statement, so concurrent opens cannot be given the same number.
191 async fn next_number(&self, repo_id: &str) -> Result<u32> {
192 let row = self
193 .db
194 .prepare(
195 "INSERT INTO counters (repo_id, last) VALUES (?, 1)
196 ON CONFLICT (repo_id) DO UPDATE SET last = last + 1
197 RETURNING last AS n",
198 )
199 .bind(&[repo_id.into()])?
200 .first::<NumberRow>(None)
201 .await?;
202 row.map(|row| row.n)
203 .ok_or_else(|| worker::Error::RustError("no number was assigned".into()))
204 }
205
206 async fn issue(&self, repo_id: &str, number: u32) -> Result<Option<Issue>> {
207 Ok(self
208 .db
209 .prepare(format!(
210 "SELECT {ISSUE_COLUMNS} FROM issues WHERE repo_id = ? AND number = ?"
211 ))
212 .bind(&[repo_id.into(), number.into()])?
213 .first::<IssueRow>(None)
214 .await?
215 .map(Issue::from))
216 }
217
218 async fn pull(&self, repo_id: &str, number: u32) -> Result<Option<Pull>> {
219 Ok(self
220 .db
221 .prepare("SELECT * FROM pulls WHERE repo_id = ? AND number = ?")
222 .bind(&[repo_id.into(), number.into()])?
223 .first::<PullRow>(None)
224 .await?
225 .map(Pull::from))
226 }
227
228 async fn comments(&self, repo_id: &str, number: u32) -> Result<Vec<Comment>> {
229 let rows = self
230 .db
231 .prepare(
232 "SELECT * FROM comments WHERE repo_id = ? AND number = ? ORDER BY id LIMIT 500",
233 )
234 .bind(&[repo_id.into(), number.into()])?
235 .all()
236 .await?
237 .results::<CommentRow>()?;
238 Ok(rows.into_iter().map(Comment::from).collect())
239 }
240
241 /// The repository and one of its issues, as seen by `viewer`.
242 async fn issue_at(
243 &self,
244 path: &RepoPath,
245 number: u32,
246 viewer: &Viewer,
247 ) -> Result<Outcome<(Repo, Issue)>> {
248 let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
249 return Ok(no_issue());
250 };
251 Ok(match self.issue(&repo.id, number).await? {
252 Some(issue) => Outcome::Ok((repo, issue)),
253 None => no_issue(),
254 })
255 }
256
257 /// The repository and one of its pull requests, as seen by `viewer`.
258 async fn pull_at(
259 &self,
260 path: &RepoPath,
261 number: u32,
262 viewer: &Viewer,
263 ) -> Result<Outcome<(Repo, Pull)>> {
264 let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
265 return Ok(no_pull());
266 };
267 Ok(match self.pull(&repo.id, number).await? {
268 Some(pull) => Outcome::Ok((repo, pull)),
269 None => no_pull(),
270 })
271 }
272
273 /// Records something that happened to an issue or a pull request, so
274 /// that it shows in the conversation where it happened. `text` is what
275 /// `author` did, as the rest of a sentence starting with their name.
276 pub(crate) async fn note(
277 &self,
278 repo_id: &str,
279 number: u32,
280 author: (&str, &str),
281 text: &str,
282 ) -> Result<()> {
283 let now = now_ms();
284 self.db
285 .prepare(
286 "INSERT INTO comments
287 (id, repo_id, number, author_id, author_name, body, kind, created_at)
288 VALUES (?, ?, ?, ?, ?, ?, 'event', ?)",
289 )
290 .bind(&[
291 new_id("cmt", now).into(),
292 repo_id.into(),
293 number.into(),
294 author.0.into(),
295 author.1.into(),
296 text.into(),
297 rfc3339(now).into(),
298 ])?
299 .run()
300 .await?;
301 Ok(())
302 }
303
304 /// Notes who was added to and removed from a list of people, such as
305 /// "assigned ana" or "requested a review from g1t-agent".
306 async fn note_changes(
307 &self,
308 repo_id: &str,
309 number: u32,
310 actor: &User,
311 before: &[String],
312 after: &[String],
313 (added, removed): (&str, &str),
314 ) -> Result<()> {
315 let joined = |names: Vec<&String>| {
316 names
317 .into_iter()
318 .map(String::as_str)
319 .collect::<Vec<_>>()
320 .join(", ")
321 };
322 let new: Vec<&String> = after.iter().filter(|name| !before.contains(name)).collect();
323 let gone: Vec<&String> = before.iter().filter(|name| !after.contains(name)).collect();
324 let who = (actor.id.as_str(), actor.username.as_str());
325 if !new.is_empty() {
326 // Taking something on oneself reads better said that way.
327 let text = if added == "assigned" && new == [&actor.username] {
328 "self-assigned this".to_owned()
329 } else {
330 format!("{added} {}", joined(new))
331 };
332 self.note(repo_id, number, who, &text).await?;
333 }
334 if !gone.is_empty() {
335 self.note(repo_id, number, who, &format!("{removed} {}", joined(gone)))
336 .await?;
337 }
338 Ok(())
339 }
340
341 fn issue_event(issue: &Issue) -> IssueEvent {
342 IssueEvent {
343 issue_id: issue.id.clone(),
344 repo_id: issue.repo_id.clone(),
345 number: issue.number,
346 ..IssueEvent::default()
347 }
348 }
349
350 /// The commit a pull request's change is at in git right now: its
351 /// fork's default branch, or its branch.
352 async fn live_head(&self, pull: &Pull) -> Result<Option<String>> {
353 g1t_kit::call(
354 &self.repos,
355 "head",
356 &HeadArgs {
357 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| pull.repo_id.clone()),
358 branch: pull.branch.clone().unwrap_or_default(),
359 },
360 )
361 .await
362 }
363
364 fn pull_event(pull: &Pull) -> PullEvent {
365 PullEvent {
366 pull_id: pull.id.clone(),
367 repo_id: pull.repo_id.clone(),
368 number: pull.number,
369 issue: pull.issue,
370 ..PullEvent::default()
371 }
372 }
373
374 // --- Issues ------------------------------------------------------------
375
376 async fn open_issue(&self, a: OpenIssueArgs) -> Result<Outcome<Issue>> {
377 if !a.actor.verified {
378 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
379 }
380 let title = match valid_title(&a.title) {
381 Ok(title) => title,
382 Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
383 };
384 let Some(labels) = normalize_labels(&a.labels) else {
385 return Ok(Outcome::fail(
386 FailureCode::Invalid,
387 "An issue can have up to 10 labels of up to 40 characters each.",
388 ));
389 };
390 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
391 let checks: Vec<&str> = a
392 .checks
393 .iter()
394 .map(|check| check.trim())
395 .filter(|check| !check.is_empty())
396 .collect();
397
398 let now = now_ms();
399 let id = new_id("iss", now);
400 let number = self.next_number(&repo.id).await?;
401 let timestamp = rfc3339(now);
402 self.db
403 .prepare(
404 "INSERT INTO issues
405 (id, repo_id, number, title, body, labels, checks, author_id, author_name,
406 created_at, updated_at)
407 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
408 )
409 .bind(&[
410 id.as_str().into(),
411 repo.id.as_str().into(),
412 number.into(),
413 title.into(),
414 a.body.trim().into(),
415 serde_json::to_string(&labels)?.into(),
416 serde_json::to_string(&checks)?.into(),
417 a.actor.id.as_str().into(),
418 a.actor.username.as_str().into(),
419 timestamp.as_str().into(),
420 timestamp.as_str().into(),
421 ])?
422 .run()
423 .await?;
424 let Some(issue) = self.issue(&repo.id, number).await? else {
425 return Ok(no_issue());
426 };
427 self.apply_label_rule(&a.actor, &issue, &[]).await?;
428 self.publish(
429 "issue.opened",
430 &repo.id,
431 &a.actor,
432 IssueEvent {
433 title: Some(issue.title.clone()),
434 ..Self::issue_event(&issue)
435 },
436 )
437 .await?;
438 Ok(Outcome::Ok(issue))
439 }
440
441 async fn list_issues(&self, a: ListIssuesArgs) -> Result<Outcome<Vec<Issue>>> {
442 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
443 let state = state_name(a.state).map_or(JsValue::NULL, JsValue::from);
444 let label = a
445 .label
446 .map(|label| label.trim().to_lowercase())
447 .filter(|label| !label.is_empty());
448 let rows = self
449 .db
450 .prepare(format!(
451 "SELECT {ISSUE_COLUMNS} FROM issues
452 WHERE repo_id = ? AND (? IS NULL OR state = ?)
453 AND (? IS NULL OR EXISTS
454 (SELECT 1 FROM json_each(issues.labels) WHERE json_each.value = ?))
455 ORDER BY number DESC LIMIT ?"
456 ))
457 .bind(&[
458 repo.id.into(),
459 state.clone(),
460 state,
461 optional(&label),
462 optional(&label),
463 LIST_PAGE.into(),
464 ])?
465 .all()
466 .await?
467 .results::<IssueRow>()?;
468 Ok(Outcome::Ok(rows.into_iter().map(Issue::from).collect()))
469 }
470
471 async fn get_issue(&self, a: ViewArgs) -> Result<Outcome<IssueDetail>> {
472 let (repo, issue) = check!(self.issue_at(&a.repo, a.number, &a.viewer).await?);
473 let pulls = self
474 .db
475 .prepare("SELECT * FROM pulls WHERE issue_id = ? ORDER BY number")
476 .bind(&[issue.id.as_str().into()])?
477 .all()
478 .await?
479 .results::<PullRow>()?;
480 Ok(Outcome::Ok(IssueDetail {
481 comments: self.comments(&repo.id, issue.number).await?,
482 pulls: pulls.into_iter().map(Pull::from).collect(),
483 issue,
484 }))
485 }
486
487 /// The issue, if `actor` wrote it or belongs to the repository's workspace.
488 async fn manageable_issue(
489 &self,
490 actor: &User,
491 path: &RepoPath,
492 number: u32,
493 ) -> Result<Outcome<Issue>> {
494 let (repo, issue) = check!(self.issue_at(path, number, &Some(actor.clone())).await?);
495 if issue.author.id != actor.id && !actor.is_member(&repo.namespace) {
496 return Ok(Outcome::fail(
497 FailureCode::Forbidden,
498 "Only the author or a member of the workspace can change an issue.",
499 ));
500 }
501 Ok(Outcome::Ok(issue))
502 }
503
504 async fn update_issue(&self, a: UpdateIssueArgs) -> Result<Outcome<Issue>> {
505 let issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
506 let title = match a.title.as_deref().map(valid_title) {
507 Some(Err(message)) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
508 Some(Ok(title)) => Some(title.to_owned()),
509 None => None,
510 };
511 let labels = match a.labels.as_deref().map(normalize_labels) {
512 Some(None) => {
513 return Ok(Outcome::fail(
514 FailureCode::Invalid,
515 "An issue can have up to 10 labels of up to 40 characters each.",
516 ));
517 }
518 Some(Some(labels)) => Some(serde_json::to_string(&labels)?),
519 None => None,
520 };
521 let assignees = match a.assignees {
522 Some(names) => Some(check!(self.valid_assignees(names).await?)),
523 None => None,
524 };
525 let assigned = assignees.as_ref().map(serde_json::to_string).transpose()?;
526 let body = a.body.map(|body| body.trim().to_owned());
527 self.db
528 .prepare(
529 "UPDATE issues
530 SET title = COALESCE(?, title), body = COALESCE(?, body),
531 labels = COALESCE(?, labels), assignees = COALESCE(?, assignees),
532 updated_at = ?
533 WHERE id = ?",
534 )
535 .bind(&[
536 optional(&title),
537 optional(&body),
538 optional(&labels),
539 optional(&assigned),
540 rfc3339(now_ms()).into(),
541 issue.id.as_str().into(),
542 ])?
543 .run()
544 .await?;
545 let before = issue.assignees.clone();
546 let labels_before = issue.labels.clone();
547 let Some(issue) = self.issue(&issue.repo_id, issue.number).await? else {
548 return Ok(no_issue());
549 };
550 self.apply_label_rule(&a.actor, &issue, &labels_before).await?;
551 self.publish(
552 "issue.updated",
553 &issue.repo_id,
554 &a.actor,
555 Self::issue_event(&issue),
556 )
557 .await?;
558 if let Some(assignees) = assignees {
559 self.note_changes(
560 &issue.repo_id,
561 issue.number,
562 &a.actor,
563 &before,
564 &assignees,
565 ("assigned", "unassigned"),
566 )
567 .await?;
568 self.publish(
569 "issue.assigned",
570 &issue.repo_id,
571 &a.actor,
572 IssueEvent {
573 assignees: Some(assignees),
574 ..Self::issue_event(&issue)
575 },
576 )
577 .await?;
578 }
579 Ok(Outcome::Ok(issue))
580 }
581
582 /// Usernames as given, tidied, if each names an account.
583 async fn valid_assignees(&self, names: Vec<String>) -> Result<Outcome<Vec<String>>> {
584 let mut assignees: Vec<String> = Vec::new();
585 for name in names {
586 let name = name.trim().trim_start_matches('@').to_lowercase();
587 if name.is_empty() || assignees.contains(&name) {
588 continue;
589 }
590 if assignees.len() == MAX_ASSIGNEES {
591 return Ok(Outcome::fail(
592 FailureCode::Invalid,
593 format!("An issue can be assigned to at most {MAX_ASSIGNEES} people."),
594 ));
595 }
596 let account: Viewer = g1t_kit::call(
597 &self.identity,
598 "user_by_username",
599 &UsernameArgs {
600 username: name.clone(),
601 },
602 )
603 .await?;
604 if account.is_none() {
605 return Ok(Outcome::fail(
606 FailureCode::Invalid,
607 format!("There is no account named {name}."),
608 ));
609 }
610 assignees.push(name);
611 }
612 Ok(Outcome::Ok(assignees))
613 }
614
615 /// Open issues assigned to the viewer, in every repository. Callers
616 /// show only those in repositories the viewer can still see.
617 async fn list_assigned_issues(&self, a: ViewerArgs) -> Result<Vec<Issue>> {
618 let Some(viewer) = a.viewer else {
619 return Ok(Vec::new());
620 };
621 let rows = self
622 .db
623 .prepare(format!(
624 "SELECT {ISSUE_COLUMNS} FROM issues
625 WHERE state = 'open' AND EXISTS (
626 SELECT 1 FROM json_each(issues.assignees) WHERE json_each.value = ?)
627 ORDER BY updated_at DESC LIMIT 50"
628 ))
629 .bind(&[viewer.username.into()])?
630 .all()
631 .await?
632 .results::<IssueRow>()?;
633 Ok(rows.into_iter().map(Issue::from).collect())
634 }
635
636 async fn close_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::Closed {
639 return Ok(Outcome::fail(
640 FailureCode::Conflict,
641 "This issue is already closed.",
642 ));
643 }
644 let reason = a.reason.unwrap_or(IssueReason::Completed);
645 let now = rfc3339(now_ms());
646 self.db
647 .prepare(
648 "UPDATE issues SET state = 'closed', reason = ?, closed_at = ?, updated_at = ?
649 WHERE id = ?",
650 )
651 .bind(&[
652 reason.as_str().into(),
653 now.as_str().into(),
654 now.as_str().into(),
655 issue.id.as_str().into(),
656 ])?
657 .run()
658 .await?;
659 self.publish(
660 "issue.closed",
661 &issue.repo_id,
662 &a.actor,
663 IssueEvent {
664 reason: Some(reason.as_str()),
665 ..Self::issue_event(&issue)
666 },
667 )
668 .await?;
669 self.note(
670 &issue.repo_id,
671 issue.number,
672 (&a.actor.id, &a.actor.username),
673 match reason {
674 IssueReason::Completed => "closed this as completed",
675 IssueReason::NotPlanned => "closed this as not planned",
676 },
677 )
678 .await?;
679 issue.state = State::Closed;
680 issue.reason = Some(reason);
681 issue.closed_at = Some(now.clone());
682 issue.updated_at = now;
683 Ok(Outcome::Ok(issue))
684 }
685
686 async fn reopen_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
687 let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
688 if issue.state == State::Open {
689 return Ok(Outcome::fail(
690 FailureCode::Conflict,
691 "This issue is already open.",
692 ));
693 }
694 let now = rfc3339(now_ms());
695 self.db
696 .prepare(
697 "UPDATE issues
698 SET state = 'open', reason = NULL, resolved_by = NULL, closed_at = NULL,
699 updated_at = ?
700 WHERE id = ?",
701 )
702 .bind(&[now.as_str().into(), issue.id.as_str().into()])?
703 .run()
704 .await?;
705 self.publish(
706 "issue.reopened",
707 &issue.repo_id,
708 &a.actor,
709 Self::issue_event(&issue),
710 )
711 .await?;
712 self.note(
713 &issue.repo_id,
714 issue.number,
715 (&a.actor.id, &a.actor.username),
716 "reopened this",
717 )
718 .await?;
719 issue.state = State::Open;
720 issue.reason = None;
721 issue.resolved_by = None;
722 issue.closed_at = None;
723 issue.updated_at = now;
724 Ok(Outcome::Ok(issue))
725 }
726
727 /// The default labels, then every other label in use on the repository.
728 async fn list_labels(&self, a: ViewArgs) -> Result<Outcome<Vec<String>>> {
729 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
730 let used = self
731 .db
732 .prepare(
733 "SELECT DISTINCT json_each.value AS value
734 FROM issues, json_each(issues.labels)
735 WHERE issues.repo_id = ? ORDER BY 1 LIMIT 200",
736 )
737 .bind(&[repo.id.into()])?
738 .all()
739 .await?
740 .results::<ValueRow>()?;
741 let mut labels: Vec<String> = DEFAULT_LABELS.iter().map(|label| (*label).into()).collect();
742 for row in used {
743 if !labels.contains(&row.value) {
744 labels.push(row.value);
745 }
746 }
747 Ok(Outcome::Ok(labels))
748 }
749
750 async fn counts(&self, a: ViewArgs) -> Result<Outcome<Counts>> {
751 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
752 let counts = self
753 .db
754 .prepare(
755 "SELECT
756 (SELECT count(*) FROM issues WHERE repo_id = ? AND state = 'open') AS issues,
757 (SELECT count(*) FROM pulls
758 WHERE repo_id = ? AND status IN ('draft', 'open')) AS pulls",
759 )
760 .bind(&[repo.id.as_str().into(), repo.id.as_str().into()])?
761 .first::<Counts>(None)
762 .await?;
763 Ok(Outcome::Ok(counts.unwrap_or(Counts {
764 issues: 0,
765 pulls: 0,
766 })))
767 }
768
769 // --- Comments ----------------------------------------------------------
770
771 async fn add_comment(&self, a: AddCommentArgs) -> Result<Outcome<Comment>> {
772 if !a.actor.verified {
773 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
774 }
775 let body = a.body.trim();
776 // An approval speaks for itself; anything else has to say something.
777 if body.is_empty() && a.verdict != Some(Verdict::Approve) {
778 return Ok(Outcome::fail(
779 FailureCode::Invalid,
780 "A comment cannot be empty.",
781 ));
782 }
783 let path = a
784 .path
785 .as_deref()
786 .map(str::trim)
787 .filter(|path| !path.is_empty());
788 let line = a.line.filter(|line| *line > 0 && path.is_some());
789 if body.chars().count() > MAX_ENTRY_CHARS {
790 return Ok(Outcome::fail(
791 FailureCode::Invalid,
792 "That comment is too long.",
793 ));
794 }
795 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
796 // The number names an issue or a pull request, never both.
797 let mut pull_id = None;
798 let table = if self.issue(&repo.id, a.number).await?.is_some() {
799 if path.is_some() || a.verdict.is_some() {
800 return Ok(Outcome::fail(
801 FailureCode::Invalid,
802 "Only a pull request can be reviewed or commented on by line.",
803 ));
804 }
805 "issues"
806 } else if let Some(pull) = self.pull(&repo.id, a.number).await? {
807 if a.verdict.is_some() && pull.author.id == a.actor.id {
808 return Ok(Outcome::fail(
809 FailureCode::Forbidden,
810 "You cannot approve or request changes on your own pull request.",
811 ));
812 }
813 pull_id = Some(pull.id.clone());
814 "pulls"
815 } else {
816 return Ok(Outcome::fail(
817 FailureCode::NotFound,
818 "No issue or pull request has that number.",
819 ));
820 };
821
822 let now = now_ms();
823 let comment = Comment {
824 kind: CommentKind::Comment,
825 id: new_id("cmt", now),
826 author: a.actor.clone(),
827 body: body.to_owned(),
828 path: path.map(str::to_owned),
829 line,
830 verdict: a.verdict,
831 created_at: rfc3339(now),
832 };
833 self.db
834 .batch(vec![
835 self.db
836 .prepare(
837 "INSERT INTO comments
838 (id, repo_id, number, author_id, author_name, body, path, line,
839 verdict, created_at)
840 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
841 )
842 .bind(&[
843 comment.id.as_str().into(),
844 repo.id.as_str().into(),
845 a.number.into(),
846 a.actor.id.as_str().into(),
847 a.actor.username.as_str().into(),
848 body.into(),
849 optional(&comment.path),
850 optional_number(line),
851 a.verdict
852 .map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
853 comment.created_at.as_str().into(),
854 ])?,
855 self.db
856 .prepare(format!(
857 "UPDATE {table} SET updated_at = ? WHERE repo_id = ? AND number = ?"
858 ))
859 .bind(&[
860 comment.created_at.as_str().into(),
861 repo.id.as_str().into(),
862 a.number.into(),
863 ])?,
864 ])
865 .await?;
866 self.note_mention(&a.actor, &repo, a.number, &comment, pull_id.as_deref()).await?;
867 self.publish(
868 "comment.created",
869 &repo.id,
870 &a.actor,
871 CommentCreated {
872 comment_id: comment.id.clone(),
873 repo_id: repo.id.clone(),
874 number: a.number,
875 pull_id,
876 verdict: a.verdict,
877 },
878 )
879 .await?;
880 Ok(Outcome::Ok(comment))
881 }
882
883 // --- Pull requests -----------------------------------------------------
884
885 async fn open_pull(&self, a: OpenPullArgs) -> Result<Outcome<Pull>> {
886 if !a.actor.verified {
887 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
888 }
889 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
890 let issue = match a.issue {
891 Some(number) => match self.issue(&repo.id, number).await? {
892 Some(issue) if issue.state == State::Open => Some(issue),
893 Some(_) => {
894 return Ok(Outcome::fail(
895 FailureCode::Conflict,
896 "This issue is closed.",
897 ));
898 }
899 None => return Ok(no_issue()),
900 },
901 None => None,
902 };
903 // A pull request for an issue takes the issue's title unless given one.
904 let title = match (a.title.trim(), &issue) {
905 ("", Some(issue)) => issue.title.clone(),
906 (title, _) => match valid_title(title) {
907 Ok(title) => title.to_owned(),
908 Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
909 },
910 };
911 let agent = match a.agent.trim() {
912 "" => "agent",
913 agent => agent,
914 };
915 let runtime = match a.runtime {
916 Runtime::Hosted => "hosted",
917 Runtime::External => "external",
918 };
919
920 let now = now_ms();
921 let id = new_id("pr", now);
922 let branch = a
923 .branch
924 .as_deref()
925 .map(str::trim)
926 .filter(|branch| !branch.is_empty());
927 // The change is on a branch already pushed to the repository, or
928 // will be made in a fork created for this pull request.
929 let (fork, head) = match branch {
930 Some(branch) => {
931 if branch == repo.default_branch {
932 return Ok(Outcome::fail(
933 FailureCode::Invalid,
934 format!("Choose a branch other than {branch}."),
935 ));
936 }
937 let head: Option<String> = g1t_kit::call(
938 &self.repos,
939 "head",
940 &HeadArgs {
941 repo_id: repo.id.clone(),
942 branch: branch.to_owned(),
943 },
944 )
945 .await?;
946 let Some(head) = head else {
947 return Ok(Outcome::fail(
948 FailureCode::NotFound,
949 format!("There is no branch named {branch}. Push it first."),
950 ));
951 };
952 let existing = self
953 .db
954 .prepare(
955 "SELECT number AS n FROM pulls
956 WHERE repo_id = ? AND source_branch = ? AND status IN ('draft', 'open')",
957 )
958 .bind(&[repo.id.as_str().into(), branch.into()])?
959 .first::<NumberRow>(None)
960 .await?;
961 if let Some(existing) = existing {
962 return Ok(Outcome::fail(
963 FailureCode::Conflict,
964 format!("Pull request #{} is already open for {branch}.", existing.n),
965 ));
966 }
967 (None, Some(head))
968 }
969 None => {
970 let fork: Outcome<Repo> = g1t_kit::call(
971 &self.repos,
972 "fork_for_pull",
973 &ForkArgs {
974 source_id: repo.id.clone(),
975 pull_id: id.clone(),
976 actor: a.actor.clone(),
977 },
978 )
979 .await?;
980 (Some(check!(fork)), None)
981 }
982 };
983 // A branch already holds the work, so its pull request is ready for
984 // review from the start; one with a fork starts as a draft.
985 let status = if branch.is_some() { "open" } else { "draft" };
986 let body = Some(a.body.trim().to_owned()).filter(|body| !body.is_empty());
987
988 let number = self.next_number(&repo.id).await?;
989 let timestamp = rfc3339(now);
990 self.db
991 .prepare(
992 "INSERT INTO pulls
993 (id, repo_id, number, issue_id, issue_number, title, body, agent, runtime,
994 status, fork_repo_id, fork_namespace, fork_name, source_branch, head_commit,
995 author_id, author_name, created_at, updated_at)
996 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
997 )
998 .bind(&[
999 id.as_str().into(),
1000 repo.id.as_str().into(),
1001 number.into(),
1002 optional(&issue.as_ref().map(|issue| issue.id.clone())),
1003 optional_number(issue.as_ref().map(|issue| issue.number)),
1004 title.into(),
1005 optional(&body),
1006 agent.into(),
1007 runtime.into(),
1008 status.into(),
1009 optional(&fork.as_ref().map(|fork| fork.id.clone())),
1010 optional(&fork.as_ref().map(|fork| fork.namespace.clone())),
1011 optional(&fork.as_ref().map(|fork| fork.name.clone())),
1012 optional(&branch.map(str::to_owned)),
1013 optional(&head),
1014 a.actor.id.as_str().into(),
1015 a.actor.username.as_str().into(),
1016 timestamp.as_str().into(),
1017 timestamp.as_str().into(),
1018 ])?
1019 .run()
1020 .await?;
1021 let Some(pull) = self.pull(&repo.id, number).await? else {
1022 return Ok(no_pull());
1023 };
1024 self.manage(&pull).await?;
1025 // Someone is on it now, so it is no longer waiting for an agent.
1026 if let Some(issue) = pull.issue {
1027 self.db
1028 .prepare("UPDATE issues SET queued_by = NULL WHERE repo_id = ? AND number = ?")
1029 .bind(&[repo.id.as_str().into(), issue.into()])?
1030 .run()
1031 .await?;
1032 }
1033 if let Some(issue) = pull.issue {
1034 let text = if lifecycle::made_by_g1t(&pull) {
1035 format!("assigned this to g1t-agent, which opened #{}", pull.number)
1036 } else {
1037 format!("opened #{} for this", pull.number)
1038 };
1039 self.note(&repo.id, issue, (&a.actor.id, &a.actor.username), &text)
1040 .await?;
1041 }
1042 self.publish(
1043 "pull.opened",
1044 &repo.id,
1045 &a.actor,
1046 PullEvent {
1047 agent: Some(pull.agent.clone()),
1048 ..Self::pull_event(&pull)
1049 },
1050 )
1051 .await?;
1052 Ok(Outcome::Ok(pull))
1053 }
1054
1055 async fn list_pulls(&self, a: ListPullsArgs) -> Result<Outcome<Vec<Pull>>> {
1056 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
1057 let filter = match a.state {
1058 Some(State::Open) => "AND status IN ('draft', 'open')",
1059 Some(State::Closed) => "AND status IN ('merged', 'closed')",
1060 None => "",
1061 };
1062 let rows = self
1063 .db
1064 .prepare(format!(
1065 "SELECT * FROM pulls WHERE repo_id = ? {filter} ORDER BY number DESC LIMIT ?"
1066 ))
1067 .bind(&[repo.id.into(), LIST_PAGE.into()])?
1068 .all()
1069 .await?
1070 .results::<PullRow>()?;
1071 Ok(Outcome::Ok(rows.into_iter().map(Pull::from).collect()))
1072 }
1073
1074 async fn get_pull(&self, a: ViewArgs) -> Result<Outcome<PullDetail>> {
1075 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
1076 let issue = match pull.issue {
1077 Some(number) => self.issue(&repo.id, number).await?,
1078 None => None,
1079 };
1080 let mut pull = pull;
1081 // Worked out on each push; this covers a pull request from before
1082 // that was recorded.
1083 if pull.files.is_empty() && pull.head_commit.is_some() {
1084 pull.files = self.refresh_files(&pull).await?;
1085 }
1086 // Everything else at once: none of it depends on the rest, and each
1087 // is a round trip of its own.
1088 let standing = async {
1089 // Mergeability first: where g1t sees a pull request through, a
1090 // conflict decides its next step.
1091 let (merge, behind) =
1092 try_join(self.mergeability(&pull), self.is_behind(&repo.id, &pull)).await?;
1093 let lifecycle = self
1094 .assess(&pull, &issue, behind)
1095 .await?
1096 .map(|(lifecycle, _)| lifecycle);
1097 Ok::<_, worker::Error>((behind, lifecycle, merge))
1098 };
1099 let (((behind, lifecycle, (mergeable, conflicts)), (landing, stalled), comments), (checks, overlaps, review_pending)) =
1100 try_join(
1101 try_join3(standing, self.landing_state(&pull.id), self.comments(&repo.id, pull.number)),
1102 try_join3(
1103 self.latest_checks(&pull.id),
1104 self.overlaps(&pull),
1105 self.review_pending(&pull.id),
1106 ),
1107 )
1108 .await?;
1109 Ok(Outcome::Ok(PullDetail {
1110 comments,
1111 checks,
1112 overlaps,
1113 behind,
1114 review_pending,
1115 lifecycle,
1116 landing,
1117 stalled,
1118 messages: self.messages(&pull.id).await?,
1119 statuses: self.statuses(&repo.id, pull.head_commit.as_deref()).await?,
1120 mergeable,
1121 conflicts,
1122 earlier_checks: self.earlier_checks(&pull.id).await?,
1123 issue,
1124 pull,
1125 }))
1126 }
1127
1128 /// The pull request, if it is still active and `actor` opened it or
1129 /// belongs to the repository's workspace.
1130 async fn manageable_pull(
1131 &self,
1132 actor: &User,
1133 path: &RepoPath,
1134 number: u32,
1135 ) -> Result<Outcome<Pull>> {
1136 let (repo, pull) = check!(self.pull_at(path, number, &Some(actor.clone())).await?);
1137 if pull.author.id != actor.id && !actor.is_member(&repo.namespace) {
1138 return Ok(Outcome::fail(
1139 FailureCode::Forbidden,
1140 "Only whoever opened a pull request, or a member of the workspace, can change it.",
1141 ));
1142 }
1143 if !pull.status.is_active() {
1144 return Ok(Outcome::fail(
1145 FailureCode::Conflict,
1146 format!("This pull request is already {}.", pull.status.as_str()),
1147 ));
1148 }
1149 Ok(Outcome::Ok(pull))
1150 }
1151
1152 /// Brings a pull request up to date with the default branch without a
1153 /// sandbox, where the repos service can do that safely. Whoever could
1154 /// have pushed the merge themselves may ask: whoever opened it, for a
1155 /// fork; any member, for a branch of the repository. When it needs a
1156 /// real merge, says so, naming the conflicting files if a probe found
1157 /// them, and pushes nothing.
1158 async fn catch_up_pull(&self, a: PullActionArgs) -> Result<Outcome<PullBranchUpdate>> {
1159 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &Some(a.actor.clone())).await?);
1160 if !pull.status.is_active() {
1161 return Ok(Outcome::fail(
1162 FailureCode::Conflict,
1163 format!("This pull request is already {}.", pull.status.as_str()),
1164 ));
1165 }
1166 let allowed = if pull.fork_repo_id.is_some() {
1167 pull.author.id == a.actor.id
1168 } else {
1169 a.actor.is_member(&repo.namespace)
1170 };
1171 if !allowed {
1172 return Ok(Outcome::fail(
1173 FailureCode::Forbidden,
1174 if pull.fork_repo_id.is_some() {
1175 "Only whoever opened this pull request can update it."
1176 } else {
1177 "Only members of the workspace can update this pull request."
1178 },
1179 ));
1180 }
1181 let updated: Outcome<PullBranchUpdate> = g1t_kit::call(
1182 &self.repos,
1183 "update_pull_branch",
1184 &UpdatePullBranchArgs {
1185 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
1186 branch: pull.branch.clone(),
1187 number: pull.number,
1188 actor: a.actor,
1189 },
1190 )
1191 .await?;
1192 // A probe that found conflicts says more than "both changed it".
1193 if let Outcome::Ok(PullBranchUpdate::NeedsAgent { .. }) = &updated
1194 && let Some(files) = self.conflicting_files(&pull).await?
1195 && !files.is_empty()
1196 {
1197 return Ok(Outcome::Ok(PullBranchUpdate::NeedsAgent {
1198 reason: NeedsAgentReason::Conflicting,
1199 detail: "Merging it conflicts.".to_owned(),
1200 paths: files,
1201 }));
1202 }
1203 Ok(updated)
1204 }
1205
1206 async fn update_pull(&self, a: UpdatePullArgs) -> Result<Outcome<Pull>> {
1207 let pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1208 let assignees = match a.assignees {
1209 Some(names) => Some(check!(self.valid_assignees(names).await?)),
1210 None => None,
1211 };
1212 let reviewers = match a.reviewers {
1213 Some(names) => {
1214 // A g1t agent is not an account; everyone else has to be.
1215 let agent = names
1216 .iter()
1217 .any(|name| name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME));
1218 let people = names
1219 .into_iter()
1220 .filter(|name| !name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME))
1221 .collect();
1222 let mut reviewers = check!(self.valid_assignees(people).await?);
1223 reviewers.retain(|name| *name != pull.author.username);
1224 if agent {
1225 reviewers.insert(0, reviews::AGENT_NAME.to_owned());
1226 }
1227 Some(reviewers)
1228 }
1229 None => None,
1230 };
1231 self.db
1232 .prepare(
1233 "UPDATE pulls
1234 SET assignees = COALESCE(?, assignees), reviewers = COALESCE(?, reviewers),
1235 updated_at = ?
1236 WHERE id = ?",
1237 )
1238 .bind(&[
1239 optional(&assignees.as_ref().map(serde_json::to_string).transpose()?),
1240 optional(&reviewers.as_ref().map(serde_json::to_string).transpose()?),
1241 rfc3339(now_ms()).into(),
1242 pull.id.as_str().into(),
1243 ])?
1244 .run()
1245 .await?;
1246 if let Some(assignees) = &assignees {
1247 self.note_changes(
1248 &pull.repo_id,
1249 pull.number,
1250 &a.actor,
1251 &pull.assignees,
1252 assignees,
1253 ("assigned", "unassigned"),
1254 )
1255 .await?;
1256 }
1257 if let Some(reviewers) = &reviewers {
1258 self.note_changes(
1259 &pull.repo_id,
1260 pull.number,
1261 &a.actor,
1262 &pull.reviewers,
1263 reviewers,
1264 (
1265 "requested a review from",
1266 "withdrew the request for a review from",
1267 ),
1268 )
1269 .await?;
1270 }
1271 Ok(match self.pull(&pull.repo_id, pull.number).await? {
1272 Some(pull) => Outcome::Ok(pull),
1273 None => no_pull(),
1274 })
1275 }
1276
1277 /// Marks a draft ready for review, or updates the description of one
1278 /// that already is.
1279 async fn ready_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1280 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1281 let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
1282 let now = rfc3339(now_ms());
1283 self.db
1284 .prepare(
1285 "UPDATE pulls SET status = 'open', body = COALESCE(?, body), updated_at = ?
1286 WHERE id = ?",
1287 )
1288 .bind(&[
1289 optional(&summary),
1290 now.as_str().into(),
1291 pull.id.as_str().into(),
1292 ])?
1293 .run()
1294 .await?;
1295 if pull.status == PullStatus::Draft {
1296 // The head as it is now: the push that came just before may not
1297 // have reached `head_commit` yet, and workflows run on it.
1298 let commit = self.live_head(&pull).await?.or_else(|| pull.head_commit.clone());
1299 self.publish(
1300 "pull.ready",
1301 &pull.repo_id,
1302 &a.actor,
1303 PullEvent {
1304 commit,
1305 ..Self::pull_event(&pull)
1306 },
1307 )
1308 .await?;
1309 }
1310 if pull.status == PullStatus::Draft {
1311 self.note(
1312 &pull.repo_id,
1313 pull.number,
1314 (&a.actor.id, &a.actor.username),
1315 "marked this ready for review",
1316 )
1317 .await?;
1318 }
1319 pull.status = PullStatus::Open;
1320 pull.body = summary.or(pull.body);
1321 pull.updated_at = now;
1322 Ok(Outcome::Ok(pull))
1323 }
1324
1325 async fn close_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1326 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1327 let now = rfc3339(now_ms());
1328 self.db
1329 .prepare("UPDATE pulls SET status = 'closed', updated_at = ? WHERE id = ?")
1330 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
1331 .run()
1332 .await?;
1333 self.publish(
1334 "pull.closed",
1335 &pull.repo_id,
1336 &a.actor,
1337 Self::pull_event(&pull),
1338 )
1339 .await?;
1340 self.note(
1341 &pull.repo_id,
1342 pull.number,
1343 (&a.actor.id, &a.actor.username),
1344 "closed this",
1345 )
1346 .await?;
1347 // A closed pull request leaves the merge queue.
1348 if self
1349 .leave(&pull.repo_id, &pull, QueueState::Removed, Some("It was closed."))
1350 .await?
1351 {
1352 self.publish_as(
1353 "queue.changed",
1354 &pull.repo_id,
1355 None,
1356 g1t_contracts::events::QueueChanged {
1357 repo_id: pull.repo_id.clone(),
1358 },
1359 )
1360 .await?;
1361 }
1362 pull.status = PullStatus::Closed;
1363 pull.updated_at = now;
1364 Ok(Outcome::Ok(pull))
1365 }
1366
1367 /// Lands the pull request on the repository's default branch. Unless
1368 /// told to keep it open, that resolves the issue it was for: the issue
1369 /// closes naming this pull request, and the others still in progress
1370 /// for it close as superseded.
1371 async fn merge_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1372 let viewer = Some(a.actor.clone());
1373 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
1374 match pull.status {
1375 PullStatus::Open => {}
1376 PullStatus::Draft => {
1377 return Ok(Outcome::fail(
1378 FailureCode::Conflict,
1379 "This pull request is still a draft. Mark it ready for review first.",
1380 ));
1381 }
1382 status => {
1383 return Ok(Outcome::fail(
1384 FailureCode::Conflict,
1385 format!("This pull request is already {}.", status.as_str()),
1386 ));
1387 }
1388 }
1389 let settings = self.settings(&repo.id).await?;
1390 // Where the repository does not allow it, asking to ignore the
1391 // checks changes nothing.
1392 if !a.ignore_checks || !settings.allow_ignoring_checks {
1393 let waiting = match pull.check_status {
1394 Some(CheckStatus::Queued | CheckStatus::Running) => {
1395 Some("The acceptance checks are still running.")
1396 }
1397 Some(CheckStatus::Failed) => Some("The acceptance checks did not pass."),
1398 Some(CheckStatus::Errored) => Some("The acceptance checks could not be run."),
1399 Some(CheckStatus::Passed) | None => None,
1400 };
1401 // Workflows run on its head count as checks too.
1402 let workflows = statuses::WorkflowFacts::of(&self.statuses(&repo.id, pull.head_commit.as_deref()).await?).refusal();
1403 let waiting = waiting.map(str::to_owned).or(workflows);
1404 if let Some(reason) = waiting {
1405 let remedy = if settings.allow_ignoring_checks {
1406 "Wait or fix them, or merge anyway by ignoring the checks."
1407 } else {
1408 "This repository only merges pull requests whose checks pass."
1409 };
1410 return Ok(Outcome::fail(
1411 FailureCode::Conflict,
1412 format!("{reason} {remedy}"),
1413 ));
1414 }
1415 }
1416 if let Some(missing) = self.approvals_gap(&settings, &pull).await? {
1417 return Ok(Outcome::fail(FailureCode::Conflict, missing));
1418 }
1419 // Known ahead of time to conflict: neither a merge nor the queue
1420 // would get through, so say what has to be resolved now.
1421 if let Some(files) = self.conflicting_files(&pull).await? {
1422 let named = if files.is_empty() {
1423 String::new()
1424 } else {
1425 format!(" in {}", files.join(", "))
1426 };
1427 return Ok(Outcome::fail(
1428 FailureCode::Conflict,
1429 format!(
1430 "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.",
1431 repo.default_branch
1432 ),
1433 ));
1434 }
1435
1436 // A repository that merges through a queue: it joins the queue, and
1437 // lands once its state together with everything ahead has passed.
1438 if settings.merge_queue {
1439 if !a.actor.verified {
1440 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1441 }
1442 if !a.actor.is_member(&repo.namespace) {
1443 return Ok(Outcome::fail(
1444 FailureCode::Forbidden,
1445 "Only members of the repository's workspace can merge a pull request.",
1446 ));
1447 }
1448 return self.enqueue(&repo, &pull, &a.actor, a.keep_issue_open).await;
1449 }
1450
1451 // The default branch has moved under it. Unless the repository
1452 // insists on that being dealt with first, bring it up to date and
1453 // land it when that is done.
1454 if self.is_behind(&repo.id, &pull).await? {
1455 if settings.require_up_to_date {
1456 return Ok(Outcome::fail(
1457 FailureCode::Conflict,
1458 format!(
1459 "{} 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.",
1460 repo.default_branch
1461 ),
1462 ));
1463 }
1464 if !a.actor.verified {
1465 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1466 }
1467 if !a.actor.is_member(&repo.namespace) {
1468 return Ok(Outcome::fail(
1469 FailureCode::Forbidden,
1470 "Only members of the repository's workspace can merge a pull request.",
1471 ));
1472 }
1473 self.request_landing(&pull, &a.actor, a.keep_issue_open)
1474 .await?;
1475 return Ok(Outcome::Ok(pull));
1476 }
1477
1478 // Whether the actor may write to the repository is decided by repos.
1479 let landed: Outcome<Landed> = g1t_kit::call(
1480 &self.repos,
1481 "land",
1482 &LandArgs {
1483 // A pull request from a branch lands from the repository itself.
1484 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
1485 branch: pull.branch.clone(),
1486 actor: a.actor.clone(),
1487 },
1488 )
1489 .await?;
1490 let landed = check!(landed);
1491 Ok(Outcome::Ok(
1492 self.record_merge(&repo, pull, &a.actor, a.keep_issue_open, landed)
1493 .await?,
1494 ))
1495 }
1496
1497 /// Records a pull request as merged once the default branch holds it:
1498 /// closes its issue, supersedes the others for it, and says so.
1499 pub(crate) async fn record_merge(
1500 &self,
1501 repo: &Repo,
1502 mut pull: Pull,
1503 actor: &User,
1504 keep_issue_open: bool,
1505 landed: Landed,
1506 ) -> Result<Pull> {
1507 let issue = match pull.issue {
1508 Some(number) if !keep_issue_open => self
1509 .issue(&repo.id, number)
1510 .await?
1511 .filter(|issue| issue.state == State::Open),
1512 _ => None,
1513 };
1514 let now = rfc3339(now_ms());
1515 let mut statements = vec![
1516 self.db
1517 .prepare(
1518 "UPDATE pulls
1519 SET status = 'merged', head_commit = ?, merge_base = ?, merged_by = ?,
1520 merged_at = ?, updated_at = ?
1521 WHERE id = ?",
1522 )
1523 .bind(&[
1524 landed.commit.as_str().into(),
1525 optional(&landed.previous),
1526 actor.username.as_str().into(),
1527 now.as_str().into(),
1528 now.as_str().into(),
1529 pull.id.as_str().into(),
1530 ])?,
1531 ];
1532 if let Some(issue) = &issue {
1533 statements.push(
1534 self.db
1535 .prepare(
1536 "UPDATE issues
1537 SET state = 'closed', reason = 'completed', resolved_by = ?,
1538 closed_at = ?, updated_at = ?
1539 WHERE id = ?",
1540 )
1541 .bind(&[
1542 pull.number.into(),
1543 now.as_str().into(),
1544 now.as_str().into(),
1545 issue.id.as_str().into(),
1546 ])?,
1547 );
1548 statements.push(
1549 self.db
1550 .prepare(
1551 "UPDATE pulls SET status = 'closed', superseded_by = ?, updated_at = ?
1552 WHERE issue_id = ? AND id != ? AND status IN ('draft', 'open')",
1553 )
1554 .bind(&[
1555 pull.number.into(),
1556 now.as_str().into(),
1557 issue.id.as_str().into(),
1558 pull.id.as_str().into(),
1559 ])?,
1560 );
1561 }
1562 self.db.batch(statements).await?;
1563
1564 self.publish(
1565 "pull.merged",
1566 &repo.id,
1567 actor,
1568 PullEvent {
1569 commit: Some(landed.commit.clone()),
1570 ..Self::pull_event(&pull)
1571 },
1572 )
1573 .await?;
1574 if let Some(issue) = &issue {
1575 self.publish(
1576 "issue.closed",
1577 &repo.id,
1578 actor,
1579 IssueEvent {
1580 reason: Some(IssueReason::Completed.as_str()),
1581 resolved_by: Some(pull.number),
1582 ..Self::issue_event(issue)
1583 },
1584 )
1585 .await?;
1586 }
1587
1588 let who = (actor.id.as_str(), actor.username.as_str());
1589 self.note(&repo.id, pull.number, who, "merged this").await?;
1590 if let Some(issue) = &issue {
1591 self.note(
1592 &repo.id,
1593 issue.number,
1594 who,
1595 &format!("closed this by merging #{}", pull.number),
1596 )
1597 .await?;
1598 }
1599 pull.status = PullStatus::Merged;
1600 pull.head_commit = Some(landed.commit.clone());
1601 pull.merge_base = landed.previous;
1602 pull.merged_by = Some(actor.username.clone());
1603 pull.merged_at = Some(now.clone());
1604 pull.updated_at = now;
1605 Ok(pull)
1606 }
1607
1608 async fn list_active_pulls(&self, a: ViewerArgs) -> Result<Vec<ActivePull>> {
1609 let Some(viewer) = a.viewer else {
1610 return Ok(Vec::new());
1611 };
1612 let found = self
1613 .db
1614 .prepare(
1615 "SELECT * FROM pulls
1616 WHERE author_id = ? AND status IN ('draft', 'open')
1617 ORDER BY updated_at DESC LIMIT 50",
1618 )
1619 .bind(&[viewer.id.into()])?
1620 .all()
1621 .await?;
1622 let snapshots = found.results::<Snapshot>()?;
1623 let pulls: Vec<Pull> = found.results::<PullRow>()?.into_iter().map(Pull::from).collect();
1624 // Each one at once: its issue, and where it stands. That is the
1625 // remembered assessment when there is one, and worked out otherwise.
1626 try_join_all(pulls.into_iter().zip(snapshots).map(|(pull, snapshot)| async move {
1627 let issue = match pull.issue {
1628 Some(number) => self.issue(&pull.repo_id, number).await?,
1629 None => None,
1630 };
1631 // Only a pull request g1t is seeing through has a lifecycle.
1632 let lifecycle = if !lifecycle::made_by_g1t(&pull) || snapshot.managed == 0 {
1633 None
1634 } else if let (Some(stage), Some(detail)) = (snapshot.stage, snapshot.stage_detail) {
1635 Some(Lifecycle {
1636 stage,
1637 detail,
1638 revisions: snapshot.revisions,
1639 })
1640 } else {
1641 let behind = self.is_behind(&pull.repo_id, &pull).await?;
1642 self.assess(&pull, &issue, behind)
1643 .await?
1644 .map(|(lifecycle, _)| lifecycle)
1645 };
1646 Ok::<_, worker::Error>(ActivePull {
1647 pull,
1648 issue,
1649 lifecycle,
1650 })
1651 }))
1652 .await
1653 }
1654
1655 // --- Sessions ----------------------------------------------------------
1656
1657 async fn append_session(&self, a: AppendSessionArgs) -> Result<Outcome<Appended>> {
1658 if a.entries.is_empty() {
1659 return Ok(Outcome::Ok(Appended { count: 0 }));
1660 }
1661 if a.entries.len() > MAX_ENTRY_BATCH {
1662 return Ok(Outcome::fail(
1663 FailureCode::Invalid,
1664 format!("Send at most {MAX_ENTRY_BATCH} entries at a time."),
1665 ));
1666 }
1667 let viewer = Some(a.actor.clone());
1668 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
1669 if pull.author.id != a.actor.id {
1670 return Ok(Outcome::fail(
1671 FailureCode::Forbidden,
1672 "Only whoever opened a pull request can record its session.",
1673 ));
1674 }
1675
1676 let now = rfc3339(now_ms());
1677 let count = a.entries.len() as u32;
1678 let mut statements = Vec::with_capacity(a.entries.len() + 1);
1679 for entry in a.entries {
1680 let kind = serde_json::to_value(entry.kind)?;
1681 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
1682 // Each insert takes the next sequence number itself, so two
1683 // writers appending at once cannot collide.
1684 statements.push(
1685 self.db
1686 .prepare(
1687 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
1688 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
1689 FROM session_entries WHERE pull_id = ?",
1690 )
1691 .bind(&[
1692 pull.id.as_str().into(),
1693 kind.as_str().unwrap_or("note").into(),
1694 text.into(),
1695 optional(&entry.tool),
1696 optional(&entry.commit.or_else(|| pull.head_commit.clone())),
1697 now.as_str().into(),
1698 pull.id.as_str().into(),
1699 ])?,
1700 );
1701 }
1702 statements.push(
1703 self.db
1704 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
1705 .bind(&[now.as_str().into(), pull.id.as_str().into()])?,
1706 );
1707 self.db.batch(statements).await?;
1708 self.publish(
1709 "session.appended",
1710 &pull.repo_id,
1711 &a.actor,
1712 SessionAppended {
1713 pull_id: pull.id.clone(),
1714 repo_id: pull.repo_id.clone(),
1715 number: pull.number,
1716 count,
1717 },
1718 )
1719 .await?;
1720 Ok(Outcome::Ok(Appended { count }))
1721 }
1722
1723 /// Adds entries to a pull request's session, each taking the next
1724 /// sequence number, without announcing it.
1725 pub(crate) async fn append_entries(&self, pull: &Pull, entries: &[NewSessionEntry]) -> Result<()> {
1726 let now = rfc3339(now_ms());
1727 let mut statements = Vec::with_capacity(entries.len());
1728 for entry in entries {
1729 let kind = serde_json::to_value(entry.kind)?;
1730 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
1731 statements.push(
1732 self.db
1733 .prepare(
1734 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
1735 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
1736 FROM session_entries WHERE pull_id = ?",
1737 )
1738 .bind(&[
1739 pull.id.as_str().into(),
1740 kind.as_str().unwrap_or("note").into(),
1741 text.into(),
1742 optional(&entry.tool),
1743 optional(&entry.commit.clone().or_else(|| pull.head_commit.clone())),
1744 now.as_str().into(),
1745 pull.id.as_str().into(),
1746 ])?,
1747 );
1748 }
1749 self.db.batch(statements).await?;
1750 Ok(())
1751 }
1752
1753 async fn read_session(&self, a: ViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
1754 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
1755 let rows = self
1756 .db
1757 .prepare(
1758 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
1759 WHERE pull_id = ? AND seq > ? ORDER BY seq LIMIT ?",
1760 )
1761 .bind(&[pull.id.into(), a.after_seq.into(), SESSION_PAGE.into()])?
1762 .all()
1763 .await?
1764 .results::<SessionRow>()?;
1765 Ok(Outcome::Ok(
1766 rows.into_iter().map(SessionEntry::from).collect(),
1767 ))
1768 }
1769
1770 /// A push moves the head of the pull request it concerns: the one whose
1771 /// fork was pushed to, or the one opened from the branch that moved.
1772 async fn on_event(&self, event: &Event) -> Result<()> {
1773 if event.kind != "git.push" {
1774 return Ok(());
1775 }
1776 let (Some(repo_id), Some(after), Some(git_ref)) = (
1777 event.repo_id.as_deref(),
1778 event.data["after"].as_str(),
1779 event.data["ref"].as_str(),
1780 ) else {
1781 return Ok(());
1782 };
1783 let now = rfc3339(now_ms());
1784 // The head moved, so whatever the checks said no longer applies, and
1785 // whatever step g1t was waiting on has been taken.
1786 let moved = "UPDATE pulls
1787 SET head_commit = ?, updated_at = ?, check_status = NULL, check_run_id = NULL,
1788 working_on = NULL, working_until = NULL, stalled = NULL";
1789 let active = "status IN ('draft', 'open') AND head_commit IS NOT ?";
1790 let returning = "RETURNING id, repo_id, number, issue_number, status";
1791 let mut pulls: Vec<MovedRow> = Vec::new();
1792 // A fork carries its pull request on its default branch.
1793 if event.data["defaultBranch"].as_bool() == Some(true) {
1794 pulls.extend(
1795 self.db
1796 .prepare(format!(
1797 "{moved} WHERE fork_repo_id = ? AND {active} {returning}"
1798 ))
1799 .bind(&[
1800 after.into(),
1801 now.as_str().into(),
1802 repo_id.into(),
1803 after.into(),
1804 ])?
1805 .all()
1806 .await?
1807 .results::<MovedRow>()?,
1808 );
1809 }
1810 if let Some(branch) = git_ref.strip_prefix("refs/heads/") {
1811 pulls.extend(
1812 self.db
1813 .prepare(format!(
1814 "{moved} WHERE repo_id = ? AND source_branch = ? AND {active} {returning}"
1815 ))
1816 .bind(&[
1817 after.into(),
1818 now.as_str().into(),
1819 repo_id.into(),
1820 branch.into(),
1821 after.into(),
1822 ])?
1823 .all()
1824 .await?
1825 .results::<MovedRow>()?,
1826 );
1827 }
1828 // What each now changes, so overlaps show while the work is under way.
1829 for moved in &pulls {
1830 if let Some(pull) = self.pull_by_id(&moved.id).await? {
1831 self.refresh_files(&pull).await?;
1832 }
1833 }
1834 // A merge that was waiting for this push to bring it up to date.
1835 for moved in &pulls {
1836 self.land_if_requested(&moved.id).await?;
1837 }
1838 // Whether each still merges cleanly, and, when a default branch
1839 // moved, every open pull request into it.
1840 let moved_ids: Vec<String> = pulls.iter().map(|pull| pull.id.clone()).collect();
1841 self.after_push(repo_id, event.data["defaultBranch"].as_bool() == Some(true), &moved_ids)
1842 .await;
1843 // A draft is announced when it is marked ready instead.
1844 for pull in pulls
1845 .into_iter()
1846 .filter(|pull| pull.status == PullStatus::Open)
1847 {
1848 self.publish_as(
1849 "pull.updated",
1850 &pull.repo_id,
1851 event.actor.clone(),
1852 PullEvent {
1853 pull_id: pull.id,
1854 repo_id: pull.repo_id.clone(),
1855 number: pull.number,
1856 issue: pull.issue_number,
1857 commit: Some(after.to_owned()),
1858 ..PullEvent::default()
1859 },
1860 )
1861 .await?;
1862 }
1863 Ok(())
1864 }
1865}
1866
1867fn service(env: &Env) -> Result<Work> {
1868 Ok(Work {
1869 db: env.d1("DB")?,
1870 identity: env.service("IDENTITY")?,
1871 repos: env.service("REPOS")?,
1872 events: env.service("EVENTS")?,
1873 actions: env.service("ACTIONS")?,
1874 })
1875}
1876
1877#[event(fetch)]
1878async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
1879 let Some(method) = rpc_method(&request) else {
1880 return Response::error("Not found", 404);
1881 };
1882 let body: serde_json::Value = request.json().await?;
1883 let work = service(&env)?;
1884
1885 match method.as_str() {
1886 "open_issue" => reply(&work.open_issue(args(body)?).await?),
1887 "list_issues" => reply(&work.list_issues(args(body)?).await?),
1888 "get_issue" => reply(&work.get_issue(args(body)?).await?),
1889 "update_issue" => reply(&work.update_issue(args(body)?).await?),
1890 "close_issue" => reply(&work.close_issue(args(body)?).await?),
1891 "reopen_issue" => reply(&work.reopen_issue(args(body)?).await?),
1892 "list_labels" => reply(&work.list_labels(args(body)?).await?),
1893 "counts" => reply(&work.counts(args(body)?).await?),
1894 "add_comment" => reply(&work.add_comment(args(body)?).await?),
1895 "start_checks" => reply(&work.start_checks(args(body)?).await?),
1896 "report_checks" => reply(&work.report_checks(args(body)?).await?),
1897 "set_commit_status" => reply(&work.set_commit_status(args(body)?).await?),
1898 "start_review" => reply(&work.start_review(args(body)?).await?),
1899 "advance" => reply(&work.advance(args(body)?).await?),
1900 "stall" => reply(&work.stall(args(body)?).await?),
1901 "managed_pulls" => reply(&work.managed_pulls(args(body)?).await?),
1902 "queue" => reply(&work.queue(args(body)?).await?),
1903 "queue_build" => reply(&work.queue_build(args(body)?).await?),
1904 "report_queue" => reply(&work.report_queue(args(body)?).await?),
1905 "remove_from_queue" => reply(&work.remove_from_queue(args(body)?).await?),
1906 "message_agent" => reply(&work.message_agent(args(body)?).await?),
1907 "locate_pull" => reply(&work.locate_pull(args(body)?).await?),
1908 "answer_message" => reply(&work.answer_message(args(body)?).await?),
1909 "take_messages" => reply(&work.take_messages(args(body)?).await?),
1910 "wake_for_messages" => reply(&work.wake_for_messages(args(body)?).await?),
1911 "catch_up_job" => reply(&work.catch_up_job(args(body)?).await?),
1912 "get_settings" => reply(&work.get_settings(args(body)?).await?),
1913 "update_settings" => reply(&work.update_settings(args(body)?).await?),
1914 "report_review" => reply(&work.report_review(args(body)?).await?),
1915 "open_pull" => reply(&work.open_pull(args(body)?).await?),
1916 "list_pulls" => reply(&work.list_pulls(args(body)?).await?),
1917 "get_pull" => reply(&work.get_pull(args(body)?).await?),
1918 "update_pull" => reply(&work.update_pull(args(body)?).await?),
1919 "catch_up_pull" => reply(&work.catch_up_pull(args(body)?).await?),
1920 "ready_pull" => reply(&work.ready_pull(args(body)?).await?),
1921 "close_pull" => reply(&work.close_pull(args(body)?).await?),
1922 "merge_pull" => reply(&work.merge_pull(args(body)?).await?),
1923 "list_active_pulls" => reply(&work.list_active_pulls(args(body)?).await?),
1924 "by_author" => reply(&work.by_author(args(body)?).await?),
1925 "start_plan" => reply(&work.start_plan(args(body)?).await?),
1926 "report_plan" => reply(&work.report_plan(args(body)?).await?),
1927 "get_plan" => reply(&work.get_plan(args(body)?).await?),
1928 "list_plans" => reply(&work.list_plans(args(body)?).await?),
1929 "apply_plan" => reply(&work.apply_plan(args(body)?).await?),
1930 "queue_issue" => reply(&work.queue_issue(args(body)?).await?),
1931 "ready_issues" => reply(&work.ready_issues(args(body)?).await?),
1932 "list_assigned_issues" => reply(&work.list_assigned_issues(args(body)?).await?),
1933 "append_session" => reply(&work.append_session(args(body)?).await?),
1934 "read_session" => reply(&work.read_session(args(body)?).await?),
1935 // Agents at work, their sessions, and memory (runs.rs, memory.rs).
1936 "open_run" => reply(&work.open_run(args(body)?).await?),
1937 "report_run" => reply(&work.report_run(args(body)?).await?),
1938 "stop_run" => reply(&work.stop_run(args(body)?).await?),
1939 "list_runs" => reply(&work.list_runs(args(body)?).await?),
1940 "get_run" => reply(&work.get_run(args(body)?).await?),
1941 "list_sessions" => reply(&work.list_sessions(args(body)?).await?),
1942 "get_session" => reply(&work.get_session(args(body)?).await?),
1943 "list_memories" => reply(&work.list_memories(args(body)?).await?),
1944 "add_memory" => reply(&work.add_memory(args(body)?).await?),
1945 "update_memory" => reply(&work.update_memory(args(body)?).await?),
1946 "delete_memory" => reply(&work.delete_memory(args(body)?).await?),
1947 "recall" => reply(&work.recall(args(body)?).await?),
1948 "memory_context" => reply(&work.memory_context(args(body)?).await?),
1949 // What agents may do in a sandbox (guardrails.rs).
1950 "get_guardrails" => reply(&work.get_guardrails(args(body)?).await?),
1951 "update_guardrails" => reply(&work.update_guardrails(args(body)?).await?),
1952 "run_guardrails" => reply(&work.run_guardrails(args(body)?).await?),
1953 "start_mergecheck" => reply(&work.start_mergecheck(args(body)?).await?),
1954 "report_mergecheck" => reply(&work.report_mergecheck(args(body)?).await?),
1955 // Memory that fills itself, and its review queue (capture.rs).
1956 method if capture::METHODS.contains(&method) => capture::dispatch(&work, method, body).await,
1957 // @g1t-agent in comments, and the label rule (mentions.rs).
1958 "take_mention" => reply(&work.take_mention(args(body)?).await?),
1959 "mention_revision" => reply(&work.mention_revision(args(body)?).await?),
1960 "reply_mention" => reply(&work.reply_mention(args(body)?).await?),
1961 "get_agent_rules" => reply(&work.get_agent_rules(args(body)?).await?),
1962 "set_agent_rules" => reply(&work.set_agent_rules(args(body)?).await?),
1963 _ => Response::error("Unknown method", 404),
1964 }
1965}
1966
1967/// Events from the bus, delivered on this service's own queue.
1968#[event(queue)]
1969async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
1970 let work = service(&env)?;
1971 for message in batch.messages()? {
1972 // A workspace renamed: its agent runs and memory move to the slug it has now.
1973 if g1t_kit::rename::on_event(&env, &env.d1("DB")?, message.body(), &[memory::RENAMED, guardrails::RENAMED].concat()).await? {
1974 message.ack();
1975 continue;
1976 }
1977 capture::on_event(&work, message.body()).await;
1978 work.on_event(message.body()).await?;
1979 message.ack();
1980 }
1981 Ok(())
1982}