flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/work/src/lib.rs

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