Skip to content

g1t/services/work/src/lib.rs

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