Skip to content

g1t/services/work/src/lib.rs

2,628 lines109,745 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 // Unnamed, the change is its author's, unless an agent opened it.
1058 let agent = match a.agent.trim() {
1059 "" if g1t_contracts::rules::is_agent(&a.actor) => "agent",
1060 "" => a.actor.username.as_str(),
1061 agent => agent,
1062 };
1063 let runtime = match a.runtime {
1064 Runtime::Hosted => "hosted",
1065 Runtime::External => "external",
1066 };
1067
1068 let now = now_ms();
1069 let id = new_id("pr", now);
1070 let branch = a
1071 .branch
1072 .as_deref()
1073 .map(str::trim)
1074 .filter(|branch| !branch.is_empty());
1075 // The branch it merges into: the default branch unless another is
1076 // asked for, which has to exist.
1077 let base = a.base.as_deref().and_then(|base| stored_base(base, &repo));
1078 if let Some(base) = &base {
1079 if branch == Some(base.as_str()) {
1080 return Ok(Outcome::fail(
1081 FailureCode::Invalid,
1082 format!("A pull request cannot merge {base} into itself. Choose another base."),
1083 ));
1084 }
1085 let exists: Option<String> = g1t_kit::call(
1086 &self.repos,
1087 "head",
1088 &HeadArgs { repo_id: repo.id.clone(), branch: base.clone() },
1089 )
1090 .await?;
1091 if exists.is_none() {
1092 return Ok(Outcome::fail(FailureCode::NotFound, format!("There is no branch named {base} to merge into.")));
1093 }
1094 }
1095 let base_name = base.clone().unwrap_or_else(|| repo.default_branch.clone());
1096 // The change is on a branch already pushed to the repository, or
1097 // will be made in a fork created for this pull request.
1098 let (fork, head) = match branch {
1099 Some(branch) => {
1100 if branch == base_name {
1101 return Ok(Outcome::fail(
1102 FailureCode::Invalid,
1103 format!("Choose a branch other than {branch}."),
1104 ));
1105 }
1106 let head: Option<String> = g1t_kit::call(
1107 &self.repos,
1108 "head",
1109 &HeadArgs {
1110 repo_id: repo.id.clone(),
1111 branch: branch.to_owned(),
1112 },
1113 )
1114 .await?;
1115 let Some(head) = head else {
1116 return Ok(Outcome::fail(
1117 FailureCode::NotFound,
1118 format!("There is no branch named {branch}. Push it first."),
1119 ));
1120 };
1121 // One open pull request for each branch and base.
1122 let existing = self
1123 .db
1124 .prepare(
1125 "SELECT number AS n FROM pulls
1126 WHERE repo_id = ? AND source_branch = ? AND status IN ('draft', 'open')
1127 AND base_branch IS ?",
1128 )
1129 .bind(&[repo.id.as_str().into(), branch.into(), optional(&base)])?
1130 .first::<NumberRow>(None)
1131 .await?;
1132 if let Some(existing) = existing {
1133 return Ok(Outcome::fail(
1134 FailureCode::Conflict,
1135 format!("Pull request #{} is already open from {branch} into {base_name}.", existing.n),
1136 ));
1137 }
1138 (None, Some(head))
1139 }
1140 None => {
1141 let fork: Outcome<Repo> = g1t_kit::call(
1142 &self.repos,
1143 "fork_for_pull",
1144 &ForkArgs {
1145 source_id: repo.id.clone(),
1146 pull_id: id.clone(),
1147 actor: a.actor.clone(),
1148 },
1149 )
1150 .await?;
1151 (Some(check!(fork)), None)
1152 }
1153 };
1154 // A branch already holds the work, so its pull request is ready for
1155 // review from the start; one with a fork starts as a draft.
1156 let status = if branch.is_some() { "open" } else { "draft" };
1157 let body = Some(a.body.trim().to_owned()).filter(|body| !body.is_empty());
1158
1159 let number = self.next_number(&repo.id).await?;
1160 let timestamp = rfc3339(now);
1161 // A change g1t makes is g1t's, for whoever asked for it. This is
1162 // what lifecycle::made_by_g1t reads back.
1163 let by_g1t = matches!(a.runtime, Runtime::Hosted) && agent == reviews::AGENT_NAME && fork.is_some();
1164 let (author, requested_by) = authorship(&a.actor, by_g1t);
1165 self.db
1166 .prepare(
1167 "INSERT INTO pulls
1168 (id, repo_id, number, issue_id, issue_number, title, body, agent, runtime,
1169 status, fork_repo_id, fork_namespace, fork_name, source_branch, head_commit,
1170 author_id, author_name, requested_by_id, requested_by_name, created_at, updated_at,
1171 base_branch)
1172 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
1173 )
1174 .bind(&[
1175 id.as_str().into(),
1176 repo.id.as_str().into(),
1177 number.into(),
1178 optional(&issue.as_ref().map(|issue| issue.id.clone())),
1179 optional_number(issue.as_ref().map(|issue| issue.number)),
1180 title.into(),
1181 optional(&body),
1182 agent.into(),
1183 runtime.into(),
1184 status.into(),
1185 optional(&fork.as_ref().map(|fork| fork.id.clone())),
1186 optional(&fork.as_ref().map(|fork| fork.namespace.clone())),
1187 optional(&fork.as_ref().map(|fork| fork.name.clone())),
1188 optional(&branch.map(str::to_owned)),
1189 optional(&head),
1190 author.id.as_str().into(),
1191 author.username.as_str().into(),
1192 optional(&requested_by.as_ref().map(|user| user.id.clone())),
1193 optional(&requested_by.as_ref().map(|user| user.username.clone())),
1194 timestamp.as_str().into(),
1195 timestamp.as_str().into(),
1196 optional(&base),
1197 ])?
1198 .run()
1199 .await?;
1200 let Some(mut pull) = self.pull(&repo.id, number).await? else {
1201 return Ok(no_pull());
1202 };
1203 fill_base(&mut pull, &repo);
1204 self.manage(&pull).await?;
1205 // Someone is on it now, so it is no longer waiting for an agent.
1206 if let Some(issue) = pull.issue {
1207 self.db
1208 .prepare("UPDATE issues SET queued_by = NULL WHERE repo_id = ? AND number = ?")
1209 .bind(&[repo.id.as_str().into(), issue.into()])?
1210 .run()
1211 .await?;
1212 }
1213 if let Some(issue) = pull.issue {
1214 let text = if lifecycle::made_by_g1t(&pull) {
1215 format!("assigned this to g1t, which opened #{}", pull.number)
1216 } else {
1217 format!("opened #{} for this", pull.number)
1218 };
1219 self.note(&repo.id, issue, (&a.actor.id, &a.actor.username), &text)
1220 .await?;
1221 }
1222 self.publish(
1223 "pull.opened",
1224 &repo.id,
1225 &a.actor,
1226 PullEvent {
1227 agent: Some(pull.agent.clone()),
1228 base: pull.base.clone(),
1229 ..Self::pull_event(&pull)
1230 },
1231 )
1232 .await?;
1233 // Its code owners asked to review (codeowners.rs).
1234 self.refresh_code_owners(&pull).await;
1235 let pull = self.pull(&repo.id, number).await?.unwrap_or(pull);
1236 Ok(Outcome::Ok(pull))
1237 }
1238
1239 async fn list_pulls(&self, a: ListPullsArgs) -> Result<Outcome<Vec<Pull>>> {
1240 let filter = match a.state {
1241 Some(State::Open) => "AND status IN ('draft', 'open')",
1242 Some(State::Closed) => "AND status IN ('merged', 'closed')",
1243 None => "",
1244 };
1245 let label = a
1246 .label
1247 .map(|label| label.split_whitespace().collect::<Vec<_>>().join(" ").to_lowercase())
1248 .filter(|label| !label.is_empty());
1249 let milestone = a.milestone;
1250 let read = |repo_id: String| {
1251 let label = label.clone();
1252 async move {
1253 let query = self
1254 .db
1255 .prepare(format!(
1256 "SELECT {PULL_COLUMNS} FROM pulls WHERE repo_id = ? {filter}
1257 AND (? IS NULL OR EXISTS
1258 (SELECT 1 FROM json_each(pulls.labels) WHERE json_each.value = ?))
1259 AND (? IS NULL OR milestone = ?)
1260 ORDER BY number DESC LIMIT ?"
1261 ))
1262 .bind(&[
1263 repo_id.into(),
1264 optional(&label),
1265 optional(&label),
1266 optional_number(milestone),
1267 optional_number(milestone),
1268 LIST_PAGE.into(),
1269 ])?;
1270 self.timing.db(1, query.all()).await?.results::<PullRow>()
1271 }
1272 };
1273 let (repo, rows) = check!(self.repo_then(&a.repo, &a.viewer, read).await?);
1274 let base = a.base.map(|base| base.trim().to_owned()).filter(|base| !base.is_empty());
1275 Ok(Outcome::Ok(
1276 rows.into_iter()
1277 .map(|row| {
1278 let mut pull = Pull::from(row);
1279 fill_base(&mut pull, &repo);
1280 pull
1281 })
1282 .filter(|pull| base.as_deref().is_none_or(|base| pull.base.as_deref() == Some(base)))
1283 .collect(),
1284 ))
1285 }
1286
1287 /// `pulls_for_repos`: what `list_pulls` gives, open and closed, for many
1288 /// repositories at once: one access check with repos for all of them
1289 /// and one query, instead of two of each per repository.
1290 async fn pulls_for_repos(&self, a: PullsForReposArgs) -> Result<Vec<RepoPulls>> {
1291 let ids: Vec<String> = a.repo_ids.into_iter().take(MAX_PULLS_FOR_REPOS).collect();
1292 if ids.is_empty() {
1293 return Ok(Vec::new());
1294 }
1295 let limit = a.limit.clamp(1, LIST_PAGE);
1296 // The rows are read beside the access check, for every id asked
1297 // about; those of repositories the viewer cannot read are dropped.
1298 let asked = serde_json::to_string(&ids)?;
1299 let check = ReadableArgs { ids, viewer: a.viewer };
1300 let readable = self.timing.rpc(g1t_kit::call::<_, Vec<Repo>>(&self.repos, "readable", &check));
1301 let (readable, rows) = try_join(readable, self.timing.db(1, self.newest_pulls(asked, limit))).await?;
1302 if readable.is_empty() {
1303 return Ok(Vec::new());
1304 }
1305 let mut answer: Vec<RepoPulls> = readable
1306 .iter()
1307 .map(|repo| RepoPulls { repo_id: repo.id.clone(), open: Vec::new(), closed: Vec::new() })
1308 .collect();
1309 for pull in rows.into_iter().map(Pull::from) {
1310 let Some(entry) = answer.iter_mut().find(|entry| entry.repo_id == pull.repo_id) else {
1311 continue;
1312 };
1313 match pull.status {
1314 PullStatus::Draft | PullStatus::Open => entry.open.push(pull),
1315 PullStatus::Merged | PullStatus::Closed => entry.closed.push(pull),
1316 }
1317 }
1318 for entry in &mut answer {
1319 entry.open.sort_by_key(|pull| std::cmp::Reverse(pull.number));
1320 entry.closed.sort_by_key(|pull| std::cmp::Reverse(pull.number));
1321 }
1322 Ok(answer)
1323 }
1324
1325 /// The newest `limit` of each repository's open (draft or open) and
1326 /// closed (merged or closed) pull requests, for the ids in `ids` (JSON).
1327 async fn newest_pulls(&self, ids: String, limit: u32) -> Result<Vec<PullRow>> {
1328 self.db
1329 .prepare(format!(
1330 "SELECT * FROM (
1331 SELECT {PULL_COLUMNS}, ROW_NUMBER() OVER (
1332 PARTITION BY pulls.repo_id, pulls.status IN ('draft', 'open') ORDER BY pulls.number DESC
1333 ) AS place
1334 FROM pulls WHERE pulls.repo_id IN (SELECT value FROM json_each(?1))
1335 ) WHERE place <= ?2"
1336 ))
1337 .bind(&[ids.into(), limit.into()])?
1338 .all()
1339 .await?
1340 .results::<PullRow>()
1341 }
1342
1343 async fn get_pull(&self, a: ViewArgs) -> Result<Outcome<PullDetail>> {
1344 let number = a.number;
1345 // Every row the page and the lifecycle read, in one batch started
1346 // beside the access check; the helpers below read from it.
1347 let namespace = a.repo.namespace.clone();
1348 let read = |repo_id: String| self.prefetch_pull(repo_id, namespace.clone(), number);
1349 let Outcome::Ok((repo, Some(found))) = self.repo_then(&a.repo, &a.viewer, read).await? else {
1350 return Ok(no_pull());
1351 };
1352 let Some(row) = found.first::<PullRow>(prefetch::Slot::Pull)? else {
1353 return Ok(no_pull());
1354 };
1355 let stored = found.first::<StoredBehind>(prefetch::Slot::Pull)?;
1356 let issue = found.first::<IssueRow>(prefetch::Slot::Issue)?.map(Issue::from);
1357 let comments: Vec<Comment> =
1358 found.rows::<CommentRow>(prefetch::Slot::Comments)?.into_iter().map(Comment::from).collect();
1359 self.keep_prefetched(Some(found));
1360 let default_branch = repo.default_branch.clone();
1361 let detail = self.pull_detail(repo, Pull::from(row), issue, comments, stored, &a.viewer).await;
1362 self.keep_prefetched(None);
1363 // The branch it merges into, named, for whoever reads it.
1364 Ok(match detail? {
1365 Outcome::Ok(mut detail) => {
1366 if detail.pull.base.as_deref().is_none_or(str::is_empty) {
1367 detail.pull.base = Some(default_branch);
1368 }
1369 Outcome::Ok(detail)
1370 }
1371 failed => failed,
1372 })
1373 }
1374
1375 async fn pull_detail(
1376 &self,
1377 repo: Repo,
1378 mut pull: Pull,
1379 issue: Option<Issue>,
1380 comments: Vec<Comment>,
1381 stored: Option<StoredBehind>,
1382 viewer: &Viewer,
1383 ) -> Result<Outcome<PullDetail>> {
1384 // Whether it is behind, as worked out with its mergeability on the
1385 // last push to either side (mergeability.rs), when that was for
1386 // its head as it is now; otherwise asked of the repos service.
1387 let known_behind = stored.and_then(|stored| stored.for_head(pull.head_commit.as_deref()));
1388 // Worked out on each push; this covers a pull request from before
1389 // that was recorded.
1390 if pull.files.is_empty() && pull.head_commit.is_some() {
1391 pull.files = self.refresh_files(&pull).await?;
1392 }
1393 // Everything else at once: none of it depends on the rest, and each
1394 // is a round trip of its own.
1395 let standing = async {
1396 // Mergeability first: where g1t sees a pull request through, a
1397 // conflict decides its next step.
1398 let behind = async {
1399 match known_behind {
1400 Some(behind) => Ok(behind),
1401 None => {
1402 let behind = self.is_behind(&repo.id, &pull).await?;
1403 // Kept for the next view when the mergeability on
1404 // record is for this head: a pull request from
1405 // before `behind` was kept asks once.
1406 if let Some(head) = pull.head_commit.as_deref()
1407 && pull.status.is_active()
1408 {
1409 self.db
1410 .prepare(
1411 "UPDATE pulls SET behind = ?1
1412 WHERE id = ?2 AND behind IS NULL AND mergeable_key LIKE ?3 || '..%'",
1413 )
1414 .bind(&[u32::from(behind).into(), pull.id.as_str().into(), head.into()])?
1415 .run()
1416 .await?;
1417 }
1418 Ok(behind)
1419 }
1420 }
1421 };
1422 let (merge, behind) = try_join(self.mergeability(&pull), behind).await?;
1423 let assessed = self.assess_with_confidence(&pull, &issue, behind).await?;
1424 let confidence = assessed.as_ref().and_then(|(_, _, confidence)| confidence.clone());
1425 let lifecycle = assessed.map(|(lifecycle, _, _)| lifecycle);
1426 Ok::<_, worker::Error>((behind, (lifecycle, confidence), merge))
1427 };
1428 let (((behind, (lifecycle, confidence), (mergeable, conflicts)), (landing, stalled), comments), (checks, overlaps, review_pending)) =
1429 try_join(
1430 try_join3(standing, self.landing_state(&pull.id), async { Ok(comments) }),
1431 try_join3(
1432 self.latest_checks(&pull.id),
1433 self.overlaps(&pull),
1434 self.review_pending(&pull.id),
1435 ),
1436 )
1437 .await?;
1438 // As just worked out, rather than as it was read.
1439 if confidence.is_some() {
1440 pull.confidence = confidence;
1441 }
1442 // The rules of the branch it merges into, as they stack, and which
1443 // of them it does not meet yet, for whoever is looking.
1444 let (statuses, settings, gate) = try_join3(
1445 self.statuses(&repo.id, pull.head_commit.as_deref()),
1446 self.settings_on(&repo, &pull),
1447 async {
1448 if pull.status.is_active() {
1449 self.merge_gate(&repo, &pull, viewer.as_ref(), false, true).await.map(Some)
1450 } else {
1451 Ok(None)
1452 }
1453 },
1454 )
1455 .await?;
1456 let code_owners = self.pull_code_owners(&pull, &comments, &settings).await?;
1457 Ok(Outcome::Ok(PullDetail {
1458 required_checks: required_checks(&settings.required_checks, &statuses),
1459 rules: gate.map(|gate| rulesets::merge_rules(&gate.judged, &gate.requirements, pull.base_branch(&repo.default_branch) == repo.default_branch)),
1460 code_owners,
1461 comments,
1462 checks,
1463 overlaps,
1464 behind,
1465 review_pending,
1466 lifecycle,
1467 landing,
1468 stalled,
1469 messages: self.messages(&pull.id).await?,
1470 statuses,
1471 mergeable,
1472 conflicts,
1473 earlier_checks: self.earlier_checks(&pull.id).await?,
1474 issue,
1475 pull,
1476 }))
1477 }
1478
1479 /// The pull request, if it is still active and `actor` opened it or
1480 /// may triage the repository's pull requests.
1481 async fn manageable_pull(
1482 &self,
1483 actor: &User,
1484 path: &RepoPath,
1485 number: u32,
1486 ) -> Result<Outcome<Pull>> {
1487 let (repo, pull) = check!(self.pull_at(path, number, &Some(actor.clone())).await?);
1488 check!(writable(&repo));
1489 if !pull.is_owned_by(&actor.id) {
1490 check!(allowed(Some(actor), &repo, Capability::Triage));
1491 }
1492 if !pull.status.is_active() {
1493 return Ok(Outcome::fail(
1494 FailureCode::Conflict,
1495 format!("This pull request is already {}.", pull.status.as_str()),
1496 ));
1497 }
1498 Ok(Outcome::Ok(pull))
1499 }
1500
1501 /// Brings a pull request up to date with the default branch without a
1502 /// sandbox, where the repos service can do that safely. Whoever could
1503 /// have pushed the merge themselves may ask: whoever opened it (or asked
1504 /// g1t for it), for a fork; anyone who may push, for a branch of the
1505 /// repository. When it needs a
1506 /// real merge, says so, naming the conflicting files if a probe found
1507 /// them, and pushes nothing.
1508 async fn catch_up_pull(&self, a: PullActionArgs) -> Result<Outcome<PullBranchUpdate>> {
1509 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &Some(a.actor.clone())).await?);
1510 check!(writable(&repo));
1511 if !pull.status.is_active() {
1512 return Ok(Outcome::fail(
1513 FailureCode::Conflict,
1514 format!("This pull request is already {}.", pull.status.as_str()),
1515 ));
1516 }
1517 if pull.fork_repo_id.is_some() {
1518 if !pull.is_owned_by(&a.actor.id) {
1519 return Ok(Outcome::fail(
1520 FailureCode::Forbidden,
1521 "Only whoever opened this pull request, or asked g1t for it, can update it.",
1522 ));
1523 }
1524 } else {
1525 check!(allowed(Some(&a.actor), &repo, Capability::Push));
1526 }
1527 let updated: Outcome<PullBranchUpdate> = g1t_kit::call(
1528 &self.repos,
1529 "update_pull_branch",
1530 &UpdatePullBranchArgs {
1531 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
1532 branch: pull.branch.clone(),
1533 number: pull.number,
1534 actor: a.actor,
1535 target_branch: pull.base.clone(),
1536 },
1537 )
1538 .await?;
1539 // A probe that found conflicts says more than "both changed it".
1540 if let Outcome::Ok(PullBranchUpdate::NeedsAgent { .. }) = &updated
1541 && let Some(files) = self.conflicting_files(&pull).await?
1542 && !files.is_empty()
1543 {
1544 return Ok(Outcome::Ok(PullBranchUpdate::NeedsAgent {
1545 reason: NeedsAgentReason::Conflicting,
1546 detail: "Merging it conflicts.".to_owned(),
1547 paths: files,
1548 }));
1549 }
1550 Ok(updated)
1551 }
1552
1553 async fn update_pull(&self, a: UpdatePullArgs) -> Result<Outcome<Pull>> {
1554 let pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1555 // The branch it merges into, its milestone and its labels first:
1556 // each can be refused, and then nothing else changes.
1557 if a.base.is_some() || a.milestone.is_some() || a.labels.is_some() {
1558 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
1559 if let Some(base) = &a.base {
1560 check!(self.change_base(&a.actor, &repo, &pull, base).await?);
1561 }
1562 if let Some(number) = a.milestone {
1563 check!(self.set_milestone(&a.actor, &repo, &labels::Item::Pull(pull.clone()), number).await?);
1564 }
1565 if let Some(labels) = &a.labels {
1566 check!(self.relabel(&a.actor, &repo, &labels::Item::Pull(pull.clone()), labels).await?);
1567 }
1568 }
1569 let assignees = match a.assignees {
1570 Some(names) => Some(check!(self.valid_assignees(names).await?)),
1571 None => None,
1572 };
1573 // Teams, named `workspace/team`, apart from the people.
1574 let (team_names, a_reviewers) = match a.reviewers {
1575 Some(names) => {
1576 let (teams, people): (Vec<String>, Vec<String>) =
1577 names.into_iter().partition(|name| team_reviews::team_name(name).is_some());
1578 (Some(teams), Some(people))
1579 }
1580 None => (None, None),
1581 };
1582 let teams = match team_names {
1583 Some(names) => Some(check!(self.valid_team_reviewers(&a.actor, &a.repo, &pull, names).await?)),
1584 None => None,
1585 };
1586 let reviewers = match a_reviewers {
1587 Some(names) => {
1588 // g1t is not an account; everyone else has to be.
1589 let agent = names
1590 .iter()
1591 .any(|name| name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME));
1592 let people = names
1593 .into_iter()
1594 .filter(|name| !name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME))
1595 .collect();
1596 let mut reviewers = check!(self.valid_assignees(people).await?);
1597 // Nobody is asked to review their own, nor what they had g1t make.
1598 reviewers.retain(|name| *name != pull.owner().username);
1599 if agent {
1600 reviewers.insert(0, reviews::AGENT_NAME.to_owned());
1601 }
1602 Some(reviewers)
1603 }
1604 None => None,
1605 };
1606 self.db
1607 .prepare(
1608 "UPDATE pulls
1609 SET assignees = COALESCE(?, assignees), updated_at = ?
1610 WHERE id = ?",
1611 )
1612 .bind(&[
1613 optional(&assignees.as_ref().map(serde_json::to_string).transpose()?),
1614 rfc3339(now_ms()).into(),
1615 pull.id.as_str().into(),
1616 ])?
1617 .run()
1618 .await?;
1619 if let Some(assignees) = &assignees {
1620 self.note_changes(
1621 &pull.repo_id,
1622 pull.number,
1623 &a.actor,
1624 &pull.assignees,
1625 assignees,
1626 ("assigned", "unassigned"),
1627 )
1628 .await?;
1629 }
1630 // Who was newly assigned or asked to review, and whose request was
1631 // withdrawn: the inbox tells them, and webhooks say so.
1632 let newly = |after: &[String], before: &[String]| -> Vec<String> {
1633 after.iter().filter(|name| !before.contains(name)).cloned().collect()
1634 };
1635 if let Some(assignees) = &assignees {
1636 let added = newly(assignees, &pull.assignees);
1637 if !added.is_empty() {
1638 self.publish(
1639 "pull.assigned",
1640 &pull.repo_id,
1641 &a.actor,
1642 PullEvent {
1643 assignees: Some(assignees.clone()),
1644 added: Some(added),
1645 ..Self::pull_event(&pull)
1646 },
1647 )
1648 .await?;
1649 }
1650 }
1651 if reviewers.is_some() || teams.is_some() {
1652 let people = reviewers.unwrap_or_else(|| pull.reviewers.clone());
1653 let teams = teams.unwrap_or_else(|| pull.team_reviewers.clone());
1654 self.set_reviewers(&pull, people, teams, Some(&a.actor), false).await?;
1655 }
1656 Ok(match self.pull(&pull.repo_id, pull.number).await? {
1657 Some(pull) => Outcome::Ok(pull),
1658 None => no_pull(),
1659 })
1660 }
1661
1662 /// Points an open pull request at another branch to merge into. Needs
1663 /// the Write role. What it would merge, whether it is behind, its
1664 /// mergeability and its checks are all worked out against the new base.
1665 async fn change_base(&self, actor: &User, repo: &Repo, pull: &Pull, base: &str) -> Result<Outcome<()>> {
1666 check!(allowed(Some(actor), repo, Capability::Push));
1667 let base = base.trim();
1668 if base.is_empty() {
1669 return Ok(Outcome::fail(FailureCode::Invalid, "Name the branch it should merge into."));
1670 }
1671 let before = pull.base_branch(&repo.default_branch).to_owned();
1672 if base == before {
1673 return Ok(Outcome::Ok(()));
1674 }
1675 if pull.fork_repo_id.is_none() && pull.branch.as_deref() == Some(base) {
1676 return Ok(Outcome::fail(
1677 FailureCode::Invalid,
1678 format!("A pull request cannot merge {base} into itself. Choose another base."),
1679 ));
1680 }
1681 let exists: Option<String> = g1t_kit::call(
1682 &self.repos,
1683 "head",
1684 &HeadArgs { repo_id: repo.id.clone(), branch: base.to_owned() },
1685 )
1686 .await?;
1687 if exists.is_none() {
1688 return Ok(Outcome::fail(FailureCode::NotFound, format!("There is no branch named {base} to merge into.")));
1689 }
1690 // In the merge queue it was headed for the default branch; it
1691 // leaves the queue for another base.
1692 let left_queue = self
1693 .leave(&pull.repo_id, pull, QueueState::Removed, Some("Its base branch changed."))
1694 .await?;
1695 let stored = stored_base(base, repo);
1696 self.db
1697 .prepare(
1698 "UPDATE pulls
1699 SET base_branch = ?, updated_at = ?, land_requested = NULL, land_requested_at = NULL, behind = NULL,
1700 mergeable = NULL, mergeable_key = NULL, conflicts = NULL
1701 WHERE id = ?",
1702 )
1703 .bind(&[optional(&stored), rfc3339(now_ms()).into(), pull.id.as_str().into()])?
1704 .run()
1705 .await?;
1706 self.note(
1707 &pull.repo_id,
1708 pull.number,
1709 (&actor.id, &actor.username),
1710 &format!("changed the base branch from `{before}` to `{base}`"),
1711 )
1712 .await?;
1713 self.publish(
1714 "pull.base_changed",
1715 &pull.repo_id,
1716 actor,
1717 PullEvent { base: Some(base.to_owned()), ..Self::pull_event(pull) },
1718 )
1719 .await?;
1720 if left_queue {
1721 self.publish_as(
1722 "queue.changed",
1723 &pull.repo_id,
1724 None,
1725 g1t_contracts::events::QueueChanged { repo_id: pull.repo_id.clone() },
1726 )
1727 .await?;
1728 }
1729 // Whether it merges cleanly into the new base.
1730 if let Some(moved) = self.pull_by_id(&pull.id).await?
1731 && let Err(error) = self.assess_mergeability(&moved).await
1732 {
1733 worker::console_warn!("mergeability of {}: {error}", pull.id);
1734 }
1735 Ok(Outcome::Ok(()))
1736 }
1737
1738 /// Marks a draft ready for review, or updates the description of one
1739 /// that already is.
1740 async fn ready_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1741 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1742 let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
1743 let now = rfc3339(now_ms());
1744 self.db
1745 .prepare(
1746 "UPDATE pulls SET status = 'open', body = COALESCE(?, body), updated_at = ?
1747 WHERE id = ?",
1748 )
1749 .bind(&[
1750 optional(&summary),
1751 now.as_str().into(),
1752 pull.id.as_str().into(),
1753 ])?
1754 .run()
1755 .await?;
1756 if pull.status == PullStatus::Draft {
1757 // The head as it is now: the push that came just before may not
1758 // have reached `head_commit` yet, and workflows run on it.
1759 let commit = self.live_head(&pull).await?.or_else(|| pull.head_commit.clone());
1760 self.publish(
1761 "pull.ready",
1762 &pull.repo_id,
1763 &a.actor,
1764 PullEvent {
1765 commit,
1766 ..Self::pull_event(&pull)
1767 },
1768 )
1769 .await?;
1770 }
1771 if pull.status == PullStatus::Draft {
1772 self.note(
1773 &pull.repo_id,
1774 pull.number,
1775 (&a.actor.id, &a.actor.username),
1776 "marked this ready for review",
1777 )
1778 .await?;
1779 }
1780 pull.status = PullStatus::Open;
1781 pull.body = summary.or(pull.body);
1782 pull.updated_at = now;
1783 // A draft's code owners are asked once it is ready.
1784 self.refresh_code_owners(&pull).await;
1785 if let Some(fresh) = self.pull(&pull.repo_id, pull.number).await? {
1786 pull.reviewers = fresh.reviewers;
1787 pull.team_reviewers = fresh.team_reviewers;
1788 }
1789 Ok(Outcome::Ok(pull))
1790 }
1791
1792 async fn close_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1793 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1794 let now = rfc3339(now_ms());
1795 self.db
1796 .prepare("UPDATE pulls SET status = 'closed', updated_at = ? WHERE id = ?")
1797 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
1798 .run()
1799 .await?;
1800 self.publish(
1801 "pull.closed",
1802 &pull.repo_id,
1803 &a.actor,
1804 Self::pull_event(&pull),
1805 )
1806 .await?;
1807 self.note(
1808 &pull.repo_id,
1809 pull.number,
1810 (&a.actor.id, &a.actor.username),
1811 "closed this",
1812 )
1813 .await?;
1814 // A closed pull request leaves the merge queue.
1815 if self
1816 .leave(&pull.repo_id, &pull, QueueState::Removed, Some("It was closed."))
1817 .await?
1818 {
1819 self.publish_as(
1820 "queue.changed",
1821 &pull.repo_id,
1822 None,
1823 g1t_contracts::events::QueueChanged {
1824 repo_id: pull.repo_id.clone(),
1825 },
1826 )
1827 .await?;
1828 }
1829 pull.status = PullStatus::Closed;
1830 pull.updated_at = now;
1831 Ok(Outcome::Ok(pull))
1832 }
1833
1834 /// Lands the pull request on the repository's default branch. Unless
1835 /// told to keep it open, that resolves the issue it was for: the issue
1836 /// closes naming this pull request, and the others still in progress
1837 /// for it close as superseded.
1838 async fn merge_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1839 let viewer = Some(a.actor.clone());
1840 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
1841 check!(writable(&repo));
1842 match pull.status {
1843 PullStatus::Open => {}
1844 PullStatus::Draft => {
1845 return Ok(Outcome::fail(
1846 FailureCode::Conflict,
1847 "This pull request is still a draft. Mark it ready for review first.",
1848 ));
1849 }
1850 status => {
1851 return Ok(Outcome::fail(
1852 FailureCode::Conflict,
1853 format!("This pull request is already {}.", status.as_str()),
1854 ));
1855 }
1856 }
1857 let base = pull.base_branch(&repo.default_branch).to_owned();
1858 // The rules of the branch it merges into, as they stack: the one
1859 // gate for the merge button, the API, MCP, auto-merge and the queue,
1860 // for a person's pull request and an agent's alike. A bypass counts
1861 // when the merger asks for it, or for g1t when a ruleset lists it.
1862 let gate = self.merge_gate(&repo, &pull, Some(&a.actor), a.ignore_checks, true).await?;
1863 let bypassable = gate.bypassable();
1864 let gate = if a.bypass_rules || a.actor.is_system() { gate } else { gate.without_bypass() };
1865 let settings = rulesets::overlay(self.settings(&repo.id).await?, &gate.requirements, base == repo.default_branch);
1866 if pull.check_status == Some(CheckStatus::Failed) && !(a.ignore_checks && settings.allow_ignoring_checks) {
1867 return Ok(Outcome::fail(
1868 FailureCode::Conflict,
1869 "It failed in the merge queue; push a fix to try again.",
1870 ));
1871 }
1872 if let Some(refusal) = gate.refusal() {
1873 if access::can(Some(&a.actor), &repo, Capability::Merge) {
1874 self.record_merge_evaluations(&repo, &pull, &gate).await;
1875 }
1876 let offer = if bypassable && !a.bypass_rules {
1877 " You may bypass these rules: merge again and ask to bypass them (bypass_rules)."
1878 } else {
1879 ""
1880 };
1881 return Ok(Outcome::fail(FailureCode::Conflict, format!("{refusal}{offer}")));
1882 }
1883 // Known ahead of time to conflict: neither a merge nor the queue
1884 // would get through, so say what has to be resolved now.
1885 if let Some(files) = self.conflicting_files(&pull).await? {
1886 let named = if files.is_empty() {
1887 String::new()
1888 } else {
1889 format!(" in {}", files.join(", "))
1890 };
1891 return Ok(Outcome::fail(
1892 FailureCode::Conflict,
1893 format!(
1894 "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."
1895 ),
1896 ));
1897 }
1898
1899 // A repository that merges through a queue: it joins the queue, and
1900 // lands once its state together with everything ahead has passed.
1901 if settings.merge_queue {
1902 if !a.actor.verified {
1903 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1904 }
1905 check!(allowed(Some(&a.actor), &repo, Capability::Merge));
1906 self.record_merge_evaluations(&repo, &pull, &gate).await;
1907 return self.enqueue(&repo, &pull, &a.actor, a.keep_issue_open).await;
1908 }
1909
1910 // The default branch has moved under it. Unless the repository
1911 // insists on that being dealt with first, bring it up to date and
1912 // land it when that is done.
1913 if self.is_behind(&repo.id, &pull).await? {
1914 if settings.require_up_to_date {
1915 return Ok(Outcome::fail(
1916 FailureCode::Conflict,
1917 format!(
1918 "{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."
1919 ),
1920 ));
1921 }
1922 if !a.actor.verified {
1923 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1924 }
1925 check!(allowed(Some(&a.actor), &repo, Capability::Merge));
1926 self.record_merge_evaluations(&repo, &pull, &gate).await;
1927 self.request_landing(&pull, &a.actor, a.keep_issue_open)
1928 .await?;
1929 return Ok(Outcome::Ok(pull));
1930 }
1931
1932 // Whether the actor may write to the repository is decided by repos.
1933 let landed: Outcome<Landed> = g1t_kit::call(
1934 &self.repos,
1935 "land",
1936 &LandArgs {
1937 // A pull request from a branch lands from the repository itself.
1938 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
1939 branch: pull.branch.clone(),
1940 actor: a.actor.clone(),
1941 target_branch: Some(base.clone()),
1942 },
1943 )
1944 .await?;
1945 let landed = check!(landed);
1946 self.record_merge_evaluations(&repo, &pull, &gate).await;
1947 Ok(Outcome::Ok(
1948 self.record_merge(&repo, pull, &a.actor, a.keep_issue_open, landed)
1949 .await?,
1950 ))
1951 }
1952
1953 /// Records a pull request as merged once the default branch holds it:
1954 /// closes its issue, supersedes the others for it, and says so.
1955 pub(crate) async fn record_merge(
1956 &self,
1957 repo: &Repo,
1958 mut pull: Pull,
1959 actor: &User,
1960 keep_issue_open: bool,
1961 landed: Landed,
1962 ) -> Result<Pull> {
1963 // Only a merge into the default branch resolves the issue: into
1964 // another branch, the work has not landed yet.
1965 let keep_issue_open = keep_issue_open || !pull.targets_default(&repo.default_branch);
1966 let issue = match pull.issue {
1967 Some(number) if !keep_issue_open => self
1968 .issue(&repo.id, number)
1969 .await?
1970 .filter(|issue| issue.state == State::Open),
1971 _ => None,
1972 };
1973 let now = rfc3339(now_ms());
1974 let mut statements = vec![
1975 self.db
1976 .prepare(
1977 "UPDATE pulls
1978 SET status = 'merged', head_commit = ?, merge_base = ?, merged_by = ?,
1979 merged_at = ?, updated_at = ?
1980 WHERE id = ?",
1981 )
1982 .bind(&[
1983 landed.commit.as_str().into(),
1984 optional(&landed.previous),
1985 actor.username.as_str().into(),
1986 now.as_str().into(),
1987 now.as_str().into(),
1988 pull.id.as_str().into(),
1989 ])?,
1990 ];
1991 if let Some(issue) = &issue {
1992 statements.push(
1993 self.db
1994 .prepare(
1995 "UPDATE issues
1996 SET state = 'closed', reason = 'completed', resolved_by = ?,
1997 closed_at = ?, updated_at = ?
1998 WHERE id = ?",
1999 )
2000 .bind(&[
2001 pull.number.into(),
2002 now.as_str().into(),
2003 now.as_str().into(),
2004 issue.id.as_str().into(),
2005 ])?,
2006 );
2007 statements.push(
2008 self.db
2009 .prepare(
2010 "UPDATE pulls SET status = 'closed', superseded_by = ?, updated_at = ?
2011 WHERE issue_id = ? AND id != ? AND status IN ('draft', 'open')",
2012 )
2013 .bind(&[
2014 pull.number.into(),
2015 now.as_str().into(),
2016 issue.id.as_str().into(),
2017 pull.id.as_str().into(),
2018 ])?,
2019 );
2020 }
2021 self.db.batch(statements).await?;
2022
2023 self.publish(
2024 "pull.merged",
2025 &repo.id,
2026 actor,
2027 PullEvent {
2028 commit: Some(landed.commit.clone()),
2029 ..Self::pull_event(&pull)
2030 },
2031 )
2032 .await?;
2033 if let Some(issue) = &issue {
2034 self.publish(
2035 "issue.closed",
2036 &repo.id,
2037 actor,
2038 IssueEvent {
2039 reason: Some(IssueReason::Completed.as_str()),
2040 resolved_by: Some(pull.number),
2041 ..Self::issue_event(issue)
2042 },
2043 )
2044 .await?;
2045 }
2046
2047 let who = (actor.id.as_str(), actor.username.as_str());
2048 self.note(&repo.id, pull.number, who, "merged this").await?;
2049 if let Some(issue) = &issue {
2050 self.note(
2051 &repo.id,
2052 issue.number,
2053 who,
2054 &format!("closed this by merging #{}", pull.number),
2055 )
2056 .await?;
2057 }
2058 pull.status = PullStatus::Merged;
2059 pull.head_commit = Some(landed.commit.clone());
2060 pull.merge_base = landed.previous;
2061 pull.merged_by = Some(actor.username.clone());
2062 pull.merged_at = Some(now.clone());
2063 pull.updated_at = now;
2064 Ok(pull)
2065 }
2066
2067 async fn list_active_pulls(&self, a: ViewerArgs) -> Result<Vec<ActivePull>> {
2068 let Some(viewer) = a.viewer else {
2069 return Ok(Vec::new());
2070 };
2071 // The pull requests and their issues, in one round trip: their own,
2072 // and those g1t made for them (Pull::owner).
2073 let author = [JsValue::from(viewer.id.as_str())];
2074 let found = self
2075 .timing
2076 .db(
2077 2,
2078 self.db.batch(vec![
2079 self.db
2080 .prepare(format!(
2081 "SELECT {PULL_COLUMNS} FROM pulls
2082 WHERE COALESCE(requested_by_id, author_id) = ?1 AND status IN ('draft', 'open')
2083 ORDER BY updated_at DESC LIMIT 50"
2084 ))
2085 .bind(&author)?,
2086 self.db
2087 .prepare(format!(
2088 "SELECT {ISSUE_COLUMNS} FROM issues WHERE issues.id IN (
2089 SELECT issue_id FROM pulls
2090 WHERE COALESCE(requested_by_id, author_id) = ?1 AND status IN ('draft', 'open') AND issue_id IS NOT NULL
2091 ORDER BY updated_at DESC LIMIT 50)"
2092 ))
2093 .bind(&author)?,
2094 ]),
2095 )
2096 .await?;
2097 let (Some(found), Some(issues)) = (found.first(), found.get(1)) else {
2098 return Ok(Vec::new());
2099 };
2100 let snapshots = found.results::<Snapshot>()?;
2101 let pulls: Vec<Pull> = found.results::<PullRow>()?.into_iter().map(Pull::from).collect();
2102 let issues: Vec<Issue> = issues.results::<IssueRow>()?.into_iter().map(Issue::from).collect();
2103 let issues = &issues;
2104 // Where each stands: the remembered assessment when there is one,
2105 // and worked out otherwise.
2106 try_join_all(pulls.into_iter().zip(snapshots).map(|(pull, snapshot)| async move {
2107 let issue = pull.issue.and_then(|number| {
2108 issues
2109 .iter()
2110 .find(|issue| issue.repo_id == pull.repo_id && issue.number == number)
2111 .cloned()
2112 });
2113 // Only a pull request g1t is seeing through has a lifecycle.
2114 let lifecycle = if !lifecycle::made_by_g1t(&pull) || snapshot.managed == 0 {
2115 None
2116 } else if let (Some(stage), Some(detail)) = (snapshot.stage, snapshot.stage_detail) {
2117 Some(Lifecycle {
2118 stage,
2119 detail,
2120 revisions: snapshot.revisions,
2121 })
2122 } else {
2123 let behind = self.is_behind(&pull.repo_id, &pull).await?;
2124 self.assess(&pull, &issue, behind)
2125 .await?
2126 .map(|(lifecycle, _)| lifecycle)
2127 };
2128 Ok::<_, worker::Error>(ActivePull {
2129 pull,
2130 issue,
2131 lifecycle,
2132 })
2133 }))
2134 .await
2135 }
2136
2137 // --- Sessions ----------------------------------------------------------
2138
2139 async fn append_session(&self, a: AppendSessionArgs) -> Result<Outcome<Appended>> {
2140 if a.entries.is_empty() {
2141 return Ok(Outcome::Ok(Appended { count: 0 }));
2142 }
2143 if a.entries.len() > MAX_ENTRY_BATCH {
2144 return Ok(Outcome::fail(
2145 FailureCode::Invalid,
2146 format!("Send at most {MAX_ENTRY_BATCH} entries at a time."),
2147 ));
2148 }
2149 let viewer = Some(a.actor.clone());
2150 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
2151 if !pull.is_owned_by(&a.actor.id) {
2152 return Ok(Outcome::fail(
2153 FailureCode::Forbidden,
2154 "Only whoever opened a pull request, or asked g1t for it, can record its session.",
2155 ));
2156 }
2157
2158 let now = rfc3339(now_ms());
2159 let count = a.entries.len() as u32;
2160 let mut statements = Vec::with_capacity(a.entries.len() + 1);
2161 for entry in a.entries {
2162 let kind = serde_json::to_value(entry.kind)?;
2163 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
2164 // Each insert takes the next sequence number itself, so two
2165 // writers appending at once cannot collide.
2166 statements.push(
2167 self.db
2168 .prepare(
2169 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
2170 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
2171 FROM session_entries WHERE pull_id = ?",
2172 )
2173 .bind(&[
2174 pull.id.as_str().into(),
2175 kind.as_str().unwrap_or("note").into(),
2176 text.into(),
2177 optional(&entry.tool),
2178 optional(&entry.commit.or_else(|| pull.head_commit.clone())),
2179 now.as_str().into(),
2180 pull.id.as_str().into(),
2181 ])?,
2182 );
2183 }
2184 statements.push(
2185 self.db
2186 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
2187 .bind(&[now.as_str().into(), pull.id.as_str().into()])?,
2188 );
2189 self.db.batch(statements).await?;
2190 self.publish(
2191 "session.appended",
2192 &pull.repo_id,
2193 &a.actor,
2194 SessionAppended {
2195 pull_id: pull.id.clone(),
2196 repo_id: pull.repo_id.clone(),
2197 number: pull.number,
2198 count,
2199 },
2200 )
2201 .await?;
2202 Ok(Outcome::Ok(Appended { count }))
2203 }
2204
2205 /// Adds entries to a pull request's session, each taking the next
2206 /// sequence number, without announcing it.
2207 pub(crate) async fn append_entries(&self, pull: &Pull, entries: &[NewSessionEntry]) -> Result<()> {
2208 let now = rfc3339(now_ms());
2209 let mut statements = Vec::with_capacity(entries.len());
2210 for entry in entries {
2211 let kind = serde_json::to_value(entry.kind)?;
2212 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
2213 statements.push(
2214 self.db
2215 .prepare(
2216 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
2217 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
2218 FROM session_entries WHERE pull_id = ?",
2219 )
2220 .bind(&[
2221 pull.id.as_str().into(),
2222 kind.as_str().unwrap_or("note").into(),
2223 text.into(),
2224 optional(&entry.tool),
2225 optional(&entry.commit.clone().or_else(|| pull.head_commit.clone())),
2226 now.as_str().into(),
2227 pull.id.as_str().into(),
2228 ])?,
2229 );
2230 }
2231 self.db.batch(statements).await?;
2232 Ok(())
2233 }
2234
2235 async fn read_session(&self, a: ViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
2236 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
2237 let rows = self
2238 .db
2239 .prepare(
2240 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
2241 WHERE pull_id = ? AND seq > ? ORDER BY seq LIMIT ?",
2242 )
2243 .bind(&[pull.id.into(), a.after_seq.into(), SESSION_PAGE.into()])?
2244 .all()
2245 .await?
2246 .results::<SessionRow>()?;
2247 Ok(Outcome::Ok(
2248 rows.into_iter().map(SessionEntry::from).collect(),
2249 ))
2250 }
2251
2252 /// A push moves the head of the pull request it concerns: the one whose
2253 /// fork was pushed to, or the one opened from the branch that moved.
2254 async fn on_event(&self, event: &Event) -> Result<()> {
2255 // A new repository starts with the default labels.
2256 if event.kind == "repo.created"
2257 && let Some(repo_id) = event.repo_id.as_deref()
2258 {
2259 return self.seed_labels(repo_id).await;
2260 }
2261 // Pull requests into the branch that became the default merge into
2262 // the default branch, which is stored as none.
2263 if event.kind == "repo.default_branch_changed"
2264 && let (Some(repo_id), Some(to)) = (event.repo_id.as_deref(), event.data["to"].as_str())
2265 {
2266 self.db
2267 .prepare(
2268 "UPDATE pulls SET base_branch = NULL
2269 WHERE repo_id = ? AND base_branch = ? AND status IN ('draft', 'open')",
2270 )
2271 .bind(&[repo_id.into(), to.into()])?
2272 .run()
2273 .await?;
2274 return Ok(());
2275 }
2276 if event.kind != "git.push" {
2277 return Ok(());
2278 }
2279 let (Some(repo_id), Some(after), Some(git_ref)) = (
2280 event.repo_id.as_deref(),
2281 event.data["after"].as_str(),
2282 event.data["ref"].as_str(),
2283 ) else {
2284 return Ok(());
2285 };
2286 let now = rfc3339(now_ms());
2287 // Who moved it, for rules about the most recent push.
2288 let pusher: JsValue = event.actor.as_deref().map_or(JsValue::NULL, Into::into);
2289 // The head moved, so whatever the checks said no longer applies, and
2290 // whatever step g1t was waiting on has been taken.
2291 let moved = "UPDATE pulls
2292 SET head_commit = ?, updated_at = ?, head_pushed_by = ?, head_pushed_at = ?, check_status = NULL, check_run_id = NULL,
2293 working_on = NULL, working_until = NULL, stalled = NULL";
2294 let active = "status IN ('draft', 'open') AND head_commit IS NOT ?";
2295 let returning =
2296 "RETURNING id, repo_id, number, issue_number, status, author_id, author_name, requested_by_id, requested_by_name";
2297 let mut pulls: Vec<MovedRow> = Vec::new();
2298 // A fork carries its pull request on its default branch.
2299 if event.data["defaultBranch"].as_bool() == Some(true) {
2300 pulls.extend(
2301 self.db
2302 .prepare(format!(
2303 "{moved} WHERE fork_repo_id = ? AND {active} {returning}"
2304 ))
2305 .bind(&[
2306 after.into(),
2307 now.as_str().into(),
2308 pusher.clone(),
2309 now.as_str().into(),
2310 repo_id.into(),
2311 after.into(),
2312 ])?
2313 .all()
2314 .await?
2315 .results::<MovedRow>()?,
2316 );
2317 }
2318 if let Some(branch) = git_ref.strip_prefix("refs/heads/") {
2319 pulls.extend(
2320 self.db
2321 .prepare(format!(
2322 "{moved} WHERE repo_id = ? AND source_branch = ? AND {active} {returning}"
2323 ))
2324 .bind(&[
2325 after.into(),
2326 now.as_str().into(),
2327 pusher.clone(),
2328 now.as_str().into(),
2329 repo_id.into(),
2330 branch.into(),
2331 after.into(),
2332 ])?
2333 .all()
2334 .await?
2335 .results::<MovedRow>()?,
2336 );
2337 }
2338 // What each now changes, so overlaps show while the work is under way.
2339 for moved in &pulls {
2340 if let Some(mut pull) = self.pull_by_id(&moved.id).await? {
2341 pull.files = self.refresh_files(&pull).await?;
2342 // Owners of files it now changes are asked too.
2343 self.refresh_code_owners(&pull).await;
2344 }
2345 }
2346 // A merge that was waiting for this push to bring it up to date.
2347 for moved in &pulls {
2348 self.land_if_requested(&moved.id).await?;
2349 }
2350 // Whether each still merges cleanly, and, when a default branch
2351 // moved, every open pull request into it.
2352 let moved_ids: Vec<String> = pulls.iter().map(|pull| pull.id.clone()).collect();
2353 let into = match event.data["defaultBranch"].as_bool() {
2354 Some(true) => mergeability::Moved::DefaultBranch,
2355 _ => match git_ref.strip_prefix("refs/heads/") {
2356 Some(branch) => mergeability::Moved::Branch(branch),
2357 None => mergeability::Moved::Nothing,
2358 },
2359 };
2360 self.after_push(repo_id, into, &moved_ids).await;
2361 // A draft is announced when it is marked ready instead.
2362 for pull in pulls
2363 .into_iter()
2364 .filter(|pull| pull.status == PullStatus::Open)
2365 {
2366 self.publish_as(
2367 "pull.updated",
2368 &pull.repo_id,
2369 event.actor.clone(),
2370 PullEvent {
2371 author: Some(g1t_contracts::credentials::Principal { id: pull.author_id, username: pull.author_name }),
2372 requested_by: pull
2373 .requested_by_id
2374 .zip(pull.requested_by_name)
2375 .map(|(id, username)| g1t_contracts::credentials::Principal { id, username }),
2376 pull_id: pull.id,
2377 repo_id: pull.repo_id.clone(),
2378 number: pull.number,
2379 issue: pull.issue_number,
2380 commit: Some(after.to_owned()),
2381 ..PullEvent::default()
2382 },
2383 )
2384 .await?;
2385 }
2386 Ok(())
2387 }
2388}
2389
2390fn service(env: &Env) -> Result<Work> {
2391 Ok(Work {
2392 db: env.d1("DB")?,
2393 identity: env.service("IDENTITY")?,
2394 repos: env.service("REPOS")?,
2395 events: env.service("EVENTS")?,
2396 actions: env.service("ACTIONS")?,
2397 timing: g1t_kit::d1::Timing::default(),
2398 prefetched: std::cell::RefCell::new(None),
2399 known_repos: std::cell::RefCell::new(std::collections::HashMap::new()),
2400 })
2401}
2402
2403#[event(fetch)]
2404async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
2405 let Some(method) = rpc_method(&request) else {
2406 return Response::error("Not found", 404);
2407 };
2408 // A replica near the caller when it asks for one (crates/kit/src/d1.rs).
2409 let (db, served) = g1t_kit::d1::open(&env, "DB", &request)?;
2410 let body: serde_json::Value = request.json().await?;
2411 let mut work = service(&env)?;
2412 work.db = db;
2413
2414 let answered = match method.as_str() {
2415 "open_issue" => reply(&work.open_issue(args(body)?).await?),
2416 "delegate_issue" => reply(&work.delegate_issue(args(body)?).await?),
2417 "report_confidence" => reply(&work.report_confidence(args(body)?).await?),
2418 "list_issues" => reply(&work.list_issues(args(body)?).await?),
2419 "get_issue" => reply(&work.get_issue(args(body)?).await?),
2420 "update_issue" => reply(&work.update_issue(args(body)?).await?),
2421 "close_issue" => reply(&work.close_issue(args(body)?).await?),
2422 "reopen_issue" => reply(&work.reopen_issue(args(body)?).await?),
2423 "list_labels" => reply(&work.list_labels(args(body)?).await?),
2424 "save_label" => reply(&work.save_label(args(body)?).await?),
2425 "delete_label" => reply(&work.delete_label(args(body)?).await?),
2426 "add_default_labels" => reply(&work.add_default_labels(args(body)?).await?),
2427 "set_labels" => reply(&work.set_labels(args(body)?).await?),
2428 "list_milestones" => reply(&work.list_milestones(args(body)?).await?),
2429 "get_milestone" => reply(&work.get_milestone(args(body)?).await?),
2430 "save_milestone" => reply(&work.save_milestone(args(body)?).await?),
2431 "delete_milestone" => reply(&work.delete_milestone(args(body)?).await?),
2432 "counts" => reply(&work.counts(args(body)?).await?),
2433 "add_comment" => reply(&work.add_comment(args(body)?).await?),
2434 "start_checks" => reply(&work.start_checks(args(body)?).await?),
2435 "seen_checks" => reply(&work.seen_checks(args(body)?).await?),
2436 "report_checks" => reply(&work.report_checks(args(body)?).await?),
2437 "set_commit_status" => reply(&work.set_commit_status(args(body)?).await?),
2438 "start_review" => reply(&work.start_review(args(body)?).await?),
2439 "advance" => reply(&work.advance(args(body)?).await?),
2440 "stall" => reply(&work.stall(args(body)?).await?),
2441 "managed_pulls" => reply(&work.managed_pulls(args(body)?).await?),
2442 "queue" => reply(&work.queue(args(body)?).await?),
2443 "queue_build" => reply(&work.queue_build(args(body)?).await?),
2444 "report_queue" => reply(&work.report_queue(args(body)?).await?),
2445 "remove_from_queue" => reply(&work.remove_from_queue(args(body)?).await?),
2446 "message_agent" => reply(&work.message_agent(args(body)?).await?),
2447 "locate_pull" => reply(&work.locate_pull(args(body)?).await?),
2448 "inbox_subject" => reply(&work.inbox_subject(args(body)?).await?),
2449 "answer_message" => reply(&work.answer_message(args(body)?).await?),
2450 "take_messages" => reply(&work.take_messages(args(body)?).await?),
2451 "wake_for_messages" => reply(&work.wake_for_messages(args(body)?).await?),
2452 "catch_up_job" => reply(&work.catch_up_job(args(body)?).await?),
2453 "get_settings" => reply(&work.get_settings(args(body)?).await?),
2454 "codeowners_errors" => reply(&work.codeowners_errors(args(body)?).await?),
2455 "update_settings" => reply(&work.update_settings(args(body)?).await?),
2456 "report_review" => reply(&work.report_review(args(body)?).await?),
2457 "open_pull" => reply(&work.open_pull(args(body)?).await?),
2458 "list_pulls" => reply(&work.list_pulls(args(body)?).await?),
2459 "pulls_for_repos" => reply(&work.pulls_for_repos(args(body)?).await?),
2460 "get_pull" => reply(&work.get_pull(args(body)?).await?),
2461 "update_pull" => reply(&work.update_pull(args(body)?).await?),
2462 "catch_up_pull" => reply(&work.catch_up_pull(args(body)?).await?),
2463 "ready_pull" => reply(&work.ready_pull(args(body)?).await?),
2464 "close_pull" => reply(&work.close_pull(args(body)?).await?),
2465 "merge_pull" => reply(&work.merge_pull(args(body)?).await?),
2466 "list_active_pulls" => reply(&work.list_active_pulls(args(body)?).await?),
2467 "by_author" => reply(&work.by_author(args(body)?).await?),
2468 "start_plan" => reply(&work.start_plan(args(body)?).await?),
2469 "report_plan" => reply(&work.report_plan(args(body)?).await?),
2470 "get_plan" => reply(&work.get_plan(args(body)?).await?),
2471 "list_plans" => reply(&work.list_plans(args(body)?).await?),
2472 "apply_plan" => reply(&work.apply_plan(args(body)?).await?),
2473 "queue_issue" => reply(&work.queue_issue(args(body)?).await?),
2474 "ready_issues" => reply(&work.ready_issues(args(body)?).await?),
2475 "list_assigned_issues" => reply(&work.list_assigned_issues(args(body)?).await?),
2476 "append_session" => reply(&work.append_session(args(body)?).await?),
2477 "read_session" => reply(&work.read_session(args(body)?).await?),
2478 // Agents at work, their sessions, and memory (runs.rs, memory.rs).
2479 "open_run" => reply(&work.open_run(args(body)?).await?),
2480 "report_run" => reply(&work.report_run(args(body)?).await?),
2481 "stop_run" => reply(&work.stop_run(args(body)?).await?),
2482 "list_runs" => reply(&work.list_runs(args(body)?).await?),
2483 "get_run" => reply(&work.get_run(args(body)?).await?),
2484 "list_sessions" => reply(&work.list_sessions(args(body)?).await?),
2485 "get_session" => reply(&work.get_session(args(body)?).await?),
2486 "list_memories" => reply(&work.list_memories(args(body)?).await?),
2487 "add_memory" => reply(&work.add_memory(args(body)?).await?),
2488 "update_memory" => reply(&work.update_memory(args(body)?).await?),
2489 "delete_memory" => reply(&work.delete_memory(args(body)?).await?),
2490 "recall" => reply(&work.recall(args(body)?).await?),
2491 "memory_context" => reply(&work.memory_context(args(body)?).await?),
2492 // What agents may do in a sandbox (guardrails.rs).
2493 "get_guardrails" => reply(&work.get_guardrails(args(body)?).await?),
2494 "update_guardrails" => reply(&work.update_guardrails(args(body)?).await?),
2495 "run_guardrails" => reply(&work.run_guardrails(args(body)?).await?),
2496 // Plan caps the runner applies (compute.rs).
2497 "active_agents" => reply(&work.active_agents(args(body)?).await?),
2498 "issue_spend" => reply(&work.issue_spend(args(body)?).await?),
2499 "wait_for_slot" => reply(&work.wait_for_slot(args(body)?).await?),
2500 "agent_comment" => reply(&work.agent_comment(args(body)?).await?),
2501 "add_wait" => reply(&work.add_wait(args(body)?).await?),
2502 "waiting_workspaces" => reply(&work.waiting_workspaces(args(body)?).await?),
2503 "take_wait" => reply(&work.take_wait(args(body)?).await?),
2504 // The runs whose sandboxes stop with their repository (retired.rs).
2505 "runs_in_repo" => reply(&work.runs_in_repo(args(body)?).await?),
2506 "run_cost" => reply(&work.run_cost(args(body)?).await?),
2507 "start_mergecheck" => reply(&work.start_mergecheck(args(body)?).await?),
2508 "report_mergecheck" => reply(&work.report_mergecheck(args(body)?).await?),
2509 // Memory that fills itself, and its review queue (capture.rs).
2510 method if capture::METHODS.contains(&method) => capture::dispatch(&work, method, body).await,
2511 // @g1t in comments, and the label rule (mentions.rs).
2512 "take_mention" => reply(&work.take_mention(args(body)?).await?),
2513 "mention_revision" => reply(&work.mention_revision(args(body)?).await?),
2514 "reply_mention" => reply(&work.reply_mention(args(body)?).await?),
2515 "get_agent_rules" => reply(&work.get_agent_rules(args(body)?).await?),
2516 "set_agent_rules" => reply(&work.set_agent_rules(args(body)?).await?),
2517 // Rulesets (rulesets.rs): kept here, enforced here on merge and by
2518 // repos on push.
2519 "list_rulesets" => reply(&work.list_rulesets(args(body)?).await?),
2520 "get_ruleset" => reply(&work.get_ruleset(args(body)?).await?),
2521 "save_ruleset" => reply(&work.save_ruleset(args(body)?).await?),
2522 "delete_ruleset" => reply(&work.delete_ruleset(args(body)?).await?),
2523 "effective_rules" => reply(&work.effective_rules(args(body)?).await?),
2524 "rule_evaluations" => reply(&work.rule_evaluations(args(body)?).await?),
2525 "ref_rules" => reply(&work.ref_rules(args(body)?).await?),
2526 "record_evaluations" => reply(&work.record_evaluations(args(body)?).await?),
2527 "set_requires_pull_request" => reply(&work.set_requires_pull_request(args(body)?).await?),
2528 _ => Response::error("Unknown method", 404),
2529 };
2530 served.finish_timed(answered, &work.timing)
2531}
2532
2533/// Events from the bus, delivered on this service's own queue.
2534#[event(queue)]
2535async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
2536 let work = service(&env)?;
2537 for message in batch.messages()? {
2538 // A workspace renamed: its agent runs and memory move to the slug it has now.
2539 if g1t_kit::rename::on_event(&env, &env.d1("DB")?, message.body(), &[memory::RENAMED, guardrails::RENAMED, rulesets::RENAMED].concat()).await? {
2540 message.ack();
2541 continue;
2542 }
2543 // A repository renamed or transferred: its runs, memory, guardrails
2544 // and runs waiting for a slot follow.
2545 if g1t_kit::transfer::on_event(&env, &env.d1("DB")?, message.body(), &[memory::TRANSFERRED, guardrails::TRANSFERRED, retired::WAITS_MOVED].concat()).await? {
2546 message.ack();
2547 continue;
2548 }
2549 // A workspace deleted: what it kept for itself goes.
2550 if g1t_kit::deleted::on_event(&env.d1("DB")?, message.body(), memory::DELETED).await? {
2551 message.ack();
2552 continue;
2553 }
2554 // A repository deleted, archived or purged, or a branch renamed (retired.rs).
2555 if work.on_retired(message.body()).await? {
2556 message.ack();
2557 continue;
2558 }
2559 capture::on_event(&work, message.body()).await;
2560 work.on_event(message.body()).await?;
2561 message.ack();
2562 }
2563 Ok(())
2564}
2565
2566/// The rules that once read a pull request's author read its owner now:
2567/// whoever asked g1t for it, or its author. For each, the person who asked
2568/// is held to what an author was, and g1t's agent (a token it works with)
2569/// gains nothing by being the author.
2570#[cfg(test)]
2571mod owner_rules {
2572 use super::*;
2573 use crate::rows::stored::{ASKER, G1T, pull};
2574 use g1t_contracts::identity::AGENT_ID;
2575
2576 const SOMEONE: &str = "usr_2";
2577
2578 #[test]
2579 fn no_self_approval() {
2580 // add_comment refuses a verdict on one that is theirs.
2581 let made = pull(G1T, Some(ASKER));
2582 assert!(made.is_owned_by(ASKER.0), "the person who asked cannot approve it");
2583 assert!(!made.is_owned_by(AGENT_ID), "g1t's review agent still gives its verdict");
2584 assert!(!made.is_owned_by(SOMEONE));
2585 }
2586
2587 #[test]
2588 fn what_an_author_could_do_without_a_role() {
2589 // manageable_pull (update, ready, close), catch_up_pull on a fork,
2590 // append_session, and steering with message_agent: theirs to do.
2591 let made = pull(G1T, Some(ASKER));
2592 assert!(made.is_owned_by(ASKER.0));
2593 assert!(!made.is_owned_by(SOMEONE), "anyone else still needs the role");
2594 assert!(!made.is_owned_by(AGENT_ID), "being its author gives g1t's tokens nothing more");
2595 }
2596
2597 #[test]
2598 fn nobody_is_asked_to_review_what_they_asked_for() {
2599 // update_pull drops the owner from the reviewers asked.
2600 let made = pull(G1T, Some(ASKER));
2601 let mut reviewers = vec!["syntaqx".to_owned(), "ana".to_owned()];
2602 reviewers.retain(|name| *name != made.owner().username);
2603 assert_eq!(reviewers, ["ana"]);
2604 }
2605
2606 #[test]
2607 fn sandboxes_act_as_whoever_asked() {
2608 // LifecycleJob, ReviewJob, MergecheckJob and the merge queue's job
2609 // carry who the sandbox's credential acts for: a real account.
2610 let made = pull(G1T, Some(ASKER));
2611 assert_eq!(made.owner().id, ASKER.0);
2612 let acts_as = made.requested_by.unwrap_or(made.author);
2613 assert_eq!(acts_as.id, ASKER.0);
2614 // g1t's own work, which nobody asked for, acts as g1t, as before.
2615 let own = pull(("g1t", "g1t"), None);
2616 assert_eq!(own.requested_by.unwrap_or(own.author).id, "g1t");
2617 }
2618
2619 #[test]
2620 fn events_name_g1t_and_whoever_asked() {
2621 let made = pull(G1T, Some(ASKER));
2622 let event = serde_json::to_value(Work::pull_event(&made)).unwrap();
2623 assert_eq!(event["author"], serde_json::json!({ "id": AGENT_ID, "username": "g1t" }));
2624 assert_eq!(event["requestedBy"], serde_json::json!({ "id": "usr_1", "username": "syntaqx" }));
2625 let own = serde_json::to_value(Work::pull_event(&pull(ASKER, None))).unwrap();
2626 assert!(own.get("requestedBy").is_none());
2627 }
2628}