pr_01m47d15m3e54sn21z27rpy5n9/services/work/src/lib.rs

1,293 lines47,178 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 rows;
9
10use g1t_contracts::events::{
11 CommentCreated, Event, IssueEvent, NewEvent, Publish, PullEvent, SessionAppended,
12};
13use g1t_contracts::repos::{ForkArgs, GetArgs, HeadArgs, LandArgs, Landed, Repo, RepoPath};
14use g1t_contracts::time::rfc3339;
15use g1t_contracts::work::*;
16use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
17use g1t_kit::{args, now_ms, reply, rpc_method};
18use serde::Serialize;
19use worker::wasm_bindgen::JsValue;
20use worker::{
21 Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event,
22};
23
24use rows::{CommentRow, IssueRow, MovedRow, NumberRow, PullRow, SessionRow, ValueRow};
25
26const SOURCE: &str = "work";
27const MAX_ENTRY_BATCH: usize = 200;
28const MAX_ENTRY_CHARS: usize = 64_000;
29const MAX_TITLE_CHARS: usize = 200;
30const SESSION_PAGE: u32 = 500;
31const LIST_PAGE: u32 = 100;
32const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
33
34const ISSUE_COLUMNS: &str = "issues.*,
35 (SELECT count(*) FROM pulls WHERE pulls.issue_id = issues.id) AS pull_count,
36 (SELECT count(*) FROM comments
37 WHERE comments.repo_id = issues.repo_id AND comments.number = issues.number) AS comment_count";
38
39fn no_issue<T>() -> Outcome<T> {
40 Outcome::fail(FailureCode::NotFound, "Issue not found.")
41}
42
43fn no_pull<T>() -> Outcome<T> {
44 Outcome::fail(FailureCode::NotFound, "Pull request not found.")
45}
46
47fn optional(value: &Option<String>) -> JsValue {
48 value.as_deref().map_or(JsValue::NULL, JsValue::from)
49}
50
51fn optional_number(value: Option<u32>) -> JsValue {
52 value.map_or(JsValue::NULL, JsValue::from)
53}
54
55/// The lowercase name a `State` is stored and sent as.
56fn state_name(state: Option<State>) -> Option<&'static str> {
57 state.map(|state| match state {
58 State::Open => "open",
59 State::Closed => "closed",
60 })
61}
62
63/// A trimmed title, or why it cannot be used.
64fn valid_title(title: &str) -> std::result::Result<&str, &'static str> {
65 let title = title.trim();
66 if title.is_empty() {
67 Err("A title is required.")
68 } else if title.chars().count() > MAX_TITLE_CHARS {
69 Err("That title is too long.")
70 } else {
71 Ok(title)
72 }
73}
74
75/// Unwraps an `Outcome`, returning its failure from the enclosing method.
76macro_rules! check {
77 ($outcome:expr) => {
78 match $outcome {
79 Outcome::Ok(value) => value,
80 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
81 }
82 };
83}
84
85struct Work {
86 db: D1Database,
87 repos: Fetcher,
88 events: Fetcher,
89}
90
91impl Work {
92 async fn publish<T: Serialize>(
93 &self,
94 kind: &'static str,
95 repo_id: &str,
96 actor: &User,
97 data: T,
98 ) -> Result<()> {
99 self.publish_as(kind, repo_id, Some(actor.id.clone()), data)
100 .await
101 }
102
103 /// Publishes an event caused by `actor`, or by g1t itself.
104 async fn publish_as<T: Serialize>(
105 &self,
106 kind: &'static str,
107 repo_id: &str,
108 actor: Option<String>,
109 data: T,
110 ) -> Result<()> {
111 let event = NewEvent {
112 kind,
113 source: SOURCE,
114 repo_id: Some(repo_id.to_owned()),
115 actor,
116 data,
117 };
118 g1t_kit::call(
119 &self.events,
120 "publish",
121 &Publish {
122 events: vec![event],
123 },
124 )
125 .await
126 }
127
128 /// The repository, if the viewer may see it. Whether they may is
129 /// decided by the repos service.
130 async fn repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Outcome<Repo>> {
131 g1t_kit::call(
132 &self.repos,
133 "get",
134 &GetArgs {
135 path: path.clone(),
136 viewer: viewer.clone(),
137 },
138 )
139 .await
140 }
141
142 /// The next number in the repository's sequence. Taking it is one
143 /// statement, so concurrent opens cannot be given the same number.
144 async fn next_number(&self, repo_id: &str) -> Result<u32> {
145 let row = self
146 .db
147 .prepare(
148 "INSERT INTO counters (repo_id, last) VALUES (?, 1)
149 ON CONFLICT (repo_id) DO UPDATE SET last = last + 1
150 RETURNING last AS n",
151 )
152 .bind(&[repo_id.into()])?
153 .first::<NumberRow>(None)
154 .await?;
155 row.map(|row| row.n)
156 .ok_or_else(|| worker::Error::RustError("no number was assigned".into()))
157 }
158
159 async fn issue(&self, repo_id: &str, number: u32) -> Result<Option<Issue>> {
160 Ok(self
161 .db
162 .prepare(format!(
163 "SELECT {ISSUE_COLUMNS} FROM issues WHERE repo_id = ? AND number = ?"
164 ))
165 .bind(&[repo_id.into(), number.into()])?
166 .first::<IssueRow>(None)
167 .await?
168 .map(Issue::from))
169 }
170
171 async fn pull(&self, repo_id: &str, number: u32) -> Result<Option<Pull>> {
172 Ok(self
173 .db
174 .prepare("SELECT * FROM pulls WHERE repo_id = ? AND number = ?")
175 .bind(&[repo_id.into(), number.into()])?
176 .first::<PullRow>(None)
177 .await?
178 .map(Pull::from))
179 }
180
181 async fn comments(&self, repo_id: &str, number: u32) -> Result<Vec<Comment>> {
182 let rows = self
183 .db
184 .prepare(
185 "SELECT * FROM comments WHERE repo_id = ? AND number = ? ORDER BY id LIMIT 500",
186 )
187 .bind(&[repo_id.into(), number.into()])?
188 .all()
189 .await?
190 .results::<CommentRow>()?;
191 Ok(rows.into_iter().map(Comment::from).collect())
192 }
193
194 /// The repository and one of its issues, as seen by `viewer`.
195 async fn issue_at(
196 &self,
197 path: &RepoPath,
198 number: u32,
199 viewer: &Viewer,
200 ) -> Result<Outcome<(Repo, Issue)>> {
201 let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
202 return Ok(no_issue());
203 };
204 Ok(match self.issue(&repo.id, number).await? {
205 Some(issue) => Outcome::Ok((repo, issue)),
206 None => no_issue(),
207 })
208 }
209
210 /// The repository and one of its pull requests, as seen by `viewer`.
211 async fn pull_at(
212 &self,
213 path: &RepoPath,
214 number: u32,
215 viewer: &Viewer,
216 ) -> Result<Outcome<(Repo, Pull)>> {
217 let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
218 return Ok(no_pull());
219 };
220 Ok(match self.pull(&repo.id, number).await? {
221 Some(pull) => Outcome::Ok((repo, pull)),
222 None => no_pull(),
223 })
224 }
225
226 fn issue_event(issue: &Issue) -> IssueEvent {
227 IssueEvent {
228 issue_id: issue.id.clone(),
229 repo_id: issue.repo_id.clone(),
230 number: issue.number,
231 ..IssueEvent::default()
232 }
233 }
234
235 fn pull_event(pull: &Pull) -> PullEvent {
236 PullEvent {
237 pull_id: pull.id.clone(),
238 repo_id: pull.repo_id.clone(),
239 number: pull.number,
240 issue: pull.issue,
241 ..PullEvent::default()
242 }
243 }
244
245 // --- Issues ------------------------------------------------------------
246
247 async fn open_issue(&self, a: OpenIssueArgs) -> Result<Outcome<Issue>> {
248 if !a.actor.verified {
249 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
250 }
251 let title = match valid_title(&a.title) {
252 Ok(title) => title,
253 Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
254 };
255 let Some(labels) = normalize_labels(&a.labels) else {
256 return Ok(Outcome::fail(
257 FailureCode::Invalid,
258 "An issue can have up to 10 labels of up to 40 characters each.",
259 ));
260 };
261 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
262 let checks: Vec<&str> = a
263 .checks
264 .iter()
265 .map(|check| check.trim())
266 .filter(|check| !check.is_empty())
267 .collect();
268
269 let now = now_ms();
270 let id = new_id("iss", now);
271 let number = self.next_number(&repo.id).await?;
272 let timestamp = rfc3339(now);
273 self.db
274 .prepare(
275 "INSERT INTO issues
276 (id, repo_id, number, title, body, labels, checks, author_id, author_name,
277 created_at, updated_at)
278 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
279 )
280 .bind(&[
281 id.as_str().into(),
282 repo.id.as_str().into(),
283 number.into(),
284 title.into(),
285 a.body.trim().into(),
286 serde_json::to_string(&labels)?.into(),
287 serde_json::to_string(&checks)?.into(),
288 a.actor.id.as_str().into(),
289 a.actor.username.as_str().into(),
290 timestamp.as_str().into(),
291 timestamp.as_str().into(),
292 ])?
293 .run()
294 .await?;
295 let Some(issue) = self.issue(&repo.id, number).await? else {
296 return Ok(no_issue());
297 };
298 self.publish(
299 "issue.opened",
300 &repo.id,
301 &a.actor,
302 IssueEvent {
303 title: Some(issue.title.clone()),
304 ..Self::issue_event(&issue)
305 },
306 )
307 .await?;
308 Ok(Outcome::Ok(issue))
309 }
310
311 async fn list_issues(&self, a: ListIssuesArgs) -> Result<Outcome<Vec<Issue>>> {
312 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
313 let state = state_name(a.state).map_or(JsValue::NULL, JsValue::from);
314 let label = a
315 .label
316 .map(|label| label.trim().to_lowercase())
317 .filter(|label| !label.is_empty());
318 let rows = self
319 .db
320 .prepare(format!(
321 "SELECT {ISSUE_COLUMNS} FROM issues
322 WHERE repo_id = ? AND (? IS NULL OR state = ?)
323 AND (? IS NULL OR EXISTS
324 (SELECT 1 FROM json_each(issues.labels) WHERE json_each.value = ?))
325 ORDER BY number DESC LIMIT ?"
326 ))
327 .bind(&[
328 repo.id.into(),
329 state.clone(),
330 state,
331 optional(&label),
332 optional(&label),
333 LIST_PAGE.into(),
334 ])?
335 .all()
336 .await?
337 .results::<IssueRow>()?;
338 Ok(Outcome::Ok(rows.into_iter().map(Issue::from).collect()))
339 }
340
341 async fn get_issue(&self, a: ViewArgs) -> Result<Outcome<IssueDetail>> {
342 let (repo, issue) = check!(self.issue_at(&a.repo, a.number, &a.viewer).await?);
343 let pulls = self
344 .db
345 .prepare("SELECT * FROM pulls WHERE issue_id = ? ORDER BY number")
346 .bind(&[issue.id.as_str().into()])?
347 .all()
348 .await?
349 .results::<PullRow>()?;
350 Ok(Outcome::Ok(IssueDetail {
351 comments: self.comments(&repo.id, issue.number).await?,
352 pulls: pulls.into_iter().map(Pull::from).collect(),
353 issue,
354 }))
355 }
356
357 /// The issue, if `actor` wrote it or belongs to the repository's workspace.
358 async fn manageable_issue(
359 &self,
360 actor: &User,
361 path: &RepoPath,
362 number: u32,
363 ) -> Result<Outcome<Issue>> {
364 let (repo, issue) = check!(self.issue_at(path, number, &Some(actor.clone())).await?);
365 if issue.author.id != actor.id && !actor.is_member(&repo.namespace) {
366 return Ok(Outcome::fail(
367 FailureCode::Forbidden,
368 "Only the author or a member of the workspace can change an issue.",
369 ));
370 }
371 Ok(Outcome::Ok(issue))
372 }
373
374 async fn update_issue(&self, a: UpdateIssueArgs) -> Result<Outcome<Issue>> {
375 let issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
376 let title = match a.title.as_deref().map(valid_title) {
377 Some(Err(message)) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
378 Some(Ok(title)) => Some(title.to_owned()),
379 None => None,
380 };
381 let labels = match a.labels.as_deref().map(normalize_labels) {
382 Some(None) => {
383 return Ok(Outcome::fail(
384 FailureCode::Invalid,
385 "An issue can have up to 10 labels of up to 40 characters each.",
386 ));
387 }
388 Some(Some(labels)) => Some(serde_json::to_string(&labels)?),
389 None => None,
390 };
391 let body = a.body.map(|body| body.trim().to_owned());
392 self.db
393 .prepare(
394 "UPDATE issues
395 SET title = COALESCE(?, title), body = COALESCE(?, body),
396 labels = COALESCE(?, labels), updated_at = ?
397 WHERE id = ?",
398 )
399 .bind(&[
400 optional(&title),
401 optional(&body),
402 optional(&labels),
403 rfc3339(now_ms()).into(),
404 issue.id.as_str().into(),
405 ])?
406 .run()
407 .await?;
408 let Some(issue) = self.issue(&issue.repo_id, issue.number).await? else {
409 return Ok(no_issue());
410 };
411 self.publish(
412 "issue.updated",
413 &issue.repo_id,
414 &a.actor,
415 Self::issue_event(&issue),
416 )
417 .await?;
418 Ok(Outcome::Ok(issue))
419 }
420
421 async fn close_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
422 let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
423 if issue.state == State::Closed {
424 return Ok(Outcome::fail(
425 FailureCode::Conflict,
426 "This issue is already closed.",
427 ));
428 }
429 let reason = a.reason.unwrap_or(IssueReason::Completed);
430 let now = rfc3339(now_ms());
431 self.db
432 .prepare(
433 "UPDATE issues SET state = 'closed', reason = ?, closed_at = ?, updated_at = ?
434 WHERE id = ?",
435 )
436 .bind(&[
437 reason.as_str().into(),
438 now.as_str().into(),
439 now.as_str().into(),
440 issue.id.as_str().into(),
441 ])?
442 .run()
443 .await?;
444 self.publish(
445 "issue.closed",
446 &issue.repo_id,
447 &a.actor,
448 IssueEvent {
449 reason: Some(reason.as_str()),
450 ..Self::issue_event(&issue)
451 },
452 )
453 .await?;
454 issue.state = State::Closed;
455 issue.reason = Some(reason);
456 issue.closed_at = Some(now.clone());
457 issue.updated_at = now;
458 Ok(Outcome::Ok(issue))
459 }
460
461 async fn reopen_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
462 let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
463 if issue.state == State::Open {
464 return Ok(Outcome::fail(
465 FailureCode::Conflict,
466 "This issue is already open.",
467 ));
468 }
469 let now = rfc3339(now_ms());
470 self.db
471 .prepare(
472 "UPDATE issues
473 SET state = 'open', reason = NULL, resolved_by = NULL, closed_at = NULL,
474 updated_at = ?
475 WHERE id = ?",
476 )
477 .bind(&[now.as_str().into(), issue.id.as_str().into()])?
478 .run()
479 .await?;
480 self.publish(
481 "issue.reopened",
482 &issue.repo_id,
483 &a.actor,
484 Self::issue_event(&issue),
485 )
486 .await?;
487 issue.state = State::Open;
488 issue.reason = None;
489 issue.resolved_by = None;
490 issue.closed_at = None;
491 issue.updated_at = now;
492 Ok(Outcome::Ok(issue))
493 }
494
495 /// The default labels, then every other label in use on the repository.
496 async fn list_labels(&self, a: ViewArgs) -> Result<Outcome<Vec<String>>> {
497 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
498 let used = self
499 .db
500 .prepare(
501 "SELECT DISTINCT json_each.value AS value
502 FROM issues, json_each(issues.labels)
503 WHERE issues.repo_id = ? ORDER BY 1 LIMIT 200",
504 )
505 .bind(&[repo.id.into()])?
506 .all()
507 .await?
508 .results::<ValueRow>()?;
509 let mut labels: Vec<String> = DEFAULT_LABELS.iter().map(|label| (*label).into()).collect();
510 for row in used {
511 if !labels.contains(&row.value) {
512 labels.push(row.value);
513 }
514 }
515 Ok(Outcome::Ok(labels))
516 }
517
518 async fn counts(&self, a: ViewArgs) -> Result<Outcome<Counts>> {
519 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
520 let counts = self
521 .db
522 .prepare(
523 "SELECT
524 (SELECT count(*) FROM issues WHERE repo_id = ? AND state = 'open') AS issues,
525 (SELECT count(*) FROM pulls
526 WHERE repo_id = ? AND status IN ('draft', 'open')) AS pulls",
527 )
528 .bind(&[repo.id.as_str().into(), repo.id.as_str().into()])?
529 .first::<Counts>(None)
530 .await?;
531 Ok(Outcome::Ok(counts.unwrap_or(Counts {
532 issues: 0,
533 pulls: 0,
534 })))
535 }
536
537 // --- Comments ----------------------------------------------------------
538
539 async fn add_comment(&self, a: AddCommentArgs) -> Result<Outcome<Comment>> {
540 if !a.actor.verified {
541 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
542 }
543 let body = a.body.trim();
544 // An approval speaks for itself; anything else has to say something.
545 if body.is_empty() && a.verdict != Some(Verdict::Approve) {
546 return Ok(Outcome::fail(
547 FailureCode::Invalid,
548 "A comment cannot be empty.",
549 ));
550 }
551 let path = a
552 .path
553 .as_deref()
554 .map(str::trim)
555 .filter(|path| !path.is_empty());
556 let line = a.line.filter(|line| *line > 0 && path.is_some());
557 if body.chars().count() > MAX_ENTRY_CHARS {
558 return Ok(Outcome::fail(
559 FailureCode::Invalid,
560 "That comment is too long.",
561 ));
562 }
563 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
564 // The number names an issue or a pull request, never both.
565 let table = if self.issue(&repo.id, a.number).await?.is_some() {
566 if path.is_some() || a.verdict.is_some() {
567 return Ok(Outcome::fail(
568 FailureCode::Invalid,
569 "Only a pull request can be reviewed or commented on by line.",
570 ));
571 }
572 "issues"
573 } else if let Some(pull) = self.pull(&repo.id, a.number).await? {
574 if a.verdict.is_some() && pull.author.id == a.actor.id {
575 return Ok(Outcome::fail(
576 FailureCode::Forbidden,
577 "You cannot approve or request changes on your own pull request.",
578 ));
579 }
580 "pulls"
581 } else {
582 return Ok(Outcome::fail(
583 FailureCode::NotFound,
584 "No issue or pull request has that number.",
585 ));
586 };
587
588 let now = now_ms();
589 let comment = Comment {
590 id: new_id("cmt", now),
591 author: a.actor.clone(),
592 body: body.to_owned(),
593 path: path.map(str::to_owned),
594 line,
595 verdict: a.verdict,
596 created_at: rfc3339(now),
597 };
598 self.db
599 .batch(vec![
600 self.db
601 .prepare(
602 "INSERT INTO comments
603 (id, repo_id, number, author_id, author_name, body, path, line,
604 verdict, created_at)
605 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
606 )
607 .bind(&[
608 comment.id.as_str().into(),
609 repo.id.as_str().into(),
610 a.number.into(),
611 a.actor.id.as_str().into(),
612 a.actor.username.as_str().into(),
613 body.into(),
614 optional(&comment.path),
615 optional_number(line),
616 a.verdict
617 .map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
618 comment.created_at.as_str().into(),
619 ])?,
620 self.db
621 .prepare(format!(
622 "UPDATE {table} SET updated_at = ? WHERE repo_id = ? AND number = ?"
623 ))
624 .bind(&[
625 comment.created_at.as_str().into(),
626 repo.id.as_str().into(),
627 a.number.into(),
628 ])?,
629 ])
630 .await?;
631 self.publish(
632 "comment.created",
633 &repo.id,
634 &a.actor,
635 CommentCreated {
636 comment_id: comment.id.clone(),
637 repo_id: repo.id.clone(),
638 number: a.number,
639 },
640 )
641 .await?;
642 Ok(Outcome::Ok(comment))
643 }
644
645 // --- Pull requests -----------------------------------------------------
646
647 async fn open_pull(&self, a: OpenPullArgs) -> Result<Outcome<Pull>> {
648 if !a.actor.verified {
649 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
650 }
651 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
652 let issue = match a.issue {
653 Some(number) => match self.issue(&repo.id, number).await? {
654 Some(issue) if issue.state == State::Open => Some(issue),
655 Some(_) => {
656 return Ok(Outcome::fail(
657 FailureCode::Conflict,
658 "This issue is closed.",
659 ));
660 }
661 None => return Ok(no_issue()),
662 },
663 None => None,
664 };
665 // A pull request for an issue takes the issue's title unless given one.
666 let title = match (a.title.trim(), &issue) {
667 ("", Some(issue)) => issue.title.clone(),
668 (title, _) => match valid_title(title) {
669 Ok(title) => title.to_owned(),
670 Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
671 },
672 };
673 let agent = match a.agent.trim() {
674 "" => "agent",
675 agent => agent,
676 };
677 let runtime = match a.runtime {
678 Runtime::Hosted => "hosted",
679 Runtime::External => "external",
680 };
681
682 let now = now_ms();
683 let id = new_id("pr", now);
684 let branch = a
685 .branch
686 .as_deref()
687 .map(str::trim)
688 .filter(|branch| !branch.is_empty());
689 // The change is on a branch already pushed to the repository, or
690 // will be made in a fork created for this pull request.
691 let (fork, head) = match branch {
692 Some(branch) => {
693 if branch == repo.default_branch {
694 return Ok(Outcome::fail(
695 FailureCode::Invalid,
696 format!("Choose a branch other than {branch}."),
697 ));
698 }
699 let head: Option<String> = g1t_kit::call(
700 &self.repos,
701 "head",
702 &HeadArgs {
703 repo_id: repo.id.clone(),
704 branch: branch.to_owned(),
705 },
706 )
707 .await?;
708 let Some(head) = head else {
709 return Ok(Outcome::fail(
710 FailureCode::NotFound,
711 format!("There is no branch named {branch}. Push it first."),
712 ));
713 };
714 let existing = self
715 .db
716 .prepare(
717 "SELECT number AS n FROM pulls
718 WHERE repo_id = ? AND source_branch = ? AND status IN ('draft', 'open')",
719 )
720 .bind(&[repo.id.as_str().into(), branch.into()])?
721 .first::<NumberRow>(None)
722 .await?;
723 if let Some(existing) = existing {
724 return Ok(Outcome::fail(
725 FailureCode::Conflict,
726 format!("Pull request #{} is already open for {branch}.", existing.n),
727 ));
728 }
729 (None, Some(head))
730 }
731 None => {
732 let fork: Outcome<Repo> = g1t_kit::call(
733 &self.repos,
734 "fork_for_pull",
735 &ForkArgs {
736 source_id: repo.id.clone(),
737 pull_id: id.clone(),
738 actor: a.actor.clone(),
739 },
740 )
741 .await?;
742 (Some(check!(fork)), None)
743 }
744 };
745 // A branch already holds the work, so its pull request is ready for
746 // review from the start; one with a fork starts as a draft.
747 let status = if branch.is_some() { "open" } else { "draft" };
748 let body = Some(a.body.trim().to_owned()).filter(|body| !body.is_empty());
749
750 let number = self.next_number(&repo.id).await?;
751 let timestamp = rfc3339(now);
752 self.db
753 .prepare(
754 "INSERT INTO pulls
755 (id, repo_id, number, issue_id, issue_number, title, body, agent, runtime,
756 status, fork_repo_id, fork_namespace, fork_name, source_branch, head_commit,
757 author_id, author_name, created_at, updated_at)
758 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
759 )
760 .bind(&[
761 id.as_str().into(),
762 repo.id.as_str().into(),
763 number.into(),
764 optional(&issue.as_ref().map(|issue| issue.id.clone())),
765 optional_number(issue.as_ref().map(|issue| issue.number)),
766 title.into(),
767 optional(&body),
768 agent.into(),
769 runtime.into(),
770 status.into(),
771 optional(&fork.as_ref().map(|fork| fork.id.clone())),
772 optional(&fork.as_ref().map(|fork| fork.namespace.clone())),
773 optional(&fork.as_ref().map(|fork| fork.name.clone())),
774 optional(&branch.map(str::to_owned)),
775 optional(&head),
776 a.actor.id.as_str().into(),
777 a.actor.username.as_str().into(),
778 timestamp.as_str().into(),
779 timestamp.as_str().into(),
780 ])?
781 .run()
782 .await?;
783 let Some(pull) = self.pull(&repo.id, number).await? else {
784 return Ok(no_pull());
785 };
786 self.publish(
787 "pull.opened",
788 &repo.id,
789 &a.actor,
790 PullEvent {
791 agent: Some(pull.agent.clone()),
792 ..Self::pull_event(&pull)
793 },
794 )
795 .await?;
796 Ok(Outcome::Ok(pull))
797 }
798
799 async fn list_pulls(&self, a: ListPullsArgs) -> Result<Outcome<Vec<Pull>>> {
800 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
801 let filter = match a.state {
802 Some(State::Open) => "AND status IN ('draft', 'open')",
803 Some(State::Closed) => "AND status IN ('merged', 'closed')",
804 None => "",
805 };
806 let rows = self
807 .db
808 .prepare(format!(
809 "SELECT * FROM pulls WHERE repo_id = ? {filter} ORDER BY number DESC LIMIT ?"
810 ))
811 .bind(&[repo.id.into(), LIST_PAGE.into()])?
812 .all()
813 .await?
814 .results::<PullRow>()?;
815 Ok(Outcome::Ok(rows.into_iter().map(Pull::from).collect()))
816 }
817
818 async fn get_pull(&self, a: ViewArgs) -> Result<Outcome<PullDetail>> {
819 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
820 let issue = match pull.issue {
821 Some(number) => self.issue(&repo.id, number).await?,
822 None => None,
823 };
824 Ok(Outcome::Ok(PullDetail {
825 comments: self.comments(&repo.id, pull.number).await?,
826 checks: self.latest_checks(&pull.id).await?,
827 issue,
828 pull,
829 }))
830 }
831
832 /// The pull request, if it is still active and `actor` opened it or
833 /// belongs to the repository's workspace.
834 async fn manageable_pull(
835 &self,
836 actor: &User,
837 path: &RepoPath,
838 number: u32,
839 ) -> Result<Outcome<Pull>> {
840 let (repo, pull) = check!(self.pull_at(path, number, &Some(actor.clone())).await?);
841 if pull.author.id != actor.id && !actor.is_member(&repo.namespace) {
842 return Ok(Outcome::fail(
843 FailureCode::Forbidden,
844 "Only whoever opened a pull request, or a member of the workspace, can change it.",
845 ));
846 }
847 if !pull.status.is_active() {
848 return Ok(Outcome::fail(
849 FailureCode::Conflict,
850 format!("This pull request is already {}.", pull.status.as_str()),
851 ));
852 }
853 Ok(Outcome::Ok(pull))
854 }
855
856 /// Marks a draft ready for review, or updates the description of one
857 /// that already is.
858 async fn ready_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
859 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
860 let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
861 let now = rfc3339(now_ms());
862 self.db
863 .prepare(
864 "UPDATE pulls SET status = 'open', body = COALESCE(?, body), updated_at = ?
865 WHERE id = ?",
866 )
867 .bind(&[
868 optional(&summary),
869 now.as_str().into(),
870 pull.id.as_str().into(),
871 ])?
872 .run()
873 .await?;
874 if pull.status == PullStatus::Draft {
875 self.publish(
876 "pull.ready",
877 &pull.repo_id,
878 &a.actor,
879 Self::pull_event(&pull),
880 )
881 .await?;
882 }
883 pull.status = PullStatus::Open;
884 pull.body = summary.or(pull.body);
885 pull.updated_at = now;
886 Ok(Outcome::Ok(pull))
887 }
888
889 async fn close_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
890 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
891 let now = rfc3339(now_ms());
892 self.db
893 .prepare("UPDATE pulls SET status = 'closed', updated_at = ? WHERE id = ?")
894 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
895 .run()
896 .await?;
897 self.publish(
898 "pull.closed",
899 &pull.repo_id,
900 &a.actor,
901 Self::pull_event(&pull),
902 )
903 .await?;
904 pull.status = PullStatus::Closed;
905 pull.updated_at = now;
906 Ok(Outcome::Ok(pull))
907 }
908
909 /// Lands the pull request on the repository's default branch. Unless
910 /// told to keep it open, that resolves the issue it was for: the issue
911 /// closes naming this pull request, and the others still in progress
912 /// for it close as superseded.
913 async fn merge_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
914 let viewer = Some(a.actor.clone());
915 let (repo, mut pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
916 match pull.status {
917 PullStatus::Open => {}
918 PullStatus::Draft => {
919 return Ok(Outcome::fail(
920 FailureCode::Conflict,
921 "This pull request is still a draft. Mark it ready for review first.",
922 ));
923 }
924 status => {
925 return Ok(Outcome::fail(
926 FailureCode::Conflict,
927 format!("This pull request is already {}.", status.as_str()),
928 ));
929 }
930 }
931 if !a.ignore_checks {
932 let waiting = match pull.check_status {
933 Some(CheckStatus::Queued | CheckStatus::Running) => {
934 Some("The acceptance checks are still running.")
935 }
936 Some(CheckStatus::Failed) => Some("The acceptance checks did not pass."),
937 Some(CheckStatus::Errored) => Some("The acceptance checks could not be run."),
938 Some(CheckStatus::Passed) | None => None,
939 };
940 if let Some(reason) = waiting {
941 return Ok(Outcome::fail(
942 FailureCode::Conflict,
943 format!("{reason} Wait or fix them, or merge anyway by ignoring the checks."),
944 ));
945 }
946 }
947 let issue = match pull.issue {
948 Some(number) if !a.keep_issue_open => self
949 .issue(&repo.id, number)
950 .await?
951 .filter(|issue| issue.state == State::Open),
952 _ => None,
953 };
954
955 // Whether the actor may write to the repository is decided by repos.
956 let landed: Outcome<Landed> = g1t_kit::call(
957 &self.repos,
958 "land",
959 &LandArgs {
960 // A pull request from a branch lands from the repository itself.
961 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
962 branch: pull.branch.clone(),
963 actor: a.actor.clone(),
964 },
965 )
966 .await?;
967 let landed = check!(landed);
968
969 let now = rfc3339(now_ms());
970 let mut statements = vec![
971 self.db
972 .prepare(
973 "UPDATE pulls
974 SET status = 'merged', head_commit = ?, merge_base = ?, merged_by = ?,
975 merged_at = ?, updated_at = ?
976 WHERE id = ?",
977 )
978 .bind(&[
979 landed.commit.as_str().into(),
980 optional(&landed.previous),
981 a.actor.username.as_str().into(),
982 now.as_str().into(),
983 now.as_str().into(),
984 pull.id.as_str().into(),
985 ])?,
986 ];
987 if let Some(issue) = &issue {
988 statements.push(
989 self.db
990 .prepare(
991 "UPDATE issues
992 SET state = 'closed', reason = 'completed', resolved_by = ?,
993 closed_at = ?, updated_at = ?
994 WHERE id = ?",
995 )
996 .bind(&[
997 pull.number.into(),
998 now.as_str().into(),
999 now.as_str().into(),
1000 issue.id.as_str().into(),
1001 ])?,
1002 );
1003 statements.push(
1004 self.db
1005 .prepare(
1006 "UPDATE pulls SET status = 'closed', superseded_by = ?, updated_at = ?
1007 WHERE issue_id = ? AND id != ? AND status IN ('draft', 'open')",
1008 )
1009 .bind(&[
1010 pull.number.into(),
1011 now.as_str().into(),
1012 issue.id.as_str().into(),
1013 pull.id.as_str().into(),
1014 ])?,
1015 );
1016 }
1017 self.db.batch(statements).await?;
1018
1019 self.publish(
1020 "pull.merged",
1021 &repo.id,
1022 &a.actor,
1023 PullEvent {
1024 commit: Some(landed.commit.clone()),
1025 ..Self::pull_event(&pull)
1026 },
1027 )
1028 .await?;
1029 if let Some(issue) = &issue {
1030 self.publish(
1031 "issue.closed",
1032 &repo.id,
1033 &a.actor,
1034 IssueEvent {
1035 reason: Some(IssueReason::Completed.as_str()),
1036 resolved_by: Some(pull.number),
1037 ..Self::issue_event(issue)
1038 },
1039 )
1040 .await?;
1041 }
1042
1043 pull.status = PullStatus::Merged;
1044 pull.head_commit = Some(landed.commit);
1045 pull.merge_base = landed.previous;
1046 pull.merged_by = Some(a.actor.username);
1047 pull.merged_at = Some(now.clone());
1048 pull.updated_at = now;
1049 Ok(Outcome::Ok(pull))
1050 }
1051
1052 async fn list_active_pulls(&self, a: ViewerArgs) -> Result<Vec<ActivePull>> {
1053 let Some(viewer) = a.viewer else {
1054 return Ok(Vec::new());
1055 };
1056 let rows = self
1057 .db
1058 .prepare(
1059 "SELECT * FROM pulls
1060 WHERE author_id = ? AND status IN ('draft', 'open')
1061 ORDER BY updated_at DESC LIMIT 50",
1062 )
1063 .bind(&[viewer.id.into()])?
1064 .all()
1065 .await?
1066 .results::<PullRow>()?;
1067 let mut active = Vec::with_capacity(rows.len());
1068 for pull in rows.into_iter().map(Pull::from) {
1069 let issue = match pull.issue {
1070 Some(number) => self.issue(&pull.repo_id, number).await?,
1071 None => None,
1072 };
1073 active.push(ActivePull { pull, issue });
1074 }
1075 Ok(active)
1076 }
1077
1078 // --- Sessions ----------------------------------------------------------
1079
1080 async fn append_session(&self, a: AppendSessionArgs) -> Result<Outcome<Appended>> {
1081 if a.entries.is_empty() {
1082 return Ok(Outcome::Ok(Appended { count: 0 }));
1083 }
1084 if a.entries.len() > MAX_ENTRY_BATCH {
1085 return Ok(Outcome::fail(
1086 FailureCode::Invalid,
1087 format!("Send at most {MAX_ENTRY_BATCH} entries at a time."),
1088 ));
1089 }
1090 let viewer = Some(a.actor.clone());
1091 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
1092 if pull.author.id != a.actor.id {
1093 return Ok(Outcome::fail(
1094 FailureCode::Forbidden,
1095 "Only whoever opened a pull request can record its session.",
1096 ));
1097 }
1098
1099 let now = rfc3339(now_ms());
1100 let count = a.entries.len() as u32;
1101 let mut statements = Vec::with_capacity(a.entries.len() + 1);
1102 for entry in a.entries {
1103 let kind = serde_json::to_value(entry.kind)?;
1104 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
1105 // Each insert takes the next sequence number itself, so two
1106 // writers appending at once cannot collide.
1107 statements.push(
1108 self.db
1109 .prepare(
1110 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
1111 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
1112 FROM session_entries WHERE pull_id = ?",
1113 )
1114 .bind(&[
1115 pull.id.as_str().into(),
1116 kind.as_str().unwrap_or("note").into(),
1117 text.into(),
1118 optional(&entry.tool),
1119 optional(&entry.commit.or_else(|| pull.head_commit.clone())),
1120 now.as_str().into(),
1121 pull.id.as_str().into(),
1122 ])?,
1123 );
1124 }
1125 statements.push(
1126 self.db
1127 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
1128 .bind(&[now.as_str().into(), pull.id.as_str().into()])?,
1129 );
1130 self.db.batch(statements).await?;
1131 self.publish(
1132 "session.appended",
1133 &pull.repo_id,
1134 &a.actor,
1135 SessionAppended {
1136 pull_id: pull.id.clone(),
1137 repo_id: pull.repo_id.clone(),
1138 number: pull.number,
1139 count,
1140 },
1141 )
1142 .await?;
1143 Ok(Outcome::Ok(Appended { count }))
1144 }
1145
1146 async fn read_session(&self, a: ViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
1147 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
1148 let rows = self
1149 .db
1150 .prepare(
1151 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
1152 WHERE pull_id = ? AND seq > ? ORDER BY seq LIMIT ?",
1153 )
1154 .bind(&[pull.id.into(), a.after_seq.into(), SESSION_PAGE.into()])?
1155 .all()
1156 .await?
1157 .results::<SessionRow>()?;
1158 Ok(Outcome::Ok(
1159 rows.into_iter().map(SessionEntry::from).collect(),
1160 ))
1161 }
1162
1163 /// A push moves the head of the pull request it concerns: the one whose
1164 /// fork was pushed to, or the one opened from the branch that moved.
1165 async fn on_event(&self, event: &Event) -> Result<()> {
1166 if event.kind != "git.push" {
1167 return Ok(());
1168 }
1169 let (Some(repo_id), Some(after), Some(git_ref)) = (
1170 event.repo_id.as_deref(),
1171 event.data["after"].as_str(),
1172 event.data["ref"].as_str(),
1173 ) else {
1174 return Ok(());
1175 };
1176 let now = rfc3339(now_ms());
1177 // The head moved, so whatever the checks said no longer applies.
1178 let moved = "UPDATE pulls
1179 SET head_commit = ?, updated_at = ?, check_status = NULL, check_run_id = NULL";
1180 let active = "status IN ('draft', 'open') AND head_commit IS NOT ?";
1181 let returning = "RETURNING id, repo_id, number, issue_number, status";
1182 let mut pulls: Vec<MovedRow> = Vec::new();
1183 // A fork carries its pull request on its default branch.
1184 if event.data["defaultBranch"].as_bool() == Some(true) {
1185 pulls.extend(
1186 self.db
1187 .prepare(format!(
1188 "{moved} WHERE fork_repo_id = ? AND {active} {returning}"
1189 ))
1190 .bind(&[
1191 after.into(),
1192 now.as_str().into(),
1193 repo_id.into(),
1194 after.into(),
1195 ])?
1196 .all()
1197 .await?
1198 .results::<MovedRow>()?,
1199 );
1200 }
1201 if let Some(branch) = git_ref.strip_prefix("refs/heads/") {
1202 pulls.extend(
1203 self.db
1204 .prepare(format!(
1205 "{moved} WHERE repo_id = ? AND source_branch = ? AND {active} {returning}"
1206 ))
1207 .bind(&[
1208 after.into(),
1209 now.as_str().into(),
1210 repo_id.into(),
1211 branch.into(),
1212 after.into(),
1213 ])?
1214 .all()
1215 .await?
1216 .results::<MovedRow>()?,
1217 );
1218 }
1219 // A draft is announced when it is marked ready instead.
1220 for pull in pulls
1221 .into_iter()
1222 .filter(|pull| pull.status == PullStatus::Open)
1223 {
1224 self.publish_as(
1225 "pull.updated",
1226 &pull.repo_id,
1227 event.actor.clone(),
1228 PullEvent {
1229 pull_id: pull.id,
1230 repo_id: pull.repo_id.clone(),
1231 number: pull.number,
1232 issue: pull.issue_number,
1233 commit: Some(after.to_owned()),
1234 ..PullEvent::default()
1235 },
1236 )
1237 .await?;
1238 }
1239 Ok(())
1240 }
1241}
1242
1243fn service(env: &Env) -> Result<Work> {
1244 Ok(Work {
1245 db: env.d1("DB")?,
1246 repos: env.service("REPOS")?,
1247 events: env.service("EVENTS")?,
1248 })
1249}
1250
1251#[event(fetch)]
1252async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
1253 let Some(method) = rpc_method(&request) else {
1254 return Response::error("Not found", 404);
1255 };
1256 let body: serde_json::Value = request.json().await?;
1257 let work = service(&env)?;
1258
1259 match method.as_str() {
1260 "open_issue" => reply(&work.open_issue(args(body)?).await?),
1261 "list_issues" => reply(&work.list_issues(args(body)?).await?),
1262 "get_issue" => reply(&work.get_issue(args(body)?).await?),
1263 "update_issue" => reply(&work.update_issue(args(body)?).await?),
1264 "close_issue" => reply(&work.close_issue(args(body)?).await?),
1265 "reopen_issue" => reply(&work.reopen_issue(args(body)?).await?),
1266 "list_labels" => reply(&work.list_labels(args(body)?).await?),
1267 "counts" => reply(&work.counts(args(body)?).await?),
1268 "add_comment" => reply(&work.add_comment(args(body)?).await?),
1269 "start_checks" => reply(&work.start_checks(args(body)?).await?),
1270 "report_checks" => reply(&work.report_checks(args(body)?).await?),
1271 "open_pull" => reply(&work.open_pull(args(body)?).await?),
1272 "list_pulls" => reply(&work.list_pulls(args(body)?).await?),
1273 "get_pull" => reply(&work.get_pull(args(body)?).await?),
1274 "ready_pull" => reply(&work.ready_pull(args(body)?).await?),
1275 "close_pull" => reply(&work.close_pull(args(body)?).await?),
1276 "merge_pull" => reply(&work.merge_pull(args(body)?).await?),
1277 "list_active_pulls" => reply(&work.list_active_pulls(args(body)?).await?),
1278 "append_session" => reply(&work.append_session(args(body)?).await?),
1279 "read_session" => reply(&work.read_session(args(body)?).await?),
1280 _ => Response::error("Unknown method", 404),
1281 }
1282}
1283
1284/// Events from the bus, delivered on this service's own queue.
1285#[event(queue)]
1286async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
1287 let work = service(&env)?;
1288 for message in batch.messages()? {
1289 work.on_event(message.body()).await?;
1290 message.ack();
1291 }
1292 Ok(())
1293}