flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/work/src/lib.rs

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