Skip to content

g1t/services/work/src/lib.rs

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