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