g1t/services/work/src/lib.rs

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