Skip to content

g1t/services/work/src/lib.rs

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