pr_01m47d15m3e54sn21z27rpy5n9/services/work/src/lib.rs

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