Skip to content

g1t/services/work/src/lib.rs

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