Skip to content

g1t/services/work/src/lib.rs

2,656 lines111,485 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 codeowners;
11mod commit_checks;
12mod compute;
13mod confidence;
14mod guardrails;
15mod inbox;
16mod labels;
17mod lifecycle;
18mod memory;
19mod mentions;
20mod mergeability;
21mod milestones;
22mod plans;
23mod messages;
24mod prefetch;
25mod queue;
26mod retired;
27mod reviews;
28mod rows;
29mod rulesets;
30mod runs;
31mod settings;
32mod statuses;
33mod team_reviews;
34
35use g1t_contracts::events::{
36 CommentCreated, Event, IssueEvent, NewEvent, Publish, PullEvent, SessionAppended,
37};
38use g1t_contracts::identity::UsernameArgs;
39use g1t_contracts::repos::{
40 ForkArgs, GetArgs, HeadArgs, LandArgs, Landed, NeedsAgentReason, PullBranchUpdate, ReadableArgs, Repo, RepoPath,
41 UpdatePullBranchArgs,
42};
43use g1t_contracts::access::{self, Capability, Denied};
44use g1t_contracts::time::rfc3339;
45use g1t_contracts::work::*;
46use futures_util::future::{try_join, try_join3, try_join_all};
47use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
48use g1t_kit::{args, now_ms, reply, rpc_method};
49use serde::Serialize;
50use worker::wasm_bindgen::JsValue;
51use worker::{
52 Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event,
53};
54
55use retired::writable;
56use rows::{CommentRow, IssueRow, MovedRow, NumberRow, PULL_COLUMNS, PullRow, SessionRow, Snapshot};
57
58pub(crate) const SOURCE: &str = "work";
59const MAX_ENTRY_BATCH: usize = 200;
60const MAX_ENTRY_CHARS: usize = 64_000;
61const MAX_TITLE_CHARS: usize = 200;
62const SESSION_PAGE: u32 = 500;
63const LIST_PAGE: u32 = 100;
64const MAX_ASSIGNEES: usize = 10;
65pub(crate) const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
66
67const ISSUE_COLUMNS: &str = "issues.*,
68 (SELECT count(*) FROM pulls WHERE pulls.issue_id = issues.id) AS pull_count,
69 (SELECT agent FROM pulls
70 WHERE pulls.issue_id = issues.id AND pulls.status IN ('draft', 'open')
71 AND pulls.fork_repo_id IS NOT NULL
72 ORDER BY pulls.number DESC LIMIT 1) AS agent,
73 (SELECT count(*) FROM comments
74 WHERE comments.repo_id = issues.repo_id AND comments.number = issues.number
75 AND comments.kind = 'comment') AS comment_count,
76 (SELECT title FROM milestones
77 WHERE milestones.repo_id = issues.repo_id AND milestones.number = issues.milestone) AS milestone_title";
78
79fn no_issue<T>() -> Outcome<T> {
80 Outcome::fail(FailureCode::NotFound, "Issue not found.")
81}
82
83fn no_pull<T>() -> Outcome<T> {
84 Outcome::fail(FailureCode::NotFound, "Pull request not found.")
85}
86
87/// Whether a pull request was behind when its mergeability was last
88/// worked out, and for which pair of commits (mergeability.rs).
89#[derive(serde::Deserialize)]
90struct StoredBehind {
91 #[serde(default)]
92 behind: Option<u8>,
93 #[serde(default)]
94 mergeable_key: Option<String>,
95}
96
97impl StoredBehind {
98 /// The stored answer, if it was worked out for `head`.
99 fn for_head(&self, head: Option<&str>) -> Option<bool> {
100 let (worked_for, _) = self.mergeable_key.as_deref()?.split_once("..")?;
101 (Some(worked_for) == head).then_some(self.behind? != 0)
102 }
103}
104
105/// Refuses `actor` unless their role on `repo` has `capability`: not found
106/// when they cannot read it, forbidden with the role it needs otherwise.
107pub(crate) fn allowed(actor: Option<&User>, repo: &Repo, capability: Capability) -> Outcome<()> {
108 match access::check(actor, repo, capability) {
109 Ok(()) => Outcome::Ok(()),
110 Err(Denied::NotFound) => Outcome::fail(FailureCode::NotFound, "Repository not found."),
111 Err(Denied::Forbidden) => Outcome::fail(
112 FailureCode::Forbidden,
113 access::needs(capability, &format!("{}/{}", repo.namespace, repo.name)),
114 ),
115 }
116}
117
118fn optional(value: &Option<String>) -> JsValue {
119 value.as_deref().map_or(JsValue::NULL, JsValue::from)
120}
121
122fn optional_number(value: Option<u32>) -> JsValue {
123 value.map_or(JsValue::NULL, JsValue::from)
124}
125
126/// Names the branch a pull request merges into when it is the default
127/// branch, which is stored as none so that it follows a change of default.
128pub(crate) fn fill_base(pull: &mut Pull, repo: &Repo) {
129 if pull.base.as_deref().is_none_or(str::is_empty) {
130 pull.base = Some(repo.default_branch.clone());
131 }
132}
133
134/// The branch a pull request is stored as merging into: none for the
135/// default branch.
136fn stored_base(base: &str, repo: &Repo) -> Option<String> {
137 let base = base.trim();
138 (!base.is_empty() && base != repo.default_branch).then(|| base.to_owned())
139}
140
141/// The lowercase name a `State` is stored and sent as.
142fn state_name(state: Option<State>) -> Option<&'static str> {
143 state.map(|state| match state {
144 State::Open => "open",
145 State::Closed => "closed",
146 })
147}
148
149/// A trimmed title, or why it cannot be used.
150fn valid_title(title: &str) -> std::result::Result<&str, &'static str> {
151 let title = title.trim();
152 if title.is_empty() {
153 Err("A title is required.")
154 } else if title.chars().count() > MAX_TITLE_CHARS {
155 Err("That title is too long.")
156 } else {
157 Ok(title)
158 }
159}
160
161/// Unwraps an `Outcome`, returning its failure from the enclosing method.
162macro_rules! check {
163 ($outcome:expr) => {
164 match $outcome {
165 Outcome::Ok(value) => value,
166 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
167 }
168 };
169}
170
171struct Work {
172 db: D1Database,
173 identity: Fetcher,
174 repos: Fetcher,
175 events: Fetcher,
176 /// GitHub Actions: runs a merge queue's `merge_group` workflows.
177 actions: Fetcher,
178 /// Where this request's time went, for its `Server-Timing`.
179 timing: g1t_kit::d1::Timing,
180 /// A pull request's rows read in one batch for this request
181 /// (prefetch.rs), which the helpers below read instead of the database.
182 prefetched: std::cell::RefCell<Option<std::rc::Rc<prefetch::Prefetched>>>,
183 /// Repositories read in this request, by id, for rules that need one
184 /// where only a pull request is at hand (rulesets.rs `repo_for`).
185 known_repos: std::cell::RefCell<std::collections::HashMap<String, Repo>>,
186}
187
188impl Work {
189 async fn publish<T: Serialize>(
190 &self,
191 kind: &'static str,
192 repo_id: &str,
193 actor: &User,
194 data: T,
195 ) -> Result<()> {
196 // What a workflow job's token did is marked, so it starts no
197 // workflows (`g1t_contracts::events::CAUSED_BY_JOB`).
198 self.publish_as(kind, repo_id, Some(actor.id.clone()), g1t_contracts::events::marked(data, Some(actor)))
199 .await
200 }
201
202 /// A pull request's owner (whoever asked g1t for it, or its author) as
203 /// a viewer who can read its repository and source. Stored people carry
204 /// no memberships, so a private repository would otherwise look missing
205 /// to them. The membership given reads and nothing more: it is for
206 /// looking, never for acting.
207 pub(crate) async fn owner_viewer(&self, pull: &Pull) -> Result<Viewer> {
208 let path: Option<RepoPath> = g1t_kit::call(
209 &self.repos,
210 "path_by_id",
211 &g1t_contracts::repos::PathByIdArgs { id: pull.repo_id.clone() },
212 )
213 .await?;
214 let mut owner = pull.owner().clone();
215 if let Some(path) = path
216 && !owner.is_member(&path.namespace.to_lowercase())
217 {
218 owner.workspaces.push(g1t_contracts::Membership {
219 base_permission: Some(access::BasePermission::Read),
220 ..g1t_contracts::Membership::member(path.namespace.to_lowercase())
221 });
222 }
223 Ok(Some(owner))
224 }
225
226 /// Publishes an event caused by `actor`, or by g1t itself.
227 async fn publish_as<T: Serialize>(
228 &self,
229 kind: &'static str,
230 repo_id: &str,
231 actor: Option<String>,
232 data: T,
233 ) -> Result<()> {
234 let event = NewEvent {
235 kind,
236 source: SOURCE,
237 repo_id: Some(repo_id.to_owned()),
238 actor,
239 data,
240 };
241 g1t_kit::call(
242 &self.events,
243 "publish",
244 &Publish {
245 events: vec![event],
246 },
247 )
248 .await
249 }
250
251 /// The repository, if the viewer may see it. Whether they may is
252 /// decided by the repos service.
253 async fn repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Outcome<Repo>> {
254 let found: Outcome<Repo> = self
255 .timing
256 .rpc(g1t_kit::call(
257 &self.repos,
258 "get",
259 &GetArgs {
260 path: path.clone(),
261 viewer: viewer.clone(),
262 },
263 ))
264 .await?;
265 if let Outcome::Ok(repo) = &found {
266 self.known_repos.borrow_mut().insert(repo.id.clone(), repo.clone());
267 }
268 Ok(found)
269 }
270
271 /// The next number in the repository's sequence. Taking it is one
272 /// statement, so concurrent opens cannot be given the same number.
273 async fn next_number(&self, repo_id: &str) -> Result<u32> {
274 let row = self
275 .db
276 .prepare(
277 "INSERT INTO counters (repo_id, last) VALUES (?, 1)
278 ON CONFLICT (repo_id) DO UPDATE SET last = last + 1
279 RETURNING last AS n",
280 )
281 .bind(&[repo_id.into()])?
282 .first::<NumberRow>(None)
283 .await?;
284 row.map(|row| row.n)
285 .ok_or_else(|| worker::Error::RustError("no number was assigned".into()))
286 }
287
288 async fn issue(&self, repo_id: &str, number: u32) -> Result<Option<Issue>> {
289 Ok(self
290 .db
291 .prepare(format!(
292 "SELECT {ISSUE_COLUMNS} FROM issues WHERE repo_id = ? AND number = ?"
293 ))
294 .bind(&[repo_id.into(), number.into()])?
295 .first::<IssueRow>(None)
296 .await?
297 .map(Issue::from))
298 }
299
300 async fn pull(&self, repo_id: &str, number: u32) -> Result<Option<Pull>> {
301 Ok(self
302 .db
303 .prepare(format!("SELECT {PULL_COLUMNS} FROM pulls WHERE repo_id = ? AND number = ?"))
304 .bind(&[repo_id.into(), number.into()])?
305 .first::<PullRow>(None)
306 .await?
307 .map(Pull::from))
308 }
309
310 /// The repository and one of its issues, as seen by `viewer`.
311 async fn issue_at(
312 &self,
313 path: &RepoPath,
314 number: u32,
315 viewer: &Viewer,
316 ) -> Result<Outcome<(Repo, Issue)>> {
317 let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
318 return Ok(no_issue());
319 };
320 Ok(match self.issue(&repo.id, number).await? {
321 Some(issue) => Outcome::Ok((repo, issue)),
322 None => no_issue(),
323 })
324 }
325
326 /// The repository and one of its pull requests, as seen by `viewer`.
327 async fn pull_at(
328 &self,
329 path: &RepoPath,
330 number: u32,
331 viewer: &Viewer,
332 ) -> Result<Outcome<(Repo, Pull)>> {
333 let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
334 return Ok(no_pull());
335 };
336 Ok(match self.pull(&repo.id, number).await? {
337 Some(pull) => Outcome::Ok((repo, pull)),
338 None => no_pull(),
339 })
340 }
341
342 /// Records something that happened to an issue or a pull request, so
343 /// that it shows in the conversation where it happened. `text` is what
344 /// `author` did, as the rest of a sentence starting with their name.
345 pub(crate) async fn note(
346 &self,
347 repo_id: &str,
348 number: u32,
349 author: (&str, &str),
350 text: &str,
351 ) -> Result<()> {
352 let now = now_ms();
353 self.db
354 .prepare(
355 "INSERT INTO comments
356 (id, repo_id, number, author_id, author_name, body, kind, created_at)
357 VALUES (?, ?, ?, ?, ?, ?, 'event', ?)",
358 )
359 .bind(&[
360 new_id("cmt", now).into(),
361 repo_id.into(),
362 number.into(),
363 author.0.into(),
364 author.1.into(),
365 text.into(),
366 rfc3339(now).into(),
367 ])?
368 .run()
369 .await?;
370 Ok(())
371 }
372
373 /// Notes who was added to and removed from a list of people, such as
374 /// "assigned ana" or "requested a review from g1t".
375 async fn note_changes(
376 &self,
377 repo_id: &str,
378 number: u32,
379 actor: &User,
380 before: &[String],
381 after: &[String],
382 (added, removed): (&str, &str),
383 ) -> Result<()> {
384 let joined = |names: Vec<&String>| {
385 names
386 .into_iter()
387 .map(String::as_str)
388 .collect::<Vec<_>>()
389 .join(", ")
390 };
391 let new: Vec<&String> = after.iter().filter(|name| !before.contains(name)).collect();
392 let gone: Vec<&String> = before.iter().filter(|name| !after.contains(name)).collect();
393 let who = (actor.id.as_str(), actor.username.as_str());
394 if !new.is_empty() {
395 // Taking something on oneself reads better said that way.
396 let text = if added == "assigned" && new == [&actor.username] {
397 "self-assigned this".to_owned()
398 } else {
399 format!("{added} {}", joined(new))
400 };
401 self.note(repo_id, number, who, &text).await?;
402 }
403 if !gone.is_empty() {
404 self.note(repo_id, number, who, &format!("{removed} {}", joined(gone)))
405 .await?;
406 }
407 Ok(())
408 }
409
410 fn issue_event(issue: &Issue) -> IssueEvent {
411 IssueEvent {
412 issue_id: issue.id.clone(),
413 repo_id: issue.repo_id.clone(),
414 number: issue.number,
415 author: Some((&issue.author).into()),
416 requested_by: issue.requested_by.as_ref().map(Into::into),
417 ..IssueEvent::default()
418 }
419 }
420
421 /// The commit a pull request's change is at in git right now: its
422 /// fork's default branch, or its branch.
423 async fn live_head(&self, pull: &Pull) -> Result<Option<String>> {
424 g1t_kit::call(
425 &self.repos,
426 "head",
427 &HeadArgs {
428 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| pull.repo_id.clone()),
429 branch: pull.branch.clone().unwrap_or_default(),
430 },
431 )
432 .await
433 }
434
435 fn pull_event(pull: &Pull) -> PullEvent {
436 PullEvent {
437 pull_id: pull.id.clone(),
438 repo_id: pull.repo_id.clone(),
439 number: pull.number,
440 author: Some((&pull.author).into()),
441 requested_by: pull.requested_by.as_ref().map(Into::into),
442 issue: pull.issue,
443 confidence: pull.confidence.clone(),
444 ..PullEvent::default()
445 }
446 }
447
448 // --- Issues ------------------------------------------------------------
449
450 /// Opens an issue for g1t to take at once: refused before
451 /// anything is opened unless the actor may put agents to work here. The
452 /// runner's `delegate` starts the agent on it.
453 async fn delegate_issue(&self, a: DelegateIssueArgs) -> Result<Outcome<Issue>> {
454 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
455 check!(writable(&repo));
456 check!(allowed(Some(&a.actor), &repo, Capability::Run));
457 self.open_issue(OpenIssueArgs {
458 actor: a.actor,
459 repo: a.repo,
460 title: a.title,
461 body: a.body,
462 labels: a.labels,
463 checks: a.checks,
464 milestone: None,
465 })
466 .await
467 }
468
469 async fn open_issue(&self, a: OpenIssueArgs) -> Result<Outcome<Issue>> {
470 if !a.actor.verified {
471 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
472 }
473 let title = match valid_title(&a.title) {
474 Ok(title) => title,
475 Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
476 };
477 let Some(labels) = normalize_labels(&a.labels) else {
478 return Ok(Outcome::fail(
479 FailureCode::Invalid,
480 format!("An issue can have up to {MAX_LABELS} labels of up to {MAX_LABEL_CHARS} characters each."),
481 ));
482 };
483 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
484 check!(writable(&repo));
485 // Labels it does not have yet are made for someone who may triage;
486 // anyone else picks from those there are.
487 let colors = check!(self.ensure_labels(&a.actor, &repo, &labels).await?);
488 let milestone = match a.milestone.filter(|number| *number > 0) {
489 Some(number) => {
490 check!(allowed(Some(&a.actor), &repo, Capability::Triage));
491 check!(self.milestone_ref(&repo.id, number).await?)
492 }
493 None => None,
494 };
495 // Commands given the old way are words for the agent now: added to
496 // the body under "Definition of done". What has to pass to merge is
497 // the branch's required checks.
498 let body = with_definition_of_done(&a.body, &commands_pass(&a.checks));
499
500 let now = now_ms();
501 let id = new_id("iss", now);
502 let number = self.next_number(&repo.id).await?;
503 let timestamp = rfc3339(now);
504 // What g1t's agent files at work is g1t's, for the person it works for.
505 let (author, requested_by) = authorship(&a.actor, false);
506 self.db
507 .prepare(
508 "INSERT INTO issues
509 (id, repo_id, number, title, body, labels, checks, author_id, author_name,
510 requested_by_id, requested_by_name, created_at, updated_at, milestone)
511 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
512 )
513 .bind(&[
514 id.as_str().into(),
515 repo.id.as_str().into(),
516 number.into(),
517 title.into(),
518 body.into(),
519 serde_json::to_string(&labels)?.into(),
520 "[]".into(),
521 author.id.as_str().into(),
522 author.username.as_str().into(),
523 optional(&requested_by.as_ref().map(|user| user.id.clone())),
524 optional(&requested_by.as_ref().map(|user| user.username.clone())),
525 timestamp.as_str().into(),
526 timestamp.as_str().into(),
527 optional_number(milestone.as_ref().map(|milestone| milestone.number)),
528 ])?
529 .run()
530 .await?;
531 let Some(issue) = self.issue(&repo.id, number).await? else {
532 return Ok(no_issue());
533 };
534 self.apply_label_rule(&a.actor, &issue, &[]).await?;
535 self.publish(
536 "issue.opened",
537 &repo.id,
538 &a.actor,
539 IssueEvent {
540 title: Some(issue.title.clone()),
541 ..Self::issue_event(&issue)
542 },
543 )
544 .await?;
545 // Opened with labels and a milestone: each is said, as it would be
546 // if they were added afterwards, without notes in the conversation.
547 for label in &issue.labels {
548 let color = colors.iter().find(|(name, _)| name == label).map_or_else(|| label_color_for(label), |(_, c)| c.clone());
549 let label = Some(g1t_contracts::events::EventLabel { name: label.clone(), color });
550 self.publish("issue.labeled", &repo.id, &a.actor, IssueEvent { label, ..Self::issue_event(&issue) })
551 .await?;
552 }
553 if let Some(milestone) = milestone {
554 self.publish(
555 "issue.milestoned",
556 &repo.id,
557 &a.actor,
558 IssueEvent { milestone: Some(milestone), ..Self::issue_event(&issue) },
559 )
560 .await?;
561 }
562 Ok(Outcome::Ok(issue))
563 }
564
565 async fn list_issues(&self, a: ListIssuesArgs) -> Result<Outcome<Vec<Issue>>> {
566 let state = state_name(a.state);
567 let label = a
568 .label
569 .map(|label| label.split_whitespace().collect::<Vec<_>>().join(" ").to_lowercase())
570 .filter(|label| !label.is_empty());
571 let milestone = a.milestone;
572 let list = |repo_id: String| {
573 let label = label.clone();
574 async move {
575 let state = state.map_or(JsValue::NULL, JsValue::from);
576 let query = self
577 .db
578 .prepare(format!(
579 "SELECT {ISSUE_COLUMNS} FROM issues
580 WHERE repo_id = ? AND (? IS NULL OR state = ?)
581 AND (? IS NULL OR EXISTS
582 (SELECT 1 FROM json_each(issues.labels) WHERE json_each.value = ?))
583 AND (? IS NULL OR milestone = ?)
584 ORDER BY number DESC LIMIT ?"
585 ))
586 .bind(&[
587 repo_id.into(),
588 state.clone(),
589 state,
590 optional(&label),
591 optional(&label),
592 optional_number(milestone),
593 optional_number(milestone),
594 LIST_PAGE.into(),
595 ])?;
596 self.timing.db(1, query.all()).await?.results::<IssueRow>()
597 }
598 };
599 let (_, rows) = check!(self.repo_then(&a.repo, &a.viewer, list).await?);
600 Ok(Outcome::Ok(rows.into_iter().map(Issue::from).collect()))
601 }
602
603 /// An issue, the pull requests for it and its comments: one batch,
604 /// started beside the access check (prefetch.rs).
605 async fn get_issue(&self, a: ViewArgs) -> Result<Outcome<IssueDetail>> {
606 let number = a.number;
607 let read = |repo_id: String| async move {
608 let key = [JsValue::from(repo_id.as_str()), JsValue::from(number)];
609 let statements = vec![
610 self.db
611 .prepare(format!("SELECT {ISSUE_COLUMNS} FROM issues WHERE repo_id = ?1 AND number = ?2"))
612 .bind(&key)?,
613 self.db
614 .prepare(format!(
615 "SELECT {PULL_COLUMNS} FROM pulls
616 WHERE issue_id = (SELECT id FROM issues WHERE repo_id = ?1 AND number = ?2)
617 ORDER BY number"
618 ))
619 .bind(&key)?,
620 self.db
621 .prepare("SELECT * FROM comments WHERE repo_id = ?1 AND number = ?2 ORDER BY id LIMIT 500")
622 .bind(&key)?,
623 ];
624 let results = self.timing.db(3, self.db.batch(statements)).await?;
625 let rows = |index: usize| results.get(index).ok_or_else(|| worker::Error::RustError("short batch".into()));
626 Ok::<_, worker::Error>((
627 rows(0)?.results::<IssueRow>()?.into_iter().next().map(Issue::from),
628 rows(1)?.results::<PullRow>()?.into_iter().map(Pull::from).collect::<Vec<_>>(),
629 rows(2)?.results::<CommentRow>()?.into_iter().map(Comment::from).collect::<Vec<_>>(),
630 ))
631 };
632 let Outcome::Ok((_, (Some(issue), pulls, comments))) = self.repo_then(&a.repo, &a.viewer, read).await? else {
633 return Ok(no_issue());
634 };
635 Ok(Outcome::Ok(IssueDetail { comments, pulls, issue }))
636 }
637
638 /// The issue, if `actor` wrote it or may triage the repository's issues.
639 async fn manageable_issue(
640 &self,
641 actor: &User,
642 path: &RepoPath,
643 number: u32,
644 ) -> Result<Outcome<Issue>> {
645 let (repo, issue) = check!(self.issue_at(path, number, &Some(actor.clone())).await?);
646 check!(writable(&repo));
647 if issue.owner().id != actor.id {
648 check!(allowed(Some(actor), &repo, Capability::Triage));
649 }
650 Ok(Outcome::Ok(issue))
651 }
652
653 async fn update_issue(&self, a: UpdateIssueArgs) -> Result<Outcome<Issue>> {
654 let issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
655 let title = match a.title.as_deref().map(valid_title) {
656 Some(Err(message)) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
657 Some(Ok(title)) => Some(title.to_owned()),
658 None => None,
659 };
660 if a.labels.as_deref().is_some_and(|labels| normalize_labels(labels).is_none()) {
661 return Ok(Outcome::fail(
662 FailureCode::Invalid,
663 format!("An issue can have up to {MAX_LABELS} labels of up to {MAX_LABEL_CHARS} characters each."),
664 ));
665 }
666 let assignees = match a.assignees {
667 Some(names) => Some(check!(self.valid_assignees(names).await?)),
668 None => None,
669 };
670 // Its labels and milestone first: either can be refused, and then
671 // nothing else changes.
672 if a.milestone.is_some() || a.labels.is_some() {
673 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
674 if let Some(number) = a.milestone {
675 check!(self.set_milestone(&a.actor, &repo, &labels::Item::Issue(issue.clone()), number).await?);
676 }
677 if let Some(labels) = &a.labels {
678 check!(self.relabel(&a.actor, &repo, &labels::Item::Issue(issue.clone()), labels).await?);
679 }
680 }
681 let assigned = assignees.as_ref().map(serde_json::to_string).transpose()?;
682 let body = a.body.map(|body| body.trim().to_owned());
683 self.db
684 .prepare(
685 "UPDATE issues
686 SET title = COALESCE(?, title), body = COALESCE(?, body),
687 assignees = COALESCE(?, assignees), updated_at = ?
688 WHERE id = ?",
689 )
690 .bind(&[
691 optional(&title),
692 optional(&body),
693 optional(&assigned),
694 rfc3339(now_ms()).into(),
695 issue.id.as_str().into(),
696 ])?
697 .run()
698 .await?;
699 let before = issue.assignees.clone();
700 let labels_before = issue.labels.clone();
701 let Some(issue) = self.issue(&issue.repo_id, issue.number).await? else {
702 return Ok(no_issue());
703 };
704 self.apply_label_rule(&a.actor, &issue, &labels_before).await?;
705 self.publish(
706 "issue.updated",
707 &issue.repo_id,
708 &a.actor,
709 Self::issue_event(&issue),
710 )
711 .await?;
712 if let Some(assignees) = assignees {
713 self.note_changes(
714 &issue.repo_id,
715 issue.number,
716 &a.actor,
717 &before,
718 &assignees,
719 ("assigned", "unassigned"),
720 )
721 .await?;
722 let added: Vec<String> = assignees.iter().filter(|name| !before.contains(name)).cloned().collect();
723 self.publish(
724 "issue.assigned",
725 &issue.repo_id,
726 &a.actor,
727 IssueEvent {
728 assignees: Some(assignees),
729 added: Some(added),
730 ..Self::issue_event(&issue)
731 },
732 )
733 .await?;
734 }
735 Ok(Outcome::Ok(issue))
736 }
737
738 /// Usernames as given, tidied, if each names an account.
739 async fn valid_assignees(&self, names: Vec<String>) -> Result<Outcome<Vec<String>>> {
740 let mut assignees: Vec<String> = Vec::new();
741 for name in names {
742 let name = name.trim().trim_start_matches('@').to_lowercase();
743 if name.is_empty() || assignees.contains(&name) {
744 continue;
745 }
746 if assignees.len() == MAX_ASSIGNEES {
747 return Ok(Outcome::fail(
748 FailureCode::Invalid,
749 format!("An issue can be assigned to at most {MAX_ASSIGNEES} people."),
750 ));
751 }
752 let account: Viewer = g1t_kit::call(
753 &self.identity,
754 "user_by_username",
755 &UsernameArgs {
756 username: name.clone(),
757 },
758 )
759 .await?;
760 if account.is_none() {
761 return Ok(Outcome::fail(
762 FailureCode::Invalid,
763 format!("There is no account named {name}."),
764 ));
765 }
766 assignees.push(name);
767 }
768 Ok(Outcome::Ok(assignees))
769 }
770
771 /// Open issues assigned to the viewer, in every repository. Callers
772 /// show only those in repositories the viewer can still see.
773 async fn list_assigned_issues(&self, a: ViewerArgs) -> Result<Vec<Issue>> {
774 let Some(viewer) = a.viewer else {
775 return Ok(Vec::new());
776 };
777 let rows = self
778 .db
779 .prepare(format!(
780 "SELECT {ISSUE_COLUMNS} FROM issues
781 WHERE state = 'open' AND EXISTS (
782 SELECT 1 FROM json_each(issues.assignees) WHERE json_each.value = ?)
783 ORDER BY updated_at DESC LIMIT 50"
784 ))
785 .bind(&[viewer.username.into()])?
786 .all()
787 .await?
788 .results::<IssueRow>()?;
789 Ok(rows.into_iter().map(Issue::from).collect())
790 }
791
792 async fn close_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
793 let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
794 if issue.state == State::Closed {
795 return Ok(Outcome::fail(
796 FailureCode::Conflict,
797 "This issue is already closed.",
798 ));
799 }
800 let reason = a.reason.unwrap_or(IssueReason::Completed);
801 let now = rfc3339(now_ms());
802 self.db
803 .prepare(
804 "UPDATE issues SET state = 'closed', reason = ?, closed_at = ?, updated_at = ?
805 WHERE id = ?",
806 )
807 .bind(&[
808 reason.as_str().into(),
809 now.as_str().into(),
810 now.as_str().into(),
811 issue.id.as_str().into(),
812 ])?
813 .run()
814 .await?;
815 self.publish(
816 "issue.closed",
817 &issue.repo_id,
818 &a.actor,
819 IssueEvent {
820 reason: Some(reason.as_str()),
821 ..Self::issue_event(&issue)
822 },
823 )
824 .await?;
825 self.note(
826 &issue.repo_id,
827 issue.number,
828 (&a.actor.id, &a.actor.username),
829 match reason {
830 IssueReason::Completed => "closed this as completed",
831 IssueReason::NotPlanned => "closed this as not planned",
832 },
833 )
834 .await?;
835 issue.state = State::Closed;
836 issue.reason = Some(reason);
837 issue.closed_at = Some(now.clone());
838 issue.updated_at = now;
839 Ok(Outcome::Ok(issue))
840 }
841
842 async fn reopen_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
843 let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
844 if issue.state == State::Open {
845 return Ok(Outcome::fail(
846 FailureCode::Conflict,
847 "This issue is already open.",
848 ));
849 }
850 let now = rfc3339(now_ms());
851 self.db
852 .prepare(
853 "UPDATE issues
854 SET state = 'open', reason = NULL, resolved_by = NULL, closed_at = NULL,
855 updated_at = ?
856 WHERE id = ?",
857 )
858 .bind(&[now.as_str().into(), issue.id.as_str().into()])?
859 .run()
860 .await?;
861 self.publish(
862 "issue.reopened",
863 &issue.repo_id,
864 &a.actor,
865 Self::issue_event(&issue),
866 )
867 .await?;
868 self.note(
869 &issue.repo_id,
870 issue.number,
871 (&a.actor.id, &a.actor.username),
872 "reopened this",
873 )
874 .await?;
875 issue.state = State::Open;
876 issue.reason = None;
877 issue.resolved_by = None;
878 issue.closed_at = None;
879 issue.updated_at = now;
880 Ok(Outcome::Ok(issue))
881 }
882
883
884 async fn counts(&self, a: ViewArgs) -> Result<Outcome<Counts>> {
885 let read = |repo_id: String| async move {
886 let query = self
887 .db
888 .prepare(
889 "SELECT
890 (SELECT count(*) FROM issues WHERE repo_id = ?1 AND state = 'open') AS issues,
891 (SELECT count(*) FROM pulls
892 WHERE repo_id = ?1 AND status IN ('draft', 'open')) AS pulls",
893 )
894 .bind(&[repo_id.into()])?;
895 self.timing.db(1, query.first::<Counts>(None)).await
896 };
897 let (_, counts) = check!(self.repo_then(&a.repo, &a.viewer, read).await?);
898 Ok(Outcome::Ok(counts.unwrap_or(Counts {
899 issues: 0,
900 pulls: 0,
901 })))
902 }
903
904 // --- Comments ----------------------------------------------------------
905
906 async fn add_comment(&self, a: AddCommentArgs) -> Result<Outcome<Comment>> {
907 if !a.actor.verified {
908 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
909 }
910 let body = a.body.trim();
911 // An approval speaks for itself; anything else has to say something.
912 if body.is_empty() && a.verdict != Some(Verdict::Approve) {
913 return Ok(Outcome::fail(
914 FailureCode::Invalid,
915 "A comment cannot be empty.",
916 ));
917 }
918 let path = a
919 .path
920 .as_deref()
921 .map(str::trim)
922 .filter(|path| !path.is_empty());
923 let line = a.line.filter(|line| *line > 0 && path.is_some());
924 if body.chars().count() > MAX_ENTRY_CHARS {
925 return Ok(Outcome::fail(
926 FailureCode::Invalid,
927 "That comment is too long.",
928 ));
929 }
930 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
931 check!(writable(&repo));
932 // The number names an issue or a pull request, never both.
933 let mut pull_id = None;
934 // A command to g1t on a dependency update it opened (`@g1t rebase`)
935 // is the security service's to act on, not a mention for an agent.
936 let mut update_command = false;
937 let table = if self.issue(&repo.id, a.number).await?.is_some() {
938 if path.is_some() || a.verdict.is_some() {
939 return Ok(Outcome::fail(
940 FailureCode::Invalid,
941 "Only a pull request can be reviewed or commented on by line.",
942 ));
943 }
944 "issues"
945 } else if let Some(pull) = self.pull(&repo.id, a.number).await? {
946 if a.verdict.is_some() && pull.is_owned_by(&a.actor.id) {
947 return Ok(Outcome::fail(
948 FailureCode::Forbidden,
949 "You cannot approve or request changes on your own pull request.",
950 ));
951 }
952 pull_id = Some(pull.id.clone());
953 update_command = pull.author.is_system()
954 && g1t_contracts::updates::update_command(body).is_some();
955 "pulls"
956 } else {
957 return Ok(Outcome::fail(
958 FailureCode::NotFound,
959 "No issue or pull request has that number.",
960 ));
961 };
962
963 let now = now_ms();
964 let comment = Comment {
965 kind: CommentKind::Comment,
966 id: new_id("cmt", now),
967 author: a.actor.clone(),
968 body: body.to_owned(),
969 path: path.map(str::to_owned),
970 line,
971 verdict: a.verdict,
972 created_at: rfc3339(now),
973 };
974 self.db
975 .batch(vec![
976 self.db
977 .prepare(
978 "INSERT INTO comments
979 (id, repo_id, number, author_id, author_name, body, path, line,
980 verdict, created_at)
981 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
982 )
983 .bind(&[
984 comment.id.as_str().into(),
985 repo.id.as_str().into(),
986 a.number.into(),
987 a.actor.id.as_str().into(),
988 a.actor.username.as_str().into(),
989 body.into(),
990 optional(&comment.path),
991 optional_number(line),
992 a.verdict
993 .map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
994 comment.created_at.as_str().into(),
995 ])?,
996 self.db
997 .prepare(format!(
998 "UPDATE {table} SET updated_at = ? WHERE repo_id = ? AND number = ?"
999 ))
1000 .bind(&[
1001 comment.created_at.as_str().into(),
1002 repo.id.as_str().into(),
1003 a.number.into(),
1004 ])?,
1005 ])
1006 .await?;
1007 if !update_command {
1008 self.note_mention(&a.actor, &repo, a.number, &comment, pull_id.as_deref()).await?;
1009 }
1010 self.publish(
1011 "comment.created",
1012 &repo.id,
1013 &a.actor,
1014 CommentCreated {
1015 comment_id: comment.id.clone(),
1016 repo_id: repo.id.clone(),
1017 number: a.number,
1018 pull_id,
1019 verdict: a.verdict,
1020 },
1021 )
1022 .await?;
1023 Ok(Outcome::Ok(comment))
1024 }
1025
1026 // --- Pull requests -----------------------------------------------------
1027
1028 async fn open_pull(&self, a: OpenPullArgs) -> Result<Outcome<Pull>> {
1029 if !a.actor.verified {
1030 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1031 }
1032 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
1033 check!(writable(&repo));
1034 // g1t's own agent at work spends the workspace's compute; a pull
1035 // request anyone else's agent makes is like any other.
1036 if matches!(a.runtime, Runtime::Hosted) {
1037 check!(allowed(Some(&a.actor), &repo, Capability::Run));
1038 }
1039 let issue = match a.issue {
1040 Some(number) => match self.issue(&repo.id, number).await? {
1041 Some(issue) if issue.state == State::Open => Some(issue),
1042 Some(_) => {
1043 return Ok(Outcome::fail(
1044 FailureCode::Conflict,
1045 "This issue is closed.",
1046 ));
1047 }
1048 None => return Ok(no_issue()),
1049 },
1050 None => None,
1051 };
1052 // A pull request for an issue takes the issue's title unless given one.
1053 let title = match (a.title.trim(), &issue) {
1054 ("", Some(issue)) => issue.title.clone(),
1055 (title, _) => match valid_title(title) {
1056 Ok(title) => title.to_owned(),
1057 Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
1058 },
1059 };
1060 // Unnamed, the change is its author's, unless an agent opened it.
1061 let agent = match a.agent.trim() {
1062 "" if g1t_contracts::rules::is_agent(&a.actor) => "agent",
1063 "" => a.actor.username.as_str(),
1064 agent => agent,
1065 };
1066 let runtime = match a.runtime {
1067 Runtime::Hosted => "hosted",
1068 Runtime::External => "external",
1069 };
1070
1071 let now = now_ms();
1072 let id = new_id("pr", now);
1073 let branch = a
1074 .branch
1075 .as_deref()
1076 .map(str::trim)
1077 .filter(|branch| !branch.is_empty());
1078 // The branch it merges into: the default branch unless another is
1079 // asked for, which has to exist.
1080 let base = a.base.as_deref().and_then(|base| stored_base(base, &repo));
1081 if let Some(base) = &base {
1082 if branch == Some(base.as_str()) {
1083 return Ok(Outcome::fail(
1084 FailureCode::Invalid,
1085 format!("A pull request cannot merge {base} into itself. Choose another base."),
1086 ));
1087 }
1088 let exists: Option<String> = g1t_kit::call(
1089 &self.repos,
1090 "head",
1091 &HeadArgs { repo_id: repo.id.clone(), branch: base.clone() },
1092 )
1093 .await?;
1094 if exists.is_none() {
1095 return Ok(Outcome::fail(FailureCode::NotFound, format!("There is no branch named {base} to merge into.")));
1096 }
1097 }
1098 let base_name = base.clone().unwrap_or_else(|| repo.default_branch.clone());
1099 // The change is on a branch already pushed to the repository, or
1100 // will be made in a fork created for this pull request.
1101 let (fork, head) = match branch {
1102 Some(branch) => {
1103 if branch == base_name {
1104 return Ok(Outcome::fail(
1105 FailureCode::Invalid,
1106 format!("Choose a branch other than {branch}."),
1107 ));
1108 }
1109 let head: Option<String> = g1t_kit::call(
1110 &self.repos,
1111 "head",
1112 &HeadArgs {
1113 repo_id: repo.id.clone(),
1114 branch: branch.to_owned(),
1115 },
1116 )
1117 .await?;
1118 let Some(head) = head else {
1119 return Ok(Outcome::fail(
1120 FailureCode::NotFound,
1121 format!("There is no branch named {branch}. Push it first."),
1122 ));
1123 };
1124 // One open pull request for each branch and base.
1125 let existing = self
1126 .db
1127 .prepare(
1128 "SELECT number AS n FROM pulls
1129 WHERE repo_id = ? AND source_branch = ? AND status IN ('draft', 'open')
1130 AND base_branch IS ?",
1131 )
1132 .bind(&[repo.id.as_str().into(), branch.into(), optional(&base)])?
1133 .first::<NumberRow>(None)
1134 .await?;
1135 if let Some(existing) = existing {
1136 return Ok(Outcome::fail(
1137 FailureCode::Conflict,
1138 format!("Pull request #{} is already open from {branch} into {base_name}.", existing.n),
1139 ));
1140 }
1141 (None, Some(head))
1142 }
1143 None => {
1144 let fork: Outcome<Repo> = g1t_kit::call(
1145 &self.repos,
1146 "fork_for_pull",
1147 &ForkArgs {
1148 source_id: repo.id.clone(),
1149 pull_id: id.clone(),
1150 actor: a.actor.clone(),
1151 },
1152 )
1153 .await?;
1154 (Some(check!(fork)), None)
1155 }
1156 };
1157 // A branch already holds the work, so its pull request is ready for
1158 // review from the start; one with a fork starts as a draft.
1159 let status = if branch.is_some() { "open" } else { "draft" };
1160 let body = Some(a.body.trim().to_owned()).filter(|body| !body.is_empty());
1161
1162 let number = self.next_number(&repo.id).await?;
1163 let timestamp = rfc3339(now);
1164 // A change g1t makes is g1t's, for whoever asked for it. This is
1165 // what lifecycle::made_by_g1t reads back.
1166 let by_g1t = matches!(a.runtime, Runtime::Hosted) && agent == reviews::AGENT_NAME && fork.is_some();
1167 let (author, requested_by) = authorship(&a.actor, by_g1t);
1168 self.db
1169 .prepare(
1170 "INSERT INTO pulls
1171 (id, repo_id, number, issue_id, issue_number, title, body, agent, runtime,
1172 status, fork_repo_id, fork_namespace, fork_name, source_branch, head_commit,
1173 author_id, author_name, requested_by_id, requested_by_name, created_at, updated_at,
1174 base_branch)
1175 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
1176 )
1177 .bind(&[
1178 id.as_str().into(),
1179 repo.id.as_str().into(),
1180 number.into(),
1181 optional(&issue.as_ref().map(|issue| issue.id.clone())),
1182 optional_number(issue.as_ref().map(|issue| issue.number)),
1183 title.into(),
1184 optional(&body),
1185 agent.into(),
1186 runtime.into(),
1187 status.into(),
1188 optional(&fork.as_ref().map(|fork| fork.id.clone())),
1189 optional(&fork.as_ref().map(|fork| fork.namespace.clone())),
1190 optional(&fork.as_ref().map(|fork| fork.name.clone())),
1191 optional(&branch.map(str::to_owned)),
1192 optional(&head),
1193 author.id.as_str().into(),
1194 author.username.as_str().into(),
1195 optional(&requested_by.as_ref().map(|user| user.id.clone())),
1196 optional(&requested_by.as_ref().map(|user| user.username.clone())),
1197 timestamp.as_str().into(),
1198 timestamp.as_str().into(),
1199 optional(&base),
1200 ])?
1201 .run()
1202 .await?;
1203 let Some(mut pull) = self.pull(&repo.id, number).await? else {
1204 return Ok(no_pull());
1205 };
1206 fill_base(&mut pull, &repo);
1207 self.manage(&pull).await?;
1208 // Someone is on it now, so it is no longer waiting for an agent.
1209 if let Some(issue) = pull.issue {
1210 self.db
1211 .prepare("UPDATE issues SET queued_by = NULL WHERE repo_id = ? AND number = ?")
1212 .bind(&[repo.id.as_str().into(), issue.into()])?
1213 .run()
1214 .await?;
1215 }
1216 if let Some(issue) = pull.issue {
1217 let text = if lifecycle::made_by_g1t(&pull) {
1218 format!("assigned this to g1t, which opened #{}", pull.number)
1219 } else {
1220 format!("opened #{} for this", pull.number)
1221 };
1222 self.note(&repo.id, issue, (&a.actor.id, &a.actor.username), &text)
1223 .await?;
1224 }
1225 self.publish(
1226 "pull.opened",
1227 &repo.id,
1228 &a.actor,
1229 PullEvent {
1230 agent: Some(pull.agent.clone()),
1231 base: pull.base.clone(),
1232 ..Self::pull_event(&pull)
1233 },
1234 )
1235 .await?;
1236 // Its code owners asked to review (codeowners.rs).
1237 self.refresh_code_owners(&pull).await;
1238 let pull = self.pull(&repo.id, number).await?.unwrap_or(pull);
1239 Ok(Outcome::Ok(pull))
1240 }
1241
1242 async fn list_pulls(&self, a: ListPullsArgs) -> Result<Outcome<Vec<Pull>>> {
1243 let filter = match a.state {
1244 Some(State::Open) => "AND status IN ('draft', 'open')",
1245 Some(State::Closed) => "AND status IN ('merged', 'closed')",
1246 None => "",
1247 };
1248 let label = a
1249 .label
1250 .map(|label| label.split_whitespace().collect::<Vec<_>>().join(" ").to_lowercase())
1251 .filter(|label| !label.is_empty());
1252 let milestone = a.milestone;
1253 let read = |repo_id: String| {
1254 let label = label.clone();
1255 async move {
1256 let query = self
1257 .db
1258 .prepare(format!(
1259 "SELECT {PULL_COLUMNS} FROM pulls WHERE repo_id = ? {filter}
1260 AND (? IS NULL OR EXISTS
1261 (SELECT 1 FROM json_each(pulls.labels) WHERE json_each.value = ?))
1262 AND (? IS NULL OR milestone = ?)
1263 ORDER BY number DESC LIMIT ?"
1264 ))
1265 .bind(&[
1266 repo_id.into(),
1267 optional(&label),
1268 optional(&label),
1269 optional_number(milestone),
1270 optional_number(milestone),
1271 LIST_PAGE.into(),
1272 ])?;
1273 self.timing.db(1, query.all()).await?.results::<PullRow>()
1274 }
1275 };
1276 let (repo, rows) = check!(self.repo_then(&a.repo, &a.viewer, read).await?);
1277 let base = a.base.map(|base| base.trim().to_owned()).filter(|base| !base.is_empty());
1278 Ok(Outcome::Ok(
1279 rows.into_iter()
1280 .map(|row| {
1281 let mut pull = Pull::from(row);
1282 fill_base(&mut pull, &repo);
1283 pull
1284 })
1285 .filter(|pull| base.as_deref().is_none_or(|base| pull.base.as_deref() == Some(base)))
1286 .collect(),
1287 ))
1288 }
1289
1290 /// `pulls_for_repos`: what `list_pulls` gives, open and closed, for many
1291 /// repositories at once: one access check with repos for all of them
1292 /// and one query, instead of two of each per repository.
1293 async fn pulls_for_repos(&self, a: PullsForReposArgs) -> Result<Vec<RepoPulls>> {
1294 let ids: Vec<String> = a.repo_ids.into_iter().take(MAX_PULLS_FOR_REPOS).collect();
1295 if ids.is_empty() {
1296 return Ok(Vec::new());
1297 }
1298 let limit = a.limit.clamp(1, LIST_PAGE);
1299 // The rows are read beside the access check, for every id asked
1300 // about; those of repositories the viewer cannot read are dropped.
1301 let asked = serde_json::to_string(&ids)?;
1302 let check = ReadableArgs { ids, viewer: a.viewer };
1303 let readable = self.timing.rpc(g1t_kit::call::<_, Vec<Repo>>(&self.repos, "readable", &check));
1304 let (readable, rows) = try_join(readable, self.timing.db(1, self.newest_pulls(asked, limit))).await?;
1305 if readable.is_empty() {
1306 return Ok(Vec::new());
1307 }
1308 let mut answer: Vec<RepoPulls> = readable
1309 .iter()
1310 .map(|repo| RepoPulls { repo_id: repo.id.clone(), open: Vec::new(), closed: Vec::new() })
1311 .collect();
1312 for pull in rows.into_iter().map(Pull::from) {
1313 let Some(entry) = answer.iter_mut().find(|entry| entry.repo_id == pull.repo_id) else {
1314 continue;
1315 };
1316 match pull.status {
1317 PullStatus::Draft | PullStatus::Open => entry.open.push(pull),
1318 PullStatus::Merged | PullStatus::Closed => entry.closed.push(pull),
1319 }
1320 }
1321 for entry in &mut answer {
1322 entry.open.sort_by_key(|pull| std::cmp::Reverse(pull.number));
1323 entry.closed.sort_by_key(|pull| std::cmp::Reverse(pull.number));
1324 }
1325 Ok(answer)
1326 }
1327
1328 /// The newest `limit` of each repository's open (draft or open) and
1329 /// closed (merged or closed) pull requests, for the ids in `ids` (JSON).
1330 async fn newest_pulls(&self, ids: String, limit: u32) -> Result<Vec<PullRow>> {
1331 self.db
1332 .prepare(format!(
1333 "SELECT * FROM (
1334 SELECT {PULL_COLUMNS}, ROW_NUMBER() OVER (
1335 PARTITION BY pulls.repo_id, pulls.status IN ('draft', 'open') ORDER BY pulls.number DESC
1336 ) AS place
1337 FROM pulls WHERE pulls.repo_id IN (SELECT value FROM json_each(?1))
1338 ) WHERE place <= ?2"
1339 ))
1340 .bind(&[ids.into(), limit.into()])?
1341 .all()
1342 .await?
1343 .results::<PullRow>()
1344 }
1345
1346 async fn get_pull(&self, a: ViewArgs) -> Result<Outcome<PullDetail>> {
1347 let number = a.number;
1348 // Every row the page and the lifecycle read, in one batch started
1349 // beside the access check; the helpers below read from it.
1350 let namespace = a.repo.namespace.clone();
1351 let read = |repo_id: String| self.prefetch_pull(repo_id, namespace.clone(), number);
1352 let Outcome::Ok((repo, Some(found))) = self.repo_then(&a.repo, &a.viewer, read).await? else {
1353 return Ok(no_pull());
1354 };
1355 let Some(row) = found.first::<PullRow>(prefetch::Slot::Pull)? else {
1356 return Ok(no_pull());
1357 };
1358 let stored = found.first::<StoredBehind>(prefetch::Slot::Pull)?;
1359 let issue = found.first::<IssueRow>(prefetch::Slot::Issue)?.map(Issue::from);
1360 let comments: Vec<Comment> =
1361 found.rows::<CommentRow>(prefetch::Slot::Comments)?.into_iter().map(Comment::from).collect();
1362 self.keep_prefetched(Some(found));
1363 let default_branch = repo.default_branch.clone();
1364 let detail = self.pull_detail(repo, Pull::from(row), issue, comments, stored, &a.viewer).await;
1365 self.keep_prefetched(None);
1366 // The branch it merges into, named, for whoever reads it.
1367 Ok(match detail? {
1368 Outcome::Ok(mut detail) => {
1369 if detail.pull.base.as_deref().is_none_or(str::is_empty) {
1370 detail.pull.base = Some(default_branch);
1371 }
1372 Outcome::Ok(detail)
1373 }
1374 failed => failed,
1375 })
1376 }
1377
1378 async fn pull_detail(
1379 &self,
1380 repo: Repo,
1381 mut pull: Pull,
1382 issue: Option<Issue>,
1383 comments: Vec<Comment>,
1384 stored: Option<StoredBehind>,
1385 viewer: &Viewer,
1386 ) -> Result<Outcome<PullDetail>> {
1387 // Whether it is behind, as worked out with its mergeability on the
1388 // last push to either side (mergeability.rs), when that was for
1389 // its head as it is now; otherwise asked of the repos service.
1390 let known_behind = stored.and_then(|stored| stored.for_head(pull.head_commit.as_deref()));
1391 // Worked out on each push; this covers a pull request from before
1392 // that was recorded.
1393 if pull.files.is_empty() && pull.head_commit.is_some() {
1394 pull.files = self.refresh_files(&pull).await?;
1395 }
1396 // Everything else at once: none of it depends on the rest, and each
1397 // is a round trip of its own.
1398 let standing = async {
1399 // Mergeability first: where g1t sees a pull request through, a
1400 // conflict decides its next step.
1401 let behind = async {
1402 match known_behind {
1403 Some(behind) => Ok(behind),
1404 None => {
1405 let behind = self.is_behind(&repo.id, &pull).await?;
1406 // Kept for the next view when the mergeability on
1407 // record is for this head: a pull request from
1408 // before `behind` was kept asks once.
1409 if let Some(head) = pull.head_commit.as_deref()
1410 && pull.status.is_active()
1411 {
1412 self.db
1413 .prepare(
1414 "UPDATE pulls SET behind = ?1
1415 WHERE id = ?2 AND behind IS NULL AND mergeable_key LIKE ?3 || '..%'",
1416 )
1417 .bind(&[u32::from(behind).into(), pull.id.as_str().into(), head.into()])?
1418 .run()
1419 .await?;
1420 }
1421 Ok(behind)
1422 }
1423 }
1424 };
1425 let (merge, behind) = try_join(self.mergeability(&pull), behind).await?;
1426 let assessed = self.assess_with_confidence(&pull, &issue, behind).await?;
1427 let confidence = assessed.as_ref().and_then(|(_, _, confidence)| confidence.clone());
1428 let lifecycle = assessed.map(|(lifecycle, _, _)| lifecycle);
1429 Ok::<_, worker::Error>((behind, (lifecycle, confidence), merge))
1430 };
1431 let (((behind, (lifecycle, confidence), (mergeable, conflicts)), (landing, stalled), comments), (checks, overlaps, review_pending)) =
1432 try_join(
1433 try_join3(standing, self.landing_state(&pull.id), async { Ok(comments) }),
1434 try_join3(
1435 self.latest_checks(&pull.id),
1436 self.overlaps(&pull),
1437 self.review_pending(&pull.id),
1438 ),
1439 )
1440 .await?;
1441 // As just worked out, rather than as it was read.
1442 if confidence.is_some() {
1443 pull.confidence = confidence;
1444 }
1445 // The rules of the branch it merges into, as they stack, and which
1446 // of them it does not meet yet, for whoever is looking.
1447 let (statuses, settings, gate) = try_join3(
1448 self.statuses(&repo.id, pull.head_commit.as_deref()),
1449 self.settings_on(&repo, &pull),
1450 async {
1451 if pull.status.is_active() {
1452 self.merge_gate(&repo, &pull, viewer.as_ref(), false, true).await.map(Some)
1453 } else {
1454 Ok(None)
1455 }
1456 },
1457 )
1458 .await?;
1459 let code_owners = self.pull_code_owners(&pull, &comments, &settings).await?;
1460 Ok(Outcome::Ok(PullDetail {
1461 required_checks: required_checks(&settings.required_checks, &statuses),
1462 rules: gate.map(|gate| rulesets::merge_rules(&gate.judged, &gate.requirements, pull.base_branch(&repo.default_branch) == repo.default_branch)),
1463 code_owners,
1464 comments,
1465 checks,
1466 overlaps,
1467 behind,
1468 review_pending,
1469 lifecycle,
1470 landing,
1471 stalled,
1472 messages: self.messages(&pull.id).await?,
1473 statuses,
1474 mergeable,
1475 conflicts,
1476 earlier_checks: self.earlier_checks(&pull.id).await?,
1477 issue,
1478 pull,
1479 }))
1480 }
1481
1482 /// The pull request, if it is still active and `actor` opened it or
1483 /// may triage the repository's pull requests.
1484 async fn manageable_pull(
1485 &self,
1486 actor: &User,
1487 path: &RepoPath,
1488 number: u32,
1489 ) -> Result<Outcome<Pull>> {
1490 let (repo, pull) = check!(self.pull_at(path, number, &Some(actor.clone())).await?);
1491 check!(writable(&repo));
1492 if !pull.is_owned_by(&actor.id) {
1493 check!(allowed(Some(actor), &repo, Capability::Triage));
1494 }
1495 if !pull.status.is_active() {
1496 return Ok(Outcome::fail(
1497 FailureCode::Conflict,
1498 format!("This pull request is already {}.", pull.status.as_str()),
1499 ));
1500 }
1501 Ok(Outcome::Ok(pull))
1502 }
1503
1504 /// Brings a pull request up to date with the default branch without a
1505 /// sandbox, where the repos service can do that safely. Whoever could
1506 /// have pushed the merge themselves may ask: whoever opened it (or asked
1507 /// g1t for it), for a fork; anyone who may push, for a branch of the
1508 /// repository. When it needs a
1509 /// real merge, says so, naming the conflicting files if a probe found
1510 /// them, and pushes nothing.
1511 async fn catch_up_pull(&self, a: PullActionArgs) -> Result<Outcome<PullBranchUpdate>> {
1512 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &Some(a.actor.clone())).await?);
1513 check!(writable(&repo));
1514 if !pull.status.is_active() {
1515 return Ok(Outcome::fail(
1516 FailureCode::Conflict,
1517 format!("This pull request is already {}.", pull.status.as_str()),
1518 ));
1519 }
1520 if pull.fork_repo_id.is_some() {
1521 if !pull.is_owned_by(&a.actor.id) {
1522 return Ok(Outcome::fail(
1523 FailureCode::Forbidden,
1524 "Only whoever opened this pull request, or asked g1t for it, can update it.",
1525 ));
1526 }
1527 } else {
1528 check!(allowed(Some(&a.actor), &repo, Capability::Push));
1529 }
1530 let updated: Outcome<PullBranchUpdate> = g1t_kit::call(
1531 &self.repos,
1532 "update_pull_branch",
1533 &UpdatePullBranchArgs {
1534 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
1535 branch: pull.branch.clone(),
1536 number: pull.number,
1537 actor: a.actor,
1538 target_branch: pull.base.clone(),
1539 },
1540 )
1541 .await?;
1542 // A probe that found conflicts says more than "both changed it".
1543 if let Outcome::Ok(PullBranchUpdate::NeedsAgent { .. }) = &updated
1544 && let Some(files) = self.conflicting_files(&pull).await?
1545 && !files.is_empty()
1546 {
1547 return Ok(Outcome::Ok(PullBranchUpdate::NeedsAgent {
1548 reason: NeedsAgentReason::Conflicting,
1549 detail: "Merging it conflicts.".to_owned(),
1550 paths: files,
1551 }));
1552 }
1553 Ok(updated)
1554 }
1555
1556 async fn update_pull(&self, a: UpdatePullArgs) -> Result<Outcome<Pull>> {
1557 let pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1558 // The branch it merges into, its milestone and its labels first:
1559 // each can be refused, and then nothing else changes.
1560 if a.base.is_some() || a.milestone.is_some() || a.labels.is_some() {
1561 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
1562 if let Some(base) = &a.base {
1563 check!(self.change_base(&a.actor, &repo, &pull, base).await?);
1564 }
1565 if let Some(number) = a.milestone {
1566 check!(self.set_milestone(&a.actor, &repo, &labels::Item::Pull(pull.clone()), number).await?);
1567 }
1568 if let Some(labels) = &a.labels {
1569 check!(self.relabel(&a.actor, &repo, &labels::Item::Pull(pull.clone()), labels).await?);
1570 }
1571 }
1572 let assignees = match a.assignees {
1573 Some(names) => Some(check!(self.valid_assignees(names).await?)),
1574 None => None,
1575 };
1576 // Teams, named `workspace/team`, apart from the people.
1577 let (team_names, a_reviewers) = match a.reviewers {
1578 Some(names) => {
1579 let (teams, people): (Vec<String>, Vec<String>) =
1580 names.into_iter().partition(|name| team_reviews::team_name(name).is_some());
1581 (Some(teams), Some(people))
1582 }
1583 None => (None, None),
1584 };
1585 let teams = match team_names {
1586 Some(names) => Some(check!(self.valid_team_reviewers(&a.actor, &a.repo, &pull, names).await?)),
1587 None => None,
1588 };
1589 let reviewers = match a_reviewers {
1590 Some(names) => {
1591 // g1t is not an account; everyone else has to be.
1592 let agent = names
1593 .iter()
1594 .any(|name| name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME));
1595 let people = names
1596 .into_iter()
1597 .filter(|name| !name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME))
1598 .collect();
1599 let mut reviewers = check!(self.valid_assignees(people).await?);
1600 // Nobody is asked to review their own, nor what they had g1t make.
1601 reviewers.retain(|name| *name != pull.owner().username);
1602 if agent {
1603 reviewers.insert(0, reviews::AGENT_NAME.to_owned());
1604 }
1605 Some(reviewers)
1606 }
1607 None => None,
1608 };
1609 self.db
1610 .prepare(
1611 "UPDATE pulls
1612 SET assignees = COALESCE(?, assignees), updated_at = ?
1613 WHERE id = ?",
1614 )
1615 .bind(&[
1616 optional(&assignees.as_ref().map(serde_json::to_string).transpose()?),
1617 rfc3339(now_ms()).into(),
1618 pull.id.as_str().into(),
1619 ])?
1620 .run()
1621 .await?;
1622 if let Some(assignees) = &assignees {
1623 self.note_changes(
1624 &pull.repo_id,
1625 pull.number,
1626 &a.actor,
1627 &pull.assignees,
1628 assignees,
1629 ("assigned", "unassigned"),
1630 )
1631 .await?;
1632 }
1633 // Who was newly assigned or asked to review, and whose request was
1634 // withdrawn: the inbox tells them, and webhooks say so.
1635 let newly = |after: &[String], before: &[String]| -> Vec<String> {
1636 after.iter().filter(|name| !before.contains(name)).cloned().collect()
1637 };
1638 if let Some(assignees) = &assignees {
1639 let added = newly(assignees, &pull.assignees);
1640 if !added.is_empty() {
1641 self.publish(
1642 "pull.assigned",
1643 &pull.repo_id,
1644 &a.actor,
1645 PullEvent {
1646 assignees: Some(assignees.clone()),
1647 added: Some(added),
1648 ..Self::pull_event(&pull)
1649 },
1650 )
1651 .await?;
1652 }
1653 }
1654 if reviewers.is_some() || teams.is_some() {
1655 let people = reviewers.unwrap_or_else(|| pull.reviewers.clone());
1656 let teams = teams.unwrap_or_else(|| pull.team_reviewers.clone());
1657 self.set_reviewers(&pull, people, teams, Some(&a.actor), false).await?;
1658 }
1659 Ok(match self.pull(&pull.repo_id, pull.number).await? {
1660 Some(pull) => Outcome::Ok(pull),
1661 None => no_pull(),
1662 })
1663 }
1664
1665 /// Points an open pull request at another branch to merge into. Needs
1666 /// the Write role. What it would merge, whether it is behind, its
1667 /// mergeability and its checks are all worked out against the new base.
1668 async fn change_base(&self, actor: &User, repo: &Repo, pull: &Pull, base: &str) -> Result<Outcome<()>> {
1669 check!(allowed(Some(actor), repo, Capability::Push));
1670 let base = base.trim();
1671 if base.is_empty() {
1672 return Ok(Outcome::fail(FailureCode::Invalid, "Name the branch it should merge into."));
1673 }
1674 let before = pull.base_branch(&repo.default_branch).to_owned();
1675 if base == before {
1676 return Ok(Outcome::Ok(()));
1677 }
1678 if pull.fork_repo_id.is_none() && pull.branch.as_deref() == Some(base) {
1679 return Ok(Outcome::fail(
1680 FailureCode::Invalid,
1681 format!("A pull request cannot merge {base} into itself. Choose another base."),
1682 ));
1683 }
1684 let exists: Option<String> = g1t_kit::call(
1685 &self.repos,
1686 "head",
1687 &HeadArgs { repo_id: repo.id.clone(), branch: base.to_owned() },
1688 )
1689 .await?;
1690 if exists.is_none() {
1691 return Ok(Outcome::fail(FailureCode::NotFound, format!("There is no branch named {base} to merge into.")));
1692 }
1693 // In the merge queue it was headed for the default branch; it
1694 // leaves the queue for another base.
1695 let left_queue = self
1696 .leave(&pull.repo_id, pull, QueueState::Removed, Some("Its base branch changed."))
1697 .await?;
1698 let stored = stored_base(base, repo);
1699 self.db
1700 .prepare(
1701 "UPDATE pulls
1702 SET base_branch = ?, updated_at = ?, land_requested = NULL, land_requested_at = NULL, behind = NULL,
1703 mergeable = NULL, mergeable_key = NULL, conflicts = NULL
1704 WHERE id = ?",
1705 )
1706 .bind(&[optional(&stored), rfc3339(now_ms()).into(), pull.id.as_str().into()])?
1707 .run()
1708 .await?;
1709 self.note(
1710 &pull.repo_id,
1711 pull.number,
1712 (&actor.id, &actor.username),
1713 &format!("changed the base branch from `{before}` to `{base}`"),
1714 )
1715 .await?;
1716 self.publish(
1717 "pull.base_changed",
1718 &pull.repo_id,
1719 actor,
1720 PullEvent { base: Some(base.to_owned()), ..Self::pull_event(pull) },
1721 )
1722 .await?;
1723 if left_queue {
1724 self.publish_as(
1725 "queue.changed",
1726 &pull.repo_id,
1727 None,
1728 g1t_contracts::events::QueueChanged { repo_id: pull.repo_id.clone() },
1729 )
1730 .await?;
1731 }
1732 // Whether it merges cleanly into the new base.
1733 if let Some(moved) = self.pull_by_id(&pull.id).await?
1734 && let Err(error) = self.assess_mergeability(&moved).await
1735 {
1736 worker::console_warn!("mergeability of {}: {error}", pull.id);
1737 }
1738 Ok(Outcome::Ok(()))
1739 }
1740
1741 /// Marks a draft ready for review, or updates the description of one
1742 /// that already is.
1743 async fn ready_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1744 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1745 let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
1746 let now = rfc3339(now_ms());
1747 self.db
1748 .prepare(
1749 "UPDATE pulls SET status = 'open', body = COALESCE(?, body), updated_at = ?
1750 WHERE id = ?",
1751 )
1752 .bind(&[
1753 optional(&summary),
1754 now.as_str().into(),
1755 pull.id.as_str().into(),
1756 ])?
1757 .run()
1758 .await?;
1759 if pull.status == PullStatus::Draft {
1760 // The head as it is now: the push that came just before may not
1761 // have reached `head_commit` yet, and workflows run on it.
1762 let commit = self.live_head(&pull).await?.or_else(|| pull.head_commit.clone());
1763 self.publish(
1764 "pull.ready",
1765 &pull.repo_id,
1766 &a.actor,
1767 PullEvent {
1768 commit,
1769 ..Self::pull_event(&pull)
1770 },
1771 )
1772 .await?;
1773 }
1774 if pull.status == PullStatus::Draft {
1775 self.note(
1776 &pull.repo_id,
1777 pull.number,
1778 (&a.actor.id, &a.actor.username),
1779 "marked this ready for review",
1780 )
1781 .await?;
1782 }
1783 pull.status = PullStatus::Open;
1784 pull.body = summary.or(pull.body);
1785 pull.updated_at = now;
1786 // A draft's code owners are asked once it is ready.
1787 self.refresh_code_owners(&pull).await;
1788 if let Some(fresh) = self.pull(&pull.repo_id, pull.number).await? {
1789 pull.reviewers = fresh.reviewers;
1790 pull.team_reviewers = fresh.team_reviewers;
1791 }
1792 Ok(Outcome::Ok(pull))
1793 }
1794
1795 async fn close_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1796 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1797 let now = rfc3339(now_ms());
1798 self.db
1799 .prepare("UPDATE pulls SET status = 'closed', updated_at = ? WHERE id = ?")
1800 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
1801 .run()
1802 .await?;
1803 self.publish(
1804 "pull.closed",
1805 &pull.repo_id,
1806 &a.actor,
1807 Self::pull_event(&pull),
1808 )
1809 .await?;
1810 self.note(
1811 &pull.repo_id,
1812 pull.number,
1813 (&a.actor.id, &a.actor.username),
1814 "closed this",
1815 )
1816 .await?;
1817 // A closed pull request leaves the merge queue.
1818 if self
1819 .leave(&pull.repo_id, &pull, QueueState::Removed, Some("It was closed."))
1820 .await?
1821 {
1822 self.publish_as(
1823 "queue.changed",
1824 &pull.repo_id,
1825 None,
1826 g1t_contracts::events::QueueChanged {
1827 repo_id: pull.repo_id.clone(),
1828 },
1829 )
1830 .await?;
1831 }
1832 pull.status = PullStatus::Closed;
1833 pull.updated_at = now;
1834 Ok(Outcome::Ok(pull))
1835 }
1836
1837 /// Lands the pull request on the repository's default branch. Unless
1838 /// told to keep it open, that resolves the issue it was for: the issue
1839 /// closes naming this pull request, and the others still in progress
1840 /// for it close as superseded.
1841 async fn merge_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1842 let viewer = Some(a.actor.clone());
1843 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
1844 check!(writable(&repo));
1845 match pull.status {
1846 PullStatus::Open => {}
1847 PullStatus::Draft => {
1848 return Ok(Outcome::fail(
1849 FailureCode::Conflict,
1850 "This pull request is still a draft. Mark it ready for review first.",
1851 ));
1852 }
1853 status => {
1854 return Ok(Outcome::fail(
1855 FailureCode::Conflict,
1856 format!("This pull request is already {}.", status.as_str()),
1857 ));
1858 }
1859 }
1860 let base = pull.base_branch(&repo.default_branch).to_owned();
1861 // The rules of the branch it merges into, as they stack: the one
1862 // gate for the merge button, the API, MCP, auto-merge and the queue,
1863 // for a person's pull request and an agent's alike. A bypass counts
1864 // when the merger asks for it, or for g1t when a ruleset lists it.
1865 let gate = self.merge_gate(&repo, &pull, Some(&a.actor), a.ignore_checks, true).await?;
1866 let bypassable = gate.bypassable();
1867 let gate = if a.bypass_rules || a.actor.is_system() { gate } else { gate.without_bypass() };
1868 let settings = rulesets::overlay(self.settings(&repo.id).await?, &gate.requirements, base == repo.default_branch);
1869 if pull.check_status == Some(CheckStatus::Failed) && !(a.ignore_checks && settings.allow_ignoring_checks) {
1870 return Ok(Outcome::fail(
1871 FailureCode::Conflict,
1872 "It failed in the merge queue; push a fix to try again.",
1873 ));
1874 }
1875 if let Some(refusal) = gate.refusal() {
1876 if access::can(Some(&a.actor), &repo, Capability::Merge) {
1877 self.record_merge_evaluations(&repo, &pull, &gate).await;
1878 }
1879 let offer = if bypassable && !a.bypass_rules {
1880 " You may bypass these rules: merge again and ask to bypass them (bypass_rules)."
1881 } else {
1882 ""
1883 };
1884 return Ok(Outcome::fail(FailureCode::Conflict, format!("{refusal}{offer}")));
1885 }
1886 // Known ahead of time to conflict: neither a merge nor the queue
1887 // would get through, so say what has to be resolved now.
1888 if let Some(files) = self.conflicting_files(&pull).await? {
1889 let named = if files.is_empty() {
1890 String::new()
1891 } else {
1892 format!(" in {}", files.join(", "))
1893 };
1894 return Ok(Outcome::fail(
1895 FailureCode::Conflict,
1896 format!(
1897 "This branch has conflicts with {base}{named} that must be resolved first. Have g1t resolve them, or merge {base} into it, fix them and push."
1898 ),
1899 ));
1900 }
1901
1902 // A repository that merges through a queue: it joins the queue, and
1903 // lands once its state together with everything ahead has passed.
1904 if settings.merge_queue {
1905 if !a.actor.verified {
1906 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1907 }
1908 check!(allowed(Some(&a.actor), &repo, Capability::Merge));
1909 self.record_merge_evaluations(&repo, &pull, &gate).await;
1910 return self.enqueue(&repo, &pull, &a.actor, a.keep_issue_open).await;
1911 }
1912
1913 // The default branch has moved under it. Unless the repository
1914 // insists on that being dealt with first, bring it up to date and
1915 // land it when that is done.
1916 if self.is_behind(&repo.id, &pull).await? {
1917 if settings.require_up_to_date {
1918 return Ok(Outcome::fail(
1919 FailureCode::Conflict,
1920 format!(
1921 "{base} 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 {base} first."
1922 ),
1923 ));
1924 }
1925 if !a.actor.verified {
1926 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1927 }
1928 check!(allowed(Some(&a.actor), &repo, Capability::Merge));
1929 self.record_merge_evaluations(&repo, &pull, &gate).await;
1930 self.request_landing(&pull, &a.actor, a.keep_issue_open)
1931 .await?;
1932 return Ok(Outcome::Ok(pull));
1933 }
1934
1935 // Whether the actor may write to the repository is decided by repos.
1936 let landed: Outcome<Landed> = g1t_kit::call(
1937 &self.repos,
1938 "land",
1939 &LandArgs {
1940 // A pull request from a branch lands from the repository itself.
1941 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
1942 branch: pull.branch.clone(),
1943 actor: a.actor.clone(),
1944 target_branch: Some(base.clone()),
1945 },
1946 )
1947 .await?;
1948 let landed = check!(landed);
1949 self.record_merge_evaluations(&repo, &pull, &gate).await;
1950 Ok(Outcome::Ok(
1951 self.record_merge(&repo, pull, &a.actor, a.keep_issue_open, landed)
1952 .await?,
1953 ))
1954 }
1955
1956 /// Records a pull request as merged once the default branch holds it:
1957 /// closes its issue, supersedes the others for it, and says so.
1958 pub(crate) async fn record_merge(
1959 &self,
1960 repo: &Repo,
1961 mut pull: Pull,
1962 actor: &User,
1963 keep_issue_open: bool,
1964 landed: Landed,
1965 ) -> Result<Pull> {
1966 // Only a merge into the default branch resolves the issue: into
1967 // another branch, the work has not landed yet.
1968 let keep_issue_open = keep_issue_open || !pull.targets_default(&repo.default_branch);
1969 let issue = match pull.issue {
1970 Some(number) if !keep_issue_open => self
1971 .issue(&repo.id, number)
1972 .await?
1973 .filter(|issue| issue.state == State::Open),
1974 _ => None,
1975 };
1976 let now = rfc3339(now_ms());
1977 let mut statements = vec![
1978 self.db
1979 .prepare(
1980 "UPDATE pulls
1981 SET status = 'merged', head_commit = ?, merge_base = ?, merged_by = ?,
1982 merged_at = ?, updated_at = ?
1983 WHERE id = ?",
1984 )
1985 .bind(&[
1986 landed.commit.as_str().into(),
1987 optional(&landed.previous),
1988 actor.username.as_str().into(),
1989 now.as_str().into(),
1990 now.as_str().into(),
1991 pull.id.as_str().into(),
1992 ])?,
1993 ];
1994 if let Some(issue) = &issue {
1995 statements.push(
1996 self.db
1997 .prepare(
1998 "UPDATE issues
1999 SET state = 'closed', reason = 'completed', resolved_by = ?,
2000 closed_at = ?, updated_at = ?
2001 WHERE id = ?",
2002 )
2003 .bind(&[
2004 pull.number.into(),
2005 now.as_str().into(),
2006 now.as_str().into(),
2007 issue.id.as_str().into(),
2008 ])?,
2009 );
2010 statements.push(
2011 self.db
2012 .prepare(
2013 "UPDATE pulls SET status = 'closed', superseded_by = ?, updated_at = ?
2014 WHERE issue_id = ? AND id != ? AND status IN ('draft', 'open')",
2015 )
2016 .bind(&[
2017 pull.number.into(),
2018 now.as_str().into(),
2019 issue.id.as_str().into(),
2020 pull.id.as_str().into(),
2021 ])?,
2022 );
2023 }
2024 self.db.batch(statements).await?;
2025
2026 self.publish(
2027 "pull.merged",
2028 &repo.id,
2029 actor,
2030 PullEvent {
2031 commit: Some(landed.commit.clone()),
2032 ..Self::pull_event(&pull)
2033 },
2034 )
2035 .await?;
2036 if let Some(issue) = &issue {
2037 self.publish(
2038 "issue.closed",
2039 &repo.id,
2040 actor,
2041 IssueEvent {
2042 reason: Some(IssueReason::Completed.as_str()),
2043 resolved_by: Some(pull.number),
2044 ..Self::issue_event(issue)
2045 },
2046 )
2047 .await?;
2048 }
2049
2050 let who = (actor.id.as_str(), actor.username.as_str());
2051 self.note(&repo.id, pull.number, who, "merged this").await?;
2052 if let Some(issue) = &issue {
2053 self.note(
2054 &repo.id,
2055 issue.number,
2056 who,
2057 &format!("closed this by merging #{}", pull.number),
2058 )
2059 .await?;
2060 }
2061 pull.status = PullStatus::Merged;
2062 pull.head_commit = Some(landed.commit.clone());
2063 pull.merge_base = landed.previous;
2064 pull.merged_by = Some(actor.username.clone());
2065 pull.merged_at = Some(now.clone());
2066 pull.updated_at = now;
2067 Ok(pull)
2068 }
2069
2070 async fn list_active_pulls(&self, a: ViewerArgs) -> Result<Vec<ActivePull>> {
2071 let Some(viewer) = a.viewer else {
2072 return Ok(Vec::new());
2073 };
2074 // The pull requests and their issues, in one round trip: their own,
2075 // and those g1t made for them (Pull::owner).
2076 let author = [JsValue::from(viewer.id.as_str())];
2077 let found = self
2078 .timing
2079 .db(
2080 2,
2081 self.db.batch(vec![
2082 self.db
2083 .prepare(format!(
2084 "SELECT {PULL_COLUMNS} FROM pulls
2085 WHERE COALESCE(requested_by_id, author_id) = ?1 AND status IN ('draft', 'open')
2086 ORDER BY updated_at DESC LIMIT 50"
2087 ))
2088 .bind(&author)?,
2089 self.db
2090 .prepare(format!(
2091 "SELECT {ISSUE_COLUMNS} FROM issues WHERE issues.id IN (
2092 SELECT issue_id FROM pulls
2093 WHERE COALESCE(requested_by_id, author_id) = ?1 AND status IN ('draft', 'open') AND issue_id IS NOT NULL
2094 ORDER BY updated_at DESC LIMIT 50)"
2095 ))
2096 .bind(&author)?,
2097 ]),
2098 )
2099 .await?;
2100 let (Some(found), Some(issues)) = (found.first(), found.get(1)) else {
2101 return Ok(Vec::new());
2102 };
2103 let snapshots = found.results::<Snapshot>()?;
2104 let pulls: Vec<Pull> = found.results::<PullRow>()?.into_iter().map(Pull::from).collect();
2105 let issues: Vec<Issue> = issues.results::<IssueRow>()?.into_iter().map(Issue::from).collect();
2106 let issues = &issues;
2107 // Where each stands: the remembered assessment when there is one,
2108 // and worked out otherwise.
2109 try_join_all(pulls.into_iter().zip(snapshots).map(|(pull, snapshot)| async move {
2110 let issue = pull.issue.and_then(|number| {
2111 issues
2112 .iter()
2113 .find(|issue| issue.repo_id == pull.repo_id && issue.number == number)
2114 .cloned()
2115 });
2116 // Only a pull request g1t is seeing through has a lifecycle.
2117 let lifecycle = if !lifecycle::made_by_g1t(&pull) || snapshot.managed == 0 {
2118 None
2119 } else if let (Some(stage), Some(detail)) = (snapshot.stage, snapshot.stage_detail) {
2120 Some(Lifecycle {
2121 stage,
2122 detail,
2123 revisions: snapshot.revisions,
2124 })
2125 } else {
2126 let behind = self.is_behind(&pull.repo_id, &pull).await?;
2127 self.assess(&pull, &issue, behind)
2128 .await?
2129 .map(|(lifecycle, _)| lifecycle)
2130 };
2131 Ok::<_, worker::Error>(ActivePull {
2132 pull,
2133 issue,
2134 lifecycle,
2135 })
2136 }))
2137 .await
2138 }
2139
2140 // --- Sessions ----------------------------------------------------------
2141
2142 async fn append_session(&self, a: AppendSessionArgs) -> Result<Outcome<Appended>> {
2143 if a.entries.is_empty() {
2144 return Ok(Outcome::Ok(Appended { count: 0 }));
2145 }
2146 if a.entries.len() > MAX_ENTRY_BATCH {
2147 return Ok(Outcome::fail(
2148 FailureCode::Invalid,
2149 format!("Send at most {MAX_ENTRY_BATCH} entries at a time."),
2150 ));
2151 }
2152 let viewer = Some(a.actor.clone());
2153 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
2154 if !pull.is_owned_by(&a.actor.id) {
2155 return Ok(Outcome::fail(
2156 FailureCode::Forbidden,
2157 "Only whoever opened a pull request, or asked g1t for it, can record its session.",
2158 ));
2159 }
2160
2161 let now = rfc3339(now_ms());
2162 let count = a.entries.len() as u32;
2163 let mut statements = Vec::with_capacity(a.entries.len() + 1);
2164 for entry in a.entries {
2165 let kind = serde_json::to_value(entry.kind)?;
2166 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
2167 // Each insert takes the next sequence number itself, so two
2168 // writers appending at once cannot collide.
2169 statements.push(
2170 self.db
2171 .prepare(
2172 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
2173 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
2174 FROM session_entries WHERE pull_id = ?",
2175 )
2176 .bind(&[
2177 pull.id.as_str().into(),
2178 kind.as_str().unwrap_or("note").into(),
2179 text.into(),
2180 optional(&entry.tool),
2181 optional(&entry.commit.or_else(|| pull.head_commit.clone())),
2182 now.as_str().into(),
2183 pull.id.as_str().into(),
2184 ])?,
2185 );
2186 }
2187 statements.push(
2188 self.db
2189 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
2190 .bind(&[now.as_str().into(), pull.id.as_str().into()])?,
2191 );
2192 self.db.batch(statements).await?;
2193 self.publish(
2194 "session.appended",
2195 &pull.repo_id,
2196 &a.actor,
2197 SessionAppended {
2198 pull_id: pull.id.clone(),
2199 repo_id: pull.repo_id.clone(),
2200 number: pull.number,
2201 count,
2202 },
2203 )
2204 .await?;
2205 Ok(Outcome::Ok(Appended { count }))
2206 }
2207
2208 /// Adds entries to a pull request's session, each taking the next
2209 /// sequence number, without announcing it.
2210 pub(crate) async fn append_entries(&self, pull: &Pull, entries: &[NewSessionEntry]) -> Result<()> {
2211 let now = rfc3339(now_ms());
2212 let mut statements = Vec::with_capacity(entries.len());
2213 for entry in entries {
2214 let kind = serde_json::to_value(entry.kind)?;
2215 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
2216 statements.push(
2217 self.db
2218 .prepare(
2219 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
2220 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
2221 FROM session_entries WHERE pull_id = ?",
2222 )
2223 .bind(&[
2224 pull.id.as_str().into(),
2225 kind.as_str().unwrap_or("note").into(),
2226 text.into(),
2227 optional(&entry.tool),
2228 optional(&entry.commit.clone().or_else(|| pull.head_commit.clone())),
2229 now.as_str().into(),
2230 pull.id.as_str().into(),
2231 ])?,
2232 );
2233 }
2234 self.db.batch(statements).await?;
2235 Ok(())
2236 }
2237
2238 async fn read_session(&self, a: ViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
2239 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
2240 let rows = self
2241 .db
2242 .prepare(
2243 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
2244 WHERE pull_id = ? AND seq > ? ORDER BY seq LIMIT ?",
2245 )
2246 .bind(&[pull.id.into(), a.after_seq.into(), SESSION_PAGE.into()])?
2247 .all()
2248 .await?
2249 .results::<SessionRow>()?;
2250 Ok(Outcome::Ok(
2251 rows.into_iter().map(SessionEntry::from).collect(),
2252 ))
2253 }
2254
2255 /// A push moves the head of the pull request it concerns: the one whose
2256 /// fork was pushed to, or the one opened from the branch that moved.
2257 async fn on_event(&self, event: &Event) -> Result<()> {
2258 // A new repository starts with the default labels.
2259 if event.kind == "repo.created"
2260 && let Some(repo_id) = event.repo_id.as_deref()
2261 {
2262 return self.seed_labels(repo_id).await;
2263 }
2264 // Pull requests into the branch that became the default merge into
2265 // the default branch, which is stored as none.
2266 if event.kind == "repo.default_branch_changed"
2267 && let (Some(repo_id), Some(to)) = (event.repo_id.as_deref(), event.data["to"].as_str())
2268 {
2269 self.db
2270 .prepare(
2271 "UPDATE pulls SET base_branch = NULL
2272 WHERE repo_id = ? AND base_branch = ? AND status IN ('draft', 'open')",
2273 )
2274 .bind(&[repo_id.into(), to.into()])?
2275 .run()
2276 .await?;
2277 return Ok(());
2278 }
2279 if event.kind != "git.push" {
2280 return Ok(());
2281 }
2282 let (Some(repo_id), Some(after), Some(git_ref)) = (
2283 event.repo_id.as_deref(),
2284 event.data["after"].as_str(),
2285 event.data["ref"].as_str(),
2286 ) else {
2287 return Ok(());
2288 };
2289 let now = rfc3339(now_ms());
2290 // Who moved it, for rules about the most recent push.
2291 let pusher: JsValue = event.actor.as_deref().map_or(JsValue::NULL, Into::into);
2292 // The head moved, so whatever the checks said no longer applies, and
2293 // whatever step g1t was waiting on has been taken.
2294 let moved = "UPDATE pulls
2295 SET head_commit = ?, updated_at = ?, head_pushed_by = ?, head_pushed_at = ?, check_status = NULL, check_run_id = NULL,
2296 working_on = NULL, working_until = NULL, stalled = NULL";
2297 let active = "status IN ('draft', 'open') AND head_commit IS NOT ?";
2298 let returning =
2299 "RETURNING id, repo_id, number, issue_number, status, author_id, author_name, requested_by_id, requested_by_name";
2300 let mut pulls: Vec<MovedRow> = Vec::new();
2301 // A fork carries its pull request on its default branch.
2302 if event.data["defaultBranch"].as_bool() == Some(true) {
2303 pulls.extend(
2304 self.db
2305 .prepare(format!(
2306 "{moved} WHERE fork_repo_id = ? AND {active} {returning}"
2307 ))
2308 .bind(&[
2309 after.into(),
2310 now.as_str().into(),
2311 pusher.clone(),
2312 now.as_str().into(),
2313 repo_id.into(),
2314 after.into(),
2315 ])?
2316 .all()
2317 .await?
2318 .results::<MovedRow>()?,
2319 );
2320 }
2321 if let Some(branch) = git_ref.strip_prefix("refs/heads/") {
2322 pulls.extend(
2323 self.db
2324 .prepare(format!(
2325 "{moved} WHERE repo_id = ? AND source_branch = ? AND {active} {returning}"
2326 ))
2327 .bind(&[
2328 after.into(),
2329 now.as_str().into(),
2330 pusher.clone(),
2331 now.as_str().into(),
2332 repo_id.into(),
2333 branch.into(),
2334 after.into(),
2335 ])?
2336 .all()
2337 .await?
2338 .results::<MovedRow>()?,
2339 );
2340 }
2341 // What each now changes, so overlaps show while the work is under way.
2342 for moved in &pulls {
2343 if let Some(mut pull) = self.pull_by_id(&moved.id).await? {
2344 pull.files = self.refresh_files(&pull).await?;
2345 // Owners of files it now changes are asked too.
2346 self.refresh_code_owners(&pull).await;
2347 }
2348 }
2349 // A merge that was waiting for this push to bring it up to date.
2350 for moved in &pulls {
2351 self.land_if_requested(&moved.id).await?;
2352 }
2353 // Whether each still merges cleanly, and, when a default branch
2354 // moved, every open pull request into it.
2355 let moved_ids: Vec<String> = pulls.iter().map(|pull| pull.id.clone()).collect();
2356 let into = match event.data["defaultBranch"].as_bool() {
2357 Some(true) => mergeability::Moved::DefaultBranch,
2358 _ => match git_ref.strip_prefix("refs/heads/") {
2359 Some(branch) => mergeability::Moved::Branch(branch),
2360 None => mergeability::Moved::Nothing,
2361 },
2362 };
2363 self.after_push(repo_id, into, &moved_ids).await;
2364 // A draft is announced when it is marked ready instead.
2365 for pull in pulls
2366 .into_iter()
2367 .filter(|pull| pull.status == PullStatus::Open)
2368 {
2369 self.publish_as(
2370 "pull.updated",
2371 &pull.repo_id,
2372 event.actor.clone(),
2373 PullEvent {
2374 author: Some(g1t_contracts::credentials::Principal { id: pull.author_id, username: pull.author_name }),
2375 requested_by: pull
2376 .requested_by_id
2377 .zip(pull.requested_by_name)
2378 .map(|(id, username)| g1t_contracts::credentials::Principal { id, username }),
2379 pull_id: pull.id,
2380 repo_id: pull.repo_id.clone(),
2381 number: pull.number,
2382 issue: pull.issue_number,
2383 commit: Some(after.to_owned()),
2384 ..PullEvent::default()
2385 }
2386 .carrying(&event.data),
2387 )
2388 .await?;
2389 }
2390 Ok(())
2391 }
2392}
2393
2394/// An event's data that carries on what caused the event it follows from.
2395trait Carrying: serde::Serialize + Sized {
2396 /// As JSON, marked as a workflow job's doing when `cause` was.
2397 fn carrying(self, cause: &serde_json::Value) -> serde_json::Value {
2398 g1t_contracts::events::carried(self, cause)
2399 }
2400}
2401
2402impl Carrying for PullEvent {}
2403
2404fn service(env: &Env) -> Result<Work> {
2405 Ok(Work {
2406 db: env.d1("DB")?,
2407 identity: env.service("IDENTITY")?,
2408 repos: env.service("REPOS")?,
2409 events: env.service("EVENTS")?,
2410 actions: env.service("ACTIONS")?,
2411 timing: g1t_kit::d1::Timing::default(),
2412 prefetched: std::cell::RefCell::new(None),
2413 known_repos: std::cell::RefCell::new(std::collections::HashMap::new()),
2414 })
2415}
2416
2417#[event(fetch)]
2418async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
2419 let Some(method) = rpc_method(&request) else {
2420 return Response::error("Not found", 404);
2421 };
2422 // A replica near the caller when it asks for one (crates/kit/src/d1.rs).
2423 let (db, served) = g1t_kit::d1::open(&env, "DB", &request)?;
2424 let body: serde_json::Value = request.json().await?;
2425 let mut work = service(&env)?;
2426 work.db = db;
2427
2428 let answered = match method.as_str() {
2429 "open_issue" => reply(&work.open_issue(args(body)?).await?),
2430 "delegate_issue" => reply(&work.delegate_issue(args(body)?).await?),
2431 "report_confidence" => reply(&work.report_confidence(args(body)?).await?),
2432 "list_issues" => reply(&work.list_issues(args(body)?).await?),
2433 "get_issue" => reply(&work.get_issue(args(body)?).await?),
2434 "update_issue" => reply(&work.update_issue(args(body)?).await?),
2435 "close_issue" => reply(&work.close_issue(args(body)?).await?),
2436 "reopen_issue" => reply(&work.reopen_issue(args(body)?).await?),
2437 "list_labels" => reply(&work.list_labels(args(body)?).await?),
2438 "save_label" => reply(&work.save_label(args(body)?).await?),
2439 "delete_label" => reply(&work.delete_label(args(body)?).await?),
2440 "add_default_labels" => reply(&work.add_default_labels(args(body)?).await?),
2441 "set_labels" => reply(&work.set_labels(args(body)?).await?),
2442 "list_milestones" => reply(&work.list_milestones(args(body)?).await?),
2443 "get_milestone" => reply(&work.get_milestone(args(body)?).await?),
2444 "save_milestone" => reply(&work.save_milestone(args(body)?).await?),
2445 "delete_milestone" => reply(&work.delete_milestone(args(body)?).await?),
2446 "counts" => reply(&work.counts(args(body)?).await?),
2447 "add_comment" => reply(&work.add_comment(args(body)?).await?),
2448 "start_checks" => reply(&work.start_checks(args(body)?).await?),
2449 "seen_checks" => reply(&work.seen_checks(args(body)?).await?),
2450 "report_checks" => reply(&work.report_checks(args(body)?).await?),
2451 "set_commit_status" => reply(&work.set_commit_status(args(body)?).await?),
2452 "create_commit_status" => reply(&work.create_commit_status(args(body)?).await?),
2453 "commit_statuses" => reply(&work.commit_statuses(args(body)?).await?),
2454 "combined_status" => reply(&work.combined_status(args(body)?).await?),
2455 "create_check_run" => reply(&work.create_check_run(args(body)?).await?),
2456 "update_check_run" => reply(&work.update_check_run(args(body)?).await?),
2457 "get_check_run" => reply(&work.get_check_run(args(body)?).await?),
2458 "check_run_annotations" => reply(&work.check_run_annotations(args(body)?).await?),
2459 "ref_check_runs" => reply(&work.ref_check_runs(args(body)?).await?),
2460 "ref_check_suites" => reply(&work.ref_check_suites(args(body)?).await?),
2461 "get_check_suite" => reply(&work.get_check_suite(args(body)?).await?),
2462 "rerequest_check_run" => reply(&work.rerequest_check_run(args(body)?).await?),
2463 "rerequest_check_suite" => reply(&work.rerequest_check_suite(args(body)?).await?),
2464 "request_check_action" => reply(&work.request_check_action(args(body)?).await?),
2465 "commit_checks" => reply(&work.commit_checks(args(body)?).await?),
2466 "start_review" => reply(&work.start_review(args(body)?).await?),
2467 "advance" => reply(&work.advance(args(body)?).await?),
2468 "stall" => reply(&work.stall(args(body)?).await?),
2469 "managed_pulls" => reply(&work.managed_pulls(args(body)?).await?),
2470 "queue" => reply(&work.queue(args(body)?).await?),
2471 "queue_build" => reply(&work.queue_build(args(body)?).await?),
2472 "report_queue" => reply(&work.report_queue(args(body)?).await?),
2473 "remove_from_queue" => reply(&work.remove_from_queue(args(body)?).await?),
2474 "message_agent" => reply(&work.message_agent(args(body)?).await?),
2475 "locate_pull" => reply(&work.locate_pull(args(body)?).await?),
2476 "inbox_subject" => reply(&work.inbox_subject(args(body)?).await?),
2477 "answer_message" => reply(&work.answer_message(args(body)?).await?),
2478 "take_messages" => reply(&work.take_messages(args(body)?).await?),
2479 "wake_for_messages" => reply(&work.wake_for_messages(args(body)?).await?),
2480 "catch_up_job" => reply(&work.catch_up_job(args(body)?).await?),
2481 "get_settings" => reply(&work.get_settings(args(body)?).await?),
2482 "codeowners_errors" => reply(&work.codeowners_errors(args(body)?).await?),
2483 "update_settings" => reply(&work.update_settings(args(body)?).await?),
2484 "report_review" => reply(&work.report_review(args(body)?).await?),
2485 "open_pull" => reply(&work.open_pull(args(body)?).await?),
2486 "list_pulls" => reply(&work.list_pulls(args(body)?).await?),
2487 "pulls_for_repos" => reply(&work.pulls_for_repos(args(body)?).await?),
2488 "get_pull" => reply(&work.get_pull(args(body)?).await?),
2489 "update_pull" => reply(&work.update_pull(args(body)?).await?),
2490 "catch_up_pull" => reply(&work.catch_up_pull(args(body)?).await?),
2491 "ready_pull" => reply(&work.ready_pull(args(body)?).await?),
2492 "close_pull" => reply(&work.close_pull(args(body)?).await?),
2493 "merge_pull" => reply(&work.merge_pull(args(body)?).await?),
2494 "list_active_pulls" => reply(&work.list_active_pulls(args(body)?).await?),
2495 "by_author" => reply(&work.by_author(args(body)?).await?),
2496 "start_plan" => reply(&work.start_plan(args(body)?).await?),
2497 "report_plan" => reply(&work.report_plan(args(body)?).await?),
2498 "get_plan" => reply(&work.get_plan(args(body)?).await?),
2499 "list_plans" => reply(&work.list_plans(args(body)?).await?),
2500 "apply_plan" => reply(&work.apply_plan(args(body)?).await?),
2501 "queue_issue" => reply(&work.queue_issue(args(body)?).await?),
2502 "ready_issues" => reply(&work.ready_issues(args(body)?).await?),
2503 "list_assigned_issues" => reply(&work.list_assigned_issues(args(body)?).await?),
2504 "append_session" => reply(&work.append_session(args(body)?).await?),
2505 "read_session" => reply(&work.read_session(args(body)?).await?),
2506 // Agents at work, their sessions, and memory (runs.rs, memory.rs).
2507 "open_run" => reply(&work.open_run(args(body)?).await?),
2508 "report_run" => reply(&work.report_run(args(body)?).await?),
2509 "stop_run" => reply(&work.stop_run(args(body)?).await?),
2510 "list_runs" => reply(&work.list_runs(args(body)?).await?),
2511 "get_run" => reply(&work.get_run(args(body)?).await?),
2512 "list_sessions" => reply(&work.list_sessions(args(body)?).await?),
2513 "get_session" => reply(&work.get_session(args(body)?).await?),
2514 "list_memories" => reply(&work.list_memories(args(body)?).await?),
2515 "add_memory" => reply(&work.add_memory(args(body)?).await?),
2516 "update_memory" => reply(&work.update_memory(args(body)?).await?),
2517 "delete_memory" => reply(&work.delete_memory(args(body)?).await?),
2518 "recall" => reply(&work.recall(args(body)?).await?),
2519 "memory_context" => reply(&work.memory_context(args(body)?).await?),
2520 // What agents may do in a sandbox (guardrails.rs).
2521 "get_guardrails" => reply(&work.get_guardrails(args(body)?).await?),
2522 "update_guardrails" => reply(&work.update_guardrails(args(body)?).await?),
2523 "run_guardrails" => reply(&work.run_guardrails(args(body)?).await?),
2524 // Plan caps the runner applies (compute.rs).
2525 "active_agents" => reply(&work.active_agents(args(body)?).await?),
2526 "issue_spend" => reply(&work.issue_spend(args(body)?).await?),
2527 "wait_for_slot" => reply(&work.wait_for_slot(args(body)?).await?),
2528 "agent_comment" => reply(&work.agent_comment(args(body)?).await?),
2529 "add_wait" => reply(&work.add_wait(args(body)?).await?),
2530 "waiting_workspaces" => reply(&work.waiting_workspaces(args(body)?).await?),
2531 "take_wait" => reply(&work.take_wait(args(body)?).await?),
2532 // The runs whose sandboxes stop with their repository (retired.rs).
2533 "runs_in_repo" => reply(&work.runs_in_repo(args(body)?).await?),
2534 "run_cost" => reply(&work.run_cost(args(body)?).await?),
2535 "start_mergecheck" => reply(&work.start_mergecheck(args(body)?).await?),
2536 "report_mergecheck" => reply(&work.report_mergecheck(args(body)?).await?),
2537 // Memory that fills itself, and its review queue (capture.rs).
2538 method if capture::METHODS.contains(&method) => capture::dispatch(&work, method, body).await,
2539 // @g1t in comments, and the label rule (mentions.rs).
2540 "take_mention" => reply(&work.take_mention(args(body)?).await?),
2541 "mention_revision" => reply(&work.mention_revision(args(body)?).await?),
2542 "reply_mention" => reply(&work.reply_mention(args(body)?).await?),
2543 "get_agent_rules" => reply(&work.get_agent_rules(args(body)?).await?),
2544 "set_agent_rules" => reply(&work.set_agent_rules(args(body)?).await?),
2545 // Rulesets (rulesets.rs): kept here, enforced here on merge and by
2546 // repos on push.
2547 "list_rulesets" => reply(&work.list_rulesets(args(body)?).await?),
2548 "get_ruleset" => reply(&work.get_ruleset(args(body)?).await?),
2549 "save_ruleset" => reply(&work.save_ruleset(args(body)?).await?),
2550 "delete_ruleset" => reply(&work.delete_ruleset(args(body)?).await?),
2551 "effective_rules" => reply(&work.effective_rules(args(body)?).await?),
2552 "rule_evaluations" => reply(&work.rule_evaluations(args(body)?).await?),
2553 "ref_rules" => reply(&work.ref_rules(args(body)?).await?),
2554 "record_evaluations" => reply(&work.record_evaluations(args(body)?).await?),
2555 "set_requires_pull_request" => reply(&work.set_requires_pull_request(args(body)?).await?),
2556 _ => Response::error("Unknown method", 404),
2557 };
2558 served.finish_timed(answered, &work.timing)
2559}
2560
2561/// Events from the bus, delivered on this service's own queue.
2562#[event(queue)]
2563async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
2564 let work = service(&env)?;
2565 for message in batch.messages()? {
2566 // A workspace renamed: its agent runs and memory move to the slug it has now.
2567 if g1t_kit::rename::on_event(&env, &env.d1("DB")?, message.body(), &[memory::RENAMED, guardrails::RENAMED, rulesets::RENAMED].concat()).await? {
2568 message.ack();
2569 continue;
2570 }
2571 // A repository renamed or transferred: its runs, memory, guardrails
2572 // and runs waiting for a slot follow.
2573 if g1t_kit::transfer::on_event(&env, &env.d1("DB")?, message.body(), &[memory::TRANSFERRED, guardrails::TRANSFERRED, retired::WAITS_MOVED].concat()).await? {
2574 message.ack();
2575 continue;
2576 }
2577 // A workspace deleted: what it kept for itself goes.
2578 if g1t_kit::deleted::on_event(&env.d1("DB")?, message.body(), memory::DELETED).await? {
2579 message.ack();
2580 continue;
2581 }
2582 // A repository deleted, archived or purged, or a branch renamed (retired.rs).
2583 if work.on_retired(message.body()).await? {
2584 message.ack();
2585 continue;
2586 }
2587 capture::on_event(&work, message.body()).await;
2588 work.on_event(message.body()).await?;
2589 message.ack();
2590 }
2591 Ok(())
2592}
2593
2594/// The rules that once read a pull request's author read its owner now:
2595/// whoever asked g1t for it, or its author. For each, the person who asked
2596/// is held to what an author was, and g1t's agent (a token it works with)
2597/// gains nothing by being the author.
2598#[cfg(test)]
2599mod owner_rules {
2600 use super::*;
2601 use crate::rows::stored::{ASKER, G1T, pull};
2602 use g1t_contracts::identity::AGENT_ID;
2603
2604 const SOMEONE: &str = "usr_2";
2605
2606 #[test]
2607 fn no_self_approval() {
2608 // add_comment refuses a verdict on one that is theirs.
2609 let made = pull(G1T, Some(ASKER));
2610 assert!(made.is_owned_by(ASKER.0), "the person who asked cannot approve it");
2611 assert!(!made.is_owned_by(AGENT_ID), "g1t's review agent still gives its verdict");
2612 assert!(!made.is_owned_by(SOMEONE));
2613 }
2614
2615 #[test]
2616 fn what_an_author_could_do_without_a_role() {
2617 // manageable_pull (update, ready, close), catch_up_pull on a fork,
2618 // append_session, and steering with message_agent: theirs to do.
2619 let made = pull(G1T, Some(ASKER));
2620 assert!(made.is_owned_by(ASKER.0));
2621 assert!(!made.is_owned_by(SOMEONE), "anyone else still needs the role");
2622 assert!(!made.is_owned_by(AGENT_ID), "being its author gives g1t's tokens nothing more");
2623 }
2624
2625 #[test]
2626 fn nobody_is_asked_to_review_what_they_asked_for() {
2627 // update_pull drops the owner from the reviewers asked.
2628 let made = pull(G1T, Some(ASKER));
2629 let mut reviewers = vec!["syntaqx".to_owned(), "ana".to_owned()];
2630 reviewers.retain(|name| *name != made.owner().username);
2631 assert_eq!(reviewers, ["ana"]);
2632 }
2633
2634 #[test]
2635 fn sandboxes_act_as_whoever_asked() {
2636 // LifecycleJob, ReviewJob, MergecheckJob and the merge queue's job
2637 // carry who the sandbox's credential acts for: a real account.
2638 let made = pull(G1T, Some(ASKER));
2639 assert_eq!(made.owner().id, ASKER.0);
2640 let acts_as = made.requested_by.unwrap_or(made.author);
2641 assert_eq!(acts_as.id, ASKER.0);
2642 // g1t's own work, which nobody asked for, acts as g1t, as before.
2643 let own = pull(("g1t", "g1t"), None);
2644 assert_eq!(own.requested_by.unwrap_or(own.author).id, "g1t");
2645 }
2646
2647 #[test]
2648 fn events_name_g1t_and_whoever_asked() {
2649 let made = pull(G1T, Some(ASKER));
2650 let event = serde_json::to_value(Work::pull_event(&made)).unwrap();
2651 assert_eq!(event["author"], serde_json::json!({ "id": AGENT_ID, "username": "g1t" }));
2652 assert_eq!(event["requestedBy"], serde_json::json!({ "id": "usr_1", "username": "syntaqx" }));
2653 let own = serde_json::to_value(Work::pull_event(&pull(ASKER, None))).unwrap();
2654 assert!(own.get("requestedBy").is_none());
2655 }
2656}