g1t/services/work/src/lib.rs

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