g1t/services/work/src/lib.rs

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