Skip to content
2,480 linesCodeBlameRaw

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Inbox: the events service tells people what needs them as events arrive1//! The inbox, kept beside the event log. See `g1t_contracts::inbox`.
2//!
Inbox: threads, reasons, subscriptions and watching3//! As each batch arrives from the bus, the events that tell someone are
4//! read against what they name (the work service's `inbox_subject`), who
5//! subscribes to it and who watches its repository (subscriptions.rs), and
6//! each person told gets their thread about it brought to the top, unread,
7//! with a line added to its history. Who is told is worked out in
8//! [`notices`], from the event, its subject and that audience alone:
Inbox: the events service tells people what needs them as events arrive9//!
Inbox: threads, reasons, subscriptions and watching10//! | Event | Who | Reason | Severity |
11//! | --- | --- | --- | --- |
12//! | `agent.asked`, `pull.stalled` | the pull request's owner, and its issue's owner and assignees | agent | warning |
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar13//! | `pull.review_requested` | the reviewers asked, and the people each team asked tells | review_requested | warning |
Inbox: threads, reasons, subscriptions and watching14//! | `issue.assigned`, `pull.assigned` | the people newly assigned | assign | info |
15//! | `checks.completed`, failed or errored | the pull request's owner | ci_activity | error |
16//! | `workflow.completed`, failed | the pull request's owner, or whoever pushed | ci_activity | error |
17//! | `deployment.failed` | the pull request's owner, or whoever pushed; watchers | ci_activity | error |
18//! | `deployment.succeeded` after a failure | the same | ci_activity | success |
19//! | `review.completed` by g1t | the pull request's owner | author | success, or info for changes asked |
20//! | `pull.ready` for a change g1t made | whoever asked g1t for it | author | success |
21//! | `pull.merged`, `pull.closed`, `issue.closed`, `issue.reopened` | everyone subscribed | state_change | success for a merge, else info |
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar22//! | `comment.created` | everyone mentioned, the people of teams mentioned, then everyone subscribed | mention, team_mention, or why they are subscribed | info (success for an approval) |
23//! | `issue.opened`, `pull.opened` | whoever was assigned, asked to review or mentioned in its description (people and teams); watchers | assign, review_requested, mention, team_mention, subscribed | info |
Inbox: the events service tells people what needs them as events arrive24//!
Inbox: threads, reasons, subscriptions and watching25//! Watchers of a repository at `all` (or `custom`, for the kinds they
26//! chose) hear of every issue and pull request opened, commented on,
27//! closed, reopened or merged, and of deployments. Anyone who ignores the
28//! thread or the repository hears of nothing on it; anyone who
29//! unsubscribed hears only of what is asked of them.
30//!
31//! Nobody is told of what they did themselves, though an outcome they set
32//! off (checks, a workflow, a deployment, g1t's work) is theirs to hear of.
33//! g1t is never told. A failure here is logged and the batch goes on: the
34//! bus never waits on the inbox, so an item can be missed, but nothing else
35//! is held up.
Inbox: the events service tells people what needs them as events arrive36
37use std::collections::{HashMap, HashSet};
38
39use g1t_contracts::credentials::Principal;
40use g1t_contracts::events::{Event, WorkspaceRenamed};
41use g1t_contracts::identity::{AGENT_ID, UsernamesArgs};
42use g1t_contracts::inbox::*;
43use g1t_contracts::repos::{PathByIdArgs, ReadableArgs, Repo, RepoPath};
44use g1t_contracts::time::rfc3339;
45use g1t_contracts::{new_id, system};
46use g1t_kit::now_ms;
47use serde::Deserialize;
48use worker::wasm_bindgen::JsValue;
Inbox: threads, reasons, subscriptions and watching49use worker::{D1Database, D1PreparedStatement, Fetcher, Result};
50
51use crate::subscriptions::{self, Audience};
Inbox: the events service tells people what needs them as events arrive52
53/// Items marked done are kept this long, then removed.
54pub const DONE_DAYS: u32 = 30;
55/// No item is kept longer than this, unless it was saved.
56pub const MAX_DAYS: u32 = 180;
57/// Unread warnings shown ahead of everything else on the first page.
58const MAX_RANKED: u32 = 20;
59const MAX_TITLE: usize = 200;
60const MAX_BODY: usize = 300;
61
62/// What the inbox asks about an event before deciding who is told.
63#[derive(Debug, PartialEq, Eq)]
64pub struct Wanted {
65 pub repo_id: String,
66 /// The issue or pull request, when the event names one.
67 pub number: Option<u32>,
68 pub comment_id: Option<String>,
69}
70
71/// One person to tell, and what.
72#[derive(Debug, PartialEq, Eq)]
73pub struct Notice {
74 pub username: String,
Inbox: threads, reasons, subscriptions and watching75 pub reason: Reason,
Inbox: the events service tells people what needs them as events arrive76 pub severity: Severity,
77 pub title: String,
78 pub body: String,
79}
80
81/// Who did it, as far as is known.
82#[derive(Debug, Default)]
83pub struct Actor {
84 pub id: Option<String>,
85 pub username: Option<String>,
86}
87
88impl Actor {
89 fn is(&self, person: &Principal) -> bool {
90 self.id.as_deref().is_some_and(|id| id == person.id) || self.is_named(&person.username)
91 }
92
93 fn is_named(&self, username: &str) -> bool {
94 self.username.as_deref().is_some_and(|name| name.eq_ignore_ascii_case(username))
95 }
Inbox: threads, reasons, subscriptions and watching96
97 /// A person's name, for a title: never g1t's ids.
98 fn name(&self) -> Option<&str> {
99 self.username.as_deref().filter(|name| !is_g1t(name))
100 }
Inbox: the events service tells people what needs them as events arrive101}
102
Inbox: threads, reasons, subscriptions and watching103pub(crate) fn is_g1t(username: &str) -> bool {
Inbox: the events service tells people what needs them as events arrive104 username.eq_ignore_ascii_case(system::USERNAME) || username.eq_ignore_ascii_case("g1t-agent")
105}
106
Inbox: threads, reasons, subscriptions and watching107pub(crate) fn is_g1t_id(id: &str) -> bool {
Inbox: the events service tells people what needs them as events arrive108 system::is_system_id(id) || id == AGENT_ID
109}
110
Inbox: threads, reasons, subscriptions and watching111/// Events that end what an agent was waiting on a person for: it picked
112/// back up, its head moved, a merge was asked for, or it is over.
113const RESUMES: [&str; 5] = ["pull.resumed", "pull.updated", "pull.merge_requested", "pull.merged", "pull.closed"];
114
Inbox: the events service tells people what needs them as events arrive115/// The issue or pull request an event names, and what to read for it.
116/// None for events the inbox does not tell anyone of.
117pub fn wants(event: &Event) -> Option<Wanted> {
118 let data = &event.data;
119 let text = |key: &str| data[key].as_str().filter(|value| !value.is_empty()).map(str::to_owned);
120 let number = |key: &str| data[key].as_u64().and_then(|n| u32::try_from(n).ok());
121 let repo_id = text("repoId").or_else(|| event.repo_id.clone())?;
122 let on = |number: Option<u32>, comment_id: Option<String>| {
123 Some(Wanted {
124 repo_id: repo_id.clone(),
125 number,
126 comment_id,
127 })
128 };
129 match event.kind.as_str() {
Inbox: threads, reasons, subscriptions and watching130 "agent.asked" | "pull.stalled" | "pull.merged" | "pull.closed" | "pull.ready" | "pull.opened"
131 | "pull.review_requested" | "pull.assigned" | "issue.opened" | "issue.closed" | "issue.reopened"
132 | "issue.assigned" => on(Some(number("number")?), None),
Inbox: the events service tells people what needs them as events arrive133 "checks.completed" => match data["status"].as_str() {
134 Some("failed" | "errored") => on(Some(number("number")?), None),
135 _ => None,
136 },
137 "review.completed" => match data["verdict"].as_str() {
138 Some("approve" | "request_changes") => on(Some(number("number")?), None),
139 _ => None,
140 },
141 "workflow.completed" => match data["conclusion"].as_str() {
142 Some("failure") => on(number("pull"), None),
143 _ => None,
144 },
Inbox: threads, reasons, subscriptions and watching145 "deployment.failed" | "deployment.succeeded" => on(number("number"), None),
Merge branch 'worktree-agent-a3abfcce648e87dca'146 // A workflow run's jobs wait for their environment's reviewers.
147 "deployment.review_requested" => on(None, None),
Inbox: the events service tells people what needs them as events arrive148 "comment.created" => on(Some(number("number")?), Some(text("commentId")?)),
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar149 // Security alerts are threads of their own, not of an issue.
150 kind if SECURITY_EVENTS.contains(&kind) => on(None, None),
Inbox: the events service tells people what needs them as events arrive151 _ => None,
152 }
153}
154
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar155/// The security service's events the inbox tells people of: new alerts,
156/// and push protection bypasses asked for and decided. Fixes and
157/// dismissals are on the Security page and in webhooks.
158pub const SECURITY_EVENTS: [&str; 5] = [
159 "secret_scanning_alert.created",
160 "code_scanning_alert.created",
161 "vulnerability_alert.created",
162 "secret_scanning.bypass_requested",
163 "secret_scanning.bypass_reviewed",
164];
165
Inbox: threads, reasons, subscriptions and watching166/// Where an event's items go: which thread, what it is, and where it is.
167#[derive(Clone, Debug, PartialEq, Eq)]
168pub struct Thread {
169 pub key: String,
170 pub kind: Option<SubjectKind>,
171 pub number: Option<u32>,
172 pub run_id: Option<String>,
173 /// A path, for a subject with a page of its own.
174 pub link: Option<String>,
175}
176
177/// The thread key of an issue or pull request.
178pub fn numbered_thread(repo_id: &str, number: u32) -> String {
179 format!("{repo_id}#{number}")
180}
181
182/// The thread an event's items go to.
183pub fn thread_of(event: &Event, wanted: &Wanted, subject: Option<&InboxSubject>) -> Thread {
184 let data = &event.data;
185 let text = |key: &str| data[key].as_str().filter(|value| !value.is_empty()).map(str::to_owned);
186 match event.kind.as_str() {
187 // A project's production, or one pull request's preview.
188 "deployment.failed" | "deployment.succeeded" => {
189 let which = wanted.number.map_or_else(|| "production".to_owned(), |number| number.to_string());
190 Thread {
191 key: format!("{}/deploy/{}/{which}", wanted.repo_id, text("projectId").unwrap_or_default()),
192 kind: Some(SubjectKind::Deploy),
193 number: wanted.number,
194 run_id: text("deploymentId"),
195 link: text("path"),
196 }
197 }
Merge branch 'worktree-agent-a3abfcce648e87dca'198 // One run's deployment to one environment, waiting for review.
199 "deployment.review_requested" => Thread {
200 key: format!(
201 "{}/review/{}/{}",
202 wanted.repo_id,
203 text("runId").unwrap_or_default(),
204 text("environment").unwrap_or_default()
205 ),
206 kind: Some(SubjectKind::Run),
207 number: None,
208 run_id: text("runId"),
209 link: text("link"),
210 },
Inbox: threads, reasons, subscriptions and watching211 // A workflow on a branch: its next failure bumps the same thread.
212 "workflow.completed" if subject.is_none() => {
213 let branch = text("ref").unwrap_or_default();
214 let branch = branch.strip_prefix("refs/heads/").unwrap_or(&branch).to_owned();
215 let workflow = text("path").or_else(|| text("workflow")).unwrap_or_default();
216 Thread {
217 key: format!("{}/run/{workflow}@{branch}", wanted.repo_id),
218 kind: Some(SubjectKind::Run),
219 number: None,
220 run_id: text("runId"),
221 link: None,
222 }
223 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar224 // One alert, or one bypass request: its page is the link.
225 kind if SECURITY_EVENTS.contains(&kind) => {
226 let which = text("requestId").or_else(|| text("alertId")).unwrap_or_else(|| event.id.clone());
227 Thread {
228 key: format!("{}/security/{which}", wanted.repo_id),
229 kind: None,
230 number: None,
231 run_id: None,
232 link: text("link"),
233 }
234 }
Inbox: threads, reasons, subscriptions and watching235 _ => match wanted.number {
236 Some(number) => Thread {
237 key: numbered_thread(&wanted.repo_id, number),
238 kind: subject.and_then(|subject| subject.kind),
239 number: Some(number),
240 run_id: None,
241 link: None,
242 },
243 None => Thread {
244 key: format!("event/{}", event.id),
245 kind: None,
246 number: None,
247 run_id: None,
248 link: None,
249 },
250 },
251 }
252}
253
254/// The issue or pull request whose agent threads an event closes: what an
255/// agent was waiting on a person for is over.
256pub fn resolves(event: &Event) -> Option<String> {
257 if !RESUMES.contains(&event.kind.as_str()) {
258 return None;
259 }
260 let repo_id = event.data["repoId"].as_str().map(str::to_owned).or_else(|| event.repo_id.clone())?;
261 let number = event.data["number"].as_u64().and_then(|n| u32::try_from(n).ok())?;
262 Some(numbered_thread(&repo_id, number))
263}
264
265/// Collects who is told, each once with the most specific reason, never
266/// the actor, never g1t, and never anyone ignoring the thread.
Inbox: the events service tells people what needs them as events arrive267struct Told<'a> {
268 actor: &'a Actor,
Inbox: threads, reasons, subscriptions and watching269 audience: &'a Audience,
Inbox: the events service tells people what needs them as events arrive270 notices: Vec<Notice>,
271}
272
273impl Told<'_> {
Inbox: threads, reasons, subscriptions and watching274 /// `asked`: something asked of the person directly, told even when
275 /// they unsubscribed from the thread.
276 fn tell(&mut self, username: &str, reason: Reason, severity: Severity, title: &str, body: &str, asked: bool) {
Inbox: the events service tells people what needs them as events arrive277 let username = username.trim().trim_start_matches('@').to_lowercase();
278 if username.is_empty()
279 || is_g1t(&username)
280 || self.actor.is_named(&username)
Inbox: threads, reasons, subscriptions and watching281 || self.audience.ignores(&username)
282 || (!asked && self.audience.unsubscribed(&username))
Inbox: the events service tells people what needs them as events arrive283 {
284 return;
285 }
Inbox: threads, reasons, subscriptions and watching286 let notice = Notice {
Inbox: the events service tells people what needs them as events arrive287 username,
288 reason,
289 severity,
290 title: clip(title, MAX_TITLE),
291 body: clip(body, MAX_BODY),
Inbox: threads, reasons, subscriptions and watching292 };
293 match self.notices.iter_mut().find(|told| told.username == notice.username) {
294 Some(told) if notice.reason.rank() < told.reason.rank() => *told = notice,
295 Some(_) => {}
296 None => self.notices.push(notice),
297 }
Inbox: the events service tells people what needs them as events arrive298 }
299
300 /// A person known by id as well as name, such as an author.
Inbox: threads, reasons, subscriptions and watching301 fn tell_person(&mut self, person: &Principal, reason: Reason, severity: Severity, title: &str, body: &str, asked: bool) {
Inbox: the events service tells people what needs them as events arrive302 if is_g1t_id(&person.id) || self.actor.is(person) {
303 return;
304 }
Inbox: threads, reasons, subscriptions and watching305 self.tell(&person.username, reason, severity, title, body, asked);
306 }
307
308 /// Everyone subscribed to an issue or pull request: its owner and
309 /// author, its assignees and reviewers, and whoever subscribed by
310 /// commenting, being mentioned or by hand. `reason` overrides why
311 /// each is told, as a state change does.
312 fn tell_subscribed(&mut self, subject: &InboxSubject, reason: Option<Reason>, severity: Severity, title: &str, body: &str) {
313 let why = |own: Reason| reason.unwrap_or(own);
314 self.tell_person(subject.owner(), why(Reason::Author), severity, title, body, false);
315 self.tell_person(&subject.author, why(Reason::Author), severity, title, body, false);
316 for name in &subject.assignees {
317 self.tell(name, why(Reason::Assign), severity, title, body, false);
318 }
319 for name in &subject.reviewers {
320 self.tell(name, why(Reason::ReviewRequested), severity, title, body, false);
321 }
322 let subscribed: Vec<(String, Reason)> = self.audience.subscribed().collect();
323 for (name, own) in subscribed {
324 self.tell(&name, why(own), severity, title, body, false);
325 }
Inbox: the events service tells people what needs them as events arrive326 }
Inbox: threads, reasons, subscriptions and watching327
328 /// Everyone watching the repository for this kind of activity.
329 fn tell_watchers(&mut self, kind: &str, severity: Severity, title: &str, body: &str) {
330 let watching: Vec<String> = self.audience.watching(kind).collect();
331 for name in watching {
332 self.tell(&name, Reason::Subscribed, severity, title, body, false);
333 }
334 }
Inbox: the events service tells people what needs them as events arrive335}
336
337fn clip(text: &str, max: usize) -> String {
338 let text = text.trim();
339 if text.chars().count() <= max {
340 return text.to_owned();
341 }
342 let cut: String = text.chars().take(max - 1).collect();
343 format!("{}…", cut.trim_end())
344}
345
Inbox: threads, reasons, subscriptions and watching346/// The kind of activity a watcher chooses: `issues` or `pulls`.
347fn activity_of(subject: &InboxSubject) -> &'static str {
348 match subject.kind {
349 Some(SubjectKind::Pull) => "pulls",
350 _ => "issues",
351 }
352}
353
354/// The names in an event's list field, such as the reviewers just asked.
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar355/// The teams an event asked to review: each `workspace/slug` with the
356/// people it tells.
357fn teams_asked(data: &serde_json::Value) -> Vec<(String, Vec<String>)> {
358 data["teams"]
359 .as_array()
360 .map(|teams| {
361 teams
362 .iter()
363 .filter_map(|team| Some((team["team"].as_str()?.to_owned(), names(team, "notified"))))
364 .collect()
365 })
366 .unwrap_or_default()
367}
368
Inbox: threads, reasons, subscriptions and watching369fn names(data: &serde_json::Value, key: &str) -> Vec<String> {
370 data[key]
371 .as_array()
372 .map(|names| names.iter().filter_map(|name| name.as_str().map(str::to_owned)).collect())
373 .unwrap_or_default()
374}
375
376/// Who is told of `event`, in `repo` (`owner/name`), given what it names
377/// and who follows it. `actor` is who caused it; outcomes nobody chose
378/// (checks, workflows, deployments, a review, an agent finishing or
379/// stopping) are told whoever caused them.
380pub fn notices(event: &Event, repo: &str, actor: &Actor, subject: Option<&InboxSubject>, audience: &Audience) -> Vec<Notice> {
Inbox: the events service tells people what needs them as events arrive381 let nobody = Actor::default();
382 let data = &event.data;
383 // The issue or pull request, as titles name it: `acme/rocket#12`.
384 let at = match data["number"].as_u64().or(data["pull"].as_u64()) {
385 Some(number) if subject.is_some() => format!("{repo}#{number}"),
386 _ => repo.to_owned(),
387 };
388 let outcome = matches!(
389 event.kind.as_str(),
Inbox: threads, reasons, subscriptions and watching390 "checks.completed"
391 | "workflow.completed"
392 | "review.completed"
393 | "pull.ready"
394 | "agent.asked"
395 | "pull.stalled"
396 | "deployment.failed"
397 | "deployment.succeeded"
Merge branch 'worktree-agent-a3abfcce648e87dca'398 | "deployment.review_requested"
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar399 ) || SECURITY_EVENTS.contains(&event.kind.as_str());
Inbox: the events service tells people what needs them as events arrive400 let mut told = Told {
401 actor: if outcome { &nobody } else { actor },
Inbox: threads, reasons, subscriptions and watching402 audience,
Inbox: the events service tells people what needs them as events arrive403 notices: Vec::new(),
404 };
Inbox: threads, reasons, subscriptions and watching405 // Who did it, as titles name them.
406 let who = actor.name().unwrap_or("g1t");
Inbox: the events service tells people what needs them as events arrive407
408 match (event.kind.as_str(), subject) {
Inbox: threads, reasons, subscriptions and watching409 ("agent.asked" | "pull.stalled", Some(pull)) => {
410 let (title, body) = if event.kind == "agent.asked" {
411 (format!("An agent is waiting on {at}"), pull.title.clone())
412 } else {
413 let detail = data["detail"].as_str().map(str::trim).filter(|detail| !detail.is_empty());
414 (format!("g1t stopped on {at} and needs you"), detail.unwrap_or(&pull.title).to_owned())
415 };
416 told.tell_person(pull.owner(), Reason::Agent, Severity::Warning, &title, &body, true);
Inbox: the events service tells people what needs them as events arrive417 if let Some(issue) = &pull.issue {
Inbox: threads, reasons, subscriptions and watching418 told.tell_person(issue.owner(), Reason::Agent, Severity::Warning, &title, &body, true);
Inbox: the events service tells people what needs them as events arrive419 for name in &issue.assignees {
Inbox: threads, reasons, subscriptions and watching420 told.tell(name, Reason::Agent, Severity::Warning, &title, &body, true);
Inbox: the events service tells people what needs them as events arrive421 }
422 }
423 }
Inbox: threads, reasons, subscriptions and watching424 ("pull.review_requested", Some(pull)) => {
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar425 // Asked because the CODEOWNERS file says they own what changed.
426 let owned = data["codeOwners"].as_bool() == Some(true);
427 let title = if owned {
428 format!("{at} changes files you own")
429 } else {
430 format!("{who} asked you to review {at}")
431 };
Inbox: threads, reasons, subscriptions and watching432 for name in names(data, "reviewers") {
433 told.tell(&name, Reason::ReviewRequested, Severity::Warning, &title, &pull.title, true);
434 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar435 for (team, people) in teams_asked(data) {
436 let title = if owned {
437 format!("{at} changes files @{team} owns")
438 } else {
439 format!("{who} asked @{team} to review {at}")
440 };
441 for name in people {
442 told.tell(&name, Reason::ReviewRequested, Severity::Warning, &title, &pull.title, true);
443 }
444 }
Inbox: threads, reasons, subscriptions and watching445 }
446 ("issue.assigned" | "pull.assigned", Some(on)) => {
447 let title = format!("{who} assigned you to {at}");
448 for name in names(data, "added") {
449 told.tell(&name, Reason::Assign, Severity::Info, &title, &on.title, true);
450 }
451 }
452 ("issue.opened" | "pull.opened", Some(on)) => {
453 let title = format!("{who} opened {at}");
454 for name in &on.assignees {
455 told.tell(name, Reason::Assign, Severity::Info, &format!("{who} assigned you to {at}"), &on.title, true);
456 }
457 for name in &on.reviewers {
458 told.tell(name, Reason::ReviewRequested, Severity::Warning, &format!("{who} asked you to review {at}"), &on.title, true);
459 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar460 for name in &on.mentions {
461 told.tell(name, Reason::Mention, Severity::Info, &format!("{who} mentioned you on {at}"), &on.title, true);
462 }
463 for team in &on.team_mentions {
464 let title = format!("{who} mentioned @{} on {at}", team.team);
465 for name in &team.members {
466 told.tell(name, Reason::TeamMention, Severity::Info, &title, &on.title, true);
467 }
468 }
Inbox: threads, reasons, subscriptions and watching469 told.tell_watchers(activity_of(on), Severity::Info, &title, &on.title);
470 }
Inbox: the events service tells people what needs them as events arrive471 ("checks.completed", Some(pull)) => {
472 let title = match data["status"].as_str() {
Inbox: threads, reasons, subscriptions and watching473 Some("errored") => format!("Checks could not run on {at}"),
474 _ => format!("Checks failed on {at}"),
Inbox: the events service tells people what needs them as events arrive475 };
Inbox: threads, reasons, subscriptions and watching476 told.tell_person(pull.owner(), Reason::CiActivity, Severity::Error, &title, &pull.title, true);
Inbox: the events service tells people what needs them as events arrive477 }
478 ("workflow.completed", subject) => {
479 let workflow = data["workflow"].as_str().filter(|name| !name.is_empty()).unwrap_or("A workflow");
480 match subject {
481 Some(pull) => {
Inbox: threads, reasons, subscriptions and watching482 let title = format!("{workflow} failed on {at}");
483 told.tell_person(pull.owner(), Reason::CiActivity, Severity::Error, &title, &pull.title, true);
Inbox: the events service tells people what needs them as events arrive484 }
485 None => {
486 // Not on a pull request: whoever pushed the commit it ran on.
487 let branch = data["ref"].as_str().unwrap_or_default();
488 let branch = branch.strip_prefix("refs/heads/").unwrap_or(branch);
489 let title = format!("{workflow} failed on {branch} in {repo}");
490 let body = format!("Run {} at {}", data["number"], short(data["sha"].as_str().unwrap_or_default()));
491 if let Some(id) = &actor.id
492 && !is_g1t_id(id)
493 && let Some(name) = &actor.username
494 {
Inbox: threads, reasons, subscriptions and watching495 told.tell(name, Reason::CiActivity, Severity::Error, &title, &body, true);
496 }
497 }
498 }
499 }
500 ("deployment.failed" | "deployment.succeeded", subject) => {
501 let failed = event.kind == "deployment.failed";
502 let project = data["project"].as_str().filter(|name| !name.is_empty()).unwrap_or(repo);
503 let what = match subject {
504 Some(_) => format!("The preview of {at}"),
505 None => format!("Production of {project}"),
506 };
507 let (title, severity) = if failed {
508 (format!("{what} failed to deploy"), Severity::Error)
509 } else {
510 (format!("{what} is live"), Severity::Success)
511 };
512 let body = match (failed, data["error"].as_str().map(str::trim).filter(|error| !error.is_empty())) {
513 (true, Some(error)) => error.to_owned(),
514 _ => match subject {
515 Some(pull) => pull.title.clone(),
516 None => format!("Commit {}", short(data["commit"].as_str().unwrap_or_default())),
517 },
518 };
519 // A success is news to whoever answers for it only after a failure.
520 if failed || data["recovered"].as_bool() == Some(true) {
521 let title = if failed { title.clone() } else { format!("{what} is live again") };
522 match subject {
523 Some(pull) => told.tell_person(pull.owner(), Reason::CiActivity, severity, &title, &body, true),
524 None => {
525 let pusher = actor
526 .id
527 .as_deref()
528 .filter(|id| !is_g1t_id(id))
529 .and(actor.username.as_deref())
530 .or_else(|| data["triggeredBy"].as_str());
531 if let Some(name) = pusher {
532 told.tell(name, Reason::CiActivity, severity, &title, &body, true);
533 }
Inbox: the events service tells people what needs them as events arrive534 }
535 }
536 }
Inbox: threads, reasons, subscriptions and watching537 told.tell_watchers("deployments", severity, &title, &body);
Inbox: the events service tells people what needs them as events arrive538 }
539 ("review.completed", Some(pull)) => {
Inbox: threads, reasons, subscriptions and watching540 let (title, severity) = match data["verdict"].as_str() {
541 Some("approve") => (format!("g1t approved {at}"), Severity::Success),
542 _ => (format!("g1t asked for changes on {at}"), Severity::Info),
Inbox: the events service tells people what needs them as events arrive543 };
Inbox: threads, reasons, subscriptions and watching544 told.tell_person(pull.owner(), Reason::Author, severity, &title, &pull.title, true);
Inbox: the events service tells people what needs them as events arrive545 }
546 ("pull.ready", Some(pull)) => {
547 // A change g1t made is ready: the agent's run is over.
548 if let Some(owner) = pull.requested_by.as_ref().filter(|_| is_g1t_id(&pull.author.id) || is_g1t(&pull.author.username)) {
Inbox: threads, reasons, subscriptions and watching549 let title = format!("g1t finished {at}");
550 told.tell_person(owner, Reason::Author, Severity::Success, &title, &pull.title, true);
Inbox: the events service tells people what needs them as events arrive551 }
552 }
Inbox: threads, reasons, subscriptions and watching553 ("pull.merged" | "pull.closed" | "issue.closed" | "issue.reopened", Some(on)) => {
554 let (verb, severity) = match event.kind.as_str() {
555 "pull.merged" => ("merged", Severity::Success),
556 "issue.reopened" => ("reopened", Severity::Info),
557 _ => ("closed", Severity::Info),
558 };
559 let title = match (actor.name(), data["resolvedBy"].as_u64()) {
560 (_, Some(pull)) if event.kind == "issue.closed" => format!("{at} was closed by #{pull}"),
561 (Some(name), _) => format!("{name} {verb} {at}"),
562 (None, _) => format!("{at} was {verb}"),
Inbox: the events service tells people what needs them as events arrive563 };
Inbox: threads, reasons, subscriptions and watching564 told.tell_subscribed(on, Some(Reason::StateChange), severity, &title, &on.title);
565 told.tell_watchers(activity_of(on), severity, &title, &on.title);
Inbox: the events service tells people what needs them as events arrive566 }
Merge branch 'worktree-agent-a3abfcce648e87dca'567 ("deployment.review_requested", _) => {
568 let environment = data["environment"].as_str().unwrap_or("an environment");
569 let workflow = data["workflow"].as_str().unwrap_or("A workflow");
570 let title = format!("{workflow} is waiting for your review to deploy to {environment} in {repo}");
571 let body = data["title"].as_str().unwrap_or_default();
572 // Those the actions service names: the environment's reviewers.
573 for name in names(data, "notify") {
574 told.tell(&name, Reason::ReviewRequested, Severity::Warning, &title, body, true);
575 }
576 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar577 (kind, None) if SECURITY_EVENTS.contains(&kind) => {
578 let severity = match data["severity"].as_str() {
579 _ if kind == "secret_scanning.bypass_requested" => Severity::Warning,
580 _ if kind == "secret_scanning.bypass_reviewed" => Severity::Info,
581 Some("critical" | "high") => Severity::Error,
582 _ => Severity::Warning,
583 };
584 let what = data["title"].as_str().unwrap_or("A security alert");
585 let title = match kind {
586 "secret_scanning.bypass_requested" => format!("{who} asked to bypass push protection in {repo}"),
587 "secret_scanning.bypass_reviewed" => {
588 let verdict = data["state"].as_str().unwrap_or("reviewed");
589 format!("Your request to bypass push protection in {repo} was {verdict}")
590 }
591 _ if data["pusher"].is_string() => format!("A push to {repo} was blocked: it adds a secret"),
592 _ => format!("New security alert in {repo}"),
593 };
594 // Whoever pushed a blocked secret, and those the service names
595 // (the workspace's owners, or whoever asked for a bypass).
596 if let Some(pusher) = data["pusher"].as_str() {
597 told.tell(pusher, Reason::SecurityAlert, Severity::Error, &title, what, true);
598 }
599 for name in names(data, "notify") {
600 told.tell(&name, Reason::SecurityAlert, severity, &title, what, true);
601 }
602 // Watchers who chose security alerts, if they may see findings:
603 // the service lists who may (`members`), as findings are never
604 // shown to someone who cannot change the code.
605 if !kind.starts_with("secret_scanning.bypass") {
606 let members = names(data, "members");
607 let watching: Vec<String> = told.audience.watching("security").collect();
608 for name in watching.iter().filter(|name| members.iter().any(|member| member.eq_ignore_ascii_case(name))) {
609 told.tell(name, Reason::SecurityAlert, severity, &title, what, false);
610 }
611 }
612 }
Inbox: the events service tells people what needs them as events arrive613 ("comment.created", Some(on)) => {
614 let Some(comment) = on.comment.as_ref().filter(|comment| !comment.event) else {
615 return Vec::new();
616 };
617 // Whoever wrote it is the actor, whatever the event says.
618 let writer = Actor {
619 id: Some(comment.author.id.clone()),
620 username: Some(comment.author.username.clone()),
621 };
622 told.actor = &writer;
623 let who = &comment.author.username;
624 let body = if comment.excerpt.is_empty() { &on.title } else { &comment.excerpt };
625 for name in &comment.mentions {
Inbox: threads, reasons, subscriptions and watching626 let title = format!("{who} mentioned you on {at}");
627 told.tell(name, Reason::Mention, Severity::Info, &title, body, true);
Inbox: the events service tells people what needs them as events arrive628 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar629 for team in &comment.team_mentions {
630 let title = format!("{who} mentioned @{} on {at}", team.team);
631 for name in &team.members {
632 told.tell(name, Reason::TeamMention, Severity::Info, &title, body, true);
633 }
634 }
Inbox: threads, reasons, subscriptions and watching635 let (severity, title) = match comment.verdict.as_deref() {
636 Some("approve") => (Severity::Success, format!("{who} approved {at}")),
637 Some("request_changes") => (Severity::Info, format!("{who} asked for changes on {at}")),
638 _ => (Severity::Info, format!("{who} commented on {at}")),
Inbox: the events service tells people what needs them as events arrive639 };
Inbox: threads, reasons, subscriptions and watching640 // An approval or a request for changes is the owner's to hear
641 // of whatever they chose; the rest, as they subscribed.
642 if comment.verdict.is_some() {
643 told.tell_person(on.owner(), Reason::Author, severity, &title, body, true);
644 }
645 told.tell_subscribed(on, None, severity, &title, body);
646 told.tell_watchers(activity_of(on), severity, &title, body);
Inbox: the events service tells people what needs them as events arrive647 return told.notices;
648 }
649 _ => {}
650 }
651 told.notices
652}
653
Inbox: threads, reasons, subscriptions and watching654/// Who an event subscribes to its issue or pull request without asking,
655/// and why: whoever commented, whoever they mentioned, the people assigned
656/// and the reviewers asked. Never g1t.
657pub fn subscribes(event: &Event, subject: Option<&InboxSubject>) -> Vec<(String, Reason)> {
658 let mut people: Vec<(String, Reason)> = Vec::new();
659 let mut add = |name: &str, reason: Reason| {
660 let name = name.trim().trim_start_matches('@').to_lowercase();
661 if !name.is_empty() && !is_g1t(&name) && !people.iter().any(|(had, _)| *had == name) {
662 people.push((name, reason));
663 }
664 };
665 match event.kind.as_str() {
666 "comment.created" => {
667 if let Some(comment) = subject.and_then(|subject| subject.comment.as_ref()).filter(|comment| !comment.event) {
668 if !is_g1t_id(&comment.author.id) {
669 add(&comment.author.username, Reason::Comment);
670 }
671 for name in &comment.mentions {
672 add(name, Reason::Mention);
673 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar674 for team in &comment.team_mentions {
675 for name in &team.members {
676 add(name, Reason::TeamMention);
677 }
678 }
679 }
680 }
681 "issue.opened" | "pull.opened" => {
682 if let Some(on) = subject {
683 for name in &on.mentions {
684 add(name, Reason::Mention);
685 }
686 for team in &on.team_mentions {
687 for name in &team.members {
688 add(name, Reason::TeamMention);
689 }
690 }
Inbox: threads, reasons, subscriptions and watching691 }
692 }
693 "issue.assigned" | "pull.assigned" => {
694 for name in names(&event.data, "added") {
695 add(&name, Reason::Assign);
696 }
697 }
698 "pull.review_requested" => {
699 for name in names(&event.data, "reviewers") {
700 add(&name, Reason::ReviewRequested);
701 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar702 for (_, people) in teams_asked(&event.data) {
703 for name in people {
704 add(&name, Reason::ReviewRequested);
705 }
706 }
Inbox: threads, reasons, subscriptions and watching707 }
708 _ => {}
709 }
710 people
711}
712
Inbox: the events service tells people what needs them as events arrive713fn short(sha: &str) -> &str {
714 sha.get(..7).unwrap_or(sha)
715}
716
717/// Where an item is on g1t.sh.
Inbox: threads, reasons, subscriptions and watching718pub fn url(repo: Option<&str>, subject: Option<SubjectKind>, number: Option<u32>, run_id: Option<&str>, link: Option<&str>) -> String {
719 if let Some(link) = link.filter(|link| link.starts_with('/')) {
720 return link.to_owned();
721 }
Inbox: the events service tells people what needs them as events arrive722 let Some(repo) = repo else {
723 return "/inbox".to_owned();
724 };
725 match (subject, number, run_id) {
726 (Some(SubjectKind::Pull), Some(number), _) => format!("/{repo}/pull/{number}"),
727 (Some(SubjectKind::Issue), Some(number), _) => format!("/{repo}/issues/{number}"),
728 (Some(SubjectKind::Run), _, Some(run)) => format!("/{repo}/actions/runs/{run}"),
Inbox: threads, reasons, subscriptions and watching729 (Some(SubjectKind::Deploy), Some(number), _) => format!("/{repo}/pull/{number}"),
Inbox: the events service tells people what needs them as events arrive730 _ => format!("/{repo}"),
731 }
732}
733
734// --- Writing ---------------------------------------------------------------
735
736/// The services the inbox reads from as events arrive.
737pub struct Sources<'a> {
738 pub work: &'a Fetcher,
739 pub repos: &'a Fetcher,
740 pub identity: &'a Fetcher,
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)741 /// The notify service, told of each new item for live toasts and pushes; None, nobody is.
742 pub notify: Option<&'a Fetcher>,
Inbox: the events service tells people what needs them as events arrive743}
744
Inbox: threads, reasons, subscriptions and watching745/// One notice written, and whether it was news (not a redelivery): what
746/// may be emailed.
747struct Written {
748 notice: Notice,
749 repo_id: String,
750 url: String,
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)751 event_id: String,
Inbox: threads, reasons, subscriptions and watching752}
753
Inbox: the events service tells people what needs them as events arrive754/// Writes the items a batch from the bus calls for. Never fails the batch:
755/// what cannot be worked out is logged and left.
756pub async fn deliver(db: &D1Database, sources: &Sources<'_>, events: &[Event]) {
Inbox: threads, reasons, subscriptions and watching757 if let Err(error) = resolve(db, events).await {
758 worker::console_error!("inbox: agent threads not closed: {error}");
759 }
Fine-grained personal tokens, workspace token rules and approvals in identity760 // A workspace's own notices, about no repository (workspace_notices).
761 for event in events.iter().filter(|event| WORKSPACE_EVENTS.contains(&event.kind.as_str())) {
762 if let Err(error) = deliver_workspace(db, event).await {
763 worker::console_error!("inbox: {} {} not delivered: {error}", event.kind, event.id);
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)764 } else if let Some(notify) = sources.notify {
765 // Notify: the feed tells each person once per event, redelivered or not.
766 let (_, told) = workspace_notices(event, None);
767 let link = event.data["link"].as_str().filter(|link| link.starts_with('/')).unwrap_or("/inbox");
768 let live = told.iter().map(|notice| (notice.username.clone(), live_notification(&event.id, notice, link, &event.time))).collect();
769 tell_live(notify, live).await;
770 for notice in &told {
771 tell_inbox_count(db, notify, &notice.username).await;
772 }
Fine-grained personal tokens, workspace token rules and approvals in identity773 }
774 }
Inbox: the events service tells people what needs them as events arrive775 let wanted: Vec<(&Event, Wanted)> = events
776 .iter()
777 .filter_map(|event| wants(event).map(|wanted| (event, wanted)))
778 .collect();
Inbox: threads, reasons, subscriptions and watching779 let created: Vec<&Event> = events.iter().filter(|event| event.kind == "repo.created").collect();
780 if wanted.is_empty() && created.is_empty() {
Inbox: the events service tells people what needs them as events arrive781 return;
782 }
783 // Everyone who caused one, named in one call.
784 let ids: Vec<String> = wanted
785 .iter()
Inbox: threads, reasons, subscriptions and watching786 .map(|(event, _)| *event)
787 .chain(created.iter().copied())
788 .filter_map(|event| event.actor.clone())
Inbox: the events service tells people what needs them as events arrive789 .filter(|id| !is_g1t_id(id))
790 .collect::<HashSet<_>>()
791 .into_iter()
792 .collect();
793 let names: HashMap<String, String> = if ids.is_empty() {
794 HashMap::new()
795 } else {
796 g1t_kit::call(sources.identity, "usernames", &UsernamesArgs { ids })
797 .await
798 .unwrap_or_else(|error| {
799 worker::console_error!("inbox: could not name who acted: {error}");
800 HashMap::new()
801 })
802 };
Inbox: threads, reasons, subscriptions and watching803 for event in created {
804 if let Err(error) = subscriptions::watch_created(db, event, &names).await {
805 worker::console_error!("inbox: {} {} not watched: {error}", event.kind, event.id);
806 }
807 }
Inbox: the events service tells people what needs them as events arrive808 let mut paths: HashMap<String, Option<RepoPath>> = HashMap::new();
Inbox: threads, reasons, subscriptions and watching809 let mut watchers: HashMap<String, Vec<subscriptions::Watcher>> = HashMap::new();
810 let mut written: Vec<Written> = Vec::new();
Inbox: the events service tells people what needs them as events arrive811 for (event, wanted) in wanted {
Inbox: threads, reasons, subscriptions and watching812 match deliver_one(db, sources, &names, &mut paths, &mut watchers, event, wanted).await {
813 Ok(mut news) => written.append(&mut news),
814 Err(error) => worker::console_error!("inbox: {} {} not delivered: {error}", event.kind, event.id),
815 }
816 }
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)817 // Notify: every item that was news, live in the person's tabs (services/notify).
818 if let Some(notify) = sources.notify {
819 let live = written
820 .iter()
821 .map(|item| (item.notice.username.clone(), live_notification(&item.event_id, &item.notice, &item.url, &rfc3339(now_ms()))))
822 .collect::<Vec<_>>();
823 tell_live(notify, live).await;
824 let people: HashSet<&str> = written.iter().map(|item| item.notice.username.as_str()).collect();
825 for username in people {
826 tell_inbox_count(db, notify, username).await;
827 }
828 }
Inbox: threads, reasons, subscriptions and watching829 if let Err(error) = email(db, sources.identity, written).await {
830 worker::console_error!("inbox: emails not sent: {error}");
831 }
832}
833
Fine-grained personal tokens, workspace token rules and approvals in identity834/// Events about a workspace rather than a repository, each naming who to
835/// tell (`notify`), what to say (`title`, `body`) and where it is (`link`):
836/// identity's personal access token approvals.
837///
838/// | Event | Who | Reason | Severity |
839/// | --- | --- | --- | --- |
840/// | `token.approval_requested` | the workspace's owners | review_requested | warning |
841/// | `token.approval_reviewed` | the token's owner | author | info |
Merge workspace invitations: nobody joins a workspace without saying yes, people are found by username, your own invites can bring someone in, and nobody is left without a workspace (identity 0040)842/// | `workspace_invitation.created` | the person invited | review_requested | warning |
843/// | `workspace_invitation.accepted` | who invited them | author | success |
844/// | `workspace_invitation.declined` | who invited them | author | info |
845///
846/// An invitation's item for the person invited is done once it is
847/// accepted, declined or revoked (`workspace_invitation.revoked`, which
848/// tells nobody).
849pub const WORKSPACE_EVENTS: [&str; 5] = [
850 "token.approval_requested",
851 "token.approval_reviewed",
852 "workspace_invitation.created",
853 "workspace_invitation.accepted",
854 "workspace_invitation.declined",
855];
Fine-grained personal tokens, workspace token rules and approvals in identity856
Merge workspace invitations: nobody joins a workspace without saying yes, people are found by username, your own invites can bring someone in, and nobody is left without a workspace (identity 0040)857/// Events that close what a workspace invitation asked of its person.
858pub const INVITATION_ANSWERED: [&str; 3] =
859 ["workspace_invitation.accepted", "workspace_invitation.declined", "workspace_invitation.revoked"];
860
Fine-grained personal tokens, workspace token rules and approvals in identity861/// The notices a workspace event calls for, and the thread they go to.
862pub fn workspace_notices(event: &Event, actor: Option<&str>) -> (String, Vec<Notice>) {
863 let data = &event.data;
864 let text = |key: &str| data[key].as_str().unwrap_or_default().to_owned();
865 let (reason, severity) = match event.kind.as_str() {
Merge workspace invitations: nobody joins a workspace without saying yes, people are found by username, your own invites can bring someone in, and nobody is left without a workspace (identity 0040)866 "token.approval_requested" | "workspace_invitation.created" => (Reason::ReviewRequested, Severity::Warning),
867 "workspace_invitation.accepted" => (Reason::Author, Severity::Success),
Fine-grained personal tokens, workspace token rules and approvals in identity868 _ => (Reason::Author, Severity::Info),
869 };
Merge workspace invitations: nobody joins a workspace without saying yes, people are found by username, your own invites can bring someone in, and nobody is left without a workspace (identity 0040)870 // An invitation names its own thread; a token approval's is the token's.
871 let thread = match data["thread"].as_str().filter(|thread| !thread.is_empty()) {
872 Some(thread) => thread.to_owned(),
873 None => format!("workspace:{}/token/{}", text("workspace"), text("tokenId")),
874 };
Fine-grained personal tokens, workspace token rules and approvals in identity875 let notices = names(data, "notify")
876 .into_iter()
877 .map(|name| name.to_lowercase())
878 .filter(|name| actor.is_none_or(|actor| !actor.eq_ignore_ascii_case(name)) && !is_g1t(name))
879 .collect::<HashSet<_>>()
880 .into_iter()
881 .map(|username| Notice { username, reason, severity, title: text("title"), body: text("body") })
882 .collect();
883 (thread, notices)
884}
885
886/// Writes a workspace event's notices: a thread each, with no repository.
887async fn deliver_workspace(db: &D1Database, event: &Event) -> Result<()> {
888 let (thread, told) = workspace_notices(event, None);
889 if told.is_empty() {
890 return Ok(());
891 }
892 let now = now_ms();
893 let workspace = event.data["workspace"].as_str().unwrap_or_default().to_lowercase();
894 let link = event.data["link"].as_str().filter(|link| link.starts_with('/'));
895 let mut statements = Vec::new();
896 for notice in &told {
897 let username = notice.username.as_str();
898 statements.push(db.prepare(BUMP).bind(&[
899 new_id("ntf", now).into(),
900 username.into(),
901 thread.as_str().into(),
902 event.id.as_str().into(),
903 event.kind.as_str().into(),
904 notice.reason.as_str().into(),
905 notice.severity.as_str().into(),
906 notice.title.as_str().into(),
907 notice.body.as_str().into(),
908 workspace.as_str().into(),
909 JsValue::NULL,
910 JsValue::NULL,
911 JsValue::NULL,
912 JsValue::NULL,
913 JsValue::NULL,
914 link.map_or(JsValue::NULL, JsValue::from),
915 JsValue::NULL,
916 event.time.as_str().into(),
917 new_id("ntf", now).into(),
918 ])?);
919 statements.push(
920 db.prepare(
921 "INSERT OR IGNORE INTO inbox_activity (id, item_id, username, event_id, event, reason, severity, title, body, actor, created_at)
922 SELECT ?1, id, ?2, ?4, ?5, ?6, ?7, ?8, ?9, NULL, ?10 FROM inbox_items WHERE username = ?2 AND thread = ?3",
923 )
924 .bind(&[
925 new_id("ntf", now).into(),
926 username.into(),
927 thread.as_str().into(),
928 event.id.as_str().into(),
929 event.kind.as_str().into(),
930 notice.reason.as_str().into(),
931 notice.severity.as_str().into(),
932 notice.title.as_str().into(),
933 notice.body.as_str().into(),
934 event.time.as_str().into(),
935 ])?,
936 );
937 }
938 db.batch(statements).await?;
939 Ok(())
940}
941
Inbox: threads, reasons, subscriptions and watching942/// Closes what an agent was waiting on a person for once it is over.
943async fn resolve(db: &D1Database, events: &[Event]) -> Result<()> {
944 let now = rfc3339(now_ms());
945 let mut statements = Vec::new();
946 for thread in events.iter().filter_map(resolves) {
947 statements.push(
948 db.prepare(
949 "UPDATE inbox_items SET done_at = ?1, read_at = COALESCE(read_at, ?1)
950 WHERE thread = ?2 AND reason = 'agent' AND done_at IS NULL",
951 )
952 .bind(&[now.as_str().into(), thread.into()])?,
953 );
954 }
Merge workspace invitations: nobody joins a workspace without saying yes, people are found by username, your own invites can bring someone in, and nobody is left without a workspace (identity 0040)955 // A workspace invitation answered or revoked no longer waits on its person.
956 for event in events.iter().filter(|event| INVITATION_ANSWERED.contains(&event.kind.as_str())) {
957 let Some(thread) = event.data["thread"].as_str().filter(|thread| thread.starts_with("invitation:")) else {
958 continue;
959 };
960 statements.push(
961 db.prepare(
962 "UPDATE inbox_items SET done_at = ?1, read_at = COALESCE(read_at, ?1)
963 WHERE thread = ?2 AND reason = 'review_requested' AND done_at IS NULL",
964 )
965 .bind(&[now.as_str().into(), thread.into()])?,
966 );
967 }
Inbox: threads, reasons, subscriptions and watching968 // Reviews no longer asked for are no longer waiting.
969 for event in events.iter().filter(|event| event.kind == "pull.review_request_removed") {
970 let (Some(repo_id), Some(number)) = (
971 event.data["repoId"].as_str().map(str::to_owned).or_else(|| event.repo_id.clone()),
972 event.data["number"].as_u64(),
973 ) else {
974 continue;
975 };
976 for name in names(&event.data, "reviewers") {
977 statements.push(
978 db.prepare(
979 "UPDATE inbox_items SET done_at = ?1, read_at = COALESCE(read_at, ?1)
980 WHERE thread = ?2 AND username = ?3 AND reason = 'review_requested' AND done_at IS NULL",
981 )
982 .bind(&[
983 now.as_str().into(),
984 format!("{repo_id}#{number}").into(),
985 name.to_lowercase().into(),
986 ])?,
987 );
Inbox: the events service tells people what needs them as events arrive988 }
989 }
Inbox: threads, reasons, subscriptions and watching990 if !statements.is_empty() {
991 db.batch(statements).await?;
992 }
993 Ok(())
Inbox: the events service tells people what needs them as events arrive994}
995
996async fn deliver_one(
997 db: &D1Database,
998 sources: &Sources<'_>,
999 names: &HashMap<String, String>,
1000 paths: &mut HashMap<String, Option<RepoPath>>,
Inbox: threads, reasons, subscriptions and watching1001 watchers: &mut HashMap<String, Vec<subscriptions::Watcher>>,
Inbox: the events service tells people what needs them as events arrive1002 event: &Event,
1003 wanted: Wanted,
Inbox: threads, reasons, subscriptions and watching1004) -> Result<Vec<Written>> {
Inbox: the events service tells people what needs them as events arrive1005 if !paths.contains_key(&wanted.repo_id) {
1006 let path: Option<RepoPath> = g1t_kit::call(
1007 sources.repos,
1008 "path_by_id",
1009 &PathByIdArgs {
1010 id: wanted.repo_id.clone(),
1011 },
1012 )
1013 .await?;
1014 paths.insert(wanted.repo_id.clone(), path);
1015 }
1016 // A repository that is gone tells nobody.
1017 let Some(path) = paths.get(&wanted.repo_id).cloned().flatten() else {
Inbox: threads, reasons, subscriptions and watching1018 return Ok(Vec::new());
Inbox: the events service tells people what needs them as events arrive1019 };
1020 let subject: Option<InboxSubject> = match wanted.number {
1021 Some(number) => {
1022 let found: Option<InboxSubject> = g1t_kit::call(
1023 sources.work,
1024 "inbox_subject",
1025 &InboxSubjectArgs {
1026 repo_id: wanted.repo_id.clone(),
1027 number,
1028 comment_id: wanted.comment_id.clone(),
1029 },
1030 )
1031 .await?;
1032 if found.is_none() {
Inbox: threads, reasons, subscriptions and watching1033 return Ok(Vec::new());
Inbox: the events service tells people what needs them as events arrive1034 }
1035 found
1036 }
1037 None => None,
1038 };
1039 let actor = Actor {
1040 id: event.actor.clone(),
1041 username: match event.actor.as_deref() {
1042 Some(id) if is_g1t_id(id) => Some(system::USERNAME.to_owned()),
1043 Some(id) => names.get(id).map(|name| name.to_lowercase()),
1044 None => None,
1045 },
1046 };
Inbox: threads, reasons, subscriptions and watching1047 let thread = thread_of(event, &wanted, subject.as_ref());
1048 if !watchers.contains_key(&wanted.repo_id) {
1049 watchers.insert(wanted.repo_id.clone(), subscriptions::watchers(db, &wanted.repo_id).await?);
1050 }
1051 let audience = Audience {
1052 subscriptions: match thread.kind {
1053 Some(SubjectKind::Issue | SubjectKind::Pull) => subscriptions::of_thread(db, &thread.key).await?,
1054 _ => Vec::new(),
1055 },
1056 watchers: watchers.get(&wanted.repo_id).cloned().unwrap_or_default(),
1057 };
Inbox: the events service tells people what needs them as events arrive1058 let repo = format!("{}/{}", path.namespace, path.name).to_lowercase();
Inbox: threads, reasons, subscriptions and watching1059 let told = notices(event, &repo, &actor, subject.as_ref(), &audience);
1060 let mut statements = subscriptions::auto_subscribe(db, &thread.key, &wanted.repo_id, &subscribes(event, subject.as_ref()), &event.time)?;
Inbox: the events service tells people what needs them as events arrive1061 if told.is_empty() {
Inbox: threads, reasons, subscriptions and watching1062 if !statements.is_empty() {
1063 db.batch(statements).await?;
1064 }
1065 return Ok(Vec::new());
Inbox: the events service tells people what needs them as events arrive1066 }
1067 let now = now_ms();
Inbox: threads, reasons, subscriptions and watching1068 let link = thread.link.as_deref();
1069 let item_url = url(Some(&repo), thread.kind, thread.number, thread.run_id.as_deref(), link);
1070 // Three statements a notice: its thread brought up (unless this event
1071 // was told before), a line of history, and the history kept short.
1072 let first = statements.len();
1073 for notice in &told {
1074 let username = notice.username.as_str();
1075 let values: Vec<JsValue> = vec![
1076 new_id("ntf", now).into(),
1077 username.into(),
1078 thread.key.as_str().into(),
1079 event.id.as_str().into(),
1080 event.kind.as_str().into(),
1081 notice.reason.as_str().into(),
1082 notice.severity.as_str().into(),
1083 notice.title.as_str().into(),
1084 notice.body.as_str().into(),
1085 path.namespace.to_lowercase().into(),
1086 wanted.repo_id.as_str().into(),
1087 repo.as_str().into(),
1088 thread.kind.map_or(JsValue::NULL, |kind| kind.as_str().into()),
1089 thread.number.map_or(JsValue::NULL, JsValue::from),
1090 thread.run_id.as_deref().map_or(JsValue::NULL, JsValue::from),
1091 link.map_or(JsValue::NULL, JsValue::from),
1092 actor.username.as_deref().map_or(JsValue::NULL, JsValue::from),
1093 event.time.as_str().into(),
1094 // Same prefix as item ids, so old and new sort by time together.
1095 new_id("ntf", now).into(),
1096 ];
1097 statements.push(db.prepare(BUMP).bind(&values)?);
Inbox: the events service tells people what needs them as events arrive1098 statements.push(
1099 db.prepare(
Inbox: threads, reasons, subscriptions and watching1100 "INSERT OR IGNORE INTO inbox_activity (id, item_id, username, event_id, event, reason, severity, title, body, actor, created_at)
1101 SELECT ?1, id, ?2, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11 FROM inbox_items WHERE username = ?2 AND thread = ?3
1102 RETURNING username",
Inbox: the events service tells people what needs them as events arrive1103 )
1104 .bind(&[
1105 new_id("ntf", now).into(),
Inbox: threads, reasons, subscriptions and watching1106 username.into(),
1107 thread.key.as_str().into(),
Inbox: the events service tells people what needs them as events arrive1108 event.id.as_str().into(),
Inbox: threads, reasons, subscriptions and watching1109 event.kind.as_str().into(),
1110 notice.reason.as_str().into(),
Inbox: the events service tells people what needs them as events arrive1111 notice.severity.as_str().into(),
1112 notice.title.as_str().into(),
1113 notice.body.as_str().into(),
1114 actor.username.as_deref().map_or(JsValue::NULL, JsValue::from),
1115 event.time.as_str().into(),
1116 ])?,
1117 );
Inbox: threads, reasons, subscriptions and watching1118 statements.push(
1119 db.prepare(
1120 "DELETE FROM inbox_activity
1121 WHERE item_id = (SELECT id FROM inbox_items WHERE username = ?1 AND thread = ?2)
1122 AND id NOT IN (
1123 SELECT a.id FROM inbox_activity a
1124 WHERE a.item_id = (SELECT id FROM inbox_items WHERE username = ?1 AND thread = ?2)
1125 ORDER BY a.id DESC LIMIT ?3)",
1126 )
1127 .bind(&[username.into(), thread.key.as_str().into(), MAX_ACTIVITY.into()])?,
1128 );
1129 }
1130 let results = db.batch(statements).await?;
1131 #[derive(Deserialize)]
1132 struct Told {
1133 username: String,
Inbox: the events service tells people what needs them as events arrive1134 }
Inbox: threads, reasons, subscriptions and watching1135 // A line of history written means the event is news to that person.
1136 let mut news = HashSet::new();
1137 for (at, _) in told.iter().enumerate() {
1138 if let Some(result) = results.get(first + at * 3 + 1)
1139 && let Ok(rows) = result.results::<Told>()
1140 {
1141 news.extend(rows.into_iter().map(|row| row.username));
1142 }
1143 }
1144 Ok(told
1145 .into_iter()
1146 .filter(|notice| news.contains(&notice.username))
1147 .map(|notice| Written {
1148 notice,
1149 repo_id: wanted.repo_id.clone(),
1150 url: item_url.clone(),
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)1151 event_id: event.id.clone(),
Inbox: threads, reasons, subscriptions and watching1152 })
1153 .collect())
1154}
1155
1156/// Brings a person's thread up with new activity, or starts it. `?1` id,
1157/// `?2` username, `?3` thread, `?4` event id, `?5` event type, `?6` reason,
1158/// `?7` severity, `?8` title, `?9` body, `?10` workspace, `?11` repo id,
1159/// `?12` repo, `?13` subject, `?14` number, `?15` run id, `?16` link, `?17`
1160/// actor, `?18` time, `?19` the time-sortable id lists order by. Nothing
1161/// happens when the person was already told of this event. While the
1162/// thread is unread its most urgent severity is kept; done or not, it
1163/// comes back to the inbox. A snooze stands.
1164const BUMP: &str = "INSERT INTO inbox_items (id, username, thread, event_id, event, reason, severity, title, body,
1165 workspace, repo_id, repo, subject, number, run_id, link, actor, created_at, updated_at, bumped, activity)
1166 SELECT ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?18, ?19, 1
1167 WHERE NOT EXISTS (SELECT 1 FROM inbox_activity WHERE event_id = ?4 AND username = ?2)
1168 ON CONFLICT (username, thread) DO UPDATE SET
1169 event_id = excluded.event_id,
1170 event = excluded.event,
1171 reason = excluded.reason,
1172 severity = CASE
1173 WHEN inbox_items.read_at IS NULL AND inbox_items.done_at IS NULL
1174 AND (CASE inbox_items.severity WHEN 'warning' THEN 0 WHEN 'error' THEN 1 WHEN 'success' THEN 2 ELSE 3 END)
1175 < (CASE excluded.severity WHEN 'warning' THEN 0 WHEN 'error' THEN 1 WHEN 'success' THEN 2 ELSE 3 END)
1176 THEN inbox_items.severity ELSE excluded.severity END,
1177 title = excluded.title,
1178 body = excluded.body,
1179 workspace = excluded.workspace,
1180 repo = excluded.repo,
1181 subject = COALESCE(excluded.subject, inbox_items.subject),
1182 run_id = COALESCE(excluded.run_id, inbox_items.run_id),
1183 link = COALESCE(excluded.link, inbox_items.link),
1184 actor = excluded.actor,
1185 updated_at = excluded.updated_at,
1186 bumped = excluded.bumped,
1187 activity = inbox_items.activity + 1,
1188 read_at = NULL,
1189 done_at = NULL";
1190
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)1191/// What the notify service is told of an inbox item (`FeedNotification` in
1192/// @g1t/contracts): its kind from why the person was told, the workspace
1193/// from where it links, and the item's title and first line.
1194pub fn live_notification(event_id: &str, notice: &Notice, url: &str, at: &str) -> serde_json::Value {
1195 let kind = match notice.reason {
1196 Reason::Agent => "agent_waiting",
1197 Reason::ReviewRequested => "approval",
1198 Reason::Mention | Reason::TeamMention => "mention",
1199 _ => "inbox",
1200 };
1201 let path = url.split(['?', '#']).next().unwrap_or_default();
1202 let mut segments = path.trim_start_matches('/').split('/');
1203 let workspace = segments.next().unwrap_or_default().to_lowercase();
1204 let repo = segments.next().map(|name| format!("{workspace}/{name}"));
1205 let name = repo.unwrap_or_else(|| if workspace.is_empty() { "g1t".to_owned() } else { workspace.clone() });
1206 serde_json::json!({
1207 "id": format!("inbox:{event_id}"),
1208 "kind": kind,
1209 "workspace": workspace,
1210 "title": notice.title,
1211 "body": clip(&notice.body, 140),
1212 "href": if url.starts_with('/') { url } else { "/inbox" },
1213 "actor": { "kind": "system", "id": name, "name": name },
1214 "created_at": at,
1215 })
1216}
1217
1218/// What the notify service is told when a person's inbox count changes
1219/// (its `set_inbox`): the count as it now is, never a difference, so a
1220/// call repeated or out of order still ends right.
1221pub fn inbox_count_update(username: &str, unread: u32) -> serde_json::Value {
1222 serde_json::json!({ "username": username.to_lowercase(), "unread": unread })
1223}
1224
1225/// Tells the notify service a person's unread count as it now is, so every
1226/// tab of theirs shows it at once: after items arrive, and after marks made
1227/// anywhere (the site, the API, MCP). Logged and left when it fails.
1228pub async fn tell_inbox_count(db: &D1Database, notify: &Fetcher, username: &str) {
1229 let counted = counts(db, InboxCountsArgs { username: username.to_owned() }).await;
1230 let told: Result<serde_json::Value> = match counted {
1231 Ok(counts) => g1t_kit::call(notify, "set_inbox", &inbox_count_update(username, counts.unread)).await,
1232 Err(error) => Err(error),
1233 };
1234 if let Err(error) = told {
1235 worker::console_error!("inbox: live count not sent: {error}");
1236 }
1237}
1238
1239/// Hands items to the notify service, one call each. Logged and left when it
1240/// cannot be reached: the inbox and email do not wait on it.
1241async fn tell_live(notify: &Fetcher, items: Vec<(String, serde_json::Value)>) {
1242 #[derive(serde::Serialize)]
1243 struct NotifyArgs {
1244 username: String,
1245 notification: serde_json::Value,
1246 }
1247 for (username, notification) in items {
1248 let told: Result<serde_json::Value> = g1t_kit::call(notify, "notify", &NotifyArgs { username, notification }).await;
1249 if let Err(error) = told {
1250 worker::console_error!("inbox: live notification not sent: {error}");
1251 }
1252 }
1253}
1254
Inbox: threads, reasons, subscriptions and watching1255/// Emails what was news to people who asked to be emailed for its reason.
1256/// Identity sends each, only to a confirmed address and only if the person
1257/// can still read the repository.
1258async fn email(db: &D1Database, identity: &Fetcher, written: Vec<Written>) -> Result<()> {
1259 if written.is_empty() {
1260 return Ok(());
1261 }
1262 let people: Vec<String> = written
1263 .iter()
1264 .map(|item| item.notice.username.clone())
1265 .collect::<HashSet<_>>()
1266 .into_iter()
1267 .collect();
1268 let settings = subscriptions::settings_of(db, &people).await?;
1269 for item in written {
1270 let wants = settings.get(&item.notice.username).map_or(&DEFAULT_EMAIL[..], |settings| &settings.email[..]);
1271 if !wants.contains(&item.notice.reason) {
1272 continue;
1273 }
1274 let sent: Result<bool> = g1t_kit::call(
1275 identity,
1276 "notify_by_email",
1277 &NotifyByEmailArgs {
1278 username: item.notice.username.clone(),
1279 repo_id: item.repo_id.clone(),
1280 subject: item.notice.title.clone(),
1281 intro: item.notice.body.clone(),
1282 quote: None,
1283 path: item.url.clone(),
1284 reason: item.notice.reason,
1285 },
1286 )
1287 .await;
1288 if let Err(error) = sent {
1289 worker::console_error!("inbox: email to {} not sent: {error}", item.notice.username);
1290 }
1291 }
Inbox: the events service tells people what needs them as events arrive1292 Ok(())
1293}
1294
1295// --- Reading and changing ----------------------------------------------------
1296
1297#[derive(Deserialize)]
Inbox: threads, reasons, subscriptions and watching1298pub(crate) struct Row {
1299 pub(crate) id: String,
Inbox: the events service tells people what needs them as events arrive1300 reason: String,
1301 severity: String,
1302 title: String,
1303 body: String,
Inbox: threads, reasons, subscriptions and watching1304 event: Option<String>,
Inbox: the events service tells people what needs them as events arrive1305 workspace: Option<String>,
Inbox: threads, reasons, subscriptions and watching1306 pub(crate) repo_id: Option<String>,
Inbox: the events service tells people what needs them as events arrive1307 repo: Option<String>,
1308 subject: Option<String>,
1309 number: Option<f64>,
1310 run_id: Option<String>,
Inbox: threads, reasons, subscriptions and watching1311 link: Option<String>,
Inbox: the events service tells people what needs them as events arrive1312 actor: Option<String>,
Inbox: threads, reasons, subscriptions and watching1313 activity: Option<f64>,
Inbox: the events service tells people what needs them as events arrive1314 created_at: String,
Inbox: threads, reasons, subscriptions and watching1315 updated_at: Option<String>,
Inbox: the events service tells people what needs them as events arrive1316 read_at: Option<String>,
1317 done_at: Option<String>,
1318 saved: f64,
1319 snoozed_until: Option<String>,
Inbox: threads, reasons, subscriptions and watching1320 bumped: Option<String>,
Inbox: the events service tells people what needs them as events arrive1321}
1322
1323impl Row {
Inbox: threads, reasons, subscriptions and watching1324 pub(crate) fn into_item(self) -> InboxItem {
Inbox: the events service tells people what needs them as events arrive1325 let subject = self.subject.as_deref().and_then(SubjectKind::parse);
1326 let number = self.number.map(|n| n as u32);
1327 InboxItem {
Inbox: threads, reasons, subscriptions and watching1328 url: url(self.repo.as_deref(), subject, number, self.run_id.as_deref(), self.link.as_deref()),
Inbox: the events service tells people what needs them as events arrive1329 id: self.id,
Inbox: threads, reasons, subscriptions and watching1330 reason: Reason::parse(&self.reason).unwrap_or(Reason::Subscribed),
Inbox: the events service tells people what needs them as events arrive1331 severity: Severity::parse(&self.severity).unwrap_or(Severity::Info),
1332 title: self.title,
1333 body: self.body,
Inbox: threads, reasons, subscriptions and watching1334 event: self.event,
Inbox: the events service tells people what needs them as events arrive1335 repo: self.repo,
1336 workspace: self.workspace,
1337 subject,
1338 number,
1339 actor: self.actor,
Inbox: threads, reasons, subscriptions and watching1340 count: self.activity.map_or(1, |n| n as u32),
1341 updated_at: self.updated_at.unwrap_or_else(|| self.created_at.clone()),
Inbox: the events service tells people what needs them as events arrive1342 created_at: self.created_at,
1343 read_at: self.read_at,
1344 done_at: self.done_at,
1345 saved: self.saved != 0.0,
1346 snoozed_until: self.snoozed_until,
1347 }
1348 }
1349}
1350
Inbox: threads, reasons, subscriptions and watching1351pub(crate) const COLUMNS: &str = "id, reason, severity, title, body, event, workspace, repo_id, repo, subject, number, run_id,
1352 link, actor, activity, created_at, updated_at, read_at, done_at, saved, snoozed_until, bumped";
Inbox: the events service tells people what needs them as events arrive1353
1354/// The conditions that pick a view's items, after `username = ?1`; `?2` is now.
1355fn view_filter(view: InboxView) -> &'static str {
1356 match view {
1357 InboxView::Inbox => "done_at IS NULL AND (snoozed_until IS NULL OR snoozed_until <= ?2)",
1358 InboxView::Saved => "saved = 1",
1359 InboxView::Done => "done_at IS NOT NULL",
1360 }
1361}
1362
1363/// Whether a list ranks unread warnings first: the inbox itself, unfiltered.
1364fn ranked(a: &ListInboxArgs) -> bool {
Inbox: threads, reasons, subscriptions and watching1365 a.view == InboxView::Inbox
1366 && a.severity.is_none()
1367 && a.reason.is_none()
1368 && !a.participating
1369 && a.repo_id.is_none()
1370 && !a.unread
1371 && a.since.is_none()
1372 && a.updated_before.is_none()
Inbox: the events service tells people what needs them as events arrive1373}
1374
Inbox: threads, reasons, subscriptions and watching1375/// The conditions a list's filters add, and the values they bind, numbered
1376/// after the `bound` values already given.
1377fn filters(a: &ListInboxArgs, bound: usize) -> (Vec<String>, Vec<String>) {
1378 let mut conditions = Vec::new();
1379 let mut values: Vec<String> = Vec::new();
1380 let mut bind = |value: &str, condition: &str| {
1381 values.push(value.to_owned());
1382 conditions.push(condition.replace('?', &format!("?{}", bound + values.len())));
1383 };
1384 if let Some(severity) = a.severity {
1385 bind(severity.as_str(), "severity = ?");
1386 }
1387 if let Some(reason) = a.reason {
1388 bind(reason.as_str(), "reason = ?");
1389 }
1390 if let Some(repo_id) = &a.repo_id {
1391 bind(repo_id, "repo_id = ?");
1392 }
1393 if let Some(since) = &a.since {
1394 bind(since.trim(), "COALESCE(updated_at, created_at) >= ?");
1395 }
1396 if let Some(before) = &a.updated_before {
1397 bind(before.trim(), "COALESCE(updated_at, created_at) < ?");
1398 }
1399 if a.participating {
1400 conditions.push("reason NOT IN ('manual', 'subscribed')".to_owned());
1401 }
1402 if a.unread {
1403 conditions.push("read_at IS NULL".to_owned());
1404 }
1405 (conditions, values)
1406}
1407
Inbox: the events service tells people what needs them as events arrive1408pub async fn list(db: &D1Database, repos: &Fetcher, a: ListInboxArgs) -> Result<InboxPage> {
1409 let Some(viewer) = &a.viewer else {
1410 return Ok(InboxPage::default());
1411 };
1412 let username = viewer.username.to_lowercase();
1413 let now = rfc3339(now_ms());
1414 let limit = a.limit.unwrap_or(DEFAULT_INBOX_PAGE).clamp(1, MAX_INBOX_PAGE);
1415 let mut values: Vec<JsValue> = vec![username.as_str().into(), now.as_str().into()];
Inbox: threads, reasons, subscriptions and watching1416 let mut conditions = vec!["username = ?1".to_owned(), view_filter(a.view).to_owned()];
1417 let (filtered, bound) = filters(&a, values.len());
1418 conditions.extend(filtered);
1419 values.extend(bound.iter().map(|value| JsValue::from(value.as_str())));
Inbox: the events service tells people what needs them as events arrive1420 let ranked = ranked(&a);
1421 // Unread warnings lead the first page, and are left out of the rest.
1422 let leading = "severity = 'warning' AND read_at IS NULL";
1423 let mut rest = conditions.clone();
Inbox: threads, reasons, subscriptions and watching1424 let mut rest_values = values.clone();
Inbox: the events service tells people what needs them as events arrive1425 if ranked {
1426 rest.push(format!("NOT ({leading})"));
1427 }
1428 if let Some(before) = &a.before {
Inbox: threads, reasons, subscriptions and watching1429 rest_values.push(before.as_str().into());
1430 rest.push(format!("bumped < ?{}", rest_values.len()));
Inbox: the events service tells people what needs them as events arrive1431 }
Inbox: threads, reasons, subscriptions and watching1432 rest_values.push((limit + 1).into());
1433 let order = if a.view == InboxView::Done { "done_at DESC, bumped DESC" } else { "bumped DESC" };
Inbox: the events service tells people what needs them as events arrive1434 let mut statements = vec![
1435 db.prepare(format!(
1436 "SELECT {COLUMNS} FROM inbox_items WHERE {} ORDER BY {order} LIMIT ?{}",
1437 rest.join(" AND "),
Inbox: threads, reasons, subscriptions and watching1438 rest_values.len()
Inbox: the events service tells people what needs them as events arrive1439 ))
Inbox: threads, reasons, subscriptions and watching1440 .bind(&rest_values)?,
Inbox: the events service tells people what needs them as events arrive1441 ];
1442 if ranked && a.before.is_none() {
1443 statements.push(
1444 db.prepare(format!(
Inbox: threads, reasons, subscriptions and watching1445 "SELECT {COLUMNS} FROM inbox_items WHERE {} AND {leading} ORDER BY bumped DESC LIMIT ?3",
1446 conditions.join(" AND ")
Inbox: the events service tells people what needs them as events arrive1447 ))
1448 .bind(&[username.as_str().into(), now.as_str().into(), MAX_RANKED.into()])?,
1449 );
1450 }
1451 let results = db.batch(statements).await?;
1452 let mut rows = results.first().map(|result| result.results::<Row>()).transpose()?.unwrap_or_default();
1453 let next = if rows.len() > limit as usize {
1454 rows.truncate(limit as usize);
Inbox: threads, reasons, subscriptions and watching1455 rows.last().map(|row| row.bumped.clone().unwrap_or_else(|| row.id.clone()))
Inbox: the events service tells people what needs them as events arrive1456 } else {
1457 None
1458 };
1459 let mut items: Vec<Row> = results.get(1).map(|result| result.results::<Row>()).transpose()?.unwrap_or_default();
1460 items.extend(rows);
Inbox: threads, reasons, subscriptions and watching1461 let items = readable_only(db, repos, &a.viewer, &username, items).await?;
1462 Ok(InboxPage {
1463 items: items.into_iter().map(Row::into_item).collect(),
1464 next,
1465 })
1466}
Inbox: the events service tells people what needs them as events arrive1467
Inbox: threads, reasons, subscriptions and watching1468/// Rows about repositories the viewer can still read; the rest are
1469/// removed from their inbox as they are found.
1470pub(crate) async fn readable_only(db: &D1Database, repos: &Fetcher, viewer: &g1t_contracts::Viewer, username: &str, mut rows: Vec<Row>) -> Result<Vec<Row>> {
1471 let ids: Vec<String> = rows
Inbox: the events service tells people what needs them as events arrive1472 .iter()
1473 .filter_map(|row| row.repo_id.clone())
1474 .collect::<HashSet<_>>()
1475 .into_iter()
1476 .collect();
Inbox: threads, reasons, subscriptions and watching1477 if ids.is_empty() {
1478 return Ok(rows);
1479 }
1480 let readable: Vec<Repo> = g1t_kit::call(
1481 repos,
1482 "readable",
1483 &ReadableArgs {
1484 ids: ids.clone(),
1485 viewer: viewer.clone(),
1486 },
1487 )
1488 .await?;
1489 let readable: HashSet<String> = readable.into_iter().map(|repo| repo.id).collect();
1490 let gone: Vec<String> = ids.into_iter().filter(|id| !readable.contains(id)).collect();
1491 if !gone.is_empty() {
1492 forget(db, username, &gone).await?;
1493 rows.retain(|row| row.repo_id.as_ref().is_none_or(|id| !gone.contains(id)));
Inbox: the events service tells people what needs them as events arrive1494 }
Inbox: threads, reasons, subscriptions and watching1495 Ok(rows)
Inbox: the events service tells people what needs them as events arrive1496}
1497
1498/// Removes a person's items about repositories they cannot read.
1499async fn forget(db: &D1Database, username: &str, repo_ids: &[String]) -> Result<()> {
1500 let marks = vec!["?"; repo_ids.len()].join(", ");
1501 let mut values: Vec<JsValue> = vec![username.into()];
1502 values.extend(repo_ids.iter().map(|id| JsValue::from(id.as_str())));
Inbox: threads, reasons, subscriptions and watching1503 db.batch(vec![
1504 db.prepare(format!(
1505 "DELETE FROM inbox_activity WHERE item_id IN (SELECT id FROM inbox_items WHERE username = ? AND repo_id IN ({marks}))"
1506 ))
1507 .bind(&values)?,
1508 db.prepare(format!("DELETE FROM inbox_items WHERE username = ? AND repo_id IN ({marks})"))
1509 .bind(&values)?,
1510 ])
1511 .await?;
Inbox: the events service tells people what needs them as events arrive1512 Ok(())
1513}
1514
1515#[derive(Deserialize)]
1516struct CountRow {
1517 severity: String,
1518 n: f64,
1519}
1520
1521pub async fn counts(db: &D1Database, a: InboxCountsArgs) -> Result<InboxCounts> {
1522 let rows = db
1523 .prepare(
1524 "SELECT severity, count(*) AS n FROM inbox_items
1525 WHERE username = ? AND read_at IS NULL AND done_at IS NULL
1526 AND (snoozed_until IS NULL OR snoozed_until <= ?)
1527 GROUP BY severity",
1528 )
1529 .bind(&[a.username.to_lowercase().into(), rfc3339(now_ms()).into()])?
1530 .all()
1531 .await?
1532 .results::<CountRow>()?;
1533 let mut counts = InboxCounts::default();
1534 for row in rows {
1535 let n = row.n as u32;
1536 counts.unread += n;
1537 match Severity::parse(&row.severity) {
1538 Some(Severity::Error) => counts.error += n,
1539 Some(Severity::Warning) => counts.warning += n,
1540 Some(Severity::Success) => counts.success += n,
1541 Some(Severity::Info) | None => counts.info += n,
1542 }
1543 }
1544 Ok(counts)
1545}
1546
1547/// The change a mark makes, as a `SET` clause; `?1` is now.
1548fn mark_change(mark: InboxMark) -> &'static str {
1549 match mark {
1550 InboxMark::Read => "read_at = COALESCE(read_at, ?1)",
1551 InboxMark::Unread => "read_at = NULL",
1552 InboxMark::Done => "done_at = ?1, read_at = COALESCE(read_at, ?1)",
1553 InboxMark::Undone => "done_at = NULL",
1554 InboxMark::Save => "saved = 1",
1555 InboxMark::Unsave => "saved = 0",
1556 InboxMark::Snooze => "snoozed_until = ?2, read_at = COALESCE(read_at, ?1)",
Inbox: threads, reasons, subscriptions and watching1557 InboxMark::Unsnooze => "snoozed_until = NULL",
Inbox: the events service tells people what needs them as events arrive1558 }
1559}
1560
1561#[derive(Deserialize)]
1562struct IdRow {
1563 #[allow(dead_code)]
1564 id: String,
1565}
1566
1567/// Changes the person's own items. Returns how many changed.
1568pub async fn mark(db: &D1Database, a: MarkInboxArgs) -> Result<u32> {
1569 let now = rfc3339(now_ms());
1570 let until = match (a.mark, a.until.as_deref().map(str::trim)) {
1571 // Times compare as text (`g1t_contracts::time`); a snooze is for later.
1572 (InboxMark::Snooze, Some(until)) if until.len() == now.len() && until > now.as_str() => until.to_owned(),
1573 (InboxMark::Snooze, _) => return Ok(0),
1574 _ => String::new(),
1575 };
1576 let mut values: Vec<JsValue> = vec![now.as_str().into(), until.as_str().into(), a.username.to_lowercase().into()];
1577 let target = if a.all && a.ids.is_empty() {
1578 let mut target = "done_at IS NULL".to_owned();
Inbox: threads, reasons, subscriptions and watching1579 let mut bind = |values: &mut Vec<JsValue>, value: &str, condition: &str| {
1580 values.push(value.into());
1581 target.push_str(&format!(" AND {}", condition.replace('?', &format!("?{}", values.len()))));
1582 };
Inbox: the events service tells people what needs them as events arrive1583 if let Some(severity) = a.severity {
Inbox: threads, reasons, subscriptions and watching1584 bind(&mut values, severity.as_str(), "severity = ?");
1585 }
1586 if let Some(repo_id) = &a.repo_id {
1587 bind(&mut values, repo_id, "repo_id = ?");
1588 }
1589 if let Some(last_read_at) = a.last_read_at.as_deref().map(str::trim).filter(|at| !at.is_empty()) {
1590 bind(&mut values, last_read_at, "COALESCE(updated_at, created_at) <= ?");
Inbox: the events service tells people what needs them as events arrive1591 }
1592 target
1593 } else {
1594 let ids: Vec<&String> = a.ids.iter().take(MAX_MARK).collect();
1595 if ids.is_empty() {
1596 return Ok(0);
1597 }
1598 let first = values.len() + 1;
1599 values.extend(ids.iter().map(|id| JsValue::from(id.as_str())));
1600 let marks: Vec<String> = (first..values.len() + 1).map(|at| format!("?{at}")).collect();
1601 format!("id IN ({})", marks.join(", "))
1602 };
1603 let changed = db
1604 .prepare(format!(
1605 "UPDATE inbox_items SET {} WHERE username = ?3 AND {target} RETURNING id",
1606 mark_change(a.mark)
1607 ))
1608 .bind(&values)?
1609 .all()
1610 .await?
1611 .results::<IdRow>()?;
1612 Ok(changed.len() as u32)
1613}
1614
Inbox: threads, reasons, subscriptions and watching1615#[derive(Deserialize)]
1616struct ActivityRow {
1617 reason: String,
1618 severity: String,
1619 title: String,
1620 body: String,
1621 event: Option<String>,
1622 actor: Option<String>,
1623 created_at: String,
1624}
1625
1626/// One of the viewer's threads, with its history and their subscription.
1627pub async fn thread(db: &D1Database, repos: &Fetcher, work: &Fetcher, a: ThreadArgs) -> Result<Option<InboxThread>> {
1628 let Some(viewer) = &a.viewer else {
1629 return Ok(None);
1630 };
1631 let username = viewer.username.to_lowercase();
1632 let row = db
1633 .prepare(format!("SELECT {COLUMNS} FROM inbox_items WHERE id = ? AND username = ?"))
1634 .bind(&[a.id.as_str().into(), username.as_str().into()])?
1635 .first::<Row>(None)
1636 .await?;
1637 let Some(row) = readable_only(db, repos, &a.viewer, &username, row.into_iter().collect()).await?.pop() else {
1638 return Ok(None);
1639 };
1640 let activity = db
1641 .prepare(
1642 "SELECT reason, severity, title, body, event, actor, created_at FROM inbox_activity
1643 WHERE item_id = ? ORDER BY id DESC LIMIT ?",
1644 )
1645 .bind(&[row.id.as_str().into(), MAX_ACTIVITY.into()])?
1646 .all()
1647 .await?
1648 .results::<ActivityRow>()?
1649 .into_iter()
1650 .map(|row| InboxActivity {
1651 reason: Reason::parse(&row.reason).unwrap_or(Reason::Subscribed),
1652 severity: Severity::parse(&row.severity).unwrap_or(Severity::Info),
1653 title: row.title,
1654 body: row.body,
1655 event: row.event,
1656 actor: row.actor,
1657 created_at: row.created_at,
1658 })
1659 .collect();
1660 let item = row.into_item();
1661 let subscription = match (item.subject, item.number) {
1662 (Some(SubjectKind::Issue | SubjectKind::Pull), Some(_)) => {
1663 subscriptions::subscription(
1664 db,
1665 work,
1666 SubscriptionArgs {
1667 viewer: a.viewer.clone(),
1668 id: Some(item.id.clone()),
1669 ..SubscriptionArgs::default()
1670 },
1671 )
1672 .await?
1673 }
1674 _ => None,
1675 };
1676 Ok(Some(InboxThread {
1677 item,
1678 activity,
1679 subscription,
1680 }))
1681}
1682
Inbox: the events service tells people what needs them as events arrive1683// --- Keeping up ------------------------------------------------------------------
1684
1685/// Moves rows with renamed workspaces and repositories, and drops those
Merge account deletion: soft delete for 30 days, staff restore and purge, ghost for what remains (identity 0037)1686/// of purged repositories, deleted workspaces and purged accounts.
Inbox: the events service tells people what needs them as events arrive1687pub async fn follow(db: &D1Database, events: &[Event]) -> Result<()> {
1688 let mut statements = Vec::new();
1689 for event in events {
1690 let text = |key: &str| event.data[key].as_str().unwrap_or_default().to_lowercase();
1691 match event.kind.as_str() {
1692 "workspace.renamed" => {
1693 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
1694 continue;
1695 };
1696 let (from, to) = (renamed.from.to_lowercase(), renamed.to.to_lowercase());
1697 if from == to {
1698 continue;
1699 }
1700 statements.push(
1701 db.prepare(
Inbox: threads, reasons, subscriptions and watching1702 "UPDATE inbox_items SET repo = ?2 || substr(repo, length(?1) + 1), workspace = ?2,
1703 link = CASE WHEN link LIKE '/' || ?1 || '/%' THEN '/' || ?2 || substr(link, length(?1) + 2) ELSE link END
Inbox: the events service tells people what needs them as events arrive1704 WHERE workspace = ?1",
1705 )
1706 .bind(&[from.as_str().into(), to.as_str().into()])?,
1707 );
Inbox: threads, reasons, subscriptions and watching1708 statements.push(
1709 db.prepare("UPDATE inbox_watching SET repo = ?2 || substr(repo, length(?1) + 1) WHERE repo LIKE ?1 || '/%'")
1710 .bind(&[from.as_str().into(), to.as_str().into()])?,
1711 );
Inbox: the events service tells people what needs them as events arrive1712 }
Inbox: threads, reasons, subscriptions and watching1713 "repo.renamed" => {
1714 let path = format!("{}/{}", text("namespace"), text("to"));
1715 for table in ["inbox_items", "inbox_watching"] {
1716 statements.push(
1717 db.prepare(format!("UPDATE {table} SET repo = ? WHERE repo_id = ?"))
1718 .bind(&[path.as_str().into(), text("repoId").into()])?,
1719 );
1720 }
1721 }
1722 "repo.transferred" => {
1723 let path = format!("{}/{}", text("to"), text("name"));
1724 statements.push(
1725 db.prepare("UPDATE inbox_items SET repo = ?, workspace = ? WHERE repo_id = ?")
1726 .bind(&[path.as_str().into(), text("to").into(), text("repoId").into()])?,
1727 );
1728 statements.push(
1729 db.prepare("UPDATE inbox_watching SET repo = ? WHERE repo_id = ?")
1730 .bind(&[path.as_str().into(), text("repoId").into()])?,
1731 );
1732 }
1733 "repo.purged" => {
1734 for sql in [
1735 "DELETE FROM inbox_activity WHERE item_id IN (SELECT id FROM inbox_items WHERE repo_id = ?)",
1736 "DELETE FROM inbox_items WHERE repo_id = ?",
1737 "DELETE FROM inbox_subscriptions WHERE repo_id = ?",
1738 "DELETE FROM inbox_watching WHERE repo_id = ?",
1739 ] {
1740 statements.push(db.prepare(sql).bind(&[text("repoId").into()])?);
1741 }
1742 }
1743 "workspace.deleted" => {
1744 statements.push(
1745 db.prepare("DELETE FROM inbox_items WHERE workspace = ?")
1746 .bind(&[text("slug").into()])?,
1747 );
1748 statements.push(
1749 db.prepare("DELETE FROM inbox_watching WHERE repo LIKE ? || '/%'")
1750 .bind(&[text("slug").into()])?,
1751 );
1752 }
Merge account deletion: soft delete for 30 days, staff restore and purge, ghost for what remains (identity 0037)1753 // An account purged: its inbox, what it watched and its
1754 // settings go with it (identity's account_deletion.rs).
1755 "user.deleted" => {
1756 let username = text("username");
1757 if username.is_empty() {
1758 continue;
1759 }
1760 for sql in [
1761 "DELETE FROM inbox_activity WHERE username = ?",
1762 "DELETE FROM inbox_items WHERE username = ?",
1763 "DELETE FROM inbox_subscriptions WHERE username = ?",
1764 "DELETE FROM inbox_watching WHERE username = ?",
1765 "DELETE FROM inbox_settings WHERE username = ?",
1766 ] {
1767 statements.push(db.prepare(sql).bind(&[username.as_str().into()])?);
1768 }
1769 }
Inbox: the events service tells people what needs them as events arrive1770 _ => {}
1771 }
1772 }
1773 if !statements.is_empty() {
1774 db.batch(statements).await?;
1775 }
1776 Ok(())
1777}
1778
1779/// Removes items done more than [`DONE_DAYS`] ago, and any not saved older
Inbox: threads, reasons, subscriptions and watching1780/// than [`MAX_DAYS`], with their history. Returns how many went.
Inbox: the events service tells people what needs them as events arrive1781pub async fn purge(db: &D1Database, now: u64) -> Result<u32> {
1782 let done = crate::audit::keep_from(now, DONE_DAYS);
1783 let oldest = crate::audit::keep_from(now, MAX_DAYS);
Inbox: threads, reasons, subscriptions and watching1784 let statements: Vec<D1PreparedStatement> = vec![
1785 db.prepare("DELETE FROM inbox_items WHERE saved = 0 AND done_at < ?")
1786 .bind(&[done.into()])?,
1787 db.prepare("DELETE FROM inbox_items WHERE saved = 0 AND COALESCE(updated_at, created_at) < ?")
1788 .bind(&[oldest.into()])?,
1789 db.prepare("DELETE FROM inbox_activity WHERE NOT EXISTS (SELECT 1 FROM inbox_items WHERE inbox_items.id = inbox_activity.item_id)"),
1790 ];
1791 let results = db.batch(statements).await?;
Inbox: the events service tells people what needs them as events arrive1792 let mut removed = 0;
Inbox: threads, reasons, subscriptions and watching1793 for result in results.into_iter().take(2) {
Inbox: the events service tells people what needs them as events arrive1794 removed += result.meta()?.and_then(|meta| meta.changes).unwrap_or(0) as u32;
1795 }
1796 Ok(removed)
1797}
1798
1799#[cfg(test)]
1800mod tests {
1801 use super::*;
Inbox: threads, reasons, subscriptions and watching1802 use crate::subscriptions::{State, Subscription, Watcher};
Inbox: the events service tells people what needs them as events arrive1803 use serde_json::json;
1804
1805 fn event(kind: &str, actor: Option<&str>, data: serde_json::Value) -> Event {
1806 Event {
1807 id: "evt_1".into(),
1808 kind: kind.into(),
1809 source: "work".into(),
1810 time: "2026-10-07T12:00:00.000Z".into(),
1811 repo_id: Some("rep_1".into()),
1812 actor: actor.map(str::to_owned),
1813 data,
1814 }
1815 }
1816
1817 fn person(id: &str, username: &str) -> Principal {
1818 Principal {
1819 id: id.into(),
1820 username: username.into(),
1821 }
1822 }
1823
1824 fn actor(id: &str, username: &str) -> Actor {
1825 Actor {
1826 id: Some(id.into()),
1827 username: Some(username.into()),
1828 }
1829 }
1830
1831 /// A pull request ana opened, assigned to bo, for an issue cy filed and dee is assigned.
1832 fn pull() -> InboxSubject {
1833 InboxSubject {
1834 kind: Some(SubjectKind::Pull),
1835 title: "Add the inbox".into(),
1836 author: person("usr_ana", "ana"),
1837 assignees: vec!["bo".into()],
1838 issue: Some(Box::new(InboxSubject {
1839 kind: Some(SubjectKind::Issue),
1840 title: "An inbox".into(),
1841 author: person("usr_cy", "cy"),
1842 assignees: vec!["dee".into()],
1843 ..InboxSubject::default()
1844 })),
1845 ..InboxSubject::default()
1846 }
1847 }
1848
1849 /// A change g1t made for ana.
1850 fn g1t_pull() -> InboxSubject {
1851 InboxSubject {
1852 author: person(AGENT_ID, "g1t"),
1853 requested_by: Some(person("usr_ana", "ana")),
1854 ..pull()
1855 }
1856 }
1857
Inbox: threads, reasons, subscriptions and watching1858 fn nobody() -> Audience {
1859 Audience::default()
1860 }
1861
1862 fn told(notices: &[Notice]) -> Vec<(&str, Reason, Severity)> {
Inbox: the events service tells people what needs them as events arrive1863 notices
1864 .iter()
1865 .map(|notice| (notice.username.as_str(), notice.reason, notice.severity))
1866 .collect()
1867 }
1868
1869 fn comment(author: Principal, body: &str, mentions: &[&str]) -> InboxSubject {
1870 InboxSubject {
1871 comment: Some(InboxComment {
1872 author,
1873 excerpt: body.into(),
1874 mentions: mentions.iter().map(|name| (*name).to_owned()).collect(),
1875 ..InboxComment::default()
1876 }),
1877 ..pull()
1878 }
1879 }
1880
Inbox: threads, reasons, subscriptions and watching1881 fn subscribed(rows: &[(&str, State, Option<Reason>)]) -> Audience {
1882 Audience {
1883 subscriptions: rows
1884 .iter()
1885 .map(|(name, state, reason)| Subscription {
1886 username: (*name).into(),
1887 state: *state,
1888 reason: *reason,
1889 })
1890 .collect(),
1891 watchers: Vec::new(),
1892 }
1893 }
1894
1895 fn watched(rows: &[(&str, WatchLevel, &[&str])]) -> Audience {
1896 Audience {
1897 subscriptions: Vec::new(),
1898 watchers: rows
1899 .iter()
1900 .map(|(name, level, events)| Watcher {
1901 username: (*name).into(),
1902 level: *level,
1903 events: events.iter().map(|kind| (*kind).to_owned()).collect(),
1904 })
1905 .collect(),
1906 }
1907 }
1908
Inbox: the events service tells people what needs them as events arrive1909 #[test]
1910 fn only_events_that_tell_someone_are_read() {
1911 let asked = |kind: &str, data| wants(&event(kind, None, data));
1912 assert_eq!(
1913 asked("checks.completed", json!({ "repoId": "rep_1", "number": 4, "status": "failed" })),
1914 Some(Wanted { repo_id: "rep_1".into(), number: Some(4), comment_id: None })
1915 );
1916 assert_eq!(asked("checks.completed", json!({ "repoId": "rep_1", "number": 4, "status": "passed" })), None);
1917 assert_eq!(asked("review.completed", json!({ "repoId": "rep_1", "number": 4 })), None);
1918 assert_eq!(asked("workflow.completed", json!({ "repoId": "rep_1", "conclusion": "success" })), None);
1919 assert_eq!(
1920 asked("workflow.completed", json!({ "repoId": "rep_1", "conclusion": "failure" })),
1921 Some(Wanted { repo_id: "rep_1".into(), number: None, comment_id: None })
1922 );
1923 assert_eq!(
1924 asked("comment.created", json!({ "repoId": "rep_1", "number": 2, "commentId": "cmt_1" })).map(|w| w.comment_id),
1925 Some(Some("cmt_1".into()))
1926 );
1927 assert_eq!(asked("git.push", json!({ "repoId": "rep_1" })), None);
Inbox: threads, reasons, subscriptions and watching1928 for kind in ["issue.opened", "pull.review_requested", "pull.stalled", "issue.assigned", "pull.closed", "issue.reopened"] {
1929 assert!(asked(kind, json!({ "repoId": "rep_1", "number": 1 })).is_some(), "{kind}");
1930 }
1931 assert_eq!(asked("deployment.failed", json!({ "repoId": "rep_1" })).map(|w| w.number), Some(None));
Inbox: the events service tells people what needs them as events arrive1932 }
1933
1934 #[test]
Inbox: threads, reasons, subscriptions and watching1935 fn a_waiting_or_stopped_agent_needs_the_pull_requests_and_the_issues_people_first() {
Inbox: the events service tells people what needs them as events arrive1936 let asked = event("agent.asked", Some("usr_agent"), json!({ "number": 7 }));
Inbox: threads, reasons, subscriptions and watching1937 let notices_ = notices(&asked, "acme/rocket", &Actor::default(), Some(&pull()), &nobody());
Inbox: the events service tells people what needs them as events arrive1938 assert_eq!(
Inbox: threads, reasons, subscriptions and watching1939 told(&notices_),
Inbox: the events service tells people what needs them as events arrive1940 vec![
Inbox: threads, reasons, subscriptions and watching1941 ("ana", Reason::Agent, Severity::Warning),
1942 ("cy", Reason::Agent, Severity::Warning),
1943 ("dee", Reason::Agent, Severity::Warning),
Inbox: the events service tells people what needs them as events arrive1944 ]
1945 );
Inbox: threads, reasons, subscriptions and watching1946 assert_eq!(notices_[0].title, "An agent is waiting on acme/rocket#7");
1947 assert_eq!(notices_[0].body, "Add the inbox");
1948 // Stopping says why; even someone who unsubscribed hears of it.
1949 let stalled = event("pull.stalled", None, json!({ "number": 7, "detail": "Its checks could not be run." }));
1950 let audience = subscribed(&[("ana", State::Unsubscribed, None)]);
1951 let notices_ = notices(&stalled, "acme/rocket", &Actor::default(), Some(&g1t_pull()), &audience);
1952 assert_eq!(notices_[0].username, "ana");
1953 assert_eq!(notices_[0].title, "g1t stopped on acme/rocket#7 and needs you");
1954 assert_eq!(notices_[0].body, "Its checks could not be run.");
1955 // Ignoring the thread silences even that.
1956 let audience = subscribed(&[("ana", State::Ignored, None)]);
1957 let notices_ = notices(&stalled, "acme/rocket", &Actor::default(), Some(&g1t_pull()), &audience);
1958 assert!(!notices_.iter().any(|notice| notice.username == "ana"));
1959 }
1960
1961 #[test]
1962 fn whatever_an_agent_waits_on_closes_when_it_goes_on() {
1963 for kind in RESUMES {
1964 assert_eq!(resolves(&event(kind, None, json!({ "repoId": "rep_1", "number": 7 }))), Some("rep_1#7".into()));
1965 }
1966 assert_eq!(resolves(&event("comment.created", None, json!({ "repoId": "rep_1", "number": 7 }))), None);
Inbox: the events service tells people what needs them as events arrive1967 }
1968
1969 #[test]
Merge branch 'worktree-agent-a3abfcce648e87dca'1970 fn a_deployments_reviewers_are_told() {
1971 let asked = event(
1972 "deployment.review_requested",
1973 Some("usr_ana"),
1974 json!({
1975 "repoId": "rep_1", "runId": "run_7", "environment": "production", "workflow": "Deploy",
1976 "title": "Ship it", "notify": ["cy", "ana"], "link": "/acme/rocket/actions/runs/run_7"
1977 }),
1978 );
1979 let wanted = wants(&asked).unwrap();
1980 assert_eq!(wanted.number, None);
1981 let thread = thread_of(&asked, &wanted, None);
1982 assert_eq!(thread.key, "rep_1/review/run_7/production");
1983 assert_eq!(thread.link.as_deref(), Some("/acme/rocket/actions/runs/run_7"));
1984 let told = notices(&asked, "acme/rocket", &actor("usr_ana", "ana"), None, &nobody());
1985 let names: Vec<&str> = told.iter().map(|notice| notice.username.as_str()).collect();
1986 // Whoever started the run is told too, when they review it.
1987 assert_eq!(names, ["cy", "ana"]);
1988 assert!(told[0].title.contains("waiting for your review to deploy to production"));
1989 }
1990
1991 #[test]
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar1992 fn teams_asked_to_review_tell_the_people_they_name() {
1993 let requested = event(
1994 "pull.review_requested",
1995 Some("usr_ana"),
1996 json!({ "number": 7, "reviewers": ["cy"], "teams": [
1997 { "team": "acme/backend", "notified": ["cy", "dee", "ana"], "assigned": ["cy"] }
1998 ] }),
1999 );
2000 let notices_ = notices(&requested, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
2001 // cy was picked and is told as a reviewer; dee through the team;
2002 // never whoever asked.
2003 assert_eq!(
2004 told(&notices_),
2005 vec![("cy", Reason::ReviewRequested, Severity::Warning), ("dee", Reason::ReviewRequested, Severity::Warning)]
2006 );
2007 assert_eq!(notices_[0].title, "ana asked you to review acme/rocket#7");
2008 assert_eq!(notices_[1].title, "ana asked @acme/backend to review acme/rocket#7");
2009 assert_eq!(
2010 subscribes(&requested, Some(&pull())),
2011 vec![("cy".into(), Reason::ReviewRequested), ("dee".into(), Reason::ReviewRequested), ("ana".into(), Reason::ReviewRequested)]
2012 );
2013 // Asked by the CODEOWNERS file.
2014 let owned = event(
2015 "pull.review_requested",
2016 None,
2017 json!({ "number": 7, "reviewers": ["bo"], "codeOwners": true, "teams": [{ "team": "acme/docs", "notified": ["wren"], "assigned": [] }] }),
2018 );
2019 let notices_ = notices(&owned, "acme/rocket", &Actor::default(), Some(&pull()), &nobody());
2020 assert_eq!(notices_[0].title, "acme/rocket#7 changes files you own");
2021 assert_eq!(notices_[1].title, "acme/rocket#7 changes files @acme/docs owns");
2022 }
2023
2024 #[test]
2025 fn a_team_mention_tells_its_people_once_and_a_name_wins() {
2026 let mut on = comment(person("usr_bo", "bo"), "cc @ana @acme/backend", &["ana"]);
2027 if let Some(comment) = on.comment.as_mut() {
2028 comment.team_mentions = vec![TeamMentioned {
2029 team: "acme/backend".into(),
2030 members: vec!["ana".into(), "cy".into(), "bo".into()],
2031 }];
2032 }
2033 let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
2034 let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &nobody());
2035 let mentioned: Vec<(&str, Reason)> = notices_.iter().map(|n| (n.username.as_str(), n.reason)).filter(|(_, r)| matches!(r, Reason::Mention | Reason::TeamMention)).collect();
2036 assert_eq!(mentioned, vec![("ana", Reason::Mention), ("cy", Reason::TeamMention)]);
2037 assert_eq!(
2038 notices_.iter().find(|n| n.username == "cy").unwrap().title,
2039 "bo mentioned @acme/backend on acme/rocket#7"
2040 );
2041 let subscribed = subscribes(&created, Some(&on));
2042 assert!(subscribed.contains(&("cy".into(), Reason::TeamMention)));
2043 assert!(subscribed.contains(&("ana".into(), Reason::Mention)));
2044 }
2045
2046 #[test]
2047 fn a_description_tells_the_people_and_teams_it_mentions() {
2048 let mut on = pull();
2049 on.mentions = vec!["dee".into()];
2050 on.team_mentions = vec![TeamMentioned {
2051 team: "acme/web".into(),
2052 members: vec!["eve".into()],
2053 }];
2054 let opened = event("pull.opened", Some("usr_ana"), json!({ "number": 7 }));
2055 let notices_ = notices(&opened, "acme/rocket", &actor("usr_ana", "ana"), Some(&on), &nobody());
2056 assert!(told(&notices_).contains(&("dee", Reason::Mention, Severity::Info)));
2057 assert!(told(&notices_).contains(&("eve", Reason::TeamMention, Severity::Info)));
2058 assert_eq!(
2059 subscribes(&opened, Some(&on)),
2060 vec![("dee".into(), Reason::Mention), ("eve".into(), Reason::TeamMention)]
2061 );
2062 }
2063
2064 #[test]
Inbox: threads, reasons, subscriptions and watching2065 fn reviewers_and_assignees_asked_are_told_never_whoever_asked() {
2066 let requested = event("pull.review_requested", Some("usr_ana"), json!({ "number": 7, "reviewers": ["bo", "g1t", "ana"] }));
2067 let notices_ = notices(&requested, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
2068 assert_eq!(told(&notices_), vec![("bo", Reason::ReviewRequested, Severity::Warning)]);
2069 assert_eq!(notices_[0].title, "ana asked you to review acme/rocket#7");
2070 let assigned = event("issue.assigned", Some("usr_ana"), json!({ "number": 3, "added": ["ana", "eve"] }));
2071 let notices_ = notices(&assigned, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
2072 assert_eq!(told(&notices_), vec![("eve", Reason::Assign, Severity::Info)]);
2073 assert_eq!(notices_[0].title, "ana assigned you to acme/rocket#3");
2074 // Both subscribe whoever they name, never g1t.
2075 assert_eq!(subscribes(&requested, Some(&pull())), vec![("bo".into(), Reason::ReviewRequested), ("ana".into(), Reason::ReviewRequested)]);
2076 assert_eq!(subscribes(&assigned, Some(&pull())), vec![("ana".into(), Reason::Assign), ("eve".into(), Reason::Assign)]);
2077 }
2078
2079 #[test]
Inbox: the events service tells people what needs them as events arrive2080 fn failures_go_to_whoever_answers_for_the_change() {
2081 let failed = event("checks.completed", None, json!({ "number": 7, "status": "failed" }));
Inbox: threads, reasons, subscriptions and watching2082 assert_eq!(
2083 told(&notices(&failed, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())),
2084 vec![("ana", Reason::CiActivity, Severity::Error)]
2085 );
Inbox: the events service tells people what needs them as events arrive2086 // g1t's change is the person's who asked for it, never g1t's.
Inbox: threads, reasons, subscriptions and watching2087 let notices_ = notices(&failed, "acme/rocket", &Actor::default(), Some(&g1t_pull()), &nobody());
2088 assert_eq!(told(&notices_), vec![("ana", Reason::CiActivity, Severity::Error)]);
Inbox: the events service tells people what needs them as events arrive2089 let errored = event("checks.completed", None, json!({ "number": 7, "status": "errored" }));
Inbox: threads, reasons, subscriptions and watching2090 assert_eq!(
2091 notices(&errored, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())[0].title,
2092 "Checks could not run on acme/rocket#7"
2093 );
Inbox: the events service tells people what needs them as events arrive2094 }
2095
2096 #[test]
2097 fn a_workflow_that_fails_tells_its_pull_requests_owner_even_if_they_pushed() {
2098 let failed = event("workflow.completed", Some("usr_ana"), json!({ "pull": 7, "workflow": "CI", "conclusion": "failure" }));
Inbox: threads, reasons, subscriptions and watching2099 let notices_ = notices(&failed, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
2100 assert_eq!(told(&notices_), vec![("ana", Reason::CiActivity, Severity::Error)]);
Inbox: the events service tells people what needs them as events arrive2101 assert_eq!(notices_[0].title, "CI failed on acme/rocket#7");
2102 // On a branch: whoever pushed.
2103 let pushed = event(
2104 "workflow.completed",
2105 Some("usr_bo"),
Inbox: threads, reasons, subscriptions and watching2106 json!({ "workflow": "Deploy", "path": ".g1t/workflows/deploy.yml", "conclusion": "failure", "ref": "refs/heads/main", "number": 12, "sha": "abcdef0123", "runId": "run_9" }),
Inbox: the events service tells people what needs them as events arrive2107 );
Inbox: threads, reasons, subscriptions and watching2108 let notices_ = notices(&pushed, "acme/rocket", &actor("usr_bo", "bo"), None, &nobody());
2109 assert_eq!(told(&notices_), vec![("bo", Reason::CiActivity, Severity::Error)]);
Inbox: the events service tells people what needs them as events arrive2110 assert_eq!(notices_[0].title, "Deploy failed on main in acme/rocket");
2111 assert_eq!(notices_[0].body, "Run 12 at abcdef0");
2112 // Nobody to tell when g1t pushed.
Inbox: threads, reasons, subscriptions and watching2113 assert!(notices(&pushed, "acme/rocket", &actor("g1t", "g1t"), None, &nobody()).is_empty());
2114 // Every failure of a workflow on a branch is one thread.
2115 let wanted = wants(&pushed).unwrap();
2116 let thread = thread_of(&pushed, &wanted, None);
2117 assert_eq!(thread.key, "rep_1/run/.g1t/workflows/deploy.yml@main");
2118 assert_eq!(thread.run_id.as_deref(), Some("run_9"));
2119 }
2120
2121 #[test]
2122 fn deployments_tell_whoever_answers_for_them_and_watchers() {
2123 let failed = event(
2124 "deployment.failed",
2125 Some("usr_bo"),
2126 json!({ "repoId": "rep_1", "projectId": "prj_1", "project": "rocket", "deploymentId": "dpl_1", "error": "The build failed.", "path": "/acme/rocket/deployments/dpl_1" }),
2127 );
2128 let audience = watched(&[("cy", WatchLevel::Custom, &["deployments"]), ("dee", WatchLevel::Custom, &["issues"]), ("eve", WatchLevel::All, &[])]);
2129 let notices_ = notices(&failed, "acme/rocket", &actor("usr_bo", "bo"), None, &audience);
2130 assert_eq!(
2131 told(&notices_),
2132 vec![
2133 ("bo", Reason::CiActivity, Severity::Error),
2134 ("cy", Reason::Subscribed, Severity::Error),
2135 ("eve", Reason::Subscribed, Severity::Error),
2136 ]
2137 );
2138 assert_eq!(notices_[0].title, "Production of rocket failed to deploy");
2139 assert_eq!(notices_[0].body, "The build failed.");
2140 let thread = thread_of(&failed, &wants(&failed).unwrap(), None);
2141 assert_eq!(thread.key, "rep_1/deploy/prj_1/production");
2142 assert_eq!(url(Some("acme/rocket"), thread.kind, None, None, thread.link.as_deref()), "/acme/rocket/deployments/dpl_1");
2143 // A success is news to the owner only after a failure.
2144 let live = event("deployment.succeeded", Some("usr_bo"), json!({ "repoId": "rep_1", "projectId": "prj_1", "project": "rocket", "commit": "abcdef0123" }));
2145 assert!(notices(&live, "acme/rocket", &actor("usr_bo", "bo"), None, &nobody()).is_empty());
2146 let recovered = event("deployment.succeeded", Some("usr_bo"), json!({ "repoId": "rep_1", "projectId": "prj_1", "project": "rocket", "recovered": true }));
2147 let notices_ = notices(&recovered, "acme/rocket", &actor("usr_bo", "bo"), None, &nobody());
2148 assert_eq!(told(&notices_), vec![("bo", Reason::CiActivity, Severity::Success)]);
2149 assert_eq!(notices_[0].title, "Production of rocket is live again");
2150 // A preview's is the pull request's owner's.
2151 let preview = event("deployment.failed", None, json!({ "repoId": "rep_1", "projectId": "prj_1", "number": 7, "triggeredBy": "g1t" }));
2152 assert_eq!(told(&notices(&preview, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())), vec![("ana", Reason::CiActivity, Severity::Error)]);
2153 assert_eq!(thread_of(&preview, &wants(&preview).unwrap(), Some(&pull())).key, "rep_1/deploy/prj_1/7");
Inbox: the events service tells people what needs them as events arrive2154 }
2155
2156 #[test]
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar2157 fn security_alerts_tell_the_pusher_the_owners_and_member_watchers() {
2158 let blocked = event(
2159 "secret_scanning_alert.created",
2160 None,
2161 json!({
2162 "repoId": "rep_1", "alertId": "sec_1", "alertType": "secret_scanning", "severity": "critical",
2163 "title": "An AWS access key in config/prod.env", "link": "/acme/rocket/security/secret-scanning/sec_1",
2164 "state": "open", "pusher": "bo", "notify": ["ana"], "members": ["ana", "bo", "cy"]
2165 }),
2166 );
2167 assert_eq!(wants(&blocked), Some(Wanted { repo_id: "rep_1".into(), number: None, comment_id: None }));
2168 // cy watches for security alerts and is a member; eve watches
2169 // everything but cannot see findings; dee watches issues only.
2170 let audience = watched(&[
2171 ("cy", WatchLevel::Custom, &["security"]),
2172 ("dee", WatchLevel::Custom, &["issues"]),
2173 ("eve", WatchLevel::All, &[]),
2174 ]);
2175 let notices_ = notices(&blocked, "acme/rocket", &Actor::default(), None, &audience);
2176 assert_eq!(
2177 told(&notices_),
2178 vec![
2179 ("bo", Reason::SecurityAlert, Severity::Error),
2180 ("ana", Reason::SecurityAlert, Severity::Error),
2181 ("cy", Reason::SecurityAlert, Severity::Error),
2182 ]
2183 );
2184 assert_eq!(notices_[0].title, "A push to acme/rocket was blocked: it adds a secret");
2185 assert_eq!(notices_[0].body, "An AWS access key in config/prod.env");
2186 let thread = thread_of(&blocked, &wants(&blocked).unwrap(), None);
2187 assert_eq!(thread.key, "rep_1/security/sec_1");
2188 assert_eq!(url(Some("acme/rocket"), thread.kind, None, None, thread.link.as_deref()), "/acme/rocket/security/secret-scanning/sec_1");
2189 // A bypass request goes to the reviewers named, never to watchers.
2190 let requested = event(
2191 "secret_scanning.bypass_requested",
2192 Some("usr_bo"),
2193 json!({ "repoId": "rep_1", "alertId": "sec_1", "requestId": "byp_1", "title": "An AWS access key in a.env", "notify": ["ana"], "link": "/acme/-/security/bypass-requests" }),
2194 );
2195 let notices_ = notices(&requested, "acme/rocket", &actor("usr_bo", "bo"), None, &audience);
2196 assert_eq!(told(&notices_), vec![("ana", Reason::SecurityAlert, Severity::Warning)]);
2197 assert_eq!(notices_[0].title, "bo asked to bypass push protection in acme/rocket");
2198 // Fixes and dismissals are not news to the inbox.
2199 assert_eq!(wants(&event("code_scanning_alert.fixed", None, json!({ "repoId": "rep_1" }))), None);
2200 }
2201
2202 #[test]
Inbox: the events service tells people what needs them as events arrive2203 fn nobody_hears_of_what_they_did_themselves() {
2204 let merged = event("pull.merged", Some("usr_ana"), json!({ "number": 7 }));
Inbox: threads, reasons, subscriptions and watching2205 let by_ana = notices(&merged, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
2206 // The assignee still hears; ana, who merged, does not.
2207 assert_eq!(told(&by_ana), vec![("bo", Reason::StateChange, Severity::Success)]);
2208 let by_bo = notices(&merged, "acme/rocket", &actor("usr_bo", "bo"), Some(&pull()), &nobody());
2209 assert_eq!(told(&by_bo), vec![("ana", Reason::StateChange, Severity::Success)]);
Inbox: the events service tells people what needs them as events arrive2210 assert_eq!(by_bo[0].title, "bo merged acme/rocket#7");
Inbox: threads, reasons, subscriptions and watching2211 let by_queue = notices(&merged, "acme/rocket", &Actor::default(), Some(&pull()), &nobody());
Inbox: the events service tells people what needs them as events arrive2212 assert_eq!(by_queue[0].title, "acme/rocket#7 was merged");
2213 }
2214
2215 #[test]
Inbox: threads, reasons, subscriptions and watching2216 fn closing_and_reopening_tells_everyone_subscribed_and_watchers() {
2217 let closed = event("issue.closed", Some("usr_bo"), json!({ "number": 3, "reason": "completed" }));
2218 let issue = InboxSubject {
2219 kind: Some(SubjectKind::Issue),
2220 title: "Crash".into(),
2221 author: person("usr_cy", "cy"),
2222 assignees: vec!["bo".into()],
2223 ..InboxSubject::default()
2224 };
2225 let mut audience = subscribed(&[("eve", State::Subscribed, Some(Reason::Comment)), ("fay", State::Unsubscribed, None)]);
2226 audience.watchers = watched(&[("gus", WatchLevel::All, &[]), ("hal", WatchLevel::Custom, &["pulls"])]).watchers;
2227 let notices_ = notices(&closed, "acme/rocket", &actor("usr_bo", "bo"), Some(&issue), &audience);
2228 assert_eq!(
2229 told(&notices_),
2230 vec![
2231 ("cy", Reason::StateChange, Severity::Info),
2232 ("eve", Reason::StateChange, Severity::Info),
2233 ("gus", Reason::Subscribed, Severity::Info),
2234 ]
2235 );
2236 assert_eq!(notices_[0].title, "bo closed acme/rocket#3");
2237 let by_pull = event("issue.closed", None, json!({ "number": 3, "resolvedBy": 9 }));
2238 assert_eq!(notices(&by_pull, "acme/rocket", &Actor::default(), Some(&issue), &nobody())[0].title, "acme/rocket#3 was closed by #9");
2239 let reopened = event("issue.reopened", Some("usr_cy"), json!({ "number": 3 }));
2240 assert_eq!(told(&notices(&reopened, "acme/rocket", &actor("usr_cy", "cy"), Some(&issue), &nobody())), vec![("bo", Reason::StateChange, Severity::Info)]);
2241 }
2242
2243 #[test]
2244 fn opening_tells_who_it_names_and_watchers_of_its_kind() {
2245 let opened = event("pull.opened", Some("usr_ana"), json!({ "number": 7 }));
2246 let subject = InboxSubject { reviewers: vec!["cy".into(), "g1t".into()], ..pull() };
2247 let audience = watched(&[("bo", WatchLevel::All, &[]), ("dee", WatchLevel::Custom, &["issues"]), ("eve", WatchLevel::Custom, &["pulls"]), ("fay", WatchLevel::Participating, &[])]);
2248 let notices_ = notices(&opened, "acme/rocket", &actor("usr_ana", "ana"), Some(&subject), &audience);
2249 assert_eq!(
2250 told(&notices_),
2251 vec![
2252 ("bo", Reason::Assign, Severity::Info),
2253 ("cy", Reason::ReviewRequested, Severity::Warning),
2254 ("eve", Reason::Subscribed, Severity::Info),
2255 ]
2256 );
2257 assert_eq!(notices_[2].title, "ana opened acme/rocket#7");
2258 }
2259
2260 #[test]
2261 fn ignoring_a_repository_silences_it_and_unsubscribing_keeps_only_what_is_asked() {
2262 let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
2263 let on = comment(person("usr_bo", "bo"), "@cy have a look", &["cy"]);
2264 let audience = watched(&[("cy", WatchLevel::Ignore, &[]), ("ana", WatchLevel::All, &[])]);
2265 assert!(notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &audience).iter().all(|notice| notice.username != "cy"));
2266 let audience = subscribed(&[("cy", State::Unsubscribed, None), ("ana", State::Unsubscribed, None)]);
2267 let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &audience);
2268 // cy was mentioned, which is asked of them; ana unsubscribed from the conversation.
2269 assert_eq!(told(&notices_), vec![("cy", Reason::Mention, Severity::Info)]);
2270 }
2271
2272 #[test]
Inbox: the events service tells people what needs them as events arrive2273 fn g1t_finishing_or_reviewing_tells_the_person_it_worked_for() {
2274 // The agent acts as the person it works for: still an outcome they hear of.
2275 let ready = event("pull.ready", Some("usr_ana"), json!({ "number": 7 }));
Inbox: threads, reasons, subscriptions and watching2276 let notices_ = notices(&ready, "acme/rocket", &actor("usr_ana", "ana"), Some(&g1t_pull()), &nobody());
2277 assert_eq!(told(&notices_), vec![("ana", Reason::Author, Severity::Success)]);
Inbox: the events service tells people what needs them as events arrive2278 assert_eq!(notices_[0].title, "g1t finished acme/rocket#7");
2279 // A person's draft marked ready is not news to them.
Inbox: threads, reasons, subscriptions and watching2280 assert!(notices(&ready, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody()).is_empty());
Inbox: the events service tells people what needs them as events arrive2281
2282 let approve = event("review.completed", None, json!({ "number": 7, "verdict": "approve" }));
Inbox: threads, reasons, subscriptions and watching2283 assert_eq!(
2284 told(&notices(&approve, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())),
2285 vec![("ana", Reason::Author, Severity::Success)]
2286 );
Inbox: the events service tells people what needs them as events arrive2287 let changes = event("review.completed", None, json!({ "number": 7, "verdict": "request_changes" }));
2288 assert_eq!(
Inbox: threads, reasons, subscriptions and watching2289 told(&notices(&changes, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())),
2290 vec![("ana", Reason::Author, Severity::Info)]
Inbox: the events service tells people what needs them as events arrive2291 );
2292 }
2293
2294 #[test]
Inbox: threads, reasons, subscriptions and watching2295 fn comments_tell_those_mentioned_then_everyone_subscribed_never_the_writer() {
Inbox: the events service tells people what needs them as events arrive2296 let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
2297 let on = comment(person("usr_bo", "bo"), "@cy @bo have a look", &["cy", "bo", "g1t"]);
Inbox: threads, reasons, subscriptions and watching2298 let audience = subscribed(&[("eve", State::Subscribed, Some(Reason::Comment)), ("fay", State::Subscribed, Some(Reason::Manual))]);
2299 let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &audience);
Inbox: the events service tells people what needs them as events arrive2300 assert_eq!(
2301 told(&notices_),
Inbox: threads, reasons, subscriptions and watching2302 vec![
2303 ("cy", Reason::Mention, Severity::Info),
2304 ("ana", Reason::Author, Severity::Info),
2305 ("eve", Reason::Comment, Severity::Info),
2306 ("fay", Reason::Manual, Severity::Info),
2307 ]
Inbox: the events service tells people what needs them as events arrive2308 );
2309 assert_eq!(notices_[0].title, "bo mentioned you on acme/rocket#7");
Inbox: threads, reasons, subscriptions and watching2310 assert_eq!(notices_[1].title, "bo commented on acme/rocket#7");
Inbox: the events service tells people what needs them as events arrive2311 assert_eq!(notices_[0].body, "@cy @bo have a look");
2312 // Mentioned and the owner: told once, as mentioned.
2313 let on = comment(person("usr_bo", "bo"), "@ana", &["ana"]);
Inbox: threads, reasons, subscriptions and watching2314 assert_eq!(
2315 told(&notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &nobody())),
2316 vec![("ana", Reason::Mention, Severity::Info)]
2317 );
2318 // The owner's own comment tells the assignee, not the owner.
Inbox: the events service tells people what needs them as events arrive2319 let on = comment(person("usr_ana", "ana"), "thanks", &[]);
Inbox: threads, reasons, subscriptions and watching2320 assert_eq!(
2321 told(&notices(&created, "acme/rocket", &actor("usr_ana", "ana"), Some(&on), &nobody())),
2322 vec![("bo", Reason::Assign, Severity::Info)]
2323 );
Inbox: the events service tells people what needs them as events arrive2324 // Something that happened, not something written, tells nobody.
2325 let mut on = comment(person("usr_bo", "bo"), "assigned cy", &[]);
2326 on.comment.as_mut().unwrap().event = true;
Inbox: threads, reasons, subscriptions and watching2327 assert!(notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &nobody()).is_empty());
Inbox: the events service tells people what needs them as events arrive2328 // An approval is good news.
Inbox: threads, reasons, subscriptions and watching2329 let mut on = comment(person("usr_cy", "cy"), "", &[]);
Inbox: the events service tells people what needs them as events arrive2330 on.comment.as_mut().unwrap().verdict = Some("approve".into());
Inbox: threads, reasons, subscriptions and watching2331 let notices_ = notices(&created, "acme/rocket", &actor("usr_cy", "cy"), Some(&on), &nobody());
2332 assert_eq!(told(&notices_)[0], ("ana", Reason::Author, Severity::Success));
Inbox: the events service tells people what needs them as events arrive2333 assert_eq!(notices_[0].body, "Add the inbox");
Inbox: threads, reasons, subscriptions and watching2334 // Writing and being mentioned subscribe, never g1t.
2335 let on = comment(person("usr_bo", "bo"), "@cy", &["cy"]);
2336 assert_eq!(subscribes(&created, Some(&on)), vec![("bo".into(), Reason::Comment), ("cy".into(), Reason::Mention)]);
2337 let on = comment(person(AGENT_ID, "g1t"), "done", &[]);
2338 assert!(subscribes(&created, Some(&on)).is_empty());
2339 }
2340
2341 #[test]
2342 fn issues_and_pull_requests_are_one_thread_each() {
2343 let created = event("comment.created", Some("usr_bo"), json!({ "repoId": "rep_1", "number": 7, "commentId": "cmt_1" }));
2344 let wanted = wants(&created).unwrap();
2345 let thread = thread_of(&created, &wanted, Some(&pull()));
2346 assert_eq!(thread.key, "rep_1#7");
2347 assert_eq!(thread.kind, Some(SubjectKind::Pull));
2348 let merged = event("pull.merged", None, json!({ "repoId": "rep_1", "number": 7 }));
2349 assert_eq!(thread_of(&merged, &wants(&merged).unwrap(), Some(&pull())).key, thread.key);
Inbox: the events service tells people what needs them as events arrive2350 }
2351
2352 #[test]
2353 fn items_link_to_what_they_are_about() {
Inbox: threads, reasons, subscriptions and watching2354 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Pull), Some(7), None, None), "/acme/rocket/pull/7");
2355 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Issue), Some(3), None, None), "/acme/rocket/issues/3");
2356 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Run), None, Some("run_1"), None), "/acme/rocket/actions/runs/run_1");
2357 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Deploy), None, None, Some("/acme/site/deployments/dpl_1")), "/acme/site/deployments/dpl_1");
2358 assert_eq!(url(Some("acme/rocket"), None, None, None, None), "/acme/rocket");
2359 assert_eq!(url(None, None, None, None, None), "/inbox");
Inbox: the events service tells people what needs them as events arrive2360 }
2361
2362 #[test]
2363 fn long_titles_are_cut_to_a_line() {
2364 let long = "x".repeat(400);
2365 assert_eq!(clip(&long, MAX_TITLE).chars().count(), MAX_TITLE);
2366 assert_eq!(clip(" short ", MAX_TITLE), "short");
2367 }
2368
2369 #[test]
2370 fn marks_change_only_what_they_say() {
2371 assert_eq!(mark_change(InboxMark::Read), "read_at = COALESCE(read_at, ?1)");
2372 assert!(mark_change(InboxMark::Done).contains("done_at = ?1"));
Inbox: threads, reasons, subscriptions and watching2373 assert_eq!(mark_change(InboxMark::Unsnooze), "snoozed_until = NULL");
2374 assert!(ranked(&ListInboxArgs::default()));
2375 assert!(!ranked(&ListInboxArgs { reason: Some(Reason::Mention), ..ListInboxArgs::default() }));
2376 assert!(!ranked(&ListInboxArgs { participating: true, ..ListInboxArgs::default() }));
2377 }
2378
2379 #[test]
2380 fn filters_bind_in_order_after_what_is_bound() {
2381 let a = ListInboxArgs {
2382 reason: Some(Reason::Mention),
2383 repo_id: Some("rep_1".into()),
2384 participating: true,
2385 unread: true,
2386 ..ListInboxArgs::default()
2387 };
2388 // Two values come first: the username and now.
2389 let (conditions, values) = filters(&a, 2);
2390 assert_eq!(
2391 conditions,
2392 vec!["reason = ?3", "repo_id = ?4", "reason NOT IN ('manual', 'subscribed')", "read_at IS NULL"]
2393 );
2394 assert_eq!(values, vec!["mention", "rep_1"]);
2395 }
2396
2397 #[test]
2398 fn the_bump_keeps_whatever_is_most_urgent_while_unread() {
2399 assert!(BUMP.contains("WHERE NOT EXISTS (SELECT 1 FROM inbox_activity WHERE event_id = ?4 AND username = ?2)"));
2400 assert!(BUMP.contains("ON CONFLICT (username, thread) DO UPDATE"));
2401 assert!(BUMP.contains("activity = inbox_items.activity + 1"));
2402 // Its numbered parameters run from 1 to 19.
2403 for at in 1..=19 {
2404 assert!(BUMP.contains(&format!("?{at}")), "?{at}");
2405 }
2406 assert!(!BUMP.contains("?20"));
Inbox: the events service tells people what needs them as events arrive2407 }
Fine-grained personal tokens, workspace token rules and approvals in identity2408
2409 #[test]
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)2410 fn live_notifications_say_what_and_where() {
2411 let notice = Notice {
2412 username: "ana".into(),
2413 reason: Reason::Agent,
2414 severity: Severity::Warning,
2415 title: "g1t needs you on acme/api#4".into(),
2416 body: "Which database?".into(),
2417 };
2418 let live = live_notification("evt_1", &notice, "/Acme/api/pull/4#c1", "2026-10-08T00:00:00Z");
2419 assert_eq!(live["id"], "inbox:evt_1");
2420 assert_eq!(live["kind"], "agent_waiting");
2421 assert_eq!(live["workspace"], "acme");
2422 assert_eq!(live["href"], "/Acme/api/pull/4#c1");
2423 assert_eq!(live["actor"]["name"], "acme/api");
2424 let mention = Notice { reason: Reason::TeamMention, ..notice };
2425 assert_eq!(live_notification("evt_2", &mention, "https://elsewhere", "t")["href"], "/inbox");
2426 assert_eq!(live_notification("evt_2", &mention, "/x", "t")["kind"], "mention");
2427 }
2428
2429 #[test]
2430 fn inbox_counts_are_told_as_they_are_by_username() {
2431 assert_eq!(inbox_count_update("Ana", 3), serde_json::json!({ "username": "ana", "unread": 3 }));
2432 assert_eq!(inbox_count_update("bo", 0)["unread"], 0);
2433 }
2434
2435 #[test]
Fine-grained personal tokens, workspace token rules and approvals in identity2436 fn token_approvals_tell_the_people_named_about_no_repository() {
2437 let mut asked = event(
2438 "token.approval_requested",
2439 Some("usr_ana"),
2440 json!({ "workspace": "acme", "tokenId": "tok_1", "notify": ["Bo", "cy", "bo", "g1t"], "title": "ana asks", "body": "ci: contents: write", "link": "/acme/-/settings/tokens" }),
2441 );
2442 asked.repo_id = None;
2443 let (thread, told) = workspace_notices(&asked, None);
2444 assert_eq!(thread, "workspace:acme/token/tok_1");
2445 let mut names: Vec<&str> = told.iter().map(|notice| notice.username.as_str()).collect();
2446 names.sort();
2447 assert_eq!(names, ["bo", "cy"], "each once, lowercased, never g1t");
2448 assert!(told.iter().all(|notice| notice.reason == Reason::ReviewRequested && notice.severity == Severity::Warning));
2449 assert!(wants(&asked).is_none(), "about no repository");
2450 let reviewed = event("token.approval_reviewed", None, json!({ "workspace": "acme", "tokenId": "tok_1", "notify": ["ana"], "title": "approved" }));
2451 let (_, told) = workspace_notices(&reviewed, None);
2452 assert_eq!(told[0].reason, Reason::Author);
2453 assert!(WORKSPACE_EVENTS.contains(&"token.approval_reviewed"));
2454 }
Merge workspace invitations: nobody joins a workspace without saying yes, people are found by username, your own invites can bring someone in, and nobody is left without a workspace (identity 0040)2455
2456 #[test]
2457 fn a_workspace_invitation_asks_its_person_and_tells_whoever_sent_it_the_answer() {
2458 let mut sent = event(
2459 "workspace_invitation.created",
2460 Some("usr_owner"),
2461 json!({ "workspace": "flagon-io", "invitationId": "inv_1", "thread": "invitation:inv_1", "notify": ["daweazl"], "title": "@syntaqx invited you to join Flagon, Inc.", "link": "/invitations" }),
2462 );
2463 sent.repo_id = None;
2464 let (thread, told) = workspace_notices(&sent, None);
2465 assert_eq!(thread, "invitation:inv_1");
2466 assert_eq!(told.len(), 1);
2467 assert_eq!((told[0].username.as_str(), told[0].reason, told[0].severity), ("daweazl", Reason::ReviewRequested, Severity::Warning));
2468 let accepted = event("workspace_invitation.accepted", None, json!({ "workspace": "flagon-io", "thread": "invitation:inv_1", "notify": ["syntaqx"] }));
2469 let (thread, told) = workspace_notices(&accepted, None);
2470 assert_eq!(thread, "invitation:inv_1");
2471 assert_eq!((told[0].reason, told[0].severity), (Reason::Author, Severity::Success));
2472 let declined = event("workspace_invitation.declined", None, json!({ "workspace": "flagon-io", "thread": "invitation:inv_1", "notify": ["syntaqx"] }));
2473 assert_eq!(workspace_notices(&declined, None).1[0].severity, Severity::Info);
2474 // Answered or revoked, the person's item is done; a revocation tells nobody.
2475 for kind in ["workspace_invitation.accepted", "workspace_invitation.declined", "workspace_invitation.revoked"] {
2476 assert!(INVITATION_ANSWERED.contains(&kind));
2477 }
2478 assert!(!WORKSPACE_EVENTS.contains(&"workspace_invitation.revoked"));
2479 }
Inbox: the events service tells people what needs them as events arrive2480}

This file's history is long; its oldest lines are credited to the oldest commit read.