pr_01m47d15m3e54sn21z27rpy5n9/services/work/src/lib.rs

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