pr_01m47d24b0e6n91zwymwxg0vpx/services/work/src/lib.rs

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