g1t/services/work/src/lib.rs

2,304 lines93,424 bytesCodeBlame
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 compute;
11mod confidence;
12mod guardrails;
13mod lifecycle;
14mod memory;
15mod mentions;
16mod mergeability;
17mod plans;
18mod messages;
19mod prefetch;
20mod queue;
21mod retired;
22mod reviews;
23mod rows;
24mod runs;
25mod settings;
26mod statuses;
27
28use g1t_contracts::events::{
29 CommentCreated, Event, IssueEvent, NewEvent, Publish, PullEvent, SessionAppended,
30};
31use g1t_contracts::identity::UsernameArgs;
32use g1t_contracts::repos::{
33 ForkArgs, GetArgs, HeadArgs, LandArgs, Landed, NeedsAgentReason, PullBranchUpdate, ReadableArgs, Repo, RepoPath,
34 UpdatePullBranchArgs,
35};
36use g1t_contracts::access::{self, Capability, Denied};
37use g1t_contracts::time::rfc3339;
38use g1t_contracts::work::*;
39use futures_util::future::{try_join, try_join3, try_join_all};
40use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
41use g1t_kit::{args, now_ms, reply, rpc_method};
42use serde::Serialize;
43use worker::wasm_bindgen::JsValue;
44use worker::{
45 Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event,
46};
47
48use retired::writable;
49use rows::{CommentRow, IssueRow, MovedRow, NumberRow, PULL_COLUMNS, PullRow, SessionRow, Snapshot, ValueRow};
50
51const SOURCE: &str = "work";
52const MAX_ENTRY_BATCH: usize = 200;
53const MAX_ENTRY_CHARS: usize = 64_000;
54const MAX_TITLE_CHARS: usize = 200;
55const SESSION_PAGE: u32 = 500;
56const LIST_PAGE: u32 = 100;
57const MAX_ASSIGNEES: usize = 10;
58const UNVERIFIED: &str = "Confirm your email address first. Check your inbox, or resend the link from the banner on g1t.sh.";
59
60const ISSUE_COLUMNS: &str = "issues.*,
61 (SELECT count(*) FROM pulls WHERE pulls.issue_id = issues.id) AS pull_count,
62 (SELECT agent FROM pulls
63 WHERE pulls.issue_id = issues.id AND pulls.status IN ('draft', 'open')
64 AND pulls.fork_repo_id IS NOT NULL
65 ORDER BY pulls.number DESC LIMIT 1) AS agent,
66 (SELECT count(*) FROM comments
67 WHERE comments.repo_id = issues.repo_id AND comments.number = issues.number
68 AND comments.kind = 'comment') AS comment_count";
69
70fn no_issue<T>() -> Outcome<T> {
71 Outcome::fail(FailureCode::NotFound, "Issue not found.")
72}
73
74fn no_pull<T>() -> Outcome<T> {
75 Outcome::fail(FailureCode::NotFound, "Pull request not found.")
76}
77
78/// Whether a pull request was behind when its mergeability was last
79/// worked out, and for which pair of commits (mergeability.rs).
80#[derive(serde::Deserialize)]
81struct StoredBehind {
82 #[serde(default)]
83 behind: Option<u8>,
84 #[serde(default)]
85 mergeable_key: Option<String>,
86}
87
88impl StoredBehind {
89 /// The stored answer, if it was worked out for `head`.
90 fn for_head(&self, head: Option<&str>) -> Option<bool> {
91 let (worked_for, _) = self.mergeable_key.as_deref()?.split_once("..")?;
92 (Some(worked_for) == head).then_some(self.behind? != 0)
93 }
94}
95
96/// Refuses `actor` unless their role on `repo` has `capability`: not found
97/// when they cannot read it, forbidden with the role it needs otherwise.
98pub(crate) fn allowed(actor: Option<&User>, repo: &Repo, capability: Capability) -> Outcome<()> {
99 match access::check(actor, repo, capability) {
100 Ok(()) => Outcome::Ok(()),
101 Err(Denied::NotFound) => Outcome::fail(FailureCode::NotFound, "Repository not found."),
102 Err(Denied::Forbidden) => Outcome::fail(
103 FailureCode::Forbidden,
104 access::needs(capability, &format!("{}/{}", repo.namespace, repo.name)),
105 ),
106 }
107}
108
109fn optional(value: &Option<String>) -> JsValue {
110 value.as_deref().map_or(JsValue::NULL, JsValue::from)
111}
112
113fn optional_number(value: Option<u32>) -> JsValue {
114 value.map_or(JsValue::NULL, JsValue::from)
115}
116
117/// The lowercase name a `State` is stored and sent as.
118fn state_name(state: Option<State>) -> Option<&'static str> {
119 state.map(|state| match state {
120 State::Open => "open",
121 State::Closed => "closed",
122 })
123}
124
125/// A trimmed title, or why it cannot be used.
126fn valid_title(title: &str) -> std::result::Result<&str, &'static str> {
127 let title = title.trim();
128 if title.is_empty() {
129 Err("A title is required.")
130 } else if title.chars().count() > MAX_TITLE_CHARS {
131 Err("That title is too long.")
132 } else {
133 Ok(title)
134 }
135}
136
137/// Unwraps an `Outcome`, returning its failure from the enclosing method.
138macro_rules! check {
139 ($outcome:expr) => {
140 match $outcome {
141 Outcome::Ok(value) => value,
142 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
143 }
144 };
145}
146
147struct Work {
148 db: D1Database,
149 identity: Fetcher,
150 repos: Fetcher,
151 events: Fetcher,
152 /// GitHub Actions: runs a merge queue's `merge_group` workflows.
153 actions: Fetcher,
154 /// Where this request's time went, for its `Server-Timing`.
155 timing: g1t_kit::d1::Timing,
156 /// A pull request's rows read in one batch for this request
157 /// (prefetch.rs), which the helpers below read instead of the database.
158 prefetched: std::cell::RefCell<Option<std::rc::Rc<prefetch::Prefetched>>>,
159}
160
161impl Work {
162 async fn publish<T: Serialize>(
163 &self,
164 kind: &'static str,
165 repo_id: &str,
166 actor: &User,
167 data: T,
168 ) -> Result<()> {
169 self.publish_as(kind, repo_id, Some(actor.id.clone()), data)
170 .await
171 }
172
173 /// A pull request's owner (whoever asked g1t for it, or its author) as
174 /// a viewer who can read its repository and source. Stored people carry
175 /// no memberships, so a private repository would otherwise look missing
176 /// to them. The membership given reads and nothing more: it is for
177 /// looking, never for acting.
178 pub(crate) async fn owner_viewer(&self, pull: &Pull) -> Result<Viewer> {
179 let path: Option<RepoPath> = g1t_kit::call(
180 &self.repos,
181 "path_by_id",
182 &g1t_contracts::repos::PathByIdArgs { id: pull.repo_id.clone() },
183 )
184 .await?;
185 let mut owner = pull.owner().clone();
186 if let Some(path) = path
187 && !owner.is_member(&path.namespace.to_lowercase())
188 {
189 owner.workspaces.push(g1t_contracts::Membership {
190 base_permission: Some(access::BasePermission::Read),
191 ..g1t_contracts::Membership::member(path.namespace.to_lowercase())
192 });
193 }
194 Ok(Some(owner))
195 }
196
197 /// Publishes an event caused by `actor`, or by g1t itself.
198 async fn publish_as<T: Serialize>(
199 &self,
200 kind: &'static str,
201 repo_id: &str,
202 actor: Option<String>,
203 data: T,
204 ) -> Result<()> {
205 let event = NewEvent {
206 kind,
207 source: SOURCE,
208 repo_id: Some(repo_id.to_owned()),
209 actor,
210 data,
211 };
212 g1t_kit::call(
213 &self.events,
214 "publish",
215 &Publish {
216 events: vec![event],
217 },
218 )
219 .await
220 }
221
222 /// The repository, if the viewer may see it. Whether they may is
223 /// decided by the repos service.
224 async fn repo(&self, path: &RepoPath, viewer: &Viewer) -> Result<Outcome<Repo>> {
225 self.timing
226 .rpc(g1t_kit::call(
227 &self.repos,
228 "get",
229 &GetArgs {
230 path: path.clone(),
231 viewer: viewer.clone(),
232 },
233 ))
234 .await
235 }
236
237 /// The next number in the repository's sequence. Taking it is one
238 /// statement, so concurrent opens cannot be given the same number.
239 async fn next_number(&self, repo_id: &str) -> Result<u32> {
240 let row = self
241 .db
242 .prepare(
243 "INSERT INTO counters (repo_id, last) VALUES (?, 1)
244 ON CONFLICT (repo_id) DO UPDATE SET last = last + 1
245 RETURNING last AS n",
246 )
247 .bind(&[repo_id.into()])?
248 .first::<NumberRow>(None)
249 .await?;
250 row.map(|row| row.n)
251 .ok_or_else(|| worker::Error::RustError("no number was assigned".into()))
252 }
253
254 async fn issue(&self, repo_id: &str, number: u32) -> Result<Option<Issue>> {
255 Ok(self
256 .db
257 .prepare(format!(
258 "SELECT {ISSUE_COLUMNS} FROM issues WHERE repo_id = ? AND number = ?"
259 ))
260 .bind(&[repo_id.into(), number.into()])?
261 .first::<IssueRow>(None)
262 .await?
263 .map(Issue::from))
264 }
265
266 async fn pull(&self, repo_id: &str, number: u32) -> Result<Option<Pull>> {
267 Ok(self
268 .db
269 .prepare(format!("SELECT {PULL_COLUMNS} FROM pulls WHERE repo_id = ? AND number = ?"))
270 .bind(&[repo_id.into(), number.into()])?
271 .first::<PullRow>(None)
272 .await?
273 .map(Pull::from))
274 }
275
276 /// The repository and one of its issues, as seen by `viewer`.
277 async fn issue_at(
278 &self,
279 path: &RepoPath,
280 number: u32,
281 viewer: &Viewer,
282 ) -> Result<Outcome<(Repo, Issue)>> {
283 let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
284 return Ok(no_issue());
285 };
286 Ok(match self.issue(&repo.id, number).await? {
287 Some(issue) => Outcome::Ok((repo, issue)),
288 None => no_issue(),
289 })
290 }
291
292 /// The repository and one of its pull requests, as seen by `viewer`.
293 async fn pull_at(
294 &self,
295 path: &RepoPath,
296 number: u32,
297 viewer: &Viewer,
298 ) -> Result<Outcome<(Repo, Pull)>> {
299 let Outcome::Ok(repo) = self.repo(path, viewer).await? else {
300 return Ok(no_pull());
301 };
302 Ok(match self.pull(&repo.id, number).await? {
303 Some(pull) => Outcome::Ok((repo, pull)),
304 None => no_pull(),
305 })
306 }
307
308 /// Records something that happened to an issue or a pull request, so
309 /// that it shows in the conversation where it happened. `text` is what
310 /// `author` did, as the rest of a sentence starting with their name.
311 pub(crate) async fn note(
312 &self,
313 repo_id: &str,
314 number: u32,
315 author: (&str, &str),
316 text: &str,
317 ) -> Result<()> {
318 let now = now_ms();
319 self.db
320 .prepare(
321 "INSERT INTO comments
322 (id, repo_id, number, author_id, author_name, body, kind, created_at)
323 VALUES (?, ?, ?, ?, ?, ?, 'event', ?)",
324 )
325 .bind(&[
326 new_id("cmt", now).into(),
327 repo_id.into(),
328 number.into(),
329 author.0.into(),
330 author.1.into(),
331 text.into(),
332 rfc3339(now).into(),
333 ])?
334 .run()
335 .await?;
336 Ok(())
337 }
338
339 /// Notes who was added to and removed from a list of people, such as
340 /// "assigned ana" or "requested a review from g1t".
341 async fn note_changes(
342 &self,
343 repo_id: &str,
344 number: u32,
345 actor: &User,
346 before: &[String],
347 after: &[String],
348 (added, removed): (&str, &str),
349 ) -> Result<()> {
350 let joined = |names: Vec<&String>| {
351 names
352 .into_iter()
353 .map(String::as_str)
354 .collect::<Vec<_>>()
355 .join(", ")
356 };
357 let new: Vec<&String> = after.iter().filter(|name| !before.contains(name)).collect();
358 let gone: Vec<&String> = before.iter().filter(|name| !after.contains(name)).collect();
359 let who = (actor.id.as_str(), actor.username.as_str());
360 if !new.is_empty() {
361 // Taking something on oneself reads better said that way.
362 let text = if added == "assigned" && new == [&actor.username] {
363 "self-assigned this".to_owned()
364 } else {
365 format!("{added} {}", joined(new))
366 };
367 self.note(repo_id, number, who, &text).await?;
368 }
369 if !gone.is_empty() {
370 self.note(repo_id, number, who, &format!("{removed} {}", joined(gone)))
371 .await?;
372 }
373 Ok(())
374 }
375
376 fn issue_event(issue: &Issue) -> IssueEvent {
377 IssueEvent {
378 issue_id: issue.id.clone(),
379 repo_id: issue.repo_id.clone(),
380 number: issue.number,
381 author: Some((&issue.author).into()),
382 requested_by: issue.requested_by.as_ref().map(Into::into),
383 ..IssueEvent::default()
384 }
385 }
386
387 /// The commit a pull request's change is at in git right now: its
388 /// fork's default branch, or its branch.
389 async fn live_head(&self, pull: &Pull) -> Result<Option<String>> {
390 g1t_kit::call(
391 &self.repos,
392 "head",
393 &HeadArgs {
394 repo_id: pull.fork_repo_id.clone().unwrap_or_else(|| pull.repo_id.clone()),
395 branch: pull.branch.clone().unwrap_or_default(),
396 },
397 )
398 .await
399 }
400
401 fn pull_event(pull: &Pull) -> PullEvent {
402 PullEvent {
403 pull_id: pull.id.clone(),
404 repo_id: pull.repo_id.clone(),
405 number: pull.number,
406 author: Some((&pull.author).into()),
407 requested_by: pull.requested_by.as_ref().map(Into::into),
408 issue: pull.issue,
409 confidence: pull.confidence.clone(),
410 ..PullEvent::default()
411 }
412 }
413
414 // --- Issues ------------------------------------------------------------
415
416 /// Opens an issue for g1t to take at once: refused before
417 /// anything is opened unless the actor may put agents to work here. The
418 /// runner's `delegate` starts the agent on it.
419 async fn delegate_issue(&self, a: DelegateIssueArgs) -> Result<Outcome<Issue>> {
420 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
421 check!(writable(&repo));
422 check!(allowed(Some(&a.actor), &repo, Capability::Run));
423 self.open_issue(OpenIssueArgs {
424 actor: a.actor,
425 repo: a.repo,
426 title: a.title,
427 body: a.body,
428 labels: a.labels,
429 checks: a.checks,
430 })
431 .await
432 }
433
434 async fn open_issue(&self, a: OpenIssueArgs) -> Result<Outcome<Issue>> {
435 if !a.actor.verified {
436 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
437 }
438 let title = match valid_title(&a.title) {
439 Ok(title) => title,
440 Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
441 };
442 let Some(labels) = normalize_labels(&a.labels) else {
443 return Ok(Outcome::fail(
444 FailureCode::Invalid,
445 "An issue can have up to 10 labels of up to 40 characters each.",
446 ));
447 };
448 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
449 check!(writable(&repo));
450 // Commands given the old way are words for the agent now: added to
451 // the body under "Definition of done". What has to pass to merge is
452 // the branch's required checks.
453 let body = with_definition_of_done(&a.body, &commands_pass(&a.checks));
454
455 let now = now_ms();
456 let id = new_id("iss", now);
457 let number = self.next_number(&repo.id).await?;
458 let timestamp = rfc3339(now);
459 // What g1t's agent files at work is g1t's, for the person it works for.
460 let (author, requested_by) = authorship(&a.actor, false);
461 self.db
462 .prepare(
463 "INSERT INTO issues
464 (id, repo_id, number, title, body, labels, checks, author_id, author_name,
465 requested_by_id, requested_by_name, created_at, updated_at)
466 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
467 )
468 .bind(&[
469 id.as_str().into(),
470 repo.id.as_str().into(),
471 number.into(),
472 title.into(),
473 body.into(),
474 serde_json::to_string(&labels)?.into(),
475 "[]".into(),
476 author.id.as_str().into(),
477 author.username.as_str().into(),
478 optional(&requested_by.as_ref().map(|user| user.id.clone())),
479 optional(&requested_by.as_ref().map(|user| user.username.clone())),
480 timestamp.as_str().into(),
481 timestamp.as_str().into(),
482 ])?
483 .run()
484 .await?;
485 let Some(issue) = self.issue(&repo.id, number).await? else {
486 return Ok(no_issue());
487 };
488 self.apply_label_rule(&a.actor, &issue, &[]).await?;
489 self.publish(
490 "issue.opened",
491 &repo.id,
492 &a.actor,
493 IssueEvent {
494 title: Some(issue.title.clone()),
495 ..Self::issue_event(&issue)
496 },
497 )
498 .await?;
499 Ok(Outcome::Ok(issue))
500 }
501
502 async fn list_issues(&self, a: ListIssuesArgs) -> Result<Outcome<Vec<Issue>>> {
503 let state = state_name(a.state);
504 let label = a
505 .label
506 .map(|label| label.trim().to_lowercase())
507 .filter(|label| !label.is_empty());
508 let list = |repo_id: String| {
509 let label = label.clone();
510 async move {
511 let state = state.map_or(JsValue::NULL, JsValue::from);
512 let query = self
513 .db
514 .prepare(format!(
515 "SELECT {ISSUE_COLUMNS} FROM issues
516 WHERE repo_id = ? AND (? IS NULL OR state = ?)
517 AND (? IS NULL OR EXISTS
518 (SELECT 1 FROM json_each(issues.labels) WHERE json_each.value = ?))
519 ORDER BY number DESC LIMIT ?"
520 ))
521 .bind(&[
522 repo_id.into(),
523 state.clone(),
524 state,
525 optional(&label),
526 optional(&label),
527 LIST_PAGE.into(),
528 ])?;
529 self.timing.db(1, query.all()).await?.results::<IssueRow>()
530 }
531 };
532 let (_, rows) = check!(self.repo_then(&a.repo, &a.viewer, list).await?);
533 Ok(Outcome::Ok(rows.into_iter().map(Issue::from).collect()))
534 }
535
536 /// An issue, the pull requests for it and its comments: one batch,
537 /// started beside the access check (prefetch.rs).
538 async fn get_issue(&self, a: ViewArgs) -> Result<Outcome<IssueDetail>> {
539 let number = a.number;
540 let read = |repo_id: String| async move {
541 let key = [JsValue::from(repo_id.as_str()), JsValue::from(number)];
542 let statements = vec![
543 self.db
544 .prepare(format!("SELECT {ISSUE_COLUMNS} FROM issues WHERE repo_id = ?1 AND number = ?2"))
545 .bind(&key)?,
546 self.db
547 .prepare(format!(
548 "SELECT {PULL_COLUMNS} FROM pulls
549 WHERE issue_id = (SELECT id FROM issues WHERE repo_id = ?1 AND number = ?2)
550 ORDER BY number"
551 ))
552 .bind(&key)?,
553 self.db
554 .prepare("SELECT * FROM comments WHERE repo_id = ?1 AND number = ?2 ORDER BY id LIMIT 500")
555 .bind(&key)?,
556 ];
557 let results = self.timing.db(3, self.db.batch(statements)).await?;
558 let rows = |index: usize| results.get(index).ok_or_else(|| worker::Error::RustError("short batch".into()));
559 Ok::<_, worker::Error>((
560 rows(0)?.results::<IssueRow>()?.into_iter().next().map(Issue::from),
561 rows(1)?.results::<PullRow>()?.into_iter().map(Pull::from).collect::<Vec<_>>(),
562 rows(2)?.results::<CommentRow>()?.into_iter().map(Comment::from).collect::<Vec<_>>(),
563 ))
564 };
565 let Outcome::Ok((_, (Some(issue), pulls, comments))) = self.repo_then(&a.repo, &a.viewer, read).await? else {
566 return Ok(no_issue());
567 };
568 Ok(Outcome::Ok(IssueDetail { comments, pulls, issue }))
569 }
570
571 /// The issue, if `actor` wrote it or may triage the repository's issues.
572 async fn manageable_issue(
573 &self,
574 actor: &User,
575 path: &RepoPath,
576 number: u32,
577 ) -> Result<Outcome<Issue>> {
578 let (repo, issue) = check!(self.issue_at(path, number, &Some(actor.clone())).await?);
579 check!(writable(&repo));
580 if issue.owner().id != actor.id {
581 check!(allowed(Some(actor), &repo, Capability::Triage));
582 }
583 Ok(Outcome::Ok(issue))
584 }
585
586 async fn update_issue(&self, a: UpdateIssueArgs) -> Result<Outcome<Issue>> {
587 let issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
588 let title = match a.title.as_deref().map(valid_title) {
589 Some(Err(message)) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
590 Some(Ok(title)) => Some(title.to_owned()),
591 None => None,
592 };
593 let labels = match a.labels.as_deref().map(normalize_labels) {
594 Some(None) => {
595 return Ok(Outcome::fail(
596 FailureCode::Invalid,
597 "An issue can have up to 10 labels of up to 40 characters each.",
598 ));
599 }
600 Some(Some(labels)) => Some(serde_json::to_string(&labels)?),
601 None => None,
602 };
603 let assignees = match a.assignees {
604 Some(names) => Some(check!(self.valid_assignees(names).await?)),
605 None => None,
606 };
607 let assigned = assignees.as_ref().map(serde_json::to_string).transpose()?;
608 let body = a.body.map(|body| body.trim().to_owned());
609 self.db
610 .prepare(
611 "UPDATE issues
612 SET title = COALESCE(?, title), body = COALESCE(?, body),
613 labels = COALESCE(?, labels), assignees = COALESCE(?, assignees),
614 updated_at = ?
615 WHERE id = ?",
616 )
617 .bind(&[
618 optional(&title),
619 optional(&body),
620 optional(&labels),
621 optional(&assigned),
622 rfc3339(now_ms()).into(),
623 issue.id.as_str().into(),
624 ])?
625 .run()
626 .await?;
627 let before = issue.assignees.clone();
628 let labels_before = issue.labels.clone();
629 let Some(issue) = self.issue(&issue.repo_id, issue.number).await? else {
630 return Ok(no_issue());
631 };
632 self.apply_label_rule(&a.actor, &issue, &labels_before).await?;
633 self.publish(
634 "issue.updated",
635 &issue.repo_id,
636 &a.actor,
637 Self::issue_event(&issue),
638 )
639 .await?;
640 if let Some(assignees) = assignees {
641 self.note_changes(
642 &issue.repo_id,
643 issue.number,
644 &a.actor,
645 &before,
646 &assignees,
647 ("assigned", "unassigned"),
648 )
649 .await?;
650 self.publish(
651 "issue.assigned",
652 &issue.repo_id,
653 &a.actor,
654 IssueEvent {
655 assignees: Some(assignees),
656 ..Self::issue_event(&issue)
657 },
658 )
659 .await?;
660 }
661 Ok(Outcome::Ok(issue))
662 }
663
664 /// Usernames as given, tidied, if each names an account.
665 async fn valid_assignees(&self, names: Vec<String>) -> Result<Outcome<Vec<String>>> {
666 let mut assignees: Vec<String> = Vec::new();
667 for name in names {
668 let name = name.trim().trim_start_matches('@').to_lowercase();
669 if name.is_empty() || assignees.contains(&name) {
670 continue;
671 }
672 if assignees.len() == MAX_ASSIGNEES {
673 return Ok(Outcome::fail(
674 FailureCode::Invalid,
675 format!("An issue can be assigned to at most {MAX_ASSIGNEES} people."),
676 ));
677 }
678 let account: Viewer = g1t_kit::call(
679 &self.identity,
680 "user_by_username",
681 &UsernameArgs {
682 username: name.clone(),
683 },
684 )
685 .await?;
686 if account.is_none() {
687 return Ok(Outcome::fail(
688 FailureCode::Invalid,
689 format!("There is no account named {name}."),
690 ));
691 }
692 assignees.push(name);
693 }
694 Ok(Outcome::Ok(assignees))
695 }
696
697 /// Open issues assigned to the viewer, in every repository. Callers
698 /// show only those in repositories the viewer can still see.
699 async fn list_assigned_issues(&self, a: ViewerArgs) -> Result<Vec<Issue>> {
700 let Some(viewer) = a.viewer else {
701 return Ok(Vec::new());
702 };
703 let rows = self
704 .db
705 .prepare(format!(
706 "SELECT {ISSUE_COLUMNS} FROM issues
707 WHERE state = 'open' AND EXISTS (
708 SELECT 1 FROM json_each(issues.assignees) WHERE json_each.value = ?)
709 ORDER BY updated_at DESC LIMIT 50"
710 ))
711 .bind(&[viewer.username.into()])?
712 .all()
713 .await?
714 .results::<IssueRow>()?;
715 Ok(rows.into_iter().map(Issue::from).collect())
716 }
717
718 async fn close_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
719 let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
720 if issue.state == State::Closed {
721 return Ok(Outcome::fail(
722 FailureCode::Conflict,
723 "This issue is already closed.",
724 ));
725 }
726 let reason = a.reason.unwrap_or(IssueReason::Completed);
727 let now = rfc3339(now_ms());
728 self.db
729 .prepare(
730 "UPDATE issues SET state = 'closed', reason = ?, closed_at = ?, updated_at = ?
731 WHERE id = ?",
732 )
733 .bind(&[
734 reason.as_str().into(),
735 now.as_str().into(),
736 now.as_str().into(),
737 issue.id.as_str().into(),
738 ])?
739 .run()
740 .await?;
741 self.publish(
742 "issue.closed",
743 &issue.repo_id,
744 &a.actor,
745 IssueEvent {
746 reason: Some(reason.as_str()),
747 ..Self::issue_event(&issue)
748 },
749 )
750 .await?;
751 self.note(
752 &issue.repo_id,
753 issue.number,
754 (&a.actor.id, &a.actor.username),
755 match reason {
756 IssueReason::Completed => "closed this as completed",
757 IssueReason::NotPlanned => "closed this as not planned",
758 },
759 )
760 .await?;
761 issue.state = State::Closed;
762 issue.reason = Some(reason);
763 issue.closed_at = Some(now.clone());
764 issue.updated_at = now;
765 Ok(Outcome::Ok(issue))
766 }
767
768 async fn reopen_issue(&self, a: IssueActionArgs) -> Result<Outcome<Issue>> {
769 let mut issue = check!(self.manageable_issue(&a.actor, &a.repo, a.number).await?);
770 if issue.state == State::Open {
771 return Ok(Outcome::fail(
772 FailureCode::Conflict,
773 "This issue is already open.",
774 ));
775 }
776 let now = rfc3339(now_ms());
777 self.db
778 .prepare(
779 "UPDATE issues
780 SET state = 'open', reason = NULL, resolved_by = NULL, closed_at = NULL,
781 updated_at = ?
782 WHERE id = ?",
783 )
784 .bind(&[now.as_str().into(), issue.id.as_str().into()])?
785 .run()
786 .await?;
787 self.publish(
788 "issue.reopened",
789 &issue.repo_id,
790 &a.actor,
791 Self::issue_event(&issue),
792 )
793 .await?;
794 self.note(
795 &issue.repo_id,
796 issue.number,
797 (&a.actor.id, &a.actor.username),
798 "reopened this",
799 )
800 .await?;
801 issue.state = State::Open;
802 issue.reason = None;
803 issue.resolved_by = None;
804 issue.closed_at = None;
805 issue.updated_at = now;
806 Ok(Outcome::Ok(issue))
807 }
808
809 /// The default labels, then every other label in use on the repository.
810 async fn list_labels(&self, a: ViewArgs) -> Result<Outcome<Vec<String>>> {
811 let read = |repo_id: String| async move {
812 let query = self
813 .db
814 .prepare(
815 "SELECT DISTINCT json_each.value AS value
816 FROM issues, json_each(issues.labels)
817 WHERE issues.repo_id = ? ORDER BY 1 LIMIT 200",
818 )
819 .bind(&[repo_id.into()])?;
820 self.timing.db(1, query.all()).await?.results::<ValueRow>()
821 };
822 let (_, used) = check!(self.repo_then(&a.repo, &a.viewer, read).await?);
823 let mut labels: Vec<String> = DEFAULT_LABELS.iter().map(|label| (*label).into()).collect();
824 for row in used {
825 if !labels.contains(&row.value) {
826 labels.push(row.value);
827 }
828 }
829 Ok(Outcome::Ok(labels))
830 }
831
832 async fn counts(&self, a: ViewArgs) -> Result<Outcome<Counts>> {
833 let read = |repo_id: String| async move {
834 let query = self
835 .db
836 .prepare(
837 "SELECT
838 (SELECT count(*) FROM issues WHERE repo_id = ?1 AND state = 'open') AS issues,
839 (SELECT count(*) FROM pulls
840 WHERE repo_id = ?1 AND status IN ('draft', 'open')) AS pulls",
841 )
842 .bind(&[repo_id.into()])?;
843 self.timing.db(1, query.first::<Counts>(None)).await
844 };
845 let (_, counts) = check!(self.repo_then(&a.repo, &a.viewer, read).await?);
846 Ok(Outcome::Ok(counts.unwrap_or(Counts {
847 issues: 0,
848 pulls: 0,
849 })))
850 }
851
852 // --- Comments ----------------------------------------------------------
853
854 async fn add_comment(&self, a: AddCommentArgs) -> Result<Outcome<Comment>> {
855 if !a.actor.verified {
856 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
857 }
858 let body = a.body.trim();
859 // An approval speaks for itself; anything else has to say something.
860 if body.is_empty() && a.verdict != Some(Verdict::Approve) {
861 return Ok(Outcome::fail(
862 FailureCode::Invalid,
863 "A comment cannot be empty.",
864 ));
865 }
866 let path = a
867 .path
868 .as_deref()
869 .map(str::trim)
870 .filter(|path| !path.is_empty());
871 let line = a.line.filter(|line| *line > 0 && path.is_some());
872 if body.chars().count() > MAX_ENTRY_CHARS {
873 return Ok(Outcome::fail(
874 FailureCode::Invalid,
875 "That comment is too long.",
876 ));
877 }
878 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
879 check!(writable(&repo));
880 // The number names an issue or a pull request, never both.
881 let mut pull_id = None;
882 let table = if self.issue(&repo.id, a.number).await?.is_some() {
883 if path.is_some() || a.verdict.is_some() {
884 return Ok(Outcome::fail(
885 FailureCode::Invalid,
886 "Only a pull request can be reviewed or commented on by line.",
887 ));
888 }
889 "issues"
890 } else if let Some(pull) = self.pull(&repo.id, a.number).await? {
891 if a.verdict.is_some() && pull.is_owned_by(&a.actor.id) {
892 return Ok(Outcome::fail(
893 FailureCode::Forbidden,
894 "You cannot approve or request changes on your own pull request.",
895 ));
896 }
897 pull_id = Some(pull.id.clone());
898 "pulls"
899 } else {
900 return Ok(Outcome::fail(
901 FailureCode::NotFound,
902 "No issue or pull request has that number.",
903 ));
904 };
905
906 let now = now_ms();
907 let comment = Comment {
908 kind: CommentKind::Comment,
909 id: new_id("cmt", now),
910 author: a.actor.clone(),
911 body: body.to_owned(),
912 path: path.map(str::to_owned),
913 line,
914 verdict: a.verdict,
915 created_at: rfc3339(now),
916 };
917 self.db
918 .batch(vec![
919 self.db
920 .prepare(
921 "INSERT INTO comments
922 (id, repo_id, number, author_id, author_name, body, path, line,
923 verdict, created_at)
924 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
925 )
926 .bind(&[
927 comment.id.as_str().into(),
928 repo.id.as_str().into(),
929 a.number.into(),
930 a.actor.id.as_str().into(),
931 a.actor.username.as_str().into(),
932 body.into(),
933 optional(&comment.path),
934 optional_number(line),
935 a.verdict
936 .map_or(JsValue::NULL, |verdict| verdict.as_str().into()),
937 comment.created_at.as_str().into(),
938 ])?,
939 self.db
940 .prepare(format!(
941 "UPDATE {table} SET updated_at = ? WHERE repo_id = ? AND number = ?"
942 ))
943 .bind(&[
944 comment.created_at.as_str().into(),
945 repo.id.as_str().into(),
946 a.number.into(),
947 ])?,
948 ])
949 .await?;
950 self.note_mention(&a.actor, &repo, a.number, &comment, pull_id.as_deref()).await?;
951 self.publish(
952 "comment.created",
953 &repo.id,
954 &a.actor,
955 CommentCreated {
956 comment_id: comment.id.clone(),
957 repo_id: repo.id.clone(),
958 number: a.number,
959 pull_id,
960 verdict: a.verdict,
961 },
962 )
963 .await?;
964 Ok(Outcome::Ok(comment))
965 }
966
967 // --- Pull requests -----------------------------------------------------
968
969 async fn open_pull(&self, a: OpenPullArgs) -> Result<Outcome<Pull>> {
970 if !a.actor.verified {
971 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
972 }
973 let repo = check!(self.repo(&a.repo, &Some(a.actor.clone())).await?);
974 check!(writable(&repo));
975 // g1t's own agent at work spends the workspace's compute; a pull
976 // request anyone else's agent makes is like any other.
977 if matches!(a.runtime, Runtime::Hosted) {
978 check!(allowed(Some(&a.actor), &repo, Capability::Run));
979 }
980 let issue = match a.issue {
981 Some(number) => match self.issue(&repo.id, number).await? {
982 Some(issue) if issue.state == State::Open => Some(issue),
983 Some(_) => {
984 return Ok(Outcome::fail(
985 FailureCode::Conflict,
986 "This issue is closed.",
987 ));
988 }
989 None => return Ok(no_issue()),
990 },
991 None => None,
992 };
993 // A pull request for an issue takes the issue's title unless given one.
994 let title = match (a.title.trim(), &issue) {
995 ("", Some(issue)) => issue.title.clone(),
996 (title, _) => match valid_title(title) {
997 Ok(title) => title.to_owned(),
998 Err(message) => return Ok(Outcome::fail(FailureCode::Invalid, message)),
999 },
1000 };
1001 let agent = match a.agent.trim() {
1002 "" => "agent",
1003 agent => agent,
1004 };
1005 let runtime = match a.runtime {
1006 Runtime::Hosted => "hosted",
1007 Runtime::External => "external",
1008 };
1009
1010 let now = now_ms();
1011 let id = new_id("pr", now);
1012 let branch = a
1013 .branch
1014 .as_deref()
1015 .map(str::trim)
1016 .filter(|branch| !branch.is_empty());
1017 // The change is on a branch already pushed to the repository, or
1018 // will be made in a fork created for this pull request.
1019 let (fork, head) = match branch {
1020 Some(branch) => {
1021 if branch == repo.default_branch {
1022 return Ok(Outcome::fail(
1023 FailureCode::Invalid,
1024 format!("Choose a branch other than {branch}."),
1025 ));
1026 }
1027 let head: Option<String> = g1t_kit::call(
1028 &self.repos,
1029 "head",
1030 &HeadArgs {
1031 repo_id: repo.id.clone(),
1032 branch: branch.to_owned(),
1033 },
1034 )
1035 .await?;
1036 let Some(head) = head else {
1037 return Ok(Outcome::fail(
1038 FailureCode::NotFound,
1039 format!("There is no branch named {branch}. Push it first."),
1040 ));
1041 };
1042 let existing = self
1043 .db
1044 .prepare(
1045 "SELECT number AS n FROM pulls
1046 WHERE repo_id = ? AND source_branch = ? AND status IN ('draft', 'open')",
1047 )
1048 .bind(&[repo.id.as_str().into(), branch.into()])?
1049 .first::<NumberRow>(None)
1050 .await?;
1051 if let Some(existing) = existing {
1052 return Ok(Outcome::fail(
1053 FailureCode::Conflict,
1054 format!("Pull request #{} is already open for {branch}.", existing.n),
1055 ));
1056 }
1057 (None, Some(head))
1058 }
1059 None => {
1060 let fork: Outcome<Repo> = g1t_kit::call(
1061 &self.repos,
1062 "fork_for_pull",
1063 &ForkArgs {
1064 source_id: repo.id.clone(),
1065 pull_id: id.clone(),
1066 actor: a.actor.clone(),
1067 },
1068 )
1069 .await?;
1070 (Some(check!(fork)), None)
1071 }
1072 };
1073 // A branch already holds the work, so its pull request is ready for
1074 // review from the start; one with a fork starts as a draft.
1075 let status = if branch.is_some() { "open" } else { "draft" };
1076 let body = Some(a.body.trim().to_owned()).filter(|body| !body.is_empty());
1077
1078 let number = self.next_number(&repo.id).await?;
1079 let timestamp = rfc3339(now);
1080 // A change g1t makes is g1t's, for whoever asked for it. This is
1081 // what lifecycle::made_by_g1t reads back.
1082 let by_g1t = matches!(a.runtime, Runtime::Hosted) && agent == reviews::AGENT_NAME && fork.is_some();
1083 let (author, requested_by) = authorship(&a.actor, by_g1t);
1084 self.db
1085 .prepare(
1086 "INSERT INTO pulls
1087 (id, repo_id, number, issue_id, issue_number, title, body, agent, runtime,
1088 status, fork_repo_id, fork_namespace, fork_name, source_branch, head_commit,
1089 author_id, author_name, requested_by_id, requested_by_name, created_at, updated_at)
1090 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
1091 )
1092 .bind(&[
1093 id.as_str().into(),
1094 repo.id.as_str().into(),
1095 number.into(),
1096 optional(&issue.as_ref().map(|issue| issue.id.clone())),
1097 optional_number(issue.as_ref().map(|issue| issue.number)),
1098 title.into(),
1099 optional(&body),
1100 agent.into(),
1101 runtime.into(),
1102 status.into(),
1103 optional(&fork.as_ref().map(|fork| fork.id.clone())),
1104 optional(&fork.as_ref().map(|fork| fork.namespace.clone())),
1105 optional(&fork.as_ref().map(|fork| fork.name.clone())),
1106 optional(&branch.map(str::to_owned)),
1107 optional(&head),
1108 author.id.as_str().into(),
1109 author.username.as_str().into(),
1110 optional(&requested_by.as_ref().map(|user| user.id.clone())),
1111 optional(&requested_by.as_ref().map(|user| user.username.clone())),
1112 timestamp.as_str().into(),
1113 timestamp.as_str().into(),
1114 ])?
1115 .run()
1116 .await?;
1117 let Some(pull) = self.pull(&repo.id, number).await? else {
1118 return Ok(no_pull());
1119 };
1120 self.manage(&pull).await?;
1121 // Someone is on it now, so it is no longer waiting for an agent.
1122 if let Some(issue) = pull.issue {
1123 self.db
1124 .prepare("UPDATE issues SET queued_by = NULL WHERE repo_id = ? AND number = ?")
1125 .bind(&[repo.id.as_str().into(), issue.into()])?
1126 .run()
1127 .await?;
1128 }
1129 if let Some(issue) = pull.issue {
1130 let text = if lifecycle::made_by_g1t(&pull) {
1131 format!("assigned this to g1t, which opened #{}", pull.number)
1132 } else {
1133 format!("opened #{} for this", pull.number)
1134 };
1135 self.note(&repo.id, issue, (&a.actor.id, &a.actor.username), &text)
1136 .await?;
1137 }
1138 self.publish(
1139 "pull.opened",
1140 &repo.id,
1141 &a.actor,
1142 PullEvent {
1143 agent: Some(pull.agent.clone()),
1144 ..Self::pull_event(&pull)
1145 },
1146 )
1147 .await?;
1148 Ok(Outcome::Ok(pull))
1149 }
1150
1151 async fn list_pulls(&self, a: ListPullsArgs) -> Result<Outcome<Vec<Pull>>> {
1152 let filter = match a.state {
1153 Some(State::Open) => "AND status IN ('draft', 'open')",
1154 Some(State::Closed) => "AND status IN ('merged', 'closed')",
1155 None => "",
1156 };
1157 let read = |repo_id: String| async move {
1158 let query = self
1159 .db
1160 .prepare(format!(
1161 "SELECT {PULL_COLUMNS} FROM pulls WHERE repo_id = ? {filter} ORDER BY number DESC LIMIT ?"
1162 ))
1163 .bind(&[repo_id.into(), LIST_PAGE.into()])?;
1164 self.timing.db(1, query.all()).await?.results::<PullRow>()
1165 };
1166 let (_, rows) = check!(self.repo_then(&a.repo, &a.viewer, read).await?);
1167 Ok(Outcome::Ok(rows.into_iter().map(Pull::from).collect()))
1168 }
1169
1170 /// `pulls_for_repos`: what `list_pulls` gives, open and closed, for many
1171 /// repositories at once: one access check with repos for all of them
1172 /// and one query, instead of two of each per repository.
1173 async fn pulls_for_repos(&self, a: PullsForReposArgs) -> Result<Vec<RepoPulls>> {
1174 let ids: Vec<String> = a.repo_ids.into_iter().take(MAX_PULLS_FOR_REPOS).collect();
1175 if ids.is_empty() {
1176 return Ok(Vec::new());
1177 }
1178 let limit = a.limit.clamp(1, LIST_PAGE);
1179 // The rows are read beside the access check, for every id asked
1180 // about; those of repositories the viewer cannot read are dropped.
1181 let asked = serde_json::to_string(&ids)?;
1182 let check = ReadableArgs { ids, viewer: a.viewer };
1183 let readable = self.timing.rpc(g1t_kit::call::<_, Vec<Repo>>(&self.repos, "readable", &check));
1184 let (readable, rows) = try_join(readable, self.timing.db(1, self.newest_pulls(asked, limit))).await?;
1185 if readable.is_empty() {
1186 return Ok(Vec::new());
1187 }
1188 let mut answer: Vec<RepoPulls> = readable
1189 .iter()
1190 .map(|repo| RepoPulls { repo_id: repo.id.clone(), open: Vec::new(), closed: Vec::new() })
1191 .collect();
1192 for pull in rows.into_iter().map(Pull::from) {
1193 let Some(entry) = answer.iter_mut().find(|entry| entry.repo_id == pull.repo_id) else {
1194 continue;
1195 };
1196 match pull.status {
1197 PullStatus::Draft | PullStatus::Open => entry.open.push(pull),
1198 PullStatus::Merged | PullStatus::Closed => entry.closed.push(pull),
1199 }
1200 }
1201 for entry in &mut answer {
1202 entry.open.sort_by_key(|pull| std::cmp::Reverse(pull.number));
1203 entry.closed.sort_by_key(|pull| std::cmp::Reverse(pull.number));
1204 }
1205 Ok(answer)
1206 }
1207
1208 /// The newest `limit` of each repository's open (draft or open) and
1209 /// closed (merged or closed) pull requests, for the ids in `ids` (JSON).
1210 async fn newest_pulls(&self, ids: String, limit: u32) -> Result<Vec<PullRow>> {
1211 self.db
1212 .prepare(format!(
1213 "SELECT * FROM (
1214 SELECT {PULL_COLUMNS}, ROW_NUMBER() OVER (
1215 PARTITION BY pulls.repo_id, pulls.status IN ('draft', 'open') ORDER BY pulls.number DESC
1216 ) AS place
1217 FROM pulls WHERE pulls.repo_id IN (SELECT value FROM json_each(?1))
1218 ) WHERE place <= ?2"
1219 ))
1220 .bind(&[ids.into(), limit.into()])?
1221 .all()
1222 .await?
1223 .results::<PullRow>()
1224 }
1225
1226 async fn get_pull(&self, a: ViewArgs) -> Result<Outcome<PullDetail>> {
1227 let number = a.number;
1228 // Every row the page and the lifecycle read, in one batch started
1229 // beside the access check; the helpers below read from it.
1230 let read = |repo_id: String| self.prefetch_pull(repo_id, number);
1231 let Outcome::Ok((repo, Some(found))) = self.repo_then(&a.repo, &a.viewer, read).await? else {
1232 return Ok(no_pull());
1233 };
1234 let Some(row) = found.first::<PullRow>(prefetch::Slot::Pull)? else {
1235 return Ok(no_pull());
1236 };
1237 let stored = found.first::<StoredBehind>(prefetch::Slot::Pull)?;
1238 let issue = found.first::<IssueRow>(prefetch::Slot::Issue)?.map(Issue::from);
1239 let comments: Vec<Comment> =
1240 found.rows::<CommentRow>(prefetch::Slot::Comments)?.into_iter().map(Comment::from).collect();
1241 self.keep_prefetched(Some(found));
1242 let detail = self.pull_detail(repo, Pull::from(row), issue, comments, stored).await;
1243 self.keep_prefetched(None);
1244 detail
1245 }
1246
1247 async fn pull_detail(
1248 &self,
1249 repo: Repo,
1250 mut pull: Pull,
1251 issue: Option<Issue>,
1252 comments: Vec<Comment>,
1253 stored: Option<StoredBehind>,
1254 ) -> Result<Outcome<PullDetail>> {
1255 // Whether it is behind, as worked out with its mergeability on the
1256 // last push to either side (mergeability.rs), when that was for
1257 // its head as it is now; otherwise asked of the repos service.
1258 let known_behind = stored.and_then(|stored| stored.for_head(pull.head_commit.as_deref()));
1259 // Worked out on each push; this covers a pull request from before
1260 // that was recorded.
1261 if pull.files.is_empty() && pull.head_commit.is_some() {
1262 pull.files = self.refresh_files(&pull).await?;
1263 }
1264 // Everything else at once: none of it depends on the rest, and each
1265 // is a round trip of its own.
1266 let standing = async {
1267 // Mergeability first: where g1t sees a pull request through, a
1268 // conflict decides its next step.
1269 let behind = async {
1270 match known_behind {
1271 Some(behind) => Ok(behind),
1272 None => {
1273 let behind = self.is_behind(&repo.id, &pull).await?;
1274 // Kept for the next view when the mergeability on
1275 // record is for this head: a pull request from
1276 // before `behind` was kept asks once.
1277 if let Some(head) = pull.head_commit.as_deref()
1278 && pull.status.is_active()
1279 {
1280 self.db
1281 .prepare(
1282 "UPDATE pulls SET behind = ?1
1283 WHERE id = ?2 AND behind IS NULL AND mergeable_key LIKE ?3 || '..%'",
1284 )
1285 .bind(&[u32::from(behind).into(), pull.id.as_str().into(), head.into()])?
1286 .run()
1287 .await?;
1288 }
1289 Ok(behind)
1290 }
1291 }
1292 };
1293 let (merge, behind) = try_join(self.mergeability(&pull), behind).await?;
1294 let assessed = self.assess_with_confidence(&pull, &issue, behind).await?;
1295 let confidence = assessed.as_ref().and_then(|(_, _, confidence)| confidence.clone());
1296 let lifecycle = assessed.map(|(lifecycle, _, _)| lifecycle);
1297 Ok::<_, worker::Error>((behind, (lifecycle, confidence), merge))
1298 };
1299 let (((behind, (lifecycle, confidence), (mergeable, conflicts)), (landing, stalled), comments), (checks, overlaps, review_pending)) =
1300 try_join(
1301 try_join3(standing, self.landing_state(&pull.id), async { Ok(comments) }),
1302 try_join3(
1303 self.latest_checks(&pull.id),
1304 self.overlaps(&pull),
1305 self.review_pending(&pull.id),
1306 ),
1307 )
1308 .await?;
1309 // As just worked out, rather than as it was read.
1310 if confidence.is_some() {
1311 pull.confidence = confidence;
1312 }
1313 let (statuses, settings) =
1314 try_join(self.statuses(&repo.id, pull.head_commit.as_deref()), self.settings(&repo.id)).await?;
1315 Ok(Outcome::Ok(PullDetail {
1316 required_checks: required_checks(&settings.required_checks, &statuses),
1317 comments,
1318 checks,
1319 overlaps,
1320 behind,
1321 review_pending,
1322 lifecycle,
1323 landing,
1324 stalled,
1325 messages: self.messages(&pull.id).await?,
1326 statuses,
1327 mergeable,
1328 conflicts,
1329 earlier_checks: self.earlier_checks(&pull.id).await?,
1330 issue,
1331 pull,
1332 }))
1333 }
1334
1335 /// The pull request, if it is still active and `actor` opened it or
1336 /// may triage the repository's pull requests.
1337 async fn manageable_pull(
1338 &self,
1339 actor: &User,
1340 path: &RepoPath,
1341 number: u32,
1342 ) -> Result<Outcome<Pull>> {
1343 let (repo, pull) = check!(self.pull_at(path, number, &Some(actor.clone())).await?);
1344 check!(writable(&repo));
1345 if !pull.is_owned_by(&actor.id) {
1346 check!(allowed(Some(actor), &repo, Capability::Triage));
1347 }
1348 if !pull.status.is_active() {
1349 return Ok(Outcome::fail(
1350 FailureCode::Conflict,
1351 format!("This pull request is already {}.", pull.status.as_str()),
1352 ));
1353 }
1354 Ok(Outcome::Ok(pull))
1355 }
1356
1357 /// Brings a pull request up to date with the default branch without a
1358 /// sandbox, where the repos service can do that safely. Whoever could
1359 /// have pushed the merge themselves may ask: whoever opened it (or asked
1360 /// g1t for it), for a fork; anyone who may push, for a branch of the
1361 /// repository. When it needs a
1362 /// real merge, says so, naming the conflicting files if a probe found
1363 /// them, and pushes nothing.
1364 async fn catch_up_pull(&self, a: PullActionArgs) -> Result<Outcome<PullBranchUpdate>> {
1365 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &Some(a.actor.clone())).await?);
1366 check!(writable(&repo));
1367 if !pull.status.is_active() {
1368 return Ok(Outcome::fail(
1369 FailureCode::Conflict,
1370 format!("This pull request is already {}.", pull.status.as_str()),
1371 ));
1372 }
1373 if pull.fork_repo_id.is_some() {
1374 if !pull.is_owned_by(&a.actor.id) {
1375 return Ok(Outcome::fail(
1376 FailureCode::Forbidden,
1377 "Only whoever opened this pull request, or asked g1t for it, can update it.",
1378 ));
1379 }
1380 } else {
1381 check!(allowed(Some(&a.actor), &repo, Capability::Push));
1382 }
1383 let updated: Outcome<PullBranchUpdate> = g1t_kit::call(
1384 &self.repos,
1385 "update_pull_branch",
1386 &UpdatePullBranchArgs {
1387 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
1388 branch: pull.branch.clone(),
1389 number: pull.number,
1390 actor: a.actor,
1391 },
1392 )
1393 .await?;
1394 // A probe that found conflicts says more than "both changed it".
1395 if let Outcome::Ok(PullBranchUpdate::NeedsAgent { .. }) = &updated
1396 && let Some(files) = self.conflicting_files(&pull).await?
1397 && !files.is_empty()
1398 {
1399 return Ok(Outcome::Ok(PullBranchUpdate::NeedsAgent {
1400 reason: NeedsAgentReason::Conflicting,
1401 detail: "Merging it conflicts.".to_owned(),
1402 paths: files,
1403 }));
1404 }
1405 Ok(updated)
1406 }
1407
1408 async fn update_pull(&self, a: UpdatePullArgs) -> Result<Outcome<Pull>> {
1409 let pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1410 let assignees = match a.assignees {
1411 Some(names) => Some(check!(self.valid_assignees(names).await?)),
1412 None => None,
1413 };
1414 let reviewers = match a.reviewers {
1415 Some(names) => {
1416 // g1t is not an account; everyone else has to be.
1417 let agent = names
1418 .iter()
1419 .any(|name| name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME));
1420 let people = names
1421 .into_iter()
1422 .filter(|name| !name.trim().eq_ignore_ascii_case(reviews::AGENT_NAME))
1423 .collect();
1424 let mut reviewers = check!(self.valid_assignees(people).await?);
1425 // Nobody is asked to review their own, nor what they had g1t make.
1426 reviewers.retain(|name| *name != pull.owner().username);
1427 if agent {
1428 reviewers.insert(0, reviews::AGENT_NAME.to_owned());
1429 }
1430 Some(reviewers)
1431 }
1432 None => None,
1433 };
1434 self.db
1435 .prepare(
1436 "UPDATE pulls
1437 SET assignees = COALESCE(?, assignees), reviewers = COALESCE(?, reviewers),
1438 updated_at = ?
1439 WHERE id = ?",
1440 )
1441 .bind(&[
1442 optional(&assignees.as_ref().map(serde_json::to_string).transpose()?),
1443 optional(&reviewers.as_ref().map(serde_json::to_string).transpose()?),
1444 rfc3339(now_ms()).into(),
1445 pull.id.as_str().into(),
1446 ])?
1447 .run()
1448 .await?;
1449 if let Some(assignees) = &assignees {
1450 self.note_changes(
1451 &pull.repo_id,
1452 pull.number,
1453 &a.actor,
1454 &pull.assignees,
1455 assignees,
1456 ("assigned", "unassigned"),
1457 )
1458 .await?;
1459 }
1460 if let Some(reviewers) = &reviewers {
1461 self.note_changes(
1462 &pull.repo_id,
1463 pull.number,
1464 &a.actor,
1465 &pull.reviewers,
1466 reviewers,
1467 (
1468 "requested a review from",
1469 "withdrew the request for a review from",
1470 ),
1471 )
1472 .await?;
1473 }
1474 Ok(match self.pull(&pull.repo_id, pull.number).await? {
1475 Some(pull) => Outcome::Ok(pull),
1476 None => no_pull(),
1477 })
1478 }
1479
1480 /// Marks a draft ready for review, or updates the description of one
1481 /// that already is.
1482 async fn ready_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1483 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1484 let summary = Some(a.summary.trim().to_owned()).filter(|summary| !summary.is_empty());
1485 let now = rfc3339(now_ms());
1486 self.db
1487 .prepare(
1488 "UPDATE pulls SET status = 'open', body = COALESCE(?, body), updated_at = ?
1489 WHERE id = ?",
1490 )
1491 .bind(&[
1492 optional(&summary),
1493 now.as_str().into(),
1494 pull.id.as_str().into(),
1495 ])?
1496 .run()
1497 .await?;
1498 if pull.status == PullStatus::Draft {
1499 // The head as it is now: the push that came just before may not
1500 // have reached `head_commit` yet, and workflows run on it.
1501 let commit = self.live_head(&pull).await?.or_else(|| pull.head_commit.clone());
1502 self.publish(
1503 "pull.ready",
1504 &pull.repo_id,
1505 &a.actor,
1506 PullEvent {
1507 commit,
1508 ..Self::pull_event(&pull)
1509 },
1510 )
1511 .await?;
1512 }
1513 if pull.status == PullStatus::Draft {
1514 self.note(
1515 &pull.repo_id,
1516 pull.number,
1517 (&a.actor.id, &a.actor.username),
1518 "marked this ready for review",
1519 )
1520 .await?;
1521 }
1522 pull.status = PullStatus::Open;
1523 pull.body = summary.or(pull.body);
1524 pull.updated_at = now;
1525 Ok(Outcome::Ok(pull))
1526 }
1527
1528 async fn close_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1529 let mut pull = check!(self.manageable_pull(&a.actor, &a.repo, a.number).await?);
1530 let now = rfc3339(now_ms());
1531 self.db
1532 .prepare("UPDATE pulls SET status = 'closed', updated_at = ? WHERE id = ?")
1533 .bind(&[now.as_str().into(), pull.id.as_str().into()])?
1534 .run()
1535 .await?;
1536 self.publish(
1537 "pull.closed",
1538 &pull.repo_id,
1539 &a.actor,
1540 Self::pull_event(&pull),
1541 )
1542 .await?;
1543 self.note(
1544 &pull.repo_id,
1545 pull.number,
1546 (&a.actor.id, &a.actor.username),
1547 "closed this",
1548 )
1549 .await?;
1550 // A closed pull request leaves the merge queue.
1551 if self
1552 .leave(&pull.repo_id, &pull, QueueState::Removed, Some("It was closed."))
1553 .await?
1554 {
1555 self.publish_as(
1556 "queue.changed",
1557 &pull.repo_id,
1558 None,
1559 g1t_contracts::events::QueueChanged {
1560 repo_id: pull.repo_id.clone(),
1561 },
1562 )
1563 .await?;
1564 }
1565 pull.status = PullStatus::Closed;
1566 pull.updated_at = now;
1567 Ok(Outcome::Ok(pull))
1568 }
1569
1570 /// Lands the pull request on the repository's default branch. Unless
1571 /// told to keep it open, that resolves the issue it was for: the issue
1572 /// closes naming this pull request, and the others still in progress
1573 /// for it close as superseded.
1574 async fn merge_pull(&self, a: PullActionArgs) -> Result<Outcome<Pull>> {
1575 let viewer = Some(a.actor.clone());
1576 let (repo, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
1577 check!(writable(&repo));
1578 match pull.status {
1579 PullStatus::Open => {}
1580 PullStatus::Draft => {
1581 return Ok(Outcome::fail(
1582 FailureCode::Conflict,
1583 "This pull request is still a draft. Mark it ready for review first.",
1584 ));
1585 }
1586 status => {
1587 return Ok(Outcome::fail(
1588 FailureCode::Conflict,
1589 format!("This pull request is already {}.", status.as_str()),
1590 ));
1591 }
1592 }
1593 let settings = self.settings(&repo.id).await?;
1594 // The default branch's protection: its required checks must pass on
1595 // the head, for a person's pull request and an agent's alike. Where
1596 // the repository does not allow bypassing them, asking to bypass
1597 // them changes nothing.
1598 if !a.ignore_checks || !settings.allow_ignoring_checks {
1599 let queue = (pull.check_status == Some(CheckStatus::Failed))
1600 .then(|| "It failed in the merge queue; push a fix to try again.".to_owned());
1601 let required = statuses::WorkflowFacts::of(
1602 &self.statuses(&repo.id, pull.head_commit.as_deref()).await?,
1603 &settings.required_checks,
1604 )
1605 .refusal();
1606 if let Some(reason) = queue.or(required) {
1607 let remedy = if settings.allow_ignoring_checks {
1608 "Wait or fix them, or bypass the required checks as you merge."
1609 } else {
1610 "This repository only merges pull requests whose required checks pass."
1611 };
1612 return Ok(Outcome::fail(
1613 FailureCode::Conflict,
1614 format!("{reason} {remedy}"),
1615 ));
1616 }
1617 }
1618 if let Some(missing) = self.approvals_gap(&settings, &pull).await? {
1619 return Ok(Outcome::fail(FailureCode::Conflict, missing));
1620 }
1621 // Known ahead of time to conflict: neither a merge nor the queue
1622 // would get through, so say what has to be resolved now.
1623 if let Some(files) = self.conflicting_files(&pull).await? {
1624 let named = if files.is_empty() {
1625 String::new()
1626 } else {
1627 format!(" in {}", files.join(", "))
1628 };
1629 return Ok(Outcome::fail(
1630 FailureCode::Conflict,
1631 format!(
1632 "This branch has conflicts with {}{named} that must be resolved first. Have g1t resolve them, or merge {0} into it, fix them and push.",
1633 repo.default_branch
1634 ),
1635 ));
1636 }
1637
1638 // A repository that merges through a queue: it joins the queue, and
1639 // lands once its state together with everything ahead has passed.
1640 if settings.merge_queue {
1641 if !a.actor.verified {
1642 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1643 }
1644 check!(allowed(Some(&a.actor), &repo, Capability::Merge));
1645 return self.enqueue(&repo, &pull, &a.actor, a.keep_issue_open).await;
1646 }
1647
1648 // The default branch has moved under it. Unless the repository
1649 // insists on that being dealt with first, bring it up to date and
1650 // land it when that is done.
1651 if self.is_behind(&repo.id, &pull).await? {
1652 if settings.require_up_to_date {
1653 return Ok(Outcome::fail(
1654 FailureCode::Conflict,
1655 format!(
1656 "{} 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 {0} first.",
1657 repo.default_branch
1658 ),
1659 ));
1660 }
1661 if !a.actor.verified {
1662 return Ok(Outcome::fail(FailureCode::Forbidden, UNVERIFIED));
1663 }
1664 check!(allowed(Some(&a.actor), &repo, Capability::Merge));
1665 self.request_landing(&pull, &a.actor, a.keep_issue_open)
1666 .await?;
1667 return Ok(Outcome::Ok(pull));
1668 }
1669
1670 // Whether the actor may write to the repository is decided by repos.
1671 let landed: Outcome<Landed> = g1t_kit::call(
1672 &self.repos,
1673 "land",
1674 &LandArgs {
1675 // A pull request from a branch lands from the repository itself.
1676 source_id: pull.fork_repo_id.clone().unwrap_or_else(|| repo.id.clone()),
1677 branch: pull.branch.clone(),
1678 actor: a.actor.clone(),
1679 },
1680 )
1681 .await?;
1682 let landed = check!(landed);
1683 Ok(Outcome::Ok(
1684 self.record_merge(&repo, pull, &a.actor, a.keep_issue_open, landed)
1685 .await?,
1686 ))
1687 }
1688
1689 /// Records a pull request as merged once the default branch holds it:
1690 /// closes its issue, supersedes the others for it, and says so.
1691 pub(crate) async fn record_merge(
1692 &self,
1693 repo: &Repo,
1694 mut pull: Pull,
1695 actor: &User,
1696 keep_issue_open: bool,
1697 landed: Landed,
1698 ) -> Result<Pull> {
1699 let issue = match pull.issue {
1700 Some(number) if !keep_issue_open => self
1701 .issue(&repo.id, number)
1702 .await?
1703 .filter(|issue| issue.state == State::Open),
1704 _ => None,
1705 };
1706 let now = rfc3339(now_ms());
1707 let mut statements = vec![
1708 self.db
1709 .prepare(
1710 "UPDATE pulls
1711 SET status = 'merged', head_commit = ?, merge_base = ?, merged_by = ?,
1712 merged_at = ?, updated_at = ?
1713 WHERE id = ?",
1714 )
1715 .bind(&[
1716 landed.commit.as_str().into(),
1717 optional(&landed.previous),
1718 actor.username.as_str().into(),
1719 now.as_str().into(),
1720 now.as_str().into(),
1721 pull.id.as_str().into(),
1722 ])?,
1723 ];
1724 if let Some(issue) = &issue {
1725 statements.push(
1726 self.db
1727 .prepare(
1728 "UPDATE issues
1729 SET state = 'closed', reason = 'completed', resolved_by = ?,
1730 closed_at = ?, updated_at = ?
1731 WHERE id = ?",
1732 )
1733 .bind(&[
1734 pull.number.into(),
1735 now.as_str().into(),
1736 now.as_str().into(),
1737 issue.id.as_str().into(),
1738 ])?,
1739 );
1740 statements.push(
1741 self.db
1742 .prepare(
1743 "UPDATE pulls SET status = 'closed', superseded_by = ?, updated_at = ?
1744 WHERE issue_id = ? AND id != ? AND status IN ('draft', 'open')",
1745 )
1746 .bind(&[
1747 pull.number.into(),
1748 now.as_str().into(),
1749 issue.id.as_str().into(),
1750 pull.id.as_str().into(),
1751 ])?,
1752 );
1753 }
1754 self.db.batch(statements).await?;
1755
1756 self.publish(
1757 "pull.merged",
1758 &repo.id,
1759 actor,
1760 PullEvent {
1761 commit: Some(landed.commit.clone()),
1762 ..Self::pull_event(&pull)
1763 },
1764 )
1765 .await?;
1766 if let Some(issue) = &issue {
1767 self.publish(
1768 "issue.closed",
1769 &repo.id,
1770 actor,
1771 IssueEvent {
1772 reason: Some(IssueReason::Completed.as_str()),
1773 resolved_by: Some(pull.number),
1774 ..Self::issue_event(issue)
1775 },
1776 )
1777 .await?;
1778 }
1779
1780 let who = (actor.id.as_str(), actor.username.as_str());
1781 self.note(&repo.id, pull.number, who, "merged this").await?;
1782 if let Some(issue) = &issue {
1783 self.note(
1784 &repo.id,
1785 issue.number,
1786 who,
1787 &format!("closed this by merging #{}", pull.number),
1788 )
1789 .await?;
1790 }
1791 pull.status = PullStatus::Merged;
1792 pull.head_commit = Some(landed.commit.clone());
1793 pull.merge_base = landed.previous;
1794 pull.merged_by = Some(actor.username.clone());
1795 pull.merged_at = Some(now.clone());
1796 pull.updated_at = now;
1797 Ok(pull)
1798 }
1799
1800 async fn list_active_pulls(&self, a: ViewerArgs) -> Result<Vec<ActivePull>> {
1801 let Some(viewer) = a.viewer else {
1802 return Ok(Vec::new());
1803 };
1804 // The pull requests and their issues, in one round trip: their own,
1805 // and those g1t made for them (Pull::owner).
1806 let author = [JsValue::from(viewer.id.as_str())];
1807 let found = self
1808 .timing
1809 .db(
1810 2,
1811 self.db.batch(vec![
1812 self.db
1813 .prepare(format!(
1814 "SELECT {PULL_COLUMNS} FROM pulls
1815 WHERE COALESCE(requested_by_id, author_id) = ?1 AND status IN ('draft', 'open')
1816 ORDER BY updated_at DESC LIMIT 50"
1817 ))
1818 .bind(&author)?,
1819 self.db
1820 .prepare(format!(
1821 "SELECT {ISSUE_COLUMNS} FROM issues WHERE issues.id IN (
1822 SELECT issue_id FROM pulls
1823 WHERE COALESCE(requested_by_id, author_id) = ?1 AND status IN ('draft', 'open') AND issue_id IS NOT NULL
1824 ORDER BY updated_at DESC LIMIT 50)"
1825 ))
1826 .bind(&author)?,
1827 ]),
1828 )
1829 .await?;
1830 let (Some(found), Some(issues)) = (found.first(), found.get(1)) else {
1831 return Ok(Vec::new());
1832 };
1833 let snapshots = found.results::<Snapshot>()?;
1834 let pulls: Vec<Pull> = found.results::<PullRow>()?.into_iter().map(Pull::from).collect();
1835 let issues: Vec<Issue> = issues.results::<IssueRow>()?.into_iter().map(Issue::from).collect();
1836 let issues = &issues;
1837 // Where each stands: the remembered assessment when there is one,
1838 // and worked out otherwise.
1839 try_join_all(pulls.into_iter().zip(snapshots).map(|(pull, snapshot)| async move {
1840 let issue = pull.issue.and_then(|number| {
1841 issues
1842 .iter()
1843 .find(|issue| issue.repo_id == pull.repo_id && issue.number == number)
1844 .cloned()
1845 });
1846 // Only a pull request g1t is seeing through has a lifecycle.
1847 let lifecycle = if !lifecycle::made_by_g1t(&pull) || snapshot.managed == 0 {
1848 None
1849 } else if let (Some(stage), Some(detail)) = (snapshot.stage, snapshot.stage_detail) {
1850 Some(Lifecycle {
1851 stage,
1852 detail,
1853 revisions: snapshot.revisions,
1854 })
1855 } else {
1856 let behind = self.is_behind(&pull.repo_id, &pull).await?;
1857 self.assess(&pull, &issue, behind)
1858 .await?
1859 .map(|(lifecycle, _)| lifecycle)
1860 };
1861 Ok::<_, worker::Error>(ActivePull {
1862 pull,
1863 issue,
1864 lifecycle,
1865 })
1866 }))
1867 .await
1868 }
1869
1870 // --- Sessions ----------------------------------------------------------
1871
1872 async fn append_session(&self, a: AppendSessionArgs) -> Result<Outcome<Appended>> {
1873 if a.entries.is_empty() {
1874 return Ok(Outcome::Ok(Appended { count: 0 }));
1875 }
1876 if a.entries.len() > MAX_ENTRY_BATCH {
1877 return Ok(Outcome::fail(
1878 FailureCode::Invalid,
1879 format!("Send at most {MAX_ENTRY_BATCH} entries at a time."),
1880 ));
1881 }
1882 let viewer = Some(a.actor.clone());
1883 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &viewer).await?);
1884 if !pull.is_owned_by(&a.actor.id) {
1885 return Ok(Outcome::fail(
1886 FailureCode::Forbidden,
1887 "Only whoever opened a pull request, or asked g1t for it, can record its session.",
1888 ));
1889 }
1890
1891 let now = rfc3339(now_ms());
1892 let count = a.entries.len() as u32;
1893 let mut statements = Vec::with_capacity(a.entries.len() + 1);
1894 for entry in a.entries {
1895 let kind = serde_json::to_value(entry.kind)?;
1896 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
1897 // Each insert takes the next sequence number itself, so two
1898 // writers appending at once cannot collide.
1899 statements.push(
1900 self.db
1901 .prepare(
1902 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
1903 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
1904 FROM session_entries WHERE pull_id = ?",
1905 )
1906 .bind(&[
1907 pull.id.as_str().into(),
1908 kind.as_str().unwrap_or("note").into(),
1909 text.into(),
1910 optional(&entry.tool),
1911 optional(&entry.commit.or_else(|| pull.head_commit.clone())),
1912 now.as_str().into(),
1913 pull.id.as_str().into(),
1914 ])?,
1915 );
1916 }
1917 statements.push(
1918 self.db
1919 .prepare("UPDATE pulls SET updated_at = ? WHERE id = ?")
1920 .bind(&[now.as_str().into(), pull.id.as_str().into()])?,
1921 );
1922 self.db.batch(statements).await?;
1923 self.publish(
1924 "session.appended",
1925 &pull.repo_id,
1926 &a.actor,
1927 SessionAppended {
1928 pull_id: pull.id.clone(),
1929 repo_id: pull.repo_id.clone(),
1930 number: pull.number,
1931 count,
1932 },
1933 )
1934 .await?;
1935 Ok(Outcome::Ok(Appended { count }))
1936 }
1937
1938 /// Adds entries to a pull request's session, each taking the next
1939 /// sequence number, without announcing it.
1940 pub(crate) async fn append_entries(&self, pull: &Pull, entries: &[NewSessionEntry]) -> Result<()> {
1941 let now = rfc3339(now_ms());
1942 let mut statements = Vec::with_capacity(entries.len());
1943 for entry in entries {
1944 let kind = serde_json::to_value(entry.kind)?;
1945 let text: String = entry.text.chars().take(MAX_ENTRY_CHARS).collect();
1946 statements.push(
1947 self.db
1948 .prepare(
1949 "INSERT INTO session_entries (pull_id, seq, kind, text, tool, \"commit\", at)
1950 SELECT ?, COALESCE(MAX(seq), 0) + 1, ?, ?, ?, ?, ?
1951 FROM session_entries WHERE pull_id = ?",
1952 )
1953 .bind(&[
1954 pull.id.as_str().into(),
1955 kind.as_str().unwrap_or("note").into(),
1956 text.into(),
1957 optional(&entry.tool),
1958 optional(&entry.commit.clone().or_else(|| pull.head_commit.clone())),
1959 now.as_str().into(),
1960 pull.id.as_str().into(),
1961 ])?,
1962 );
1963 }
1964 self.db.batch(statements).await?;
1965 Ok(())
1966 }
1967
1968 async fn read_session(&self, a: ViewArgs) -> Result<Outcome<Vec<SessionEntry>>> {
1969 let (_, pull) = check!(self.pull_at(&a.repo, a.number, &a.viewer).await?);
1970 let rows = self
1971 .db
1972 .prepare(
1973 "SELECT seq, kind, text, tool, \"commit\", at FROM session_entries
1974 WHERE pull_id = ? AND seq > ? ORDER BY seq LIMIT ?",
1975 )
1976 .bind(&[pull.id.into(), a.after_seq.into(), SESSION_PAGE.into()])?
1977 .all()
1978 .await?
1979 .results::<SessionRow>()?;
1980 Ok(Outcome::Ok(
1981 rows.into_iter().map(SessionEntry::from).collect(),
1982 ))
1983 }
1984
1985 /// A push moves the head of the pull request it concerns: the one whose
1986 /// fork was pushed to, or the one opened from the branch that moved.
1987 async fn on_event(&self, event: &Event) -> Result<()> {
1988 if event.kind != "git.push" {
1989 return Ok(());
1990 }
1991 let (Some(repo_id), Some(after), Some(git_ref)) = (
1992 event.repo_id.as_deref(),
1993 event.data["after"].as_str(),
1994 event.data["ref"].as_str(),
1995 ) else {
1996 return Ok(());
1997 };
1998 let now = rfc3339(now_ms());
1999 // The head moved, so whatever the checks said no longer applies, and
2000 // whatever step g1t was waiting on has been taken.
2001 let moved = "UPDATE pulls
2002 SET head_commit = ?, updated_at = ?, check_status = NULL, check_run_id = NULL,
2003 working_on = NULL, working_until = NULL, stalled = NULL";
2004 let active = "status IN ('draft', 'open') AND head_commit IS NOT ?";
2005 let returning =
2006 "RETURNING id, repo_id, number, issue_number, status, author_id, author_name, requested_by_id, requested_by_name";
2007 let mut pulls: Vec<MovedRow> = Vec::new();
2008 // A fork carries its pull request on its default branch.
2009 if event.data["defaultBranch"].as_bool() == Some(true) {
2010 pulls.extend(
2011 self.db
2012 .prepare(format!(
2013 "{moved} WHERE fork_repo_id = ? AND {active} {returning}"
2014 ))
2015 .bind(&[
2016 after.into(),
2017 now.as_str().into(),
2018 repo_id.into(),
2019 after.into(),
2020 ])?
2021 .all()
2022 .await?
2023 .results::<MovedRow>()?,
2024 );
2025 }
2026 if let Some(branch) = git_ref.strip_prefix("refs/heads/") {
2027 pulls.extend(
2028 self.db
2029 .prepare(format!(
2030 "{moved} WHERE repo_id = ? AND source_branch = ? AND {active} {returning}"
2031 ))
2032 .bind(&[
2033 after.into(),
2034 now.as_str().into(),
2035 repo_id.into(),
2036 branch.into(),
2037 after.into(),
2038 ])?
2039 .all()
2040 .await?
2041 .results::<MovedRow>()?,
2042 );
2043 }
2044 // What each now changes, so overlaps show while the work is under way.
2045 for moved in &pulls {
2046 if let Some(pull) = self.pull_by_id(&moved.id).await? {
2047 self.refresh_files(&pull).await?;
2048 }
2049 }
2050 // A merge that was waiting for this push to bring it up to date.
2051 for moved in &pulls {
2052 self.land_if_requested(&moved.id).await?;
2053 }
2054 // Whether each still merges cleanly, and, when a default branch
2055 // moved, every open pull request into it.
2056 let moved_ids: Vec<String> = pulls.iter().map(|pull| pull.id.clone()).collect();
2057 self.after_push(repo_id, event.data["defaultBranch"].as_bool() == Some(true), &moved_ids)
2058 .await;
2059 // A draft is announced when it is marked ready instead.
2060 for pull in pulls
2061 .into_iter()
2062 .filter(|pull| pull.status == PullStatus::Open)
2063 {
2064 self.publish_as(
2065 "pull.updated",
2066 &pull.repo_id,
2067 event.actor.clone(),
2068 PullEvent {
2069 author: Some(g1t_contracts::credentials::Principal { id: pull.author_id, username: pull.author_name }),
2070 requested_by: pull
2071 .requested_by_id
2072 .zip(pull.requested_by_name)
2073 .map(|(id, username)| g1t_contracts::credentials::Principal { id, username }),
2074 pull_id: pull.id,
2075 repo_id: pull.repo_id.clone(),
2076 number: pull.number,
2077 issue: pull.issue_number,
2078 commit: Some(after.to_owned()),
2079 ..PullEvent::default()
2080 },
2081 )
2082 .await?;
2083 }
2084 Ok(())
2085 }
2086}
2087
2088fn service(env: &Env) -> Result<Work> {
2089 Ok(Work {
2090 db: env.d1("DB")?,
2091 identity: env.service("IDENTITY")?,
2092 repos: env.service("REPOS")?,
2093 events: env.service("EVENTS")?,
2094 actions: env.service("ACTIONS")?,
2095 timing: g1t_kit::d1::Timing::default(),
2096 prefetched: std::cell::RefCell::new(None),
2097 })
2098}
2099
2100#[event(fetch)]
2101async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> {
2102 let Some(method) = rpc_method(&request) else {
2103 return Response::error("Not found", 404);
2104 };
2105 // A replica near the caller when it asks for one (crates/kit/src/d1.rs).
2106 let (db, served) = g1t_kit::d1::open(&env, "DB", &request)?;
2107 let body: serde_json::Value = request.json().await?;
2108 let mut work = service(&env)?;
2109 work.db = db;
2110
2111 let answered = match method.as_str() {
2112 "open_issue" => reply(&work.open_issue(args(body)?).await?),
2113 "delegate_issue" => reply(&work.delegate_issue(args(body)?).await?),
2114 "report_confidence" => reply(&work.report_confidence(args(body)?).await?),
2115 "list_issues" => reply(&work.list_issues(args(body)?).await?),
2116 "get_issue" => reply(&work.get_issue(args(body)?).await?),
2117 "update_issue" => reply(&work.update_issue(args(body)?).await?),
2118 "close_issue" => reply(&work.close_issue(args(body)?).await?),
2119 "reopen_issue" => reply(&work.reopen_issue(args(body)?).await?),
2120 "list_labels" => reply(&work.list_labels(args(body)?).await?),
2121 "counts" => reply(&work.counts(args(body)?).await?),
2122 "add_comment" => reply(&work.add_comment(args(body)?).await?),
2123 "start_checks" => reply(&work.start_checks(args(body)?).await?),
2124 "seen_checks" => reply(&work.seen_checks(args(body)?).await?),
2125 "report_checks" => reply(&work.report_checks(args(body)?).await?),
2126 "set_commit_status" => reply(&work.set_commit_status(args(body)?).await?),
2127 "start_review" => reply(&work.start_review(args(body)?).await?),
2128 "advance" => reply(&work.advance(args(body)?).await?),
2129 "stall" => reply(&work.stall(args(body)?).await?),
2130 "managed_pulls" => reply(&work.managed_pulls(args(body)?).await?),
2131 "queue" => reply(&work.queue(args(body)?).await?),
2132 "queue_build" => reply(&work.queue_build(args(body)?).await?),
2133 "report_queue" => reply(&work.report_queue(args(body)?).await?),
2134 "remove_from_queue" => reply(&work.remove_from_queue(args(body)?).await?),
2135 "message_agent" => reply(&work.message_agent(args(body)?).await?),
2136 "locate_pull" => reply(&work.locate_pull(args(body)?).await?),
2137 "answer_message" => reply(&work.answer_message(args(body)?).await?),
2138 "take_messages" => reply(&work.take_messages(args(body)?).await?),
2139 "wake_for_messages" => reply(&work.wake_for_messages(args(body)?).await?),
2140 "catch_up_job" => reply(&work.catch_up_job(args(body)?).await?),
2141 "get_settings" => reply(&work.get_settings(args(body)?).await?),
2142 "update_settings" => reply(&work.update_settings(args(body)?).await?),
2143 "report_review" => reply(&work.report_review(args(body)?).await?),
2144 "open_pull" => reply(&work.open_pull(args(body)?).await?),
2145 "list_pulls" => reply(&work.list_pulls(args(body)?).await?),
2146 "pulls_for_repos" => reply(&work.pulls_for_repos(args(body)?).await?),
2147 "get_pull" => reply(&work.get_pull(args(body)?).await?),
2148 "update_pull" => reply(&work.update_pull(args(body)?).await?),
2149 "catch_up_pull" => reply(&work.catch_up_pull(args(body)?).await?),
2150 "ready_pull" => reply(&work.ready_pull(args(body)?).await?),
2151 "close_pull" => reply(&work.close_pull(args(body)?).await?),
2152 "merge_pull" => reply(&work.merge_pull(args(body)?).await?),
2153 "list_active_pulls" => reply(&work.list_active_pulls(args(body)?).await?),
2154 "by_author" => reply(&work.by_author(args(body)?).await?),
2155 "start_plan" => reply(&work.start_plan(args(body)?).await?),
2156 "report_plan" => reply(&work.report_plan(args(body)?).await?),
2157 "get_plan" => reply(&work.get_plan(args(body)?).await?),
2158 "list_plans" => reply(&work.list_plans(args(body)?).await?),
2159 "apply_plan" => reply(&work.apply_plan(args(body)?).await?),
2160 "queue_issue" => reply(&work.queue_issue(args(body)?).await?),
2161 "ready_issues" => reply(&work.ready_issues(args(body)?).await?),
2162 "list_assigned_issues" => reply(&work.list_assigned_issues(args(body)?).await?),
2163 "append_session" => reply(&work.append_session(args(body)?).await?),
2164 "read_session" => reply(&work.read_session(args(body)?).await?),
2165 // Agents at work, their sessions, and memory (runs.rs, memory.rs).
2166 "open_run" => reply(&work.open_run(args(body)?).await?),
2167 "report_run" => reply(&work.report_run(args(body)?).await?),
2168 "stop_run" => reply(&work.stop_run(args(body)?).await?),
2169 "list_runs" => reply(&work.list_runs(args(body)?).await?),
2170 "get_run" => reply(&work.get_run(args(body)?).await?),
2171 "list_sessions" => reply(&work.list_sessions(args(body)?).await?),
2172 "get_session" => reply(&work.get_session(args(body)?).await?),
2173 "list_memories" => reply(&work.list_memories(args(body)?).await?),
2174 "add_memory" => reply(&work.add_memory(args(body)?).await?),
2175 "update_memory" => reply(&work.update_memory(args(body)?).await?),
2176 "delete_memory" => reply(&work.delete_memory(args(body)?).await?),
2177 "recall" => reply(&work.recall(args(body)?).await?),
2178 "memory_context" => reply(&work.memory_context(args(body)?).await?),
2179 // What agents may do in a sandbox (guardrails.rs).
2180 "get_guardrails" => reply(&work.get_guardrails(args(body)?).await?),
2181 "update_guardrails" => reply(&work.update_guardrails(args(body)?).await?),
2182 "run_guardrails" => reply(&work.run_guardrails(args(body)?).await?),
2183 // Plan caps the runner applies (compute.rs).
2184 "active_agents" => reply(&work.active_agents(args(body)?).await?),
2185 "issue_spend" => reply(&work.issue_spend(args(body)?).await?),
2186 "wait_for_slot" => reply(&work.wait_for_slot(args(body)?).await?),
2187 "agent_comment" => reply(&work.agent_comment(args(body)?).await?),
2188 "add_wait" => reply(&work.add_wait(args(body)?).await?),
2189 "waiting_workspaces" => reply(&work.waiting_workspaces(args(body)?).await?),
2190 "take_wait" => reply(&work.take_wait(args(body)?).await?),
2191 // The runs whose sandboxes stop with their repository (retired.rs).
2192 "runs_in_repo" => reply(&work.runs_in_repo(args(body)?).await?),
2193 "run_cost" => reply(&work.run_cost(args(body)?).await?),
2194 "start_mergecheck" => reply(&work.start_mergecheck(args(body)?).await?),
2195 "report_mergecheck" => reply(&work.report_mergecheck(args(body)?).await?),
2196 // Memory that fills itself, and its review queue (capture.rs).
2197 method if capture::METHODS.contains(&method) => capture::dispatch(&work, method, body).await,
2198 // @g1t in comments, and the label rule (mentions.rs).
2199 "take_mention" => reply(&work.take_mention(args(body)?).await?),
2200 "mention_revision" => reply(&work.mention_revision(args(body)?).await?),
2201 "reply_mention" => reply(&work.reply_mention(args(body)?).await?),
2202 "get_agent_rules" => reply(&work.get_agent_rules(args(body)?).await?),
2203 "set_agent_rules" => reply(&work.set_agent_rules(args(body)?).await?),
2204 _ => Response::error("Unknown method", 404),
2205 };
2206 served.finish_timed(answered, &work.timing)
2207}
2208
2209/// Events from the bus, delivered on this service's own queue.
2210#[event(queue)]
2211async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
2212 let work = service(&env)?;
2213 for message in batch.messages()? {
2214 // A workspace renamed: its agent runs and memory move to the slug it has now.
2215 if g1t_kit::rename::on_event(&env, &env.d1("DB")?, message.body(), &[memory::RENAMED, guardrails::RENAMED].concat()).await? {
2216 message.ack();
2217 continue;
2218 }
2219 // A repository renamed or transferred: its runs, memory, guardrails
2220 // and runs waiting for a slot follow.
2221 if g1t_kit::transfer::on_event(&env, &env.d1("DB")?, message.body(), &[memory::TRANSFERRED, guardrails::TRANSFERRED, retired::WAITS_MOVED].concat()).await? {
2222 message.ack();
2223 continue;
2224 }
2225 // A workspace deleted: what it kept for itself goes.
2226 if g1t_kit::deleted::on_event(&env.d1("DB")?, message.body(), memory::DELETED).await? {
2227 message.ack();
2228 continue;
2229 }
2230 // A repository deleted, archived or purged, or a branch renamed (retired.rs).
2231 if work.on_retired(message.body()).await? {
2232 message.ack();
2233 continue;
2234 }
2235 capture::on_event(&work, message.body()).await;
2236 work.on_event(message.body()).await?;
2237 message.ack();
2238 }
2239 Ok(())
2240}
2241
2242/// The rules that once read a pull request's author read its owner now:
2243/// whoever asked g1t for it, or its author. For each, the person who asked
2244/// is held to what an author was, and g1t's agent (a token it works with)
2245/// gains nothing by being the author.
2246#[cfg(test)]
2247mod owner_rules {
2248 use super::*;
2249 use crate::rows::stored::{ASKER, G1T, pull};
2250 use g1t_contracts::identity::AGENT_ID;
2251
2252 const SOMEONE: &str = "usr_2";
2253
2254 #[test]
2255 fn no_self_approval() {
2256 // add_comment refuses a verdict on one that is theirs.
2257 let made = pull(G1T, Some(ASKER));
2258 assert!(made.is_owned_by(ASKER.0), "the person who asked cannot approve it");
2259 assert!(!made.is_owned_by(AGENT_ID), "g1t's review agent still gives its verdict");
2260 assert!(!made.is_owned_by(SOMEONE));
2261 }
2262
2263 #[test]
2264 fn what_an_author_could_do_without_a_role() {
2265 // manageable_pull (update, ready, close), catch_up_pull on a fork,
2266 // append_session, and steering with message_agent: theirs to do.
2267 let made = pull(G1T, Some(ASKER));
2268 assert!(made.is_owned_by(ASKER.0));
2269 assert!(!made.is_owned_by(SOMEONE), "anyone else still needs the role");
2270 assert!(!made.is_owned_by(AGENT_ID), "being its author gives g1t's tokens nothing more");
2271 }
2272
2273 #[test]
2274 fn nobody_is_asked_to_review_what_they_asked_for() {
2275 // update_pull drops the owner from the reviewers asked.
2276 let made = pull(G1T, Some(ASKER));
2277 let mut reviewers = vec!["syntaqx".to_owned(), "ana".to_owned()];
2278 reviewers.retain(|name| *name != made.owner().username);
2279 assert_eq!(reviewers, ["ana"]);
2280 }
2281
2282 #[test]
2283 fn sandboxes_act_as_whoever_asked() {
2284 // LifecycleJob, ReviewJob, MergecheckJob and the merge queue's job
2285 // carry who the sandbox's credential acts for: a real account.
2286 let made = pull(G1T, Some(ASKER));
2287 assert_eq!(made.owner().id, ASKER.0);
2288 let acts_as = made.requested_by.unwrap_or(made.author);
2289 assert_eq!(acts_as.id, ASKER.0);
2290 // g1t's own work, which nobody asked for, acts as g1t, as before.
2291 let own = pull(("g1t", "g1t"), None);
2292 assert_eq!(own.requested_by.unwrap_or(own.author).id, "g1t");
2293 }
2294
2295 #[test]
2296 fn events_name_g1t_and_whoever_asked() {
2297 let made = pull(G1T, Some(ASKER));
2298 let event = serde_json::to_value(Work::pull_event(&made)).unwrap();
2299 assert_eq!(event["author"], serde_json::json!({ "id": AGENT_ID, "username": "g1t" }));
2300 assert_eq!(event["requestedBy"], serde_json::json!({ "id": "usr_1", "username": "syntaqx" }));
2301 let own = serde_json::to_value(Work::pull_event(&pull(ASKER, None))).unwrap();
2302 assert!(own.get("requestedBy").is_none());
2303 }
2304}