Skip to content
2,964 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 let status = if branch.is_some() { "open" } else { "draft" };
1313 let body = Some(a.body.trim().to_owned()).filter(|body| !body.is_empty());
1314
1315 let number = self.next_number(&repo.id).await?;
1316 let timestamp = rfc3339(now);
1317 // A change g1t makes is g1t's, for whoever asked for it. This is
1318 // what lifecycle::made_by_g1t reads back.
1319 let by_g1t = matches!(a.runtime, Runtime::Hosted) && agent == reviews::AGENT_NAME && fork.is_some();
1320 let (author, requested_by) = authorship(&a.actor, by_g1t);
1321 self.db
1322 .prepare(
1323 "INSERT INTO pulls
1324 (id, repo_id, number, issue_id, issue_number, title, body, agent, runtime,
1325 status, fork_repo_id, fork_namespace, fork_name, source_branch, head_commit,
1326 author_id, author_name, requested_by_id, requested_by_name, created_at, updated_at,
1327 base_branch)
1328 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
1329 )
1330 .bind(&[
1331 id.as_str().into(),
1332 repo.id.as_str().into(),
1333 number.into(),
1334 optional(&issue.as_ref().map(|issue| issue.id.clone())),
1335 optional_number(issue.as_ref().map(|issue| issue.number)),
1336 title.into(),
1337 optional(&body),
1338 agent.into(),
1339 runtime.into(),
1340 status.into(),
1341 optional(&fork.as_ref().map(|fork| fork.id.clone())),
1342 optional(&fork.as_ref().map(|fork| fork.namespace.clone())),
1343 optional(&fork.as_ref().map(|fork| fork.name.clone())),
1344 optional(&branch.map(str::to_owned)),
1345 optional(&head),
1346 author.id.as_str().into(),
1347 author.username.as_str().into(),
1348 optional(&requested_by.as_ref().map(|user| user.id.clone())),
1349 optional(&requested_by.as_ref().map(|user| user.username.clone())),
1350 timestamp.as_str().into(),
1351 timestamp.as_str().into(),
1352 optional(&base),
1353 ])?
1354 .run()
1355 .await?;
1356 let Some(mut pull) = self.pull(&repo.id, number).await? else {
1357 return Ok(no_pull());
1358 };
1359 fill_base(&mut pull, &repo);
1360 self.manage(&pull).await?;
1361 // Someone is on it now, so it is no longer waiting for an agent.
1362 if let Some(issue) = pull.issue {
1363 self.db
1364 .prepare("UPDATE issues SET queued_by = NULL WHERE repo_id = ? AND number = ?")
1365 .bind(&[repo.id.as_str().into(), issue.into()])?
1366 .run()
1367 .await?;
1368 }
1369 if let Some(issue) = pull.issue {
1370 let text = if lifecycle::made_by_g1t(&pull) {
1371 format!("assigned this to g1t, which opened #{}", pull.number)
1372 } else {
1373 format!("opened #{} for this", pull.number)
1374 };
1375 self.note(&repo.id, issue, (&a.actor.id, &a.actor.username), &text)
1376 .await?;
1377 }
1378 self.publish(
1379 "pull.opened",
1380 &repo.id,
1381 &a.actor,
1382 PullEvent {
1383 agent: Some(pull.agent.clone()),
1384 base: pull.base.clone(),
1385 ..Self::pull_event(&pull)
1386 },
1387 )
1388 .await?;
1389 // Its code owners asked to review (codeowners.rs).
1390 self.refresh_code_owners(&pull).await;
1391 let pull = self.pull(&repo.id, number).await?.unwrap_or(pull);
1392 Ok(Outcome::Ok(pull))
1393 }
1394
1395 async fn list_pulls(&self, a: ListPullsArgs) -> Result<Outcome<Vec<Pull>>> {
1396 let filter = match a.state {
1397 Some(State::Open) => "AND status IN ('draft', 'open')",
1398 Some(State::Closed) => "AND status IN ('merged', 'closed')",
1399 None => "",
1400 };
1401 let label = a
1402 .label
1403 .map(|label| label.split_whitespace().collect::<Vec<_>>().join(" ").to_lowercase())
1404 .filter(|label| !label.is_empty());
1405 let milestone = a.milestone;
1406 let read = |repo_id: String| {
1407 let label = label.clone();
1408 async move {
1409 let query = self
1410 .db
1411 .prepare(format!(
1412 "SELECT {PULL_COLUMNS} FROM pulls WHERE repo_id = ? {filter}
1413 AND (? IS NULL OR EXISTS
1414 (SELECT 1 FROM json_each(pulls.labels) WHERE json_each.value = ?))
1415 AND (? IS NULL OR milestone = ?)
1416 ORDER BY number DESC LIMIT ?"
1417 ))
1418 .bind(&[
1419 repo_id.into(),
1420 optional(&label),
1421 optional(&label),
1422 optional_number(milestone),
1423 optional_number(milestone),
1424 LIST_PAGE.into(),
1425 ])?;
1426 self.timing.db(1, query.all()).await?.results::<PullRow>()
1427 }
1428 };
1429 let (repo, rows) = check!(self.repo_then(&a.repo, &a.viewer, read).await?);
1430 let base = a.base.map(|base| base.trim().to_owned()).filter(|base| !base.is_empty());
1431 Ok(Outcome::Ok(
1432 rows.into_iter()
1433 .map(|row| {
1434 let mut pull = Pull::from(row);
1435 fill_base(&mut pull, &repo);
1436 pull
1437 })
1438 .filter(|pull| base.as_deref().is_none_or(|base| pull.base.as_deref() == Some(base)))
1439 .collect(),
1440 ))
1441 }
1442
1443 /// `pulls_for_repos`: what `list_pulls` gives, open and closed, for many
1444 /// repositories at once: one access check with repos for all of them
1445 /// and one query, instead of two of each per repository.
1446 async fn pulls_for_repos(&self, a: PullsForReposArgs) -> Result<Vec<RepoPulls>> {
1447 let ids: Vec<String> = a.repo_ids.into_iter().take(MAX_PULLS_FOR_REPOS).collect();
1448 if ids.is_empty() {
1449 return Ok(Vec::new());
1450 }
1451 let limit = a.limit.clamp(1, LIST_PAGE);
1452 // The rows are read beside the access check, for every id asked
1453 // about; those of repositories the viewer cannot read are dropped.
1454 let asked = serde_json::to_string(&ids)?;
1455 let check = ReadableArgs { ids, viewer: a.viewer };
1456 let readable = self.timing.rpc(g1t_kit::call::<_, Vec<Repo>>(&self.repos, "readable", &check));
1457 let (readable, rows) = try_join(readable, self.timing.db(1, self.newest_pulls(asked, limit))).await?;
1458 if readable.is_empty() {
1459 return Ok(Vec::new());
1460 }
1461 let mut answer: Vec<RepoPulls> = readable
1462 .iter()
1463 .map(|repo| RepoPulls { repo_id: repo.id.clone(), open: Vec::new(), closed: Vec::new() })
1464 .collect();
1465 for pull in rows.into_iter().map(Pull::from) {
1466 let Some(entry) = answer.iter_mut().find(|entry| entry.repo_id == pull.repo_id) else {
1467 continue;
1468 };
1469 match pull.status {
1470 PullStatus::Draft | PullStatus::Open => entry.open.push(pull),
1471 PullStatus::Merged | PullStatus::Closed => entry.closed.push(pull),
1472 }
1473 }
1474 for entry in &mut answer {
1475 entry.open.sort_by_key(|pull| std::cmp::Reverse(pull.number));
1476 entry.closed.sort_by_key(|pull| std::cmp::Reverse(pull.number));
1477 }
1478 Ok(answer)
1479 }
1480
1481 /// The newest `limit` of each repository's open (draft or open) and
1482 /// closed (merged or closed) pull requests, for the ids in `ids` (JSON).
1483 async fn newest_pulls(&self, ids: String, limit: u32) -> Result<Vec<PullRow>> {
1484 self.db
1485 .prepare(format!(
1486 "SELECT * FROM (
1487 SELECT {PULL_COLUMNS}, ROW_NUMBER() OVER (
1488 PARTITION BY pulls.repo_id, pulls.status IN ('draft', 'open') ORDER BY pulls.number DESC
1489 ) AS place
1490 FROM pulls WHERE pulls.repo_id IN (SELECT value FROM json_each(?1))
1491 ) WHERE place <= ?2"
1492 ))
1493 .bind(&[ids.into(), limit.into()])?
1494 .all()
1495 .await?
1496 .results::<PullRow>()
1497 }
1498
1499 async fn get_pull(&self, a: ViewArgs) -> Result<Outcome<PullDetail>> {
1500 let number = a.number;
1501 // Every row the page and the lifecycle read, in one batch started
1502 // beside the access check; the helpers below read from it.
1503 let namespace = a.repo.namespace.clone();
1504 let read = |repo_id: String| self.prefetch_pull(repo_id, namespace.clone(), number);
1505 let Outcome::Ok((repo, Some(found))) = self.repo_then(&a.repo, &a.viewer, read).await? else {
1506 return Ok(no_pull());
1507 };
1508 let Some(row) = found.first::<PullRow>(prefetch::Slot::Pull)? else {
1509 return Ok(no_pull());
1510 };
1511 let stored = found.first::<StoredBehind>(prefetch::Slot::Pull)?;
1512 let issue = found.first::<IssueRow>(prefetch::Slot::Issue)?.map(Issue::from);
1513 let comments: Vec<Comment> =
1514 found.rows::<CommentRow>(prefetch::Slot::Comments)?.into_iter().map(Comment::from).collect();
1515 self.keep_prefetched(Some(found));
1516 let default_branch = repo.default_branch.clone();
1517 let detail = self.pull_detail(repo, Pull::from(row), issue, comments, stored, &a.viewer).await;
1518 self.keep_prefetched(None);
1519 // The branch it merges into, named, for whoever reads it.
1520 Ok(match detail? {
1521 Outcome::Ok(mut detail) => {
1522 if detail.pull.base.as_deref().is_none_or(str::is_empty) {
1523 detail.pull.base = Some(default_branch);
1524 }
1525 Outcome::Ok(detail)
1526 }
1527 failed => failed,
1528 })
1529 }
1530
1531 async fn pull_detail(
1532 &self,
1533 repo: Repo,
1534 mut pull: Pull,
1535 issue: Option<Issue>,
1536 comments: Vec<Comment>,
1537 stored: Option<StoredBehind>,
1538 viewer: &Viewer,
1539 ) -> Result<Outcome<PullDetail>> {
1540 // Whether it is behind, as worked out with its mergeability on the
1541 // last push to either side (mergeability.rs), when that was for
1542 // its head as it is now; otherwise asked of the repos service.
1543 let known_behind = stored.and_then(|stored| stored.for_head(pull.head_commit.as_deref()));
1544 // Worked out on each push; this covers a pull request from before
1545 // that was recorded.
1546 if pull.files.is_empty() && pull.head_commit.is_some() {
1547 pull.files = self.refresh_files(&pull).await?;
1548 }
1549 // Everything else at once: none of it depends on the rest, and each
1550 // is a round trip of its own.
1551 let standing = async {
1552 // Mergeability first: where g1t sees a pull request through, a
1553 // conflict decides its next step.
1554 let behind = async {
1555 match known_behind {
1556 Some(behind) => Ok(behind),
1557 None => {
1558 let behind = self.is_behind(&repo.id, &pull).await?;
1559 // Kept for the next view when the mergeability on
1560 // record is for this head: a pull request from
1561 // before `behind` was kept asks once.
1562 if let Some(head) = pull.head_commit.as_deref()
1563 && pull.status.is_active()
1564 {
1565 self.db
1566 .prepare(
1567 "UPDATE pulls SET behind = ?1
1568 WHERE id = ?2 AND behind IS NULL AND mergeable_key LIKE ?3 || '..%'",
1569 )
1570 .bind(&[u32::from(behind).into(), pull.id.as_str().into(), head.into()])?
1571 .run()
1572 .await?;
1573 }
1574 Ok(behind)
1575 }
1576 }
1577 };
1578 let (merge, behind) = try_join(self.mergeability(&pull), behind).await?;
1579 let assessed = self.assess_with_confidence(&pull, &issue, behind).await?;
1580 let confidence = assessed.as_ref().and_then(|(_, _, confidence)| confidence.clone());
1581 let lifecycle = assessed.map(|(lifecycle, _, _)| lifecycle);
1582 Ok::<_, worker::Error>((behind, (lifecycle, confidence), merge))
1583 };
1584 let (((behind, (lifecycle, confidence), (mergeable, conflicts)), (landing, stalled), comments), (checks, overlaps, review_pending)) =
1585 try_join(
1586 try_join3(standing, self.landing_state(&pull.id), async { Ok(comments) }),
1587 try_join3(
1588 self.latest_checks(&pull.id),
1589 self.overlaps(&pull),
1590 self.review_pending(&pull.id),
1591 ),
1592 )
1593 .await?;
1594 // As just worked out, rather than as it was read.
1595 if confidence.is_some() {
1596 pull.confidence = confidence;
1597 }
1598 // The rules of the branch it merges into, as they stack, and which
1599 // of them it does not meet yet, for whoever is looking.
1600 let (statuses, settings, gate) = try_join3(
1601 self.statuses(&repo.id, pull.head_commit.as_deref()),
1602 self.settings_on(&repo, &pull),
1603 async {
1604 if pull.status.is_active() {
1605 self.merge_gate(&repo, &pull, viewer.as_ref(), false, true).await.map(Some)
1606 } else {
1607 Ok(None)
1608 }
1609 },
1610 )
1611 .await?;
1612 let code_owners = self.pull_code_owners(&pull, &comments, &settings).await?;
1613 Ok(Outcome::Ok(PullDetail {
1614 required_checks: required_checks(&settings.required_checks, &statuses),
1615 rules: gate.map(|gate| rulesets::merge_rules(&gate.judged, &gate.requirements, pull.base_branch(&repo.default_branch) == repo.default_branch)),
1616 code_owners,
1617 comments,
1618 checks,
1619 overlaps,
1620 behind,
1621 review_pending,
1622 lifecycle,
1623 landing,
1624 stalled,
1625 messages: self.messages(&pull.id).await?,
1626 statuses,
1627 mergeable,
1628 conflicts,
1629 earlier_checks: self.earlier_checks(&pull.id).await?,
1630 issue,
1631 pull,
1632 }))
1633 }
1634
1635 /// The pull request, whatever its status, if its repository is not
1636 /// archived and `actor` opened it or may triage its pull requests.
1637 async fn managed_pull(
1638 &self,
1639 actor: &User,
1640 path: &RepoPath,
1641 number: u32,
1642 ) -> Result<Outcome<Pull>> {
1643 let (repo, pull) = check!(self.pull_at(path, number, &Some(actor.clone())).await?);
1644 check!(writable(&repo));
1645 if !pull.is_owned_by(&actor.id) {
1646 check!(allowed(Some(actor), &repo, Capability::Triage));
1647 }
1648 Ok(Outcome::Ok(pull))
1649 }
1650
1651 /// The pull request, if it is still active and `actor` opened it or
1652 /// may triage the repository's pull requests.
1653 async fn manageable_pull(
1654 &self,
1655 actor: &User,
1656 path: &RepoPath,
1657 number: u32,
1658 ) -> Result<Outcome<Pull>> {
1659 let pull = check!(self.managed_pull(actor, path, number).await?);
1660 if !pull.status.is_active() {
1661 return Ok(Outcome::fail(
1662 FailureCode::Conflict,
1663 format!("This pull request is already {}.", pull.status.as_str()),
1664 ));
1665 }
1666 Ok(Outcome::Ok(pull))
1667 }
1668
1669 /// Brings a pull request up to date with the default branch without a
1670 /// sandbox, where the repos service can do that safely. Whoever could
1671 /// have pushed the merge themselves may ask: whoever opened it (or asked
1672 /// g1t for it), for a fork; anyone who may push, for a branch of the
1673 /// repository. When it needs a
1674 /// real merge, says so, naming the conflicting files if a probe found
1675 /// them, and pushes nothing.
1676 async fn catch_up_pull(&self, a: PullActionArgs) -> Result<Outcome<PullBranchUpdate>> {
1677 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &Some(a.actor.clone())).await?);
1678 check!(writable(&repo));
1679 if !pull.status.is_active() {
1680 return Ok(Outcome::fail(
1681 FailureCode::Conflict,
1682 format!("This pull request is already {}.", pull.status.as_str()),
1683 ));
1684 }
1685 if pull.fork_repo_id.is_some() {
1686 if !pull.is_owned_by(&a.actor.id) {
1687 return Ok(Outcome::fail(
1688 FailureCode::Forbidden,
1689 "Only whoever opened this pull request, or asked g1t for it, can update it.",
1690 ));
1691 }
1692 } else {
1693 check!(allowed(Some(&a.actor), &repo, Capability::Push));
1694 }
1695 let updated: Outcome<PullBranchUpdate> = g1t_kit::call(
1696 &self.repos,
1697 "update_pull_branch",
1698 &UpdatePullBranchArgs {
1699 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
1700 branch: pull.branch.clone(),
1701 number: pull.number,
1702 actor: a.actor,
1703 target_branch: pull.base.clone(),
1704 },
1705 )
1706 .await?;
1707 // A probe that found conflicts says more than "both changed it".
1708 if let Outcome::Ok(PullBranchUpdate::NeedsAgent { .. }) = &updated
1709 && let Some(files) = self.conflicting_files(&pull).await?
1710 && !files.is_empty()
1711 {
1712 return Ok(Outcome::Ok(PullBranchUpdate::NeedsAgent {
1713 reason: NeedsAgentReason::Conflicting,
1714 detail: "Merging it conflicts.".to_owned(),
1715 paths: files,
1716 }));
1717 }
1718 Ok(updated)
1719 }
1720
1721 async fn update_pull(&self, a: UpdatePullArgs) -> Result<Outcome<Pull>> {
1722 let pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1723 // The branch it merges into, its milestone and its labels first:
1724 // each can be refused, and then nothing else changes.
1725 if a.base.is_some() || a.milestone.is_some() || a.labels.is_some() {
1726 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
1727 if let Some(base) = &a.base {
1728 check!(self.change_base(&a.actor, &repo, &pull, base).await?);
1729 }
1730 if let Some(number) = a.milestone {
1731 check!(self.set_milestone(&a.actor, &repo, &labels::Item::Pull(pull.clone()), number).await?);
1732 }
1733 if let Some(labels) = &a.labels {
1734 check!(self.relabel(&a.actor, &repo, &labels::Item::Pull(pull.clone()), labels).await?);
1735 }
1736 }
1737 let assignees = match a.assignees {
1738 Some(names) => Some(check!(self.valid_assignees(names).await?)),
1739 None => None,
1740 };
1741 // Teams, named `workspace/team`, apart from the people.
1742 let (team_names, a_reviewers) = match a.reviewers {
1743 Some(names) => {
1744 let (teams, people): (Vec<String>, Vec<String>) =
1745 names.into_iter().partition(|name| team_reviews::team_name(name).is_some());
1746 (Some(teams), Some(people))
1747 }
1748 None => (None, None),
1749 };
1750 let teams = match team_names {
1751 Some(names) => Some(check!(self.valid_team_reviewers(&a.actor, &a.repo, &pull, names).await?)),
1752 None => None,
1753 };
1754 let reviewers = match a_reviewers {
1755 Some(names) => {
1756 // g1t is not an account; everyone else has to be.
1757 let agent = names
1758 .iter()
1759 .any(|name| name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME));
1760 let people = names
1761 .into_iter()
1762 .filter(|name| !name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME))
1763 .collect();
1764 let mut reviewers = check!(self.valid_assignees(people).await?);
1765 // Nobody is asked to review their own, nor what they had g1t make.
1766 reviewers.retain(|name| *name != pull.owner().username);
1767 if agent {
1768 reviewers.insert(0, reviews::AGENT_NAME.to_owned());
1769 }
1770 Some(reviewers)
1771 }
1772 None => None,
1773 };
1774 self.db
1775 .prepare(
1776 "UPDATE pulls
1777 SET assignees = COALESCE(?, assignees), updated_at = ?
1778 WHERE id = ?",
1779 )
1780 .bind(&[
1781 optional(&assignees.as_ref().map(serde_json::to_string).transpose()?),
1782 rfc3339(now_ms()).into(),
1783 pull.id.as_str().into(),
1784 ])?
1785 .run()
1786 .await?;
1787 if let Some(assignees) = &assignees {
1788 self.note_changes(
1789 &pull.repo_id,
1790 pull.number,
1791 &a.actor,
1792 &pull.assignees,
1793 assignees,
1794 ("assigned", "unassigned"),
1795 )
1796 .await?;
1797 }
1798 // Who was newly assigned or asked to review, and whose request was
1799 // withdrawn: the inbox tells them, and webhooks say so.
1800 let newly = |after: &[String], before: &[String]| -> Vec<String> {
1801 after.iter().filter(|name| !before.contains(name)).cloned().collect()
1802 };
1803 if let Some(assignees) = &assignees {
1804 let added = newly(assignees, &pull.assignees);
1805 if !added.is_empty() {
1806 self.publish(
1807 "pull.assigned",
1808 &pull.repo_id,
1809 &a.actor,
1810 PullEvent {
1811 assignees: Some(assignees.clone()),
1812 added: Some(added),
1813 ..Self::pull_event(&pull)
1814 },
1815 )
1816 .await?;
1817 }
1818 }
1819 if reviewers.is_some() || teams.is_some() {
1820 let people = reviewers.unwrap_or_else(|| pull.reviewers.clone());
1821 let teams = teams.unwrap_or_else(|| pull.team_reviewers.clone());
1822 self.set_reviewers(&pull, people, teams, Some(&a.actor), false).await?;
1823 }
1824 Ok(match self.pull(&pull.repo_id, pull.number).await? {
1825 Some(pull) => Outcome::Ok(pull),
1826 None => no_pull(),
1827 })
1828 }
1829
1830 /// Points an open pull request at another branch to merge into. Needs
1831 /// the Write role. What it would merge, whether it is behind, its
1832 /// mergeability and its checks are all worked out against the new base.
1833 async fn change_base(&self, actor: &User, repo: &Repo, pull: &Pull, base: &str) -> Result<Outcome<()>> {
1834 check!(allowed(Some(actor), repo, Capability::Push));
1835 let base = base.trim();
1836 if base.is_empty() {
1837 return Ok(Outcome::fail(FailureCode::Invalid, "Name the branch it should merge into."));
1838 }
1839 let before = pull.base_branch(&repo.default_branch).to_owned();
1840 if base == before {
1841 return Ok(Outcome::Ok(()));
1842 }
1843 if pull.fork_repo_id.is_none() && pull.branch.as_deref() == Some(base) {
1844 return Ok(Outcome::fail(
1845 FailureCode::Invalid,
1846 format!("A pull request cannot merge {base} into itself. Choose another base."),
1847 ));
1848 }
1849 let exists: Option<String> = g1t_kit::call(
1850 &self.repos,
1851 "head",
1852 &HeadArgs { repo_id: repo.id.clone(), branch: base.to_owned() },
1853 )
1854 .await?;
1855 if exists.is_none() {
1856 return Ok(Outcome::fail(FailureCode::NotFound, format!("There is no branch named {base} to merge into.")));
1857 }
1858 // In the merge queue it was headed for the default branch; it
1859 // leaves the queue for another base.
1860 let left_queue = self
1861 .leave(&pull.repo_id, pull, QueueState::Removed, Some("Its base branch changed."))
1862 .await?;
1863 let stored = stored_base(base, repo);
1864 self.db
1865 .prepare(
1866 "UPDATE pulls
1867 SET base_branch = ?, updated_at = ?, land_requested = NULL, land_requested_at = NULL, behind = NULL,
1868 mergeable = NULL, mergeable_key = NULL, conflicts = NULL
1869 WHERE id = ?",
1870 )
1871 .bind(&[optional(&stored), rfc3339(now_ms()).into(), pull.id.as_str().into()])?
1872 .run()
1873 .await?;
1874 self.note(
1875 &pull.repo_id,
1876 pull.number,
1877 (&actor.id, &actor.username),
1878 &format!("changed the base branch from `{before}` to `{base}`"),
1879 )
1880 .await?;
1881 self.publish(
1882 "pull.base_changed",
1883 &pull.repo_id,
1884 actor,
1885 PullEvent { base: Some(base.to_owned()), ..Self::pull_event(pull) },
1886 )
1887 .await?;
1888 if left_queue {
1889 self.publish_as(
1890 "queue.changed",
1891 &pull.repo_id,
1892 None,
1893 g1t_contracts::events::QueueChanged { repo_id: pull.repo_id.clone() },
1894 )
1895 .await?;
1896 }
1897 // Whether it merges cleanly into the new base.
1898 if let Some(moved) = self.pull_by_id(&pull.id).await?
1899 && let Err(error) = self.assess_mergeability(&moved).await
1900 {
1901 worker::console_warn!("mergeability of {}: {error}", pull.id);
1902 }
1903 Ok(Outcome::Ok(()))
1904 }
1905
1906 /// Marks a draft ready for review, or updates the description of one
1907 /// that already is.
1908 async fn ready_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1909 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1910 let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
1911 let now = rfc3339(now_ms());
1912 self.db
1913 .prepare(
1914 "UPDATE pulls SET status = 'open', body = COALESCE(?, body), updated_at = ?
1915 WHERE id = ?",
1916 )
1917 .bind(&[
1918 optional(&summary),
1919 now.as_str().into(),
1920 pull.id.as_str().into(),
1921 ])?
1922 .run()
1923 .await?;
1924 if pull.status == PullStatus::Draft {
1925 // The head as it is now: the push that came just before may not
1926 // have reached `head_commit` yet, and workflows run on it.
1927 let commit = self.live_head(&pull).await?.or_else(|| pull.head_commit.clone());
1928 self.publish(
1929 "pull.ready",
1930 &pull.repo_id,
1931 &a.actor,
1932 PullEvent {
1933 commit,
1934 ..Self::pull_event(&pull)
1935 },
1936 )
1937 .await?;
1938 }
1939 if pull.status == PullStatus::Draft {
1940 self.note(
1941 &pull.repo_id,
1942 pull.number,
1943 (&a.actor.id, &a.actor.username),
1944 "marked this ready for review",
1945 )
1946 .await?;
1947 }
1948 pull.status = PullStatus::Open;
1949 pull.body = summary.or(pull.body);
1950 pull.updated_at = now;
1951 // A draft's code owners are asked once it is ready.
1952 self.refresh_code_owners(&pull).await;
1953 if let Some(fresh) = self.pull(&pull.repo_id, pull.number).await? {
1954 pull.reviewers = fresh.reviewers;
1955 pull.team_reviewers = fresh.team_reviewers;
1956 }
1957 Ok(Outcome::Ok(pull))
1958 }
1959
1960 async fn close_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1961 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1962 let now = rfc3339(now_ms());
1963 self.db
1964 // What it was closed as, so that reopening brings that back.
1965 .prepare("UPDATE pulls SET status = 'closed', closed_from = ?, updated_at = ? WHERE id = ?")
1966 .bind(&[pull.status.as_str().into(), now.as_str().into(), pull.id.as_str().into()])?
1967 .run()
1968 .await?;
1969 self.publish(
1970 "pull.closed",
1971 &pull.repo_id,
1972 &a.actor,
1973 Self::pull_event(&pull),
1974 )
1975 .await?;
1976 self.note(
1977 &pull.repo_id,
1978 pull.number,
1979 (&a.actor.id, &a.actor.username),
1980 "closed this",
1981 )
1982 .await?;
1983 // A closed pull request leaves the merge queue.
1984 if self
1985 .leave(&pull.repo_id, &pull, QueueState::Removed, Some("It was closed."))
1986 .await?
1987 {
1988 self.publish_as(
1989 "queue.changed",
1990 &pull.repo_id,
1991 None,
1992 g1t_contracts::events::QueueChanged {
1993 repo_id: pull.repo_id.clone(),
1994 },
1995 )
1996 .await?;
1997 }
1998 pull.status = PullStatus::Closed;
1999 pull.updated_at = now;
2000 Ok(Outcome::Ok(pull))
2001 }
2002
2003 /// Opens a closed pull request again: as the draft it was, if it was
2004 /// closed as one, else ready for review. A merged one stays merged.
2005 async fn reopen_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
2006 let mut pull = check!(self.managed_pull(&a.actor, &a.repo, a.number).await?);
2007 if let Some(refusal) = reopen_refusal(pull.status) {
2008 return Ok(Outcome::fail(FailureCode::Conflict, refusal));
2009 }
2010 // The head as it is now, which workflows run on. A branch of the
2011 // repository that was deleted since leaves nothing to reopen; a
2012 // fork removed after the close is made again (repos, `pull.reopened`).
2013 let head = self.live_head(&pull).await?;
2014 if head.is_none() && pull.fork_repo_id.is_none() {
2015 return Ok(Outcome::fail(
2016 FailureCode::Conflict,
2017 format!(
2018 "The branch {} no longer exists. Push it again to reopen this pull request.",
2019 pull.branch.as_deref().unwrap_or_default()
2020 ),
2021 ));
2022 }
2023 let closed_from: Option<ClosedFrom> = self
2024 .db
2025 .prepare("SELECT closed_from FROM pulls WHERE id = ?")
2026 .bind(&[pull.id.as_str().into()])?
2027 .first(None)
2028 .await?;
2029 let status = reopened_status(closed_from.and_then(|row| row.closed_from).as_deref());
2030 let now = rfc3339(now_ms());
2031 self.db
2032 .prepare("UPDATE pulls SET status = ?, closed_from = NULL, superseded_by = NULL, updated_at = ? WHERE id = ?")
2033 .bind(&[status.as_str().into(), now.as_str().into(), pull.id.as_str().into()])?
2034 .run()
2035 .await?;
2036 self.publish(
2037 "pull.reopened",
2038 &pull.repo_id,
2039 &a.actor,
2040 PullEvent {
2041 commit: head.or_else(|| pull.head_commit.clone()),
2042 ..Self::pull_event(&pull)
2043 },
2044 )
2045 .await?;
2046 self.note(
2047 &pull.repo_id,
2048 pull.number,
2049 (&a.actor.id, &a.actor.username),
2050 "reopened this",
2051 )
2052 .await?;
2053 pull.status = status;
2054 pull.superseded_by = None;
2055 pull.updated_at = now;
2056 // Whether it still merges cleanly, now that it is open again.
2057 if let Err(error) = self.assess_mergeability(&pull).await {
2058 worker::console_warn!("mergeability of {}: {error}", pull.id);
2059 }
2060 Ok(Outcome::Ok(pull))
2061 }
2062
2063 /// Turns a pull request that is ready for review back into a draft: it
2064 /// cannot be merged until it is marked ready again, and it leaves the
2065 /// merge queue and any merge that was waiting for it to catch up.
2066 async fn convert_pull_to_draft(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
2067 let mut pull = check!(self.managed_pull(&a.actor, &a.repo, a.number).await?);
2068 if let Some(refusal) = draft_refusal(pull.status) {
2069 return Ok(Outcome::fail(FailureCode::Conflict, refusal));
2070 }
2071 let now = rfc3339(now_ms());
2072 self.db
2073 .prepare(
2074 "UPDATE pulls SET status = 'draft', land_requested = NULL, land_requested_at = NULL, updated_at = ?
2075 WHERE id = ?",
2076 )
2077 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
2078 .run()
2079 .await?;
2080 self.publish(
2081 "pull.converted_to_draft",
2082 &pull.repo_id,
2083 &a.actor,
2084 Self::pull_event(&pull),
2085 )
2086 .await?;
2087 self.note(
2088 &pull.repo_id,
2089 pull.number,
2090 (&a.actor.id, &a.actor.username),
2091 "marked this as a draft",
2092 )
2093 .await?;
2094 if self
2095 .leave(&pull.repo_id, &pull, QueueState::Removed, Some("It was marked as a draft."))
2096 .await?
2097 {
2098 self.publish_as(
2099 "queue.changed",
2100 &pull.repo_id,
2101 None,
2102 g1t_contracts::events::QueueChanged {
2103 repo_id: pull.repo_id.clone(),
2104 },
2105 )
2106 .await?;
2107 }
2108 pull.status = PullStatus::Draft;
2109 pull.updated_at = now;
2110 Ok(Outcome::Ok(pull))
2111 }
2112
2113 /// Lands the pull request on the repository's default branch. Unless
2114 /// told to keep it open, that resolves the issue it was for: the issue
2115 /// closes naming this pull request, and the others still in progress
2116 /// for it close as superseded.
2117 async fn merge_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
2118 let viewer = Some(a.actor.clone());
2119 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
2120 check!(writable(&repo));
2121 match pull.status {
2122 PullStatus::Open => {}
2123 PullStatus::Draft => {
2124 return Ok(Outcome::fail(
2125 FailureCode::Conflict,
2126 "This pull request is still a draft. Mark it ready for review first.",
2127 ));
2128 }
2129 status => {
2130 return Ok(Outcome::fail(
2131 FailureCode::Conflict,
2132 format!("This pull request is already {}.", status.as_str()),
2133 ));
2134 }
2135 }
2136 let base = pull.base_branch(&repo.default_branch).to_owned();
2137 // The rules of the branch it merges into, as they stack: the one
2138 // gate for the merge button, the API, MCP, auto-merge and the queue,
2139 // for a person's pull request and an agent's alike. A bypass counts
2140 // when the merger asks for it, or for g1t when a ruleset lists it.
2141 let gate = self.merge_gate(&repo, &pull, Some(&a.actor), a.ignore_checks, true).await?;
2142 let bypassable = gate.bypassable();
2143 let gate = if a.bypass_rules || a.actor.is_system() { gate } else { gate.without_bypass() };
2144 let settings = rulesets::overlay(self.settings(&repo.id).await?, &gate.requirements, base == repo.default_branch);
2145 if pull.check_status == Some(CheckStatus::Failed) && !(a.ignore_checks && settings.allow_ignoring_checks) {
2146 return Ok(Outcome::fail(
2147 FailureCode::Conflict,
2148 "It failed in the merge queue; push a fix to try again.",
2149 ));
2150 }
2151 if let Some(refusal) = gate.refusal() {
2152 if access::can(Some(&a.actor), &repo, Capability::Merge) {
2153 self.record_merge_evaluations(&repo, &pull, &gate).await;
2154 }
2155 let offer = if bypassable && !a.bypass_rules {
2156 " You may bypass these rules: merge again and ask to bypass them (bypass_rules)."
2157 } else {
2158 ""
2159 };
2160 return Ok(Outcome::fail(FailureCode::Conflict, format!("{refusal}{offer}")));
2161 }
2162 // Known ahead of time to conflict: neither a merge nor the queue
2163 // would get through, so say what has to be resolved now.
2164 if let Some(files) = self.conflicting_files(&pull).await? {
2165 let named = if files.is_empty() {
2166 String::new()
2167 } else {
2168 format!(" in {}", files.join(", "))
2169 };
2170 return Ok(Outcome::fail(
2171 FailureCode::Conflict,
2172 format!(
2173 "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."
2174 ),
2175 ));
2176 }
2177
2178 // A repository that merges through a queue: it joins the queue, and
2179 // lands once its state together with everything ahead has passed.
2180 if settings.merge_queue {
2181 if !a.actor.verified {
2182 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
2183 }
2184 check!(allowed(Some(&a.actor), &repo, Capability::Merge));
2185 self.record_merge_evaluations(&repo, &pull, &gate).await;
2186 return self.enqueue(&repo, &pull, &a.actor, a.keep_issue_open).await;
2187 }
2188
2189 // The default branch has moved under it. Unless the repository
2190 // insists on that being dealt with first, bring it up to date and
2191 // land it when that is done.
2192 if self.is_behind(&repo.id, &pull).await? {
2193 if settings.require_up_to_date {
2194 return Ok(Outcome::fail(
2195 FailureCode::Conflict,
2196 format!(
2197 "{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."
2198 ),
2199 ));
2200 }
2201 if !a.actor.verified {
2202 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
2203 }
2204 check!(allowed(Some(&a.actor), &repo, Capability::Merge));
2205 self.record_merge_evaluations(&repo, &pull, &gate).await;
2206 self.request_landing(&pull, &a.actor, a.keep_issue_open)
2207 .await?;
2208 return Ok(Outcome::Ok(pull));
2209 }
2210
2211 // Whether the actor may write to the repository is decided by repos.
2212 let landed: Outcome<Landed> = g1t_kit::call(
2213 &self.repos,
2214 "land",
2215 &LandArgs {
2216 // A pull request from a branch lands from the repository itself.
2217 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
2218 branch: pull.branch.clone(),
2219 actor: a.actor.clone(),
2220 target_branch: Some(base.clone()),
2221 },
2222 )
2223 .await?;
2224 let landed = check!(landed);
2225 self.record_merge_evaluations(&repo, &pull, &gate).await;
2226 Ok(Outcome::Ok(
2227 self.record_merge(&repo, pull, &a.actor, a.keep_issue_open, landed)
2228 .await?,
2229 ))
2230 }
2231
2232 /// Records a pull request as merged once the default branch holds it:
2233 /// closes its issue, supersedes the others for it, and says so.
2234 pub(crate) async fn record_merge(
2235 &self,
2236 repo: &Repo,
2237 mut pull: Pull,
2238 actor: &User,
2239 keep_issue_open: bool,
2240 landed: Landed,
2241 ) -> Result<Pull> {
2242 // Only a merge into the default branch resolves the issue: into
2243 // another branch, the work has not landed yet.
2244 let keep_issue_open = keep_issue_open || !pull.targets_default(&repo.default_branch);
2245 let issue = match pull.issue {
2246 Some(number) if !keep_issue_open => self
2247 .issue(&repo.id, number)
2248 .await?
2249 .filter(|issue| issue.state == State::Open),
2250 _ => None,
2251 };
2252 let now = rfc3339(now_ms());
2253 let mut statements = vec![
2254 self.db
2255 .prepare(
2256 "UPDATE pulls
2257 SET status = 'merged', head_commit = ?, merge_base = ?, merged_by = ?,
2258 merged_at = ?, updated_at = ?
2259 WHERE id = ?",
2260 )
2261 .bind(&[
2262 landed.commit.as_str().into(),
2263 optional(&landed.previous),
2264 actor.username.as_str().into(),
2265 now.as_str().into(),
2266 now.as_str().into(),
2267 pull.id.as_str().into(),
2268 ])?,
2269 ];
2270 if let Some(issue) = &issue {
2271 statements.push(
2272 self.db
2273 .prepare(
2274 "UPDATE issues
2275 SET state = 'closed', reason = 'completed', resolved_by = ?,
2276 closed_at = ?, updated_at = ?
2277 WHERE id = ?",
2278 )
2279 .bind(&[
2280 pull.number.into(),
2281 now.as_str().into(),
2282 now.as_str().into(),
2283 issue.id.as_str().into(),
2284 ])?,
2285 );
2286 statements.push(
2287 self.db
2288 .prepare(
2289 "UPDATE pulls SET status = 'closed', superseded_by = ?, updated_at = ?
2290 WHERE issue_id = ? AND id != ? AND status IN ('draft', 'open')",
2291 )
2292 .bind(&[
2293 pull.number.into(),
2294 now.as_str().into(),
2295 issue.id.as_str().into(),
2296 pull.id.as_str().into(),
2297 ])?,
2298 );
2299 }
2300 self.db.batch(statements).await?;
2301
2302 self.publish(
2303 "pull.merged",
2304 &repo.id,
2305 actor,
2306 PullEvent {
2307 commit: Some(landed.commit.clone()),
2308 ..Self::pull_event(&pull)
2309 },
2310 )
2311 .await?;
2312 if let Some(issue) = &issue {
2313 self.publish(
2314 "issue.closed",
2315 &repo.id,
2316 actor,
2317 IssueEvent {
2318 reason: Some(IssueReason::Completed.as_str()),
2319 resolved_by: Some(pull.number),
2320 ..Self::issue_event(issue)
2321 },
2322 )
2323 .await?;
2324 }
2325
2326 let who = (actor.id.as_str(), actor.username.as_str());
2327 self.note(&repo.id, pull.number, who, "merged this").await?;
2328 if let Some(issue) = &issue {
2329 self.note(
2330 &repo.id,
2331 issue.number,
2332 who,
2333 &format!("closed this by merging #{}", pull.number),
2334 )
2335 .await?;
2336 }
2337 pull.status = PullStatus::Merged;
2338 pull.head_commit = Some(landed.commit.clone());
2339 pull.merge_base = landed.previous;
2340 pull.merged_by = Some(actor.username.clone());
2341 pull.merged_at = Some(now.clone());
2342 pull.updated_at = now;
2343 Ok(pull)
2344 }
2345
2346 async fn list_active_pulls(&self, a: ViewerArgs) -> Result<Vec<ActivePull>> {
2347 let Some(viewer) = a.viewer else {
2348 return Ok(Vec::new());
2349 };
2350 // The pull requests and their issues, in one round trip: their own,
2351 // and those g1t made for them (Pull::owner).
2352 let author = [JsValue::from(viewer.id.as_str())];
2353 let found = self
2354 .timing
2355 .db(
2356 2,
2357 self.db.batch(vec![
2358 self.db
2359 .prepare(format!(
2360 "SELECT {PULL_COLUMNS} FROM pulls
2361 WHERE COALESCE(requested_by_id, author_id) = ?1 AND status IN ('draft', 'open')
2362 ORDER BY updated_at DESC LIMIT 50"
2363 ))
2364 .bind(&author)?,
2365 self.db
2366 .prepare(format!(
2367 "SELECT {ISSUE_COLUMNS} FROM issues WHERE issues.id IN (
2368 SELECT issue_id FROM pulls
2369 WHERE COALESCE(requested_by_id, author_id) = ?1 AND status IN ('draft', 'open') AND issue_id IS NOT NULL
2370 ORDER BY updated_at DESC LIMIT 50)"
2371 ))
2372 .bind(&author)?,
2373 ]),
2374 )
2375 .await?;
2376 let (Some(found), Some(issues)) = (found.first(), found.get(1)) else {
2377 return Ok(Vec::new());
2378 };
2379 let snapshots = found.results::<Snapshot>()?;
2380 let pulls: Vec<Pull> = found.results::<PullRow>()?.into_iter().map(Pull::from).collect();
2381 let issues: Vec<Issue> = issues.results::<IssueRow>()?.into_iter().map(Issue::from).collect();
2382 let issues = &issues;
2383 // Where each stands: the remembered assessment when there is one,
2384 // and worked out otherwise.
2385 try_join_all(pulls.into_iter().zip(snapshots).map(|(pull, snapshot)| async move {
2386 let issue = pull.issue.and_then(|number| {
2387 issues
2388 .iter()
2389 .find(|issue| issue.repo_id == pull.repo_id && issue.number == number)
2390 .cloned()
2391 });
2392 // Only a pull request g1t is seeing through has a lifecycle.
2393 let lifecycle = if !lifecycle::made_by_g1t(&pull) || snapshot.managed == 0 {
2394 None
2395 } else if let (Some(stage), Some(detail)) = (snapshot.stage, snapshot.stage_detail) {
2396 Some(Lifecycle {
2397 stage,
2398 detail,
2399 revisions: snapshot.revisions,
2400 })
2401 } else {
2402 let behind = self.is_behind(&pull.repo_id, &pull).await?;
2403 self.assess(&pull, &issue, behind)
2404 .await?
2405 .map(|(lifecycle, _)| lifecycle)
2406 };
2407 Ok::<_, worker::Error>(ActivePull {
2408 pull,
2409 issue,
2410 lifecycle,
2411 })
2412 }))
2413 .await
2414 }
2415
2416 // --- Sessions ----------------------------------------------------------
2417
2418 async fn append_session(&self, a: AppendSessionArgs) -> Result<Outcome<Appended>> {
2419 if a.entries.is_empty() {
2420 return Ok(Outcome::Ok(Appended { count: 0 }));
2421 }
2422 if a.entries.len() > MAX_ENTRY_BATCH {
2423 return Ok(Outcome::fail(
2424 FailureCode::Invalid,
2425 format!("Send at most {MAX_ENTRY_BATCH} entries at a time."),
2426 ));
2427 }
2428 let viewer = Some(a.actor.clone());
2429 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
2430 if !pull.is_owned_by(&a.actor.id) {
2431 return Ok(Outcome::fail(
2432 FailureCode::Forbidden,
2433 "Only whoever opened a pull request, or asked g1t for it, can record its session.",
2434 ));
2435 }
2436
2437 let now = rfc3339(now_ms());
2438 let count = a.entries.len() as u32;
2439 let mut statements = Vec::with_capacity(a.entries.len() + 1);
2440 for entry in a.entries {
2441 let kind = serde_json::to_value(entry.kind)?;
2442 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
2443 // Each insert takes the next sequence number itself, so two
2444 // writers appending at once cannot collide.
2445 statements.push(
2446 self.db
2447 .prepare(
2448 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
2449 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
2450 FROM session_entries WHERE pull_id = ?",
2451 )
2452 .bind(&[
2453 pull.id.as_str().into(),
2454 kind.as_str().unwrap_or("note").into(),
2455 text.into(),
2456 optional(&entry.tool),
2457 optional(&entry.commit.or_else(|| pull.head_commit.clone())),
2458 now.as_str().into(),
2459 pull.id.as_str().into(),
2460 ])?,
2461 );
2462 }
2463 statements.push(
2464 self.db
2465 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
2466 .bind(&[now.as_str().into(), pull.id.as_str().into()])?,
2467 );
2468 self.db.batch(statements).await?;
2469 self.publish(
2470 "session.appended",
2471 &pull.repo_id,
2472 &a.actor,
2473 SessionAppended {
2474 pull_id: pull.id.clone(),
2475 repo_id: pull.repo_id.clone(),
2476 number: pull.number,
2477 count,
2478 },
2479 )
2480 .await?;
2481 Ok(Outcome::Ok(Appended { count }))
2482 }
2483
2484 /// Adds entries to a pull request's session, each taking the next
2485 /// sequence number, without announcing it.
2486 pub(crate) async fn append_entries(&self, pull: &Pull, entries: &[NewSessionEntry]) -> Result<()> {
2487 let now = rfc3339(now_ms());
2488 let mut statements = Vec::with_capacity(entries.len());
2489 for entry in entries {
2490 let kind = serde_json::to_value(entry.kind)?;
2491 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
2492 statements.push(
2493 self.db
2494 .prepare(
2495 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
2496 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
2497 FROM session_entries WHERE pull_id = ?",
2498 )
2499 .bind(&[
2500 pull.id.as_str().into(),
2501 kind.as_str().unwrap_or("note").into(),
2502 text.into(),
2503 optional(&entry.tool),
2504 optional(&entry.commit.clone().or_else(|| pull.head_commit.clone())),
2505 now.as_str().into(),
2506 pull.id.as_str().into(),
2507 ])?,
2508 );
2509 }
2510 self.db.batch(statements).await?;
2511 Ok(())
2512 }
2513
2514 async fn read_session(&self, a: ViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
2515 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
2516 let rows = self
2517 .db
2518 .prepare(
2519 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
2520 WHERE pull_id = ? AND seq > ? ORDER BY seq LIMIT ?",
2521 )
2522 .bind(&[pull.id.into(), a.after_seq.into(), SESSION_PAGE.into()])?
2523 .all()
2524 .await?
2525 .results::<SessionRow>()?;
2526 Ok(Outcome::Ok(
2527 rows.into_iter().map(SessionEntry::from).collect(),
2528 ))
2529 }
2530
2531 /// A push moves the head of the pull request it concerns: the one whose
2532 /// fork was pushed to, or the one opened from the branch that moved.
2533 async fn on_event(&self, event: &Event) -> Result<()> {
2534 // A new repository starts with the default labels.
2535 if event.kind == "repo.created"
2536 && let Some(repo_id) = event.repo_id.as_deref()
2537 {
2538 return self.seed_labels(repo_id).await;
2539 }
2540 // Pull requests into the branch that became the default merge into
2541 // the default branch, which is stored as none.
2542 if event.kind == "repo.default_branch_changed"
2543 && let (Some(repo_id), Some(to)) = (event.repo_id.as_deref(), event.data["to"].as_str())
2544 {
2545 self.db
2546 .prepare(
2547 "UPDATE pulls SET base_branch = NULL
2548 WHERE repo_id = ? AND base_branch = ? AND status IN ('draft', 'open')",
2549 )
2550 .bind(&[repo_id.into(), to.into()])?
2551 .run()
2552 .await?;
2553 return Ok(());
2554 }
2555 if event.kind != "git.push" {
2556 return Ok(());
2557 }
2558 let (Some(repo_id), Some(after), Some(git_ref)) = (
2559 event.repo_id.as_deref(),
2560 event.data["after"].as_str(),
2561 event.data["ref"].as_str(),
2562 ) else {
2563 return Ok(());
2564 };
2565 let now = rfc3339(now_ms());
2566 // Who moved it, for rules about the most recent push.
2567 let pusher: JsValue = event.actor.as_deref().map_or(JsValue::NULL, Into::into);
2568 // The head moved, so whatever the checks said no longer applies, and
2569 // whatever step g1t was waiting on has been taken.
2570 let moved = "UPDATE pulls
2571 SET head_commit = ?, updated_at = ?, head_pushed_by = ?, head_pushed_at = ?, check_status = NULL, check_run_id = NULL,
2572 working_on = NULL, working_until = NULL, stalled = NULL";
2573 let active = "status IN ('draft', 'open') AND head_commit IS NOT ?";
2574 let returning =
2575 "RETURNING id, repo_id, number, issue_number, status, author_id, author_name, requested_by_id, requested_by_name";
2576 let mut pulls: Vec<MovedRow> = Vec::new();
2577 // A fork carries its pull request on its default branch.
2578 if event.data["defaultBranch"].as_bool() == Some(true) {
2579 pulls.extend(
2580 self.db
2581 .prepare(format!(
2582 "{moved} WHERE fork_repo_id = ? AND {active} {returning}"
2583 ))
2584 .bind(&[
2585 after.into(),
2586 now.as_str().into(),
2587 pusher.clone(),
2588 now.as_str().into(),
2589 repo_id.into(),
2590 after.into(),
2591 ])?
2592 .all()
2593 .await?
2594 .results::<MovedRow>()?,
2595 );
2596 }
2597 if let Some(branch) = git_ref.strip_prefix("refs/heads/") {
2598 pulls.extend(
2599 self.db
2600 .prepare(format!(
2601 "{moved} WHERE repo_id = ? AND source_branch = ? AND {active} {returning}"
2602 ))
2603 .bind(&[
2604 after.into(),
2605 now.as_str().into(),
2606 pusher.clone(),
2607 now.as_str().into(),
2608 repo_id.into(),
2609 branch.into(),
2610 after.into(),
2611 ])?
2612 .all()
2613 .await?
2614 .results::<MovedRow>()?,
2615 );
2616 }
2617 // What each now changes, so overlaps show while the work is under way.
2618 for moved in &pulls {
2619 if let Some(mut pull) = self.pull_by_id(&moved.id).await? {
2620 pull.files = self.refresh_files(&pull).await?;
2621 // Owners of files it now changes are asked too.
2622 self.refresh_code_owners(&pull).await;
2623 }
2624 }
2625 // A merge that was waiting for this push to bring it up to date.
2626 for moved in &pulls {
2627 self.land_if_requested(&moved.id).await?;
2628 }
2629 // Whether each still merges cleanly, and, when a default branch
2630 // moved, every open pull request into it.
2631 let moved_ids: Vec<String> = pulls.iter().map(|pull| pull.id.clone()).collect();
2632 let into = match event.data["defaultBranch"].as_bool() {
2633 Some(true) => mergeability::Moved::DefaultBranch,
2634 _ => match git_ref.strip_prefix("refs/heads/") {
2635 Some(branch) => mergeability::Moved::Branch(branch),
2636 None => mergeability::Moved::Nothing,
2637 },
2638 };
2639 self.after_push(repo_id, into, &moved_ids).await;
2640 // A draft is announced when it is marked ready instead.
2641 for pull in pulls
2642 .into_iter()
2643 .filter(|pull| pull.status == PullStatus::Open)
2644 {
2645 self.publish_as(
2646 "pull.updated",
2647 &pull.repo_id,
2648 event.actor.clone(),
2649 PullEvent {
2650 author: Some(g1t_contracts::credentials::Principal { id: pull.author_id, username: pull.author_name }),
2651 requested_by: pull
2652 .requested_by_id
2653 .zip(pull.requested_by_name)
2654 .map(|(id, username)| g1t_contracts::credentials::Principal { id, username }),
2655 pull_id: pull.id,
2656 repo_id: pull.repo_id.clone(),
2657 number: pull.number,
2658 issue: pull.issue_number,
2659 commit: Some(after.to_owned()),
2660 ..PullEvent::default()
2661 }
2662 .carrying(&event.data),
2663 )
2664 .await?;
2665 }
2666 Ok(())
2667 }
2668}
2669
2670/// An event's data that carries on what caused the event it follows from.
2671trait Carrying: serde::Serialize + Sized {
2672 /// As JSON, marked as a workflow job's doing when `cause` was.
2673 fn carrying(self, cause: &serde_json::Value) -> serde_json::Value {
2674 g1t_contracts::events::carried(self, cause)
2675 }
2676}
2677
2678impl Carrying for PullEvent {}
2679
2680fn service(env: &Env) -> Result<Work> {
2681 Ok(Work {
2682 db: env.d1("DB")?,
2683 identity: env.service("IDENTITY")?,
2684 repos: env.service("REPOS")?,
2685 events: env.service("EVENTS")?,
2686 actions: env.service("ACTIONS")?,
2687 timing: g1t_kit::d1::Timing::default(),
2688 prefetched: std::cell::RefCell::new(None),
2689 known_repos: std::cell::RefCell::new(std::collections::HashMap::new()),
2690 })
2691}
2692
2693#[event(fetch)]
2694async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
2695 let Some(method) = rpc_method(&request) else {
2696 return Response::error("Not found", 404);
2697 };
2698 // A replica near the caller when it asks for one (crates/kit/src/d1.rs).
2699 let (db, served) = g1t_kit::d1::open(&env, "DB", &request)?;
2700 let body: serde_json::Value = request.json().await?;
2701 let mut work = service(&env)?;
2702 work.db = db;
2703
2704 let answered = match method.as_str() {
2705 "open_issue" => reply(&work.open_issue(args(body)?).await?),
2706 "delegate_issue" => reply(&work.delegate_issue(args(body)?).await?),
2707 "report_confidence" => reply(&work.report_confidence(args(body)?).await?),
2708 "list_issues" => reply(&work.list_issues(args(body)?).await?),
2709 "get_issue" => reply(&work.get_issue(args(body)?).await?),
2710 "update_issue" => reply(&work.update_issue(args(body)?).await?),
2711 "close_issue" => reply(&work.close_issue(args(body)?).await?),
2712 "reopen_issue" => reply(&work.reopen_issue(args(body)?).await?),
2713 "list_labels" => reply(&work.list_labels(args(body)?).await?),
2714 "save_label" => reply(&work.save_label(args(body)?).await?),
2715 "delete_label" => reply(&work.delete_label(args(body)?).await?),
2716 "add_default_labels" => reply(&work.add_default_labels(args(body)?).await?),
2717 "set_labels" => reply(&work.set_labels(args(body)?).await?),
2718 "list_milestones" => reply(&work.list_milestones(args(body)?).await?),
2719 "get_milestone" => reply(&work.get_milestone(args(body)?).await?),
2720 "save_milestone" => reply(&work.save_milestone(args(body)?).await?),
2721 "delete_milestone" => reply(&work.delete_milestone(args(body)?).await?),
2722 "counts" => reply(&work.counts(args(body)?).await?),
2723 "add_comment" => reply(&work.add_comment(args(body)?).await?),
2724 "edit_comment" => reply(&work.edit_comment(args(body)?).await?),
2725 "delete_comment" => reply(&work.delete_comment(args(body)?).await?),
2726 "start_checks" => reply(&work.start_checks(args(body)?).await?),
2727 "seen_checks" => reply(&work.seen_checks(args(body)?).await?),
2728 "report_checks" => reply(&work.report_checks(args(body)?).await?),
2729 "set_commit_status" => reply(&work.set_commit_status(args(body)?).await?),
2730 "create_commit_status" => reply(&work.create_commit_status(args(body)?).await?),
2731 "commit_statuses" => reply(&work.commit_statuses(args(body)?).await?),
2732 "combined_status" => reply(&work.combined_status(args(body)?).await?),
2733 "create_check_run" => reply(&work.create_check_run(args(body)?).await?),
2734 "update_check_run" => reply(&work.update_check_run(args(body)?).await?),
2735 "get_check_run" => reply(&work.get_check_run(args(body)?).await?),
2736 "check_run_annotations" => reply(&work.check_run_annotations(args(body)?).await?),
2737 "ref_check_runs" => reply(&work.ref_check_runs(args(body)?).await?),
2738 "ref_check_suites" => reply(&work.ref_check_suites(args(body)?).await?),
2739 "get_check_suite" => reply(&work.get_check_suite(args(body)?).await?),
2740 "rerequest_check_run" => reply(&work.rerequest_check_run(args(body)?).await?),
2741 "rerequest_check_suite" => reply(&work.rerequest_check_suite(args(body)?).await?),
2742 "request_check_action" => reply(&work.request_check_action(args(body)?).await?),
2743 "commit_checks" => reply(&work.commit_checks(args(body)?).await?),
2744 "start_review" => reply(&work.start_review(args(body)?).await?),
2745 "advance" => reply(&work.advance(args(body)?).await?),
2746 "stall" => reply(&work.stall(args(body)?).await?),
2747 "managed_pulls" => reply(&work.managed_pulls(args(body)?).await?),
2748 "queue" => reply(&work.queue(args(body)?).await?),
2749 "queue_build" => reply(&work.queue_build(args(body)?).await?),
2750 "report_queue" => reply(&work.report_queue(args(body)?).await?),
2751 "remove_from_queue" => reply(&work.remove_from_queue(args(body)?).await?),
2752 "message_agent" => reply(&work.message_agent(args(body)?).await?),
2753 "locate_pull" => reply(&work.locate_pull(args(body)?).await?),
2754 "inbox_subject" => reply(&work.inbox_subject(args(body)?).await?),
2755 "answer_message" => reply(&work.answer_message(args(body)?).await?),
2756 "take_messages" => reply(&work.take_messages(args(body)?).await?),
2757 "wake_for_messages" => reply(&work.wake_for_messages(args(body)?).await?),
2758 "catch_up_job" => reply(&work.catch_up_job(args(body)?).await?),
2759 "get_settings" => reply(&work.get_settings(args(body)?).await?),
2760 "codeowners_errors" => reply(&work.codeowners_errors(args(body)?).await?),
2761 "update_settings" => reply(&work.update_settings(args(body)?).await?),
2762 "report_review" => reply(&work.report_review(args(body)?).await?),
2763 "open_pull" => reply(&work.open_pull(args(body)?).await?),
2764 "list_pulls" => reply(&work.list_pulls(args(body)?).await?),
2765 "pulls_for_repos" => reply(&work.pulls_for_repos(args(body)?).await?),
2766 "get_pull" => reply(&work.get_pull(args(body)?).await?),
2767 "update_pull" => reply(&work.update_pull(args(body)?).await?),
2768 "catch_up_pull" => reply(&work.catch_up_pull(args(body)?).await?),
2769 "ready_pull" => reply(&work.ready_pull(args(body)?).await?),
2770 "close_pull" => reply(&work.close_pull(args(body)?).await?),
2771 "reopen_pull" => reply(&work.reopen_pull(args(body)?).await?),
2772 "convert_pull_to_draft" => reply(&work.convert_pull_to_draft(args(body)?).await?),
2773 "merge_pull" => reply(&work.merge_pull(args(body)?).await?),
2774 "list_active_pulls" => reply(&work.list_active_pulls(args(body)?).await?),
2775 "by_author" => reply(&work.by_author(args(body)?).await?),
2776 "contributions" => reply(&work.contributions(args(body)?).await?),
2777 "start_plan" => reply(&work.start_plan(args(body)?).await?),
2778 "report_plan" => reply(&work.report_plan(args(body)?).await?),
2779 "get_plan" => reply(&work.get_plan(args(body)?).await?),
2780 "list_plans" => reply(&work.list_plans(args(body)?).await?),
2781 "apply_plan" => reply(&work.apply_plan(args(body)?).await?),
2782 "queue_issue" => reply(&work.queue_issue(args(body)?).await?),
2783 "ready_issues" => reply(&work.ready_issues(args(body)?).await?),
2784 "list_assigned_issues" => reply(&work.list_assigned_issues(args(body)?).await?),
2785 "append_session" => reply(&work.append_session(args(body)?).await?),
2786 "read_session" => reply(&work.read_session(args(body)?).await?),
2787 // Agents at work, their sessions, and memory (runs.rs, memory.rs).
2788 "open_run" => reply(&work.open_run(args(body)?).await?),
2789 "report_run" => reply(&work.report_run(args(body)?).await?),
2790 "stop_run" => reply(&work.stop_run(args(body)?).await?),
2791 "list_runs" => reply(&work.list_runs(args(body)?).await?),
2792 "get_run" => reply(&work.get_run(args(body)?).await?),
2793 "list_sessions" => reply(&work.list_sessions(args(body)?).await?),
2794 "get_session" => reply(&work.get_session(args(body)?).await?),
2795 "list_memories" => reply(&work.list_memories(args(body)?).await?),
2796 "add_memory" => reply(&work.add_memory(args(body)?).await?),
2797 "update_memory" => reply(&work.update_memory(args(body)?).await?),
2798 "delete_memory" => reply(&work.delete_memory(args(body)?).await?),
2799 "recall" => reply(&work.recall(args(body)?).await?),
2800 "memory_context" => reply(&work.memory_context(args(body)?).await?),
2801 // What agents may do in a sandbox (guardrails.rs).
2802 "get_guardrails" => reply(&work.get_guardrails(args(body)?).await?),
2803 "update_guardrails" => reply(&work.update_guardrails(args(body)?).await?),
2804 "run_guardrails" => reply(&work.run_guardrails(args(body)?).await?),
2805 // Plan caps the runner applies (compute.rs).
2806 "active_agents" => reply(&work.active_agents(args(body)?).await?),
2807 "issue_spend" => reply(&work.issue_spend(args(body)?).await?),
2808 "wait_for_slot" => reply(&work.wait_for_slot(args(body)?).await?),
2809 "agent_comment" => reply(&work.agent_comment(args(body)?).await?),
2810 "add_wait" => reply(&work.add_wait(args(body)?).await?),
2811 "waiting_workspaces" => reply(&work.waiting_workspaces(args(body)?).await?),
2812 "take_wait" => reply(&work.take_wait(args(body)?).await?),
2813 // The runs whose sandboxes stop with their repository (retired.rs).
2814 "runs_in_repo" => reply(&work.runs_in_repo(args(body)?).await?),
2815 "run_cost" => reply(&work.run_cost(args(body)?).await?),
2816 "start_mergecheck" => reply(&work.start_mergecheck(args(body)?).await?),
2817 "report_mergecheck" => reply(&work.report_mergecheck(args(body)?).await?),
2818 // Memory that fills itself, and its review queue (capture.rs).
2819 method if capture::METHODS.contains(&method) => capture::dispatch(&work, method, body).await,
2820 // @g1t in comments, and the label rule (mentions.rs).
2821 "take_mention" => reply(&work.take_mention(args(body)?).await?),
2822 "mention_revision" => reply(&work.mention_revision(args(body)?).await?),
2823 "reply_mention" => reply(&work.reply_mention(args(body)?).await?),
2824 "get_agent_rules" => reply(&work.get_agent_rules(args(body)?).await?),
2825 "set_agent_rules" => reply(&work.set_agent_rules(args(body)?).await?),
2826 // Rulesets (rulesets.rs): kept here, enforced here on merge and by
2827 // repos on push.
2828 "list_rulesets" => reply(&work.list_rulesets(args(body)?).await?),
2829 "get_ruleset" => reply(&work.get_ruleset(args(body)?).await?),
2830 "save_ruleset" => reply(&work.save_ruleset(args(body)?).await?),
2831 "delete_ruleset" => reply(&work.delete_ruleset(args(body)?).await?),
2832 "effective_rules" => reply(&work.effective_rules(args(body)?).await?),
2833 "rule_evaluations" => reply(&work.rule_evaluations(args(body)?).await?),
2834 "ref_rules" => reply(&work.ref_rules(args(body)?).await?),
2835 "record_evaluations" => reply(&work.record_evaluations(args(body)?).await?),
2836 "set_requires_pull_request" => reply(&work.set_requires_pull_request(args(body)?).await?),
2837 _ => Response::error("Unknown method", 404),
2838 };
2839 served.finish_timed(answered, &work.timing)
2840}
2841
2842/// Events from the bus, delivered on this service's own queue.
2843#[event(queue)]
2844async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
2845 let work = service(&env)?;
2846 for message in batch.messages()? {
2847 // A workspace renamed: its agent runs and memory move to the slug it has now.
2848 if g1t_kit::rename::on_event(&env, &env.d1("DB")?, message.body(), &[memory::RENAMED, guardrails::RENAMED, rulesets::RENAMED].concat()).await? {
2849 message.ack();
2850 continue;
2851 }
2852 // A repository renamed or transferred: its runs, memory, guardrails
2853 // and runs waiting for a slot follow.
2854 if g1t_kit::transfer::on_event(&env, &env.d1("DB")?, message.body(), &[memory::TRANSFERRED, guardrails::TRANSFERRED, retired::WAITS_MOVED].concat()).await? {
2855 message.ack();
2856 continue;
2857 }
2858 // A workspace deleted: what it kept for itself goes.
2859 if g1t_kit::deleted::on_event(&env.d1("DB")?, message.body(), memory::DELETED).await? {
2860 message.ack();
2861 continue;
2862 }
2863 // An account purged: what it wrote shows as ghost (ghost.rs).
2864 let ghost = ghost::statements();
2865 let ghost: Vec<&str> = ghost.iter().map(String::as_str).collect();
2866 if g1t_kit::user_deleted::on_event(&env.d1("DB")?, message.body(), &ghost).await? {
2867 message.ack();
2868 continue;
2869 }
2870 // A repository deleted, archived or purged, or a branch renamed (retired.rs).
2871 if work.on_retired(message.body()).await? {
2872 message.ack();
2873 continue;
2874 }
2875 capture::on_event(&work, message.body()).await;
2876 work.on_event(message.body()).await?;
2877 message.ack();
2878 }
2879 Ok(())
2880}
2881
2882/// The rules that once read a pull request's author read its owner now:
2883/// whoever asked g1t for it, or its author. For each, the person who asked
2884/// is held to what an author was, and g1t's agent (a token it works with)
2885/// gains nothing by being the author.
2886#[cfg(test)]
2887mod owner_rules {
2888 use super::*;
2889 use crate::rows::stored::{ASKER, G1T, pull};
2890 use g1t_contracts::identity::AGENT_ID;
2891
2892 const SOMEONE: &str = "usr_2";
2893
2894 #[test]
2895 fn no_self_approval() {
2896 // add_comment refuses a verdict on one that is theirs.
2897 let made = pull(G1T, Some(ASKER));
2898 assert!(made.is_owned_by(ASKER.0), "the person who asked cannot approve it");
2899 assert!(!made.is_owned_by(AGENT_ID), "g1t's review agent still gives its verdict");
2900 assert!(!made.is_owned_by(SOMEONE));
2901 }
2902
2903 #[test]
2904 fn what_an_author_could_do_without_a_role() {
2905 // manageable_pull (update, ready, close), catch_up_pull on a fork,
2906 // append_session, and steering with message_agent: theirs to do.
2907 let made = pull(G1T, Some(ASKER));
2908 assert!(made.is_owned_by(ASKER.0));
2909 assert!(!made.is_owned_by(SOMEONE), "anyone else still needs the role");
2910 assert!(!made.is_owned_by(AGENT_ID), "being its author gives g1t's tokens nothing more");
2911 }
2912
2913 #[test]
2914 fn nobody_is_asked_to_review_what_they_asked_for() {
2915 // update_pull drops the owner from the reviewers asked.
2916 let made = pull(G1T, Some(ASKER));
2917 let mut reviewers = vec!["syntaqx".to_owned(), "ana".to_owned()];
2918 reviewers.retain(|name| *name != made.owner().username);
2919 assert_eq!(reviewers, ["ana"]);
2920 }
2921
2922 #[test]
2923 fn sandboxes_act_as_whoever_asked() {
2924 // LifecycleJob, ReviewJob, MergecheckJob and the merge queue's job
2925 // carry who the sandbox's credential acts for: a real account.
2926 let made = pull(G1T, Some(ASKER));
2927 assert_eq!(made.owner().id, ASKER.0);
2928 let acts_as = made.requested_by.unwrap_or(made.author);
2929 assert_eq!(acts_as.id, ASKER.0);
2930 // g1t's own work, which nobody asked for, acts as g1t, as before.
2931 let own = pull(("g1t", "g1t"), None);
2932 assert_eq!(own.requested_by.unwrap_or(own.author).id, "g1t");
2933 }
2934
2935 #[test]
2936 fn events_name_g1t_and_whoever_asked() {
2937 let made = pull(G1T, Some(ASKER));
2938 let event = serde_json::to_value(Work::pull_event(&made)).unwrap();
2939 assert_eq!(event["author"], serde_json::json!({ "id": AGENT_ID, "username": "g1t" }));
2940 assert_eq!(event["requestedBy"], serde_json::json!({ "id": "usr_1", "username": "syntaqx" }));
2941 let own = serde_json::to_value(Work::pull_event(&pull(ASKER, None))).unwrap();
2942 assert!(own.get("requestedBy").is_none());
2943 }
2944
2945 #[test]
2946 fn only_a_closed_pull_request_reopens_and_only_an_open_one_turns_draft() {
2947 assert_eq!(reopen_refusal(PullStatus::Closed), None);
2948 assert!(reopen_refusal(PullStatus::Merged).is_some(), "a merge cannot be undone");
2949 assert!(reopen_refusal(PullStatus::Open).is_some());
2950 assert!(reopen_refusal(PullStatus::Draft).is_some());
2951 assert_eq!(draft_refusal(PullStatus::Open), None);
2952 assert!(draft_refusal(PullStatus::Draft).is_some());
2953 assert!(draft_refusal(PullStatus::Closed).is_some());
2954 assert!(draft_refusal(PullStatus::Merged).is_some());
2955 }
2956
2957 #[test]
2958 fn a_pull_request_reopens_as_what_it_was_closed_as() {
2959 assert_eq!(reopened_status(Some("draft")), PullStatus::Draft);
2960 assert_eq!(reopened_status(Some("open")), PullStatus::Open);
2961 // Closed before this was recorded, or by a merge of another.
2962 assert_eq!(reopened_status(None), PullStatus::Open);
2963 }
2964}