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