g1t/services/work/src/lib.rs

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