g1t/services/work/src/lib.rs

1,204 lines43,649 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, 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, 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::{BranchRow, 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 branch = a
640 .branch
641 .as_deref()
642 .map(str::trim)
643 .filter(|branch| !branch.is_empty());
644 // The change is on a branch already pushed to the repository, or
645 // will be made in a fork created for this pull request.
646 let (fork, head) = match branch {
647 Some(branch) => {
648 if branch == repo.default_branch {
649 return Ok(Outcome::fail(
650 FailureCode::Invalid,
651 format!("Choose a branch other than {branch}."),
652 ));
653 }
654 let head: Option<String> = g1t_kit::call(
655 &self.repos,
656 "head",
657 &HeadArgs {
658 repo_id: repo.id.clone(),
659 branch: branch.to_owned(),
660 },
661 )
662 .await?;
663 let Some(head) = head else {
664 return Ok(Outcome::fail(
665 FailureCode::NotFound,
666 format!("There is no branch named {branch}. Push it first."),
667 ));
668 };
669 let existing = self
670 .db
671 .prepare(
672 "SELECT number AS n FROM pulls
673 WHERE repo_id = ? AND source_branch = ? AND status IN ('draft', 'open')",
674 )
675 .bind(&[repo.id.as_str().into(), branch.into()])?
676 .first::<NumberRow>(None)
677 .await?;
678 if let Some(existing) = existing {
679 return Ok(Outcome::fail(
680 FailureCode::Conflict,
681 format!("Pull request #{} is already open for {branch}.", existing.n),
682 ));
683 }
684 (None, Some(head))
685 }
686 None => {
687 let fork: Outcome<Repo> = g1t_kit::call(
688 &self.repos,
689 "fork_for_pull",
690 &ForkArgs {
691 source_id: repo.id.clone(),
692 pull_id: id.clone(),
693 actor: a.actor.clone(),
694 },
695 )
696 .await?;
697 (Some(check!(fork)), None)
698 }
699 };
700 // A branch already holds the work, so its pull request is ready for
701 // review from the start; one with a fork starts as a draft.
702 let status = if branch.is_some() { "open" } else { "draft" };
703 let body = Some(a.body.trim().to_owned()).filter(|body| !body.is_empty());
704
705 let number = self.next_number(&repo.id).await?;
706 let timestamp = rfc3339(now);
707 self.db
708 .prepare(
709 "INSERT INTO pulls
710 (id, repo_id, number, issue_id, issue_number, title, body, agent, runtime,
711 status, fork_repo_id, fork_namespace, fork_name, source_branch, head_commit,
712 author_id, author_name, created_at, updated_at)
713 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
714 )
715 .bind(&[
716 id.as_str().into(),
717 repo.id.as_str().into(),
718 number.into(),
719 optional(&issue.as_ref().map(|issue| issue.id.clone())),
720 optional_number(issue.as_ref().map(|issue| issue.number)),
721 title.into(),
722 optional(&body),
723 agent.into(),
724 runtime.into(),
725 status.into(),
726 optional(&fork.as_ref().map(|fork| fork.id.clone())),
727 optional(&fork.as_ref().map(|fork| fork.namespace.clone())),
728 optional(&fork.as_ref().map(|fork| fork.name.clone())),
729 optional(&branch.map(str::to_owned)),
730 optional(&head),
731 a.actor.id.as_str().into(),
732 a.actor.username.as_str().into(),
733 timestamp.as_str().into(),
734 timestamp.as_str().into(),
735 ])?
736 .run()
737 .await?;
738 let Some(pull) = self.pull(&repo.id, number).await? else {
739 return Ok(no_pull());
740 };
741 self.publish(
742 "pull.opened",
743 &repo.id,
744 &a.actor,
745 PullEvent {
746 agent: Some(pull.agent.clone()),
747 ..Self::pull_event(&pull)
748 },
749 )
750 .await?;
751 Ok(Outcome::Ok(pull))
752 }
753
754 async fn list_pulls(&self, a: ListPullsArgs) -> Result<Outcome<Vec<Pull>>> {
755 let repo = check!(self.repo(&a.repo, &a.viewer).await?);
756 let filter = match a.state {
757 Some(State::Open) => "AND status IN ('draft', 'open')",
758 Some(State::Closed) => "AND status IN ('merged', 'closed')",
759 None => "",
760 };
761 let rows = self
762 .db
763 .prepare(format!(
764 "SELECT * FROM pulls WHERE repo_id = ? {filter} ORDER BY number DESC LIMIT ?"
765 ))
766 .bind(&[repo.id.into(), LIST_PAGE.into()])?
767 .all()
768 .await?
769 .results::<PullRow>()?;
770 Ok(Outcome::Ok(rows.into_iter().map(Pull::from).collect()))
771 }
772
773 async fn get_pull(&self, a: ViewArgs) -> Result<Outcome<PullDetail>> {
774 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
775 let issue = match pull.issue {
776 Some(number) => self.issue(&repo.id, number).await?,
777 None => None,
778 };
779 Ok(Outcome::Ok(PullDetail {
780 comments: self.comments(&repo.id, pull.number).await?,
781 issue,
782 pull,
783 }))
784 }
785
786 /// The pull request, if it is still active and `actor` opened it or
787 /// belongs to the repository's workspace.
788 async fn manageable_pull(
789 &self,
790 actor: &User,
791 path: &RepoPath,
792 number: u32,
793 ) -> Result<Outcome<Pull>> {
794 let (repo, pull) = check!(self.pull_at(path, number, &Some(actor.clone())).await?);
795 if pull.author.id != actor.id && !actor.is_member(&repo.namespace) {
796 return Ok(Outcome::fail(
797 FailureCode::Forbidden,
798 "Only whoever opened a pull request, or a member of the workspace, can change it.",
799 ));
800 }
801 if !pull.status.is_active() {
802 return Ok(Outcome::fail(
803 FailureCode::Conflict,
804 format!("This pull request is already {}.", pull.status.as_str()),
805 ));
806 }
807 Ok(Outcome::Ok(pull))
808 }
809
810 /// Marks a draft ready for review, or updates the description of one
811 /// that already is.
812 async fn ready_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
813 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
814 let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
815 let now = rfc3339(now_ms());
816 self.db
817 .prepare(
818 "UPDATE pulls SET status = 'open', body = COALESCE(?, body), updated_at = ?
819 WHERE id = ?",
820 )
821 .bind(&[
822 optional(&summary),
823 now.as_str().into(),
824 pull.id.as_str().into(),
825 ])?
826 .run()
827 .await?;
828 if pull.status == PullStatus::Draft {
829 self.publish(
830 "pull.ready",
831 &pull.repo_id,
832 &a.actor,
833 Self::pull_event(&pull),
834 )
835 .await?;
836 }
837 pull.status = PullStatus::Open;
838 pull.body = summary.or(pull.body);
839 pull.updated_at = now;
840 Ok(Outcome::Ok(pull))
841 }
842
843 async fn close_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
844 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
845 let now = rfc3339(now_ms());
846 self.db
847 .prepare("UPDATE pulls SET status = 'closed', updated_at = ? WHERE id = ?")
848 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
849 .run()
850 .await?;
851 self.publish(
852 "pull.closed",
853 &pull.repo_id,
854 &a.actor,
855 Self::pull_event(&pull),
856 )
857 .await?;
858 pull.status = PullStatus::Closed;
859 pull.updated_at = now;
860 Ok(Outcome::Ok(pull))
861 }
862
863 /// Lands the pull request on the repository's default branch. Unless
864 /// told to keep it open, that resolves the issue it was for: the issue
865 /// closes naming this pull request, and the others still in progress
866 /// for it close as superseded.
867 async fn merge_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
868 let viewer = Some(a.actor.clone());
869 let (repo, mut pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
870 match pull.status {
871 PullStatus::Open => {}
872 PullStatus::Draft => {
873 return Ok(Outcome::fail(
874 FailureCode::Conflict,
875 "This pull request is still a draft. Mark it ready for review first.",
876 ));
877 }
878 status => {
879 return Ok(Outcome::fail(
880 FailureCode::Conflict,
881 format!("This pull request is already {}.", status.as_str()),
882 ));
883 }
884 }
885 let issue = match pull.issue {
886 Some(number) if !a.keep_issue_open => self
887 .issue(&repo.id, number)
888 .await?
889 .filter(|issue| issue.state == State::Open),
890 _ => None,
891 };
892
893 // Whether the actor may write to the repository is decided by repos.
894 let landed: Outcome<Landed> = g1t_kit::call(
895 &self.repos,
896 "land",
897 &LandArgs {
898 // A pull request from a branch lands from the repository itself.
899 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
900 branch: pull.branch.clone(),
901 actor: a.actor.clone(),
902 },
903 )
904 .await?;
905 let landed = check!(landed);
906
907 let now = rfc3339(now_ms());
908 let mut statements = vec![
909 self.db
910 .prepare(
911 "UPDATE pulls
912 SET status = 'merged', head_commit = ?, merge_base = ?, merged_by = ?,
913 merged_at = ?, updated_at = ?
914 WHERE id = ?",
915 )
916 .bind(&[
917 landed.commit.as_str().into(),
918 optional(&landed.previous),
919 a.actor.username.as_str().into(),
920 now.as_str().into(),
921 now.as_str().into(),
922 pull.id.as_str().into(),
923 ])?,
924 ];
925 if let Some(issue) = &issue {
926 statements.push(
927 self.db
928 .prepare(
929 "UPDATE issues
930 SET state = 'closed', reason = 'completed', resolved_by = ?,
931 closed_at = ?, updated_at = ?
932 WHERE id = ?",
933 )
934 .bind(&[
935 pull.number.into(),
936 now.as_str().into(),
937 now.as_str().into(),
938 issue.id.as_str().into(),
939 ])?,
940 );
941 statements.push(
942 self.db
943 .prepare(
944 "UPDATE pulls SET status = 'closed', superseded_by = ?, updated_at = ?
945 WHERE issue_id = ? AND id != ? AND status IN ('draft', 'open')",
946 )
947 .bind(&[
948 pull.number.into(),
949 now.as_str().into(),
950 issue.id.as_str().into(),
951 pull.id.as_str().into(),
952 ])?,
953 );
954 }
955 self.db.batch(statements).await?;
956
957 self.publish(
958 "pull.merged",
959 &repo.id,
960 &a.actor,
961 PullEvent {
962 commit: Some(landed.commit.clone()),
963 ..Self::pull_event(&pull)
964 },
965 )
966 .await?;
967 if let Some(issue) = &issue {
968 self.publish(
969 "issue.closed",
970 &repo.id,
971 &a.actor,
972 IssueEvent {
973 reason: Some(IssueReason::Completed.as_str()),
974 resolved_by: Some(pull.number),
975 ..Self::issue_event(issue)
976 },
977 )
978 .await?;
979 }
980
981 pull.status = PullStatus::Merged;
982 pull.head_commit = Some(landed.commit);
983 pull.merge_base = landed.previous;
984 pull.merged_by = Some(a.actor.username);
985 pull.merged_at = Some(now.clone());
986 pull.updated_at = now;
987 Ok(Outcome::Ok(pull))
988 }
989
990 async fn list_active_pulls(&self, a: ViewerArgs) -> Result<Vec<ActivePull>> {
991 let Some(viewer) = a.viewer else {
992 return Ok(Vec::new());
993 };
994 let rows = self
995 .db
996 .prepare(
997 "SELECT * FROM pulls
998 WHERE author_id = ? AND status IN ('draft', 'open')
999 ORDER BY updated_at DESC LIMIT 50",
1000 )
1001 .bind(&[viewer.id.into()])?
1002 .all()
1003 .await?
1004 .results::<PullRow>()?;
1005 let mut active = Vec::with_capacity(rows.len());
1006 for pull in rows.into_iter().map(Pull::from) {
1007 let issue = match pull.issue {
1008 Some(number) => self.issue(&pull.repo_id, number).await?,
1009 None => None,
1010 };
1011 active.push(ActivePull { pull, issue });
1012 }
1013 Ok(active)
1014 }
1015
1016 // --- Sessions ----------------------------------------------------------
1017
1018 async fn append_session(&self, a: AppendSessionArgs) -> Result<Outcome<Appended>> {
1019 if a.entries.is_empty() {
1020 return Ok(Outcome::Ok(Appended { count: 0 }));
1021 }
1022 if a.entries.len() > MAX_ENTRY_BATCH {
1023 return Ok(Outcome::fail(
1024 FailureCode::Invalid,
1025 format!("Send at most {MAX_ENTRY_BATCH} entries at a time."),
1026 ));
1027 }
1028 let viewer = Some(a.actor.clone());
1029 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
1030 if pull.author.id != a.actor.id {
1031 return Ok(Outcome::fail(
1032 FailureCode::Forbidden,
1033 "Only whoever opened a pull request can record its session.",
1034 ));
1035 }
1036
1037 let now = rfc3339(now_ms());
1038 let count = a.entries.len() as u32;
1039 let mut statements = Vec::with_capacity(a.entries.len() + 1);
1040 for entry in a.entries {
1041 let kind = serde_json::to_value(entry.kind)?;
1042 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
1043 // Each insert takes the next sequence number itself, so two
1044 // writers appending at once cannot collide.
1045 statements.push(
1046 self.db
1047 .prepare(
1048 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
1049 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
1050 FROM session_entries WHERE pull_id = ?",
1051 )
1052 .bind(&[
1053 pull.id.as_str().into(),
1054 kind.as_str().unwrap_or("note").into(),
1055 text.into(),
1056 optional(&entry.tool),
1057 optional(&entry.commit.or_else(|| pull.head_commit.clone())),
1058 now.as_str().into(),
1059 pull.id.as_str().into(),
1060 ])?,
1061 );
1062 }
1063 statements.push(
1064 self.db
1065 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
1066 .bind(&[now.as_str().into(), pull.id.as_str().into()])?,
1067 );
1068 self.db.batch(statements).await?;
1069 self.publish(
1070 "session.appended",
1071 &pull.repo_id,
1072 &a.actor,
1073 SessionAppended {
1074 pull_id: pull.id.clone(),
1075 repo_id: pull.repo_id.clone(),
1076 number: pull.number,
1077 count,
1078 },
1079 )
1080 .await?;
1081 Ok(Outcome::Ok(Appended { count }))
1082 }
1083
1084 async fn read_session(&self, a: ViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
1085 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
1086 let rows = self
1087 .db
1088 .prepare(
1089 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
1090 WHERE pull_id = ? AND seq > ? ORDER BY seq LIMIT ?",
1091 )
1092 .bind(&[pull.id.into(), a.after_seq.into(), SESSION_PAGE.into()])?
1093 .all()
1094 .await?
1095 .results::<SessionRow>()?;
1096 Ok(Outcome::Ok(
1097 rows.into_iter().map(SessionEntry::from).collect(),
1098 ))
1099 }
1100
1101 /// A push moves the head of the pull requests it concerns: the one
1102 /// whose fork was pushed to, or those from branches of the repository
1103 /// that was.
1104 async fn on_event(&self, event: &Delivered) -> Result<()> {
1105 if event.kind != "git.push" {
1106 return Ok(());
1107 }
1108 let (Some(repo_id), Some(after)) = (event.repo_id.as_deref(), event.data["after"].as_str())
1109 else {
1110 return Ok(());
1111 };
1112 let now = rfc3339(now_ms());
1113 self.db
1114 .prepare(
1115 "UPDATE pulls SET head_commit = ?, updated_at = ?
1116 WHERE fork_repo_id = ? AND status IN ('draft', 'open')",
1117 )
1118 .bind(&[after.into(), now.as_str().into(), repo_id.into()])?
1119 .run()
1120 .await?;
1121
1122 // The event does not say which branch moved, so each open pull
1123 // request from a branch of this repository is checked.
1124 let from_branches = self
1125 .db
1126 .prepare(
1127 "SELECT id, source_branch, head_commit FROM pulls
1128 WHERE repo_id = ? AND source_branch IS NOT NULL AND status IN ('draft', 'open')",
1129 )
1130 .bind(&[repo_id.into()])?
1131 .all()
1132 .await?
1133 .results::<BranchRow>()?;
1134 for row in from_branches {
1135 let head: Option<String> = g1t_kit::call(
1136 &self.repos,
1137 "head",
1138 &HeadArgs {
1139 repo_id: repo_id.to_owned(),
1140 branch: row.source_branch,
1141 },
1142 )
1143 .await?;
1144 if let Some(head) = head.filter(|head| Some(head) != row.head_commit.as_ref()) {
1145 self.db
1146 .prepare("UPDATE pulls SET head_commit = ?, updated_at = ? WHERE id = ?")
1147 .bind(&[head.into(), now.as_str().into(), row.id.into()])?
1148 .run()
1149 .await?;
1150 }
1151 }
1152 Ok(())
1153 }
1154}
1155
1156fn service(env: &Env) -> Result<Work> {
1157 Ok(Work {
1158 db: env.d1("DB")?,
1159 repos: env.service("REPOS")?,
1160 events: js::binding(env, "EVENTS")?,
1161 })
1162}
1163
1164#[event(fetch)]
1165async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
1166 let Some(method) = rpc_method(&request) else {
1167 return Response::error("Not found", 404);
1168 };
1169 let body: serde_json::Value = request.json().await?;
1170 let work = service(&env)?;
1171
1172 match method.as_str() {
1173 "open_issue" => reply(&work.open_issue(args(body)?).await?),
1174 "list_issues" => reply(&work.list_issues(args(body)?).await?),
1175 "get_issue" => reply(&work.get_issue(args(body)?).await?),
1176 "update_issue" => reply(&work.update_issue(args(body)?).await?),
1177 "close_issue" => reply(&work.close_issue(args(body)?).await?),
1178 "reopen_issue" => reply(&work.reopen_issue(args(body)?).await?),
1179 "list_labels" => reply(&work.list_labels(args(body)?).await?),
1180 "counts" => reply(&work.counts(args(body)?).await?),
1181 "add_comment" => reply(&work.add_comment(args(body)?).await?),
1182 "open_pull" => reply(&work.open_pull(args(body)?).await?),
1183 "list_pulls" => reply(&work.list_pulls(args(body)?).await?),
1184 "get_pull" => reply(&work.get_pull(args(body)?).await?),
1185 "ready_pull" => reply(&work.ready_pull(args(body)?).await?),
1186 "close_pull" => reply(&work.close_pull(args(body)?).await?),
1187 "merge_pull" => reply(&work.merge_pull(args(body)?).await?),
1188 "list_active_pulls" => reply(&work.list_active_pulls(args(body)?).await?),
1189 "append_session" => reply(&work.append_session(args(body)?).await?),
1190 "read_session" => reply(&work.read_session(args(body)?).await?),
1191 _ => Response::error("Unknown method", 404),
1192 }
1193}
1194
1195/// Events from the bus, delivered on this service's own queue.
1196#[event(queue)]
1197async fn queue(batch: MessageBatch<Delivered>, env: Env, _ctx: Context) -> Result<()> {
1198 let work = service(&env)?;
1199 for message in batch.messages()? {
1200 work.on_event(message.body()).await?;
1201 message.ack();
1202 }
1203 Ok(())
1204}