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