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