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