pr_01m47d15m3e54sn21z27rpy5n9/services/work/src/lib.rs

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