pr_01m47d24b0e6n91zwymwxg0vpx/services/work/src/lib.rs

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