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