flagon-io/g1t

public

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

g1t/services/work/src/lib.rs

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