Skip to content
2,362 linesCodeBlameRaw
1//! The inbox, kept beside the event log. See `g1t_contracts::inbox`.
2//!
3//! 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:
9//!
10//! | Event | Who | Reason | Severity |
11//! | --- | --- | --- | --- |
12//! | `agent.asked`, `pull.stalled` | the pull request's owner, and its issue's owner and assignees | agent | warning |
13//! | `pull.review_requested` | the reviewers asked, and the people each team asked tells | review_requested | warning |
14//! | `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 |
22//! | `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 |
24//!
25//! 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.
36
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;
49use worker::{D1Database, D1PreparedStatement, Fetcher, Result};
50
51use crate::subscriptions::{self, Audience};
52
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,
75 pub reason: Reason,
76 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 }
96
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 }
101}
102
103pub(crate) fn is_g1t(username: &str) -> bool {
104 username.eq_ignore_ascii_case(system::USERNAME) || username.eq_ignore_ascii_case("g1t-agent")
105}
106
107pub(crate) fn is_g1t_id(id: &str) -> bool {
108 system::is_system_id(id) || id == AGENT_ID
109}
110
111/// 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
115/// 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() {
130 "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),
133 "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 },
145 "deployment.failed" | "deployment.succeeded" => on(number("number"), None),
146 // A workflow run's jobs wait for their environment's reviewers.
147 "deployment.review_requested" => on(None, None),
148 "comment.created" => on(Some(number("number")?), Some(text("commentId")?)),
149 // Security alerts are threads of their own, not of an issue.
150 kind if SECURITY_EVENTS.contains(&kind) => on(None, None),
151 // So is a repository's mirror, told to whom its settings name.
152 kind if MIRROR_EVENTS.contains(&kind) => on(None, None),
153 _ => None,
154 }
155}
156
157/// The security service's events the inbox tells people of: new alerts,
158/// and push protection bypasses asked for and decided. Fixes and
159/// dismissals are on the Security page and in webhooks.
160pub const SECURITY_EVENTS: [&str; 5] = [
161 "secret_scanning_alert.created",
162 "code_scanning_alert.created",
163 "vulnerability_alert.created",
164 "secret_scanning.bypass_requested",
165 "secret_scanning.bypass_reviewed",
166];
167
168/// The integrations service's mirroring events (see
169/// `g1t_contracts::mirrors::MirrorEvent`). People are told only when the
170/// link's settings ask for the inbox, and the event names them.
171pub const MIRROR_EVENTS: [&str; 4] = ["mirror.unreachable", "mirror.reachable", "mirror.state_changed", "mirror.moved_in"];
172
173/// Where an event's items go: which thread, what it is, and where it is.
174#[derive(Clone, Debug, PartialEq, Eq)]
175pub struct Thread {
176 pub key: String,
177 pub kind: Option<SubjectKind>,
178 pub number: Option<u32>,
179 pub run_id: Option<String>,
180 /// A path, for a subject with a page of its own.
181 pub link: Option<String>,
182}
183
184/// The thread key of an issue or pull request.
185pub fn numbered_thread(repo_id: &str, number: u32) -> String {
186 format!("{repo_id}#{number}")
187}
188
189/// The thread an event's items go to.
190pub fn thread_of(event: &Event, wanted: &Wanted, subject: Option<&InboxSubject>) -> Thread {
191 let data = &event.data;
192 let text = |key: &str| data[key].as_str().filter(|value| !value.is_empty()).map(str::to_owned);
193 match event.kind.as_str() {
194 // A project's production, or one pull request's preview.
195 "deployment.failed" | "deployment.succeeded" => {
196 let which = wanted.number.map_or_else(|| "production".to_owned(), |number| number.to_string());
197 Thread {
198 key: format!("{}/deploy/{}/{which}", wanted.repo_id, text("projectId").unwrap_or_default()),
199 kind: Some(SubjectKind::Deploy),
200 number: wanted.number,
201 run_id: text("deploymentId"),
202 link: text("path"),
203 }
204 }
205 // One run's deployment to one environment, waiting for review.
206 "deployment.review_requested" => Thread {
207 key: format!(
208 "{}/review/{}/{}",
209 wanted.repo_id,
210 text("runId").unwrap_or_default(),
211 text("environment").unwrap_or_default()
212 ),
213 kind: Some(SubjectKind::Run),
214 number: None,
215 run_id: text("runId"),
216 link: text("link"),
217 },
218 // A workflow on a branch: its next failure bumps the same thread.
219 "workflow.completed" if subject.is_none() => {
220 let branch = text("ref").unwrap_or_default();
221 let branch = branch.strip_prefix("refs/heads/").unwrap_or(&branch).to_owned();
222 let workflow = text("path").or_else(|| text("workflow")).unwrap_or_default();
223 Thread {
224 key: format!("{}/run/{workflow}@{branch}", wanted.repo_id),
225 kind: Some(SubjectKind::Run),
226 number: None,
227 run_id: text("runId"),
228 link: None,
229 }
230 }
231 // A repository's mirror: one thread, its settings page the link.
232 kind if MIRROR_EVENTS.contains(&kind) => Thread {
233 key: format!("{}/mirror", wanted.repo_id),
234 kind: None,
235 number: None,
236 run_id: None,
237 link: text("link"),
238 },
239 // One alert, or one bypass request: its page is the link.
240 kind if SECURITY_EVENTS.contains(&kind) => {
241 let which = text("requestId").or_else(|| text("alertId")).unwrap_or_else(|| event.id.clone());
242 Thread {
243 key: format!("{}/security/{which}", wanted.repo_id),
244 kind: None,
245 number: None,
246 run_id: None,
247 link: text("link"),
248 }
249 }
250 _ => match wanted.number {
251 Some(number) => Thread {
252 key: numbered_thread(&wanted.repo_id, number),
253 kind: subject.and_then(|subject| subject.kind),
254 number: Some(number),
255 run_id: None,
256 link: None,
257 },
258 None => Thread {
259 key: format!("event/{}", event.id),
260 kind: None,
261 number: None,
262 run_id: None,
263 link: None,
264 },
265 },
266 }
267}
268
269/// The issue or pull request whose agent threads an event closes: what an
270/// agent was waiting on a person for is over.
271pub fn resolves(event: &Event) -> Option<String> {
272 if !RESUMES.contains(&event.kind.as_str()) {
273 return None;
274 }
275 let repo_id = event.data["repoId"].as_str().map(str::to_owned).or_else(|| event.repo_id.clone())?;
276 let number = event.data["number"].as_u64().and_then(|n| u32::try_from(n).ok())?;
277 Some(numbered_thread(&repo_id, number))
278}
279
280/// Collects who is told, each once with the most specific reason, never
281/// the actor, never g1t, and never anyone ignoring the thread.
282struct Told<'a> {
283 actor: &'a Actor,
284 audience: &'a Audience,
285 notices: Vec<Notice>,
286}
287
288impl Told<'_> {
289 /// `asked`: something asked of the person directly, told even when
290 /// they unsubscribed from the thread.
291 fn tell(&mut self, username: &str, reason: Reason, severity: Severity, title: &str, body: &str, asked: bool) {
292 let username = username.trim().trim_start_matches('@').to_lowercase();
293 if username.is_empty()
294 || is_g1t(&username)
295 || self.actor.is_named(&username)
296 || self.audience.ignores(&username)
297 || (!asked && self.audience.unsubscribed(&username))
298 {
299 return;
300 }
301 let notice = Notice {
302 username,
303 reason,
304 severity,
305 title: clip(title, MAX_TITLE),
306 body: clip(body, MAX_BODY),
307 };
308 match self.notices.iter_mut().find(|told| told.username == notice.username) {
309 Some(told) if notice.reason.rank() < told.reason.rank() => *told = notice,
310 Some(_) => {}
311 None => self.notices.push(notice),
312 }
313 }
314
315 /// A person known by id as well as name, such as an author.
316 fn tell_person(&mut self, person: &Principal, reason: Reason, severity: Severity, title: &str, body: &str, asked: bool) {
317 if is_g1t_id(&person.id) || self.actor.is(person) {
318 return;
319 }
320 self.tell(&person.username, reason, severity, title, body, asked);
321 }
322
323 /// Everyone subscribed to an issue or pull request: its owner and
324 /// author, its assignees and reviewers, and whoever subscribed by
325 /// commenting, being mentioned or by hand. `reason` overrides why
326 /// each is told, as a state change does.
327 fn tell_subscribed(&mut self, subject: &InboxSubject, reason: Option<Reason>, severity: Severity, title: &str, body: &str) {
328 let why = |own: Reason| reason.unwrap_or(own);
329 self.tell_person(subject.owner(), why(Reason::Author), severity, title, body, false);
330 self.tell_person(&subject.author, why(Reason::Author), severity, title, body, false);
331 for name in &subject.assignees {
332 self.tell(name, why(Reason::Assign), severity, title, body, false);
333 }
334 for name in &subject.reviewers {
335 self.tell(name, why(Reason::ReviewRequested), severity, title, body, false);
336 }
337 let subscribed: Vec<(String, Reason)> = self.audience.subscribed().collect();
338 for (name, own) in subscribed {
339 self.tell(&name, why(own), severity, title, body, false);
340 }
341 }
342
343 /// Everyone watching the repository for this kind of activity.
344 fn tell_watchers(&mut self, kind: &str, severity: Severity, title: &str, body: &str) {
345 let watching: Vec<String> = self.audience.watching(kind).collect();
346 for name in watching {
347 self.tell(&name, Reason::Subscribed, severity, title, body, false);
348 }
349 }
350}
351
352fn clip(text: &str, max: usize) -> String {
353 let text = text.trim();
354 if text.chars().count() <= max {
355 return text.to_owned();
356 }
357 let cut: String = text.chars().take(max - 1).collect();
358 format!("{}…", cut.trim_end())
359}
360
361/// The kind of activity a watcher chooses: `issues` or `pulls`.
362fn activity_of(subject: &InboxSubject) -> &'static str {
363 match subject.kind {
364 Some(SubjectKind::Pull) => "pulls",
365 _ => "issues",
366 }
367}
368
369/// The names in an event's list field, such as the reviewers just asked.
370/// The teams an event asked to review: each `workspace/slug` with the
371/// people it tells.
372fn teams_asked(data: &serde_json::Value) -> Vec<(String, Vec<String>)> {
373 data["teams"]
374 .as_array()
375 .map(|teams| {
376 teams
377 .iter()
378 .filter_map(|team| Some((team["team"].as_str()?.to_owned(), names(team, "notified"))))
379 .collect()
380 })
381 .unwrap_or_default()
382}
383
384fn names(data: &serde_json::Value, key: &str) -> Vec<String> {
385 data[key]
386 .as_array()
387 .map(|names| names.iter().filter_map(|name| name.as_str().map(str::to_owned)).collect())
388 .unwrap_or_default()
389}
390
391/// Who is told of `event`, in `repo` (`owner/name`), given what it names
392/// and who follows it. `actor` is who caused it; outcomes nobody chose
393/// (checks, workflows, deployments, a review, an agent finishing or
394/// stopping) are told whoever caused them.
395pub fn notices(event: &Event, repo: &str, actor: &Actor, subject: Option<&InboxSubject>, audience: &Audience) -> Vec<Notice> {
396 let nobody = Actor::default();
397 let data = &event.data;
398 // The issue or pull request, as titles name it: `acme/rocket#12`.
399 let at = match data["number"].as_u64().or(data["pull"].as_u64()) {
400 Some(number) if subject.is_some() => format!("{repo}#{number}"),
401 _ => repo.to_owned(),
402 };
403 let outcome = matches!(
404 event.kind.as_str(),
405 "checks.completed"
406 | "workflow.completed"
407 | "review.completed"
408 | "pull.ready"
409 | "agent.asked"
410 | "pull.stalled"
411 | "deployment.failed"
412 | "deployment.succeeded"
413 | "deployment.review_requested"
414 ) || SECURITY_EVENTS.contains(&event.kind.as_str())
415 || MIRROR_EVENTS.contains(&event.kind.as_str());
416 let mut told = Told {
417 actor: if outcome { &nobody } else { actor },
418 audience,
419 notices: Vec::new(),
420 };
421 // Who did it, as titles name them.
422 let who = actor.name().unwrap_or("g1t");
423
424 match (event.kind.as_str(), subject) {
425 ("agent.asked" | "pull.stalled", Some(pull)) => {
426 let (title, body) = if event.kind == "agent.asked" {
427 (format!("An agent is waiting on {at}"), pull.title.clone())
428 } else {
429 let detail = data["detail"].as_str().map(str::trim).filter(|detail| !detail.is_empty());
430 (format!("g1t stopped on {at} and needs you"), detail.unwrap_or(&pull.title).to_owned())
431 };
432 told.tell_person(pull.owner(), Reason::Agent, Severity::Warning, &title, &body, true);
433 if let Some(issue) = &pull.issue {
434 told.tell_person(issue.owner(), Reason::Agent, Severity::Warning, &title, &body, true);
435 for name in &issue.assignees {
436 told.tell(name, Reason::Agent, Severity::Warning, &title, &body, true);
437 }
438 }
439 }
440 ("pull.review_requested", Some(pull)) => {
441 // Asked because the CODEOWNERS file says they own what changed.
442 let owned = data["codeOwners"].as_bool() == Some(true);
443 let title = if owned {
444 format!("{at} changes files you own")
445 } else {
446 format!("{who} asked you to review {at}")
447 };
448 for name in names(data, "reviewers") {
449 told.tell(&name, Reason::ReviewRequested, Severity::Warning, &title, &pull.title, true);
450 }
451 for (team, people) in teams_asked(data) {
452 let title = if owned {
453 format!("{at} changes files @{team} owns")
454 } else {
455 format!("{who} asked @{team} to review {at}")
456 };
457 for name in people {
458 told.tell(&name, Reason::ReviewRequested, Severity::Warning, &title, &pull.title, true);
459 }
460 }
461 }
462 ("issue.assigned" | "pull.assigned", Some(on)) => {
463 let title = format!("{who} assigned you to {at}");
464 for name in names(data, "added") {
465 told.tell(&name, Reason::Assign, Severity::Info, &title, &on.title, true);
466 }
467 }
468 ("issue.opened" | "pull.opened", Some(on)) => {
469 let title = format!("{who} opened {at}");
470 for name in &on.assignees {
471 told.tell(name, Reason::Assign, Severity::Info, &format!("{who} assigned you to {at}"), &on.title, true);
472 }
473 for name in &on.reviewers {
474 told.tell(name, Reason::ReviewRequested, Severity::Warning, &format!("{who} asked you to review {at}"), &on.title, true);
475 }
476 for name in &on.mentions {
477 told.tell(name, Reason::Mention, Severity::Info, &format!("{who} mentioned you on {at}"), &on.title, true);
478 }
479 for team in &on.team_mentions {
480 let title = format!("{who} mentioned @{} on {at}", team.team);
481 for name in &team.members {
482 told.tell(name, Reason::TeamMention, Severity::Info, &title, &on.title, true);
483 }
484 }
485 told.tell_watchers(activity_of(on), Severity::Info, &title, &on.title);
486 }
487 ("checks.completed", Some(pull)) => {
488 let title = match data["status"].as_str() {
489 Some("errored") => format!("Checks could not run on {at}"),
490 _ => format!("Checks failed on {at}"),
491 };
492 told.tell_person(pull.owner(), Reason::CiActivity, Severity::Error, &title, &pull.title, true);
493 }
494 ("workflow.completed", subject) => {
495 let workflow = data["workflow"].as_str().filter(|name| !name.is_empty()).unwrap_or("A workflow");
496 match subject {
497 Some(pull) => {
498 let title = format!("{workflow} failed on {at}");
499 told.tell_person(pull.owner(), Reason::CiActivity, Severity::Error, &title, &pull.title, true);
500 }
501 None => {
502 // Not on a pull request: whoever pushed the commit it ran on.
503 let branch = data["ref"].as_str().unwrap_or_default();
504 let branch = branch.strip_prefix("refs/heads/").unwrap_or(branch);
505 let title = format!("{workflow} failed on {branch} in {repo}");
506 let body = format!("Run {} at {}", data["number"], short(data["sha"].as_str().unwrap_or_default()));
507 if let Some(id) = &actor.id
508 && !is_g1t_id(id)
509 && let Some(name) = &actor.username
510 {
511 told.tell(name, Reason::CiActivity, Severity::Error, &title, &body, true);
512 }
513 }
514 }
515 }
516 ("deployment.failed" | "deployment.succeeded", subject) => {
517 let failed = event.kind == "deployment.failed";
518 let project = data["project"].as_str().filter(|name| !name.is_empty()).unwrap_or(repo);
519 let what = match subject {
520 Some(_) => format!("The preview of {at}"),
521 None => format!("Production of {project}"),
522 };
523 let (title, severity) = if failed {
524 (format!("{what} failed to deploy"), Severity::Error)
525 } else {
526 (format!("{what} is live"), Severity::Success)
527 };
528 let body = match (failed, data["error"].as_str().map(str::trim).filter(|error| !error.is_empty())) {
529 (true, Some(error)) => error.to_owned(),
530 _ => match subject {
531 Some(pull) => pull.title.clone(),
532 None => format!("Commit {}", short(data["commit"].as_str().unwrap_or_default())),
533 },
534 };
535 // A success is news to whoever answers for it only after a failure.
536 if failed || data["recovered"].as_bool() == Some(true) {
537 let title = if failed { title.clone() } else { format!("{what} is live again") };
538 match subject {
539 Some(pull) => told.tell_person(pull.owner(), Reason::CiActivity, severity, &title, &body, true),
540 None => {
541 let pusher = actor
542 .id
543 .as_deref()
544 .filter(|id| !is_g1t_id(id))
545 .and(actor.username.as_deref())
546 .or_else(|| data["triggeredBy"].as_str());
547 if let Some(name) = pusher {
548 told.tell(name, Reason::CiActivity, severity, &title, &body, true);
549 }
550 }
551 }
552 }
553 told.tell_watchers("deployments", severity, &title, &body);
554 }
555 ("review.completed", Some(pull)) => {
556 let (title, severity) = match data["verdict"].as_str() {
557 Some("approve") => (format!("g1t approved {at}"), Severity::Success),
558 _ => (format!("g1t asked for changes on {at}"), Severity::Info),
559 };
560 told.tell_person(pull.owner(), Reason::Author, severity, &title, &pull.title, true);
561 }
562 ("pull.ready", Some(pull)) => {
563 // A change g1t made is ready: the agent's run is over.
564 if let Some(owner) = pull.requested_by.as_ref().filter(|_| is_g1t_id(&pull.author.id) || is_g1t(&pull.author.username)) {
565 let title = format!("g1t finished {at}");
566 told.tell_person(owner, Reason::Author, Severity::Success, &title, &pull.title, true);
567 }
568 }
569 ("pull.merged" | "pull.closed" | "issue.closed" | "issue.reopened", Some(on)) => {
570 let (verb, severity) = match event.kind.as_str() {
571 "pull.merged" => ("merged", Severity::Success),
572 "issue.reopened" => ("reopened", Severity::Info),
573 _ => ("closed", Severity::Info),
574 };
575 let title = match (actor.name(), data["resolvedBy"].as_u64()) {
576 (_, Some(pull)) if event.kind == "issue.closed" => format!("{at} was closed by #{pull}"),
577 (Some(name), _) => format!("{name} {verb} {at}"),
578 (None, _) => format!("{at} was {verb}"),
579 };
580 told.tell_subscribed(on, Some(Reason::StateChange), severity, &title, &on.title);
581 told.tell_watchers(activity_of(on), severity, &title, &on.title);
582 }
583 ("deployment.review_requested", _) => {
584 let environment = data["environment"].as_str().unwrap_or("an environment");
585 let workflow = data["workflow"].as_str().unwrap_or("A workflow");
586 let title = format!("{workflow} is waiting for your review to deploy to {environment} in {repo}");
587 let body = data["title"].as_str().unwrap_or_default();
588 // Those the actions service names: the environment's reviewers.
589 for name in names(data, "notify") {
590 told.tell(&name, Reason::ReviewRequested, Severity::Warning, &title, body, true);
591 }
592 }
593 (kind, None) if MIRROR_EVENTS.contains(&kind) => {
594 let severity = match kind {
595 "mirror.unreachable" => Severity::Warning,
596 "mirror.state_changed" if data["state"].as_str() == Some("takeover") => Severity::Warning,
597 "mirror.state_changed" if data["state"].as_str() == Some("standby") => Severity::Success,
598 _ => Severity::Info,
599 };
600 let title = data["title"].as_str().unwrap_or("A mirror changed");
601 let body = data["detail"].as_str().unwrap_or_default();
602 for name in names(data, "notify") {
603 told.tell(&name, Reason::StateChange, severity, title, body, true);
604 }
605 }
606 (kind, None) if SECURITY_EVENTS.contains(&kind) => {
607 let severity = match data["severity"].as_str() {
608 _ if kind == "secret_scanning.bypass_requested" => Severity::Warning,
609 _ if kind == "secret_scanning.bypass_reviewed" => Severity::Info,
610 Some("critical" | "high") => Severity::Error,
611 _ => Severity::Warning,
612 };
613 let what = data["title"].as_str().unwrap_or("A security alert");
614 let title = match kind {
615 "secret_scanning.bypass_requested" => format!("{who} asked to bypass push protection in {repo}"),
616 "secret_scanning.bypass_reviewed" => {
617 let verdict = data["state"].as_str().unwrap_or("reviewed");
618 format!("Your request to bypass push protection in {repo} was {verdict}")
619 }
620 _ if data["pusher"].is_string() => format!("A push to {repo} was blocked: it adds a secret"),
621 _ => format!("New security alert in {repo}"),
622 };
623 // Whoever pushed a blocked secret, and those the service names
624 // (the workspace's owners, or whoever asked for a bypass).
625 if let Some(pusher) = data["pusher"].as_str() {
626 told.tell(pusher, Reason::SecurityAlert, Severity::Error, &title, what, true);
627 }
628 for name in names(data, "notify") {
629 told.tell(&name, Reason::SecurityAlert, severity, &title, what, true);
630 }
631 // Watchers who chose security alerts, if they may see findings:
632 // the service lists who may (`members`), as findings are never
633 // shown to someone who cannot change the code.
634 if !kind.starts_with("secret_scanning.bypass") {
635 let members = names(data, "members");
636 let watching: Vec<String> = told.audience.watching("security").collect();
637 for name in watching.iter().filter(|name| members.iter().any(|member| member.eq_ignore_ascii_case(name))) {
638 told.tell(name, Reason::SecurityAlert, severity, &title, what, false);
639 }
640 }
641 }
642 ("comment.created", Some(on)) => {
643 let Some(comment) = on.comment.as_ref().filter(|comment| !comment.event) else {
644 return Vec::new();
645 };
646 // Whoever wrote it is the actor, whatever the event says.
647 let writer = Actor {
648 id: Some(comment.author.id.clone()),
649 username: Some(comment.author.username.clone()),
650 };
651 told.actor = &writer;
652 let who = &comment.author.username;
653 let body = if comment.excerpt.is_empty() { &on.title } else { &comment.excerpt };
654 for name in &comment.mentions {
655 let title = format!("{who} mentioned you on {at}");
656 told.tell(name, Reason::Mention, Severity::Info, &title, body, true);
657 }
658 for team in &comment.team_mentions {
659 let title = format!("{who} mentioned @{} on {at}", team.team);
660 for name in &team.members {
661 told.tell(name, Reason::TeamMention, Severity::Info, &title, body, true);
662 }
663 }
664 let (severity, title) = match comment.verdict.as_deref() {
665 Some("approve") => (Severity::Success, format!("{who} approved {at}")),
666 Some("request_changes") => (Severity::Info, format!("{who} asked for changes on {at}")),
667 _ => (Severity::Info, format!("{who} commented on {at}")),
668 };
669 // An approval or a request for changes is the owner's to hear
670 // of whatever they chose; the rest, as they subscribed.
671 if comment.verdict.is_some() {
672 told.tell_person(on.owner(), Reason::Author, severity, &title, body, true);
673 }
674 told.tell_subscribed(on, None, severity, &title, body);
675 told.tell_watchers(activity_of(on), severity, &title, body);
676 return told.notices;
677 }
678 _ => {}
679 }
680 told.notices
681}
682
683/// Who an event subscribes to its issue or pull request without asking,
684/// and why: whoever commented, whoever they mentioned, the people assigned
685/// and the reviewers asked. Never g1t.
686pub fn subscribes(event: &Event, subject: Option<&InboxSubject>) -> Vec<(String, Reason)> {
687 let mut people: Vec<(String, Reason)> = Vec::new();
688 let mut add = |name: &str, reason: Reason| {
689 let name = name.trim().trim_start_matches('@').to_lowercase();
690 if !name.is_empty() && !is_g1t(&name) && !people.iter().any(|(had, _)| *had == name) {
691 people.push((name, reason));
692 }
693 };
694 match event.kind.as_str() {
695 "comment.created" => {
696 if let Some(comment) = subject.and_then(|subject| subject.comment.as_ref()).filter(|comment| !comment.event) {
697 if !is_g1t_id(&comment.author.id) {
698 add(&comment.author.username, Reason::Comment);
699 }
700 for name in &comment.mentions {
701 add(name, Reason::Mention);
702 }
703 for team in &comment.team_mentions {
704 for name in &team.members {
705 add(name, Reason::TeamMention);
706 }
707 }
708 }
709 }
710 "issue.opened" | "pull.opened" => {
711 if let Some(on) = subject {
712 for name in &on.mentions {
713 add(name, Reason::Mention);
714 }
715 for team in &on.team_mentions {
716 for name in &team.members {
717 add(name, Reason::TeamMention);
718 }
719 }
720 }
721 }
722 "issue.assigned" | "pull.assigned" => {
723 for name in names(&event.data, "added") {
724 add(&name, Reason::Assign);
725 }
726 }
727 "pull.review_requested" => {
728 for name in names(&event.data, "reviewers") {
729 add(&name, Reason::ReviewRequested);
730 }
731 for (_, people) in teams_asked(&event.data) {
732 for name in people {
733 add(&name, Reason::ReviewRequested);
734 }
735 }
736 }
737 _ => {}
738 }
739 people
740}
741
742fn short(sha: &str) -> &str {
743 sha.get(..7).unwrap_or(sha)
744}
745
746/// Where an item is on g1t.sh.
747pub fn url(repo: Option<&str>, subject: Option<SubjectKind>, number: Option<u32>, run_id: Option<&str>, link: Option<&str>) -> String {
748 if let Some(link) = link.filter(|link| link.starts_with('/')) {
749 return link.to_owned();
750 }
751 let Some(repo) = repo else {
752 return "/inbox".to_owned();
753 };
754 match (subject, number, run_id) {
755 (Some(SubjectKind::Pull), Some(number), _) => format!("/{repo}/pull/{number}"),
756 (Some(SubjectKind::Issue), Some(number), _) => format!("/{repo}/issues/{number}"),
757 (Some(SubjectKind::Run), _, Some(run)) => format!("/{repo}/actions/runs/{run}"),
758 (Some(SubjectKind::Deploy), Some(number), _) => format!("/{repo}/pull/{number}"),
759 _ => format!("/{repo}"),
760 }
761}
762
763// --- Writing ---------------------------------------------------------------
764
765/// The services the inbox reads from as events arrive.
766pub struct Sources<'a> {
767 pub work: &'a Fetcher,
768 pub repos: &'a Fetcher,
769 pub identity: &'a Fetcher,
770}
771
772/// One notice written, and whether it was news (not a redelivery): what
773/// may be emailed.
774struct Written {
775 notice: Notice,
776 repo_id: String,
777 url: String,
778}
779
780/// Writes the items a batch from the bus calls for. Never fails the batch:
781/// what cannot be worked out is logged and left.
782pub async fn deliver(db: &D1Database, sources: &Sources<'_>, events: &[Event]) {
783 if let Err(error) = resolve(db, events).await {
784 worker::console_error!("inbox: agent threads not closed: {error}");
785 }
786 // A workspace's own notices, about no repository (workspace_notices).
787 for event in events.iter().filter(|event| WORKSPACE_EVENTS.contains(&event.kind.as_str())) {
788 if let Err(error) = deliver_workspace(db, event).await {
789 worker::console_error!("inbox: {} {} not delivered: {error}", event.kind, event.id);
790 }
791 }
792 let wanted: Vec<(&Event, Wanted)> = events
793 .iter()
794 .filter_map(|event| wants(event).map(|wanted| (event, wanted)))
795 .collect();
796 let created: Vec<&Event> = events.iter().filter(|event| event.kind == "repo.created").collect();
797 if wanted.is_empty() && created.is_empty() {
798 return;
799 }
800 // Everyone who caused one, named in one call.
801 let ids: Vec<String> = wanted
802 .iter()
803 .map(|(event, _)| *event)
804 .chain(created.iter().copied())
805 .filter_map(|event| event.actor.clone())
806 .filter(|id| !is_g1t_id(id))
807 .collect::<HashSet<_>>()
808 .into_iter()
809 .collect();
810 let names: HashMap<String, String> = if ids.is_empty() {
811 HashMap::new()
812 } else {
813 g1t_kit::call(sources.identity, "usernames", &UsernamesArgs { ids })
814 .await
815 .unwrap_or_else(|error| {
816 worker::console_error!("inbox: could not name who acted: {error}");
817 HashMap::new()
818 })
819 };
820 for event in created {
821 if let Err(error) = subscriptions::watch_created(db, event, &names).await {
822 worker::console_error!("inbox: {} {} not watched: {error}", event.kind, event.id);
823 }
824 }
825 let mut paths: HashMap<String, Option<RepoPath>> = HashMap::new();
826 let mut watchers: HashMap<String, Vec<subscriptions::Watcher>> = HashMap::new();
827 let mut written: Vec<Written> = Vec::new();
828 for (event, wanted) in wanted {
829 match deliver_one(db, sources, &names, &mut paths, &mut watchers, event, wanted).await {
830 Ok(mut news) => written.append(&mut news),
831 Err(error) => worker::console_error!("inbox: {} {} not delivered: {error}", event.kind, event.id),
832 }
833 }
834 if let Err(error) = email(db, sources.identity, written).await {
835 worker::console_error!("inbox: emails not sent: {error}");
836 }
837}
838
839/// Events about a workspace rather than a repository, each naming who to
840/// tell (`notify`), what to say (`title`, `body`) and where it is (`link`):
841/// identity's personal access token approvals.
842///
843/// | Event | Who | Reason | Severity |
844/// | --- | --- | --- | --- |
845/// | `token.approval_requested` | the workspace's owners | review_requested | warning |
846/// | `token.approval_reviewed` | the token's owner | author | info |
847pub const WORKSPACE_EVENTS: [&str; 2] = ["token.approval_requested", "token.approval_reviewed"];
848
849/// The notices a workspace event calls for, and the thread they go to.
850pub fn workspace_notices(event: &Event, actor: Option<&str>) -> (String, Vec<Notice>) {
851 let data = &event.data;
852 let text = |key: &str| data[key].as_str().unwrap_or_default().to_owned();
853 let (reason, severity) = match event.kind.as_str() {
854 "token.approval_requested" => (Reason::ReviewRequested, Severity::Warning),
855 _ => (Reason::Author, Severity::Info),
856 };
857 let thread = format!("workspace:{}/token/{}", text("workspace"), text("tokenId"));
858 let notices = names(data, "notify")
859 .into_iter()
860 .map(|name| name.to_lowercase())
861 .filter(|name| actor.is_none_or(|actor| !actor.eq_ignore_ascii_case(name)) && !is_g1t(name))
862 .collect::<HashSet<_>>()
863 .into_iter()
864 .map(|username| Notice { username, reason, severity, title: text("title"), body: text("body") })
865 .collect();
866 (thread, notices)
867}
868
869/// Writes a workspace event's notices: a thread each, with no repository.
870async fn deliver_workspace(db: &D1Database, event: &Event) -> Result<()> {
871 let (thread, told) = workspace_notices(event, None);
872 if told.is_empty() {
873 return Ok(());
874 }
875 let now = now_ms();
876 let workspace = event.data["workspace"].as_str().unwrap_or_default().to_lowercase();
877 let link = event.data["link"].as_str().filter(|link| link.starts_with('/'));
878 let mut statements = Vec::new();
879 for notice in &told {
880 let username = notice.username.as_str();
881 statements.push(db.prepare(BUMP).bind(&[
882 new_id("ntf", now).into(),
883 username.into(),
884 thread.as_str().into(),
885 event.id.as_str().into(),
886 event.kind.as_str().into(),
887 notice.reason.as_str().into(),
888 notice.severity.as_str().into(),
889 notice.title.as_str().into(),
890 notice.body.as_str().into(),
891 workspace.as_str().into(),
892 JsValue::NULL,
893 JsValue::NULL,
894 JsValue::NULL,
895 JsValue::NULL,
896 JsValue::NULL,
897 link.map_or(JsValue::NULL, JsValue::from),
898 JsValue::NULL,
899 event.time.as_str().into(),
900 new_id("ntf", now).into(),
901 ])?);
902 statements.push(
903 db.prepare(
904 "INSERT OR IGNORE INTO inbox_activity (id, item_id, username, event_id, event, reason, severity, title, body, actor, created_at)
905 SELECT ?1, id, ?2, ?4, ?5, ?6, ?7, ?8, ?9, NULL, ?10 FROM inbox_items WHERE username = ?2 AND thread = ?3",
906 )
907 .bind(&[
908 new_id("ntf", now).into(),
909 username.into(),
910 thread.as_str().into(),
911 event.id.as_str().into(),
912 event.kind.as_str().into(),
913 notice.reason.as_str().into(),
914 notice.severity.as_str().into(),
915 notice.title.as_str().into(),
916 notice.body.as_str().into(),
917 event.time.as_str().into(),
918 ])?,
919 );
920 }
921 db.batch(statements).await?;
922 Ok(())
923}
924
925/// Closes what an agent was waiting on a person for once it is over.
926async fn resolve(db: &D1Database, events: &[Event]) -> Result<()> {
927 let now = rfc3339(now_ms());
928 let mut statements = Vec::new();
929 for thread in events.iter().filter_map(resolves) {
930 statements.push(
931 db.prepare(
932 "UPDATE inbox_items SET done_at = ?1, read_at = COALESCE(read_at, ?1)
933 WHERE thread = ?2 AND reason = 'agent' AND done_at IS NULL",
934 )
935 .bind(&[now.as_str().into(), thread.into()])?,
936 );
937 }
938 // Reviews no longer asked for are no longer waiting.
939 for event in events.iter().filter(|event| event.kind == "pull.review_request_removed") {
940 let (Some(repo_id), Some(number)) = (
941 event.data["repoId"].as_str().map(str::to_owned).or_else(|| event.repo_id.clone()),
942 event.data["number"].as_u64(),
943 ) else {
944 continue;
945 };
946 for name in names(&event.data, "reviewers") {
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 username = ?3 AND reason = 'review_requested' AND done_at IS NULL",
951 )
952 .bind(&[
953 now.as_str().into(),
954 format!("{repo_id}#{number}").into(),
955 name.to_lowercase().into(),
956 ])?,
957 );
958 }
959 }
960 if !statements.is_empty() {
961 db.batch(statements).await?;
962 }
963 Ok(())
964}
965
966async fn deliver_one(
967 db: &D1Database,
968 sources: &Sources<'_>,
969 names: &HashMap<String, String>,
970 paths: &mut HashMap<String, Option<RepoPath>>,
971 watchers: &mut HashMap<String, Vec<subscriptions::Watcher>>,
972 event: &Event,
973 wanted: Wanted,
974) -> Result<Vec<Written>> {
975 if !paths.contains_key(&wanted.repo_id) {
976 let path: Option<RepoPath> = g1t_kit::call(
977 sources.repos,
978 "path_by_id",
979 &PathByIdArgs {
980 id: wanted.repo_id.clone(),
981 },
982 )
983 .await?;
984 paths.insert(wanted.repo_id.clone(), path);
985 }
986 // A repository that is gone tells nobody.
987 let Some(path) = paths.get(&wanted.repo_id).cloned().flatten() else {
988 return Ok(Vec::new());
989 };
990 let subject: Option<InboxSubject> = match wanted.number {
991 Some(number) => {
992 let found: Option<InboxSubject> = g1t_kit::call(
993 sources.work,
994 "inbox_subject",
995 &InboxSubjectArgs {
996 repo_id: wanted.repo_id.clone(),
997 number,
998 comment_id: wanted.comment_id.clone(),
999 },
1000 )
1001 .await?;
1002 if found.is_none() {
1003 return Ok(Vec::new());
1004 }
1005 found
1006 }
1007 None => None,
1008 };
1009 let actor = Actor {
1010 id: event.actor.clone(),
1011 username: match event.actor.as_deref() {
1012 Some(id) if is_g1t_id(id) => Some(system::USERNAME.to_owned()),
1013 Some(id) => names.get(id).map(|name| name.to_lowercase()),
1014 None => None,
1015 },
1016 };
1017 let thread = thread_of(event, &wanted, subject.as_ref());
1018 if !watchers.contains_key(&wanted.repo_id) {
1019 watchers.insert(wanted.repo_id.clone(), subscriptions::watchers(db, &wanted.repo_id).await?);
1020 }
1021 let audience = Audience {
1022 subscriptions: match thread.kind {
1023 Some(SubjectKind::Issue | SubjectKind::Pull) => subscriptions::of_thread(db, &thread.key).await?,
1024 _ => Vec::new(),
1025 },
1026 watchers: watchers.get(&wanted.repo_id).cloned().unwrap_or_default(),
1027 };
1028 let repo = format!("{}/{}", path.namespace, path.name).to_lowercase();
1029 let told = notices(event, &repo, &actor, subject.as_ref(), &audience);
1030 let mut statements = subscriptions::auto_subscribe(db, &thread.key, &wanted.repo_id, &subscribes(event, subject.as_ref()), &event.time)?;
1031 if told.is_empty() {
1032 if !statements.is_empty() {
1033 db.batch(statements).await?;
1034 }
1035 return Ok(Vec::new());
1036 }
1037 let now = now_ms();
1038 let link = thread.link.as_deref();
1039 let item_url = url(Some(&repo), thread.kind, thread.number, thread.run_id.as_deref(), link);
1040 // Three statements a notice: its thread brought up (unless this event
1041 // was told before), a line of history, and the history kept short.
1042 let first = statements.len();
1043 for notice in &told {
1044 let username = notice.username.as_str();
1045 let values: Vec<JsValue> = vec![
1046 new_id("ntf", now).into(),
1047 username.into(),
1048 thread.key.as_str().into(),
1049 event.id.as_str().into(),
1050 event.kind.as_str().into(),
1051 notice.reason.as_str().into(),
1052 notice.severity.as_str().into(),
1053 notice.title.as_str().into(),
1054 notice.body.as_str().into(),
1055 path.namespace.to_lowercase().into(),
1056 wanted.repo_id.as_str().into(),
1057 repo.as_str().into(),
1058 thread.kind.map_or(JsValue::NULL, |kind| kind.as_str().into()),
1059 thread.number.map_or(JsValue::NULL, JsValue::from),
1060 thread.run_id.as_deref().map_or(JsValue::NULL, JsValue::from),
1061 link.map_or(JsValue::NULL, JsValue::from),
1062 actor.username.as_deref().map_or(JsValue::NULL, JsValue::from),
1063 event.time.as_str().into(),
1064 // Same prefix as item ids, so old and new sort by time together.
1065 new_id("ntf", now).into(),
1066 ];
1067 statements.push(db.prepare(BUMP).bind(&values)?);
1068 statements.push(
1069 db.prepare(
1070 "INSERT OR IGNORE INTO inbox_activity (id, item_id, username, event_id, event, reason, severity, title, body, actor, created_at)
1071 SELECT ?1, id, ?2, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11 FROM inbox_items WHERE username = ?2 AND thread = ?3
1072 RETURNING username",
1073 )
1074 .bind(&[
1075 new_id("ntf", now).into(),
1076 username.into(),
1077 thread.key.as_str().into(),
1078 event.id.as_str().into(),
1079 event.kind.as_str().into(),
1080 notice.reason.as_str().into(),
1081 notice.severity.as_str().into(),
1082 notice.title.as_str().into(),
1083 notice.body.as_str().into(),
1084 actor.username.as_deref().map_or(JsValue::NULL, JsValue::from),
1085 event.time.as_str().into(),
1086 ])?,
1087 );
1088 statements.push(
1089 db.prepare(
1090 "DELETE FROM inbox_activity
1091 WHERE item_id = (SELECT id FROM inbox_items WHERE username = ?1 AND thread = ?2)
1092 AND id NOT IN (
1093 SELECT a.id FROM inbox_activity a
1094 WHERE a.item_id = (SELECT id FROM inbox_items WHERE username = ?1 AND thread = ?2)
1095 ORDER BY a.id DESC LIMIT ?3)",
1096 )
1097 .bind(&[username.into(), thread.key.as_str().into(), MAX_ACTIVITY.into()])?,
1098 );
1099 }
1100 let results = db.batch(statements).await?;
1101 #[derive(Deserialize)]
1102 struct Told {
1103 username: String,
1104 }
1105 // A line of history written means the event is news to that person.
1106 let mut news = HashSet::new();
1107 for (at, _) in told.iter().enumerate() {
1108 if let Some(result) = results.get(first + at * 3 + 1)
1109 && let Ok(rows) = result.results::<Told>()
1110 {
1111 news.extend(rows.into_iter().map(|row| row.username));
1112 }
1113 }
1114 Ok(told
1115 .into_iter()
1116 .filter(|notice| news.contains(&notice.username))
1117 .map(|notice| Written {
1118 notice,
1119 repo_id: wanted.repo_id.clone(),
1120 url: item_url.clone(),
1121 })
1122 .collect())
1123}
1124
1125/// Brings a person's thread up with new activity, or starts it. `?1` id,
1126/// `?2` username, `?3` thread, `?4` event id, `?5` event type, `?6` reason,
1127/// `?7` severity, `?8` title, `?9` body, `?10` workspace, `?11` repo id,
1128/// `?12` repo, `?13` subject, `?14` number, `?15` run id, `?16` link, `?17`
1129/// actor, `?18` time, `?19` the time-sortable id lists order by. Nothing
1130/// happens when the person was already told of this event. While the
1131/// thread is unread its most urgent severity is kept; done or not, it
1132/// comes back to the inbox. A snooze stands.
1133const BUMP: &str = "INSERT INTO inbox_items (id, username, thread, event_id, event, reason, severity, title, body,
1134 workspace, repo_id, repo, subject, number, run_id, link, actor, created_at, updated_at, bumped, activity)
1135 SELECT ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?18, ?19, 1
1136 WHERE NOT EXISTS (SELECT 1 FROM inbox_activity WHERE event_id = ?4 AND username = ?2)
1137 ON CONFLICT (username, thread) DO UPDATE SET
1138 event_id = excluded.event_id,
1139 event = excluded.event,
1140 reason = excluded.reason,
1141 severity = CASE
1142 WHEN inbox_items.read_at IS NULL AND inbox_items.done_at IS NULL
1143 AND (CASE inbox_items.severity WHEN 'warning' THEN 0 WHEN 'error' THEN 1 WHEN 'success' THEN 2 ELSE 3 END)
1144 < (CASE excluded.severity WHEN 'warning' THEN 0 WHEN 'error' THEN 1 WHEN 'success' THEN 2 ELSE 3 END)
1145 THEN inbox_items.severity ELSE excluded.severity END,
1146 title = excluded.title,
1147 body = excluded.body,
1148 workspace = excluded.workspace,
1149 repo = excluded.repo,
1150 subject = COALESCE(excluded.subject, inbox_items.subject),
1151 run_id = COALESCE(excluded.run_id, inbox_items.run_id),
1152 link = COALESCE(excluded.link, inbox_items.link),
1153 actor = excluded.actor,
1154 updated_at = excluded.updated_at,
1155 bumped = excluded.bumped,
1156 activity = inbox_items.activity + 1,
1157 read_at = NULL,
1158 done_at = NULL";
1159
1160/// Emails what was news to people who asked to be emailed for its reason.
1161/// Identity sends each, only to a confirmed address and only if the person
1162/// can still read the repository.
1163async fn email(db: &D1Database, identity: &Fetcher, written: Vec<Written>) -> Result<()> {
1164 if written.is_empty() {
1165 return Ok(());
1166 }
1167 let people: Vec<String> = written
1168 .iter()
1169 .map(|item| item.notice.username.clone())
1170 .collect::<HashSet<_>>()
1171 .into_iter()
1172 .collect();
1173 let settings = subscriptions::settings_of(db, &people).await?;
1174 for item in written {
1175 let wants = settings.get(&item.notice.username).map_or(&DEFAULT_EMAIL[..], |settings| &settings.email[..]);
1176 if !wants.contains(&item.notice.reason) {
1177 continue;
1178 }
1179 let sent: Result<bool> = g1t_kit::call(
1180 identity,
1181 "notify_by_email",
1182 &NotifyByEmailArgs {
1183 username: item.notice.username.clone(),
1184 repo_id: item.repo_id.clone(),
1185 subject: item.notice.title.clone(),
1186 intro: item.notice.body.clone(),
1187 quote: None,
1188 path: item.url.clone(),
1189 reason: item.notice.reason,
1190 },
1191 )
1192 .await;
1193 if let Err(error) = sent {
1194 worker::console_error!("inbox: email to {} not sent: {error}", item.notice.username);
1195 }
1196 }
1197 Ok(())
1198}
1199
1200// --- Reading and changing ----------------------------------------------------
1201
1202#[derive(Deserialize)]
1203pub(crate) struct Row {
1204 pub(crate) id: String,
1205 reason: String,
1206 severity: String,
1207 title: String,
1208 body: String,
1209 event: Option<String>,
1210 workspace: Option<String>,
1211 pub(crate) repo_id: Option<String>,
1212 repo: Option<String>,
1213 subject: Option<String>,
1214 number: Option<f64>,
1215 run_id: Option<String>,
1216 link: Option<String>,
1217 actor: Option<String>,
1218 activity: Option<f64>,
1219 created_at: String,
1220 updated_at: Option<String>,
1221 read_at: Option<String>,
1222 done_at: Option<String>,
1223 saved: f64,
1224 snoozed_until: Option<String>,
1225 bumped: Option<String>,
1226}
1227
1228impl Row {
1229 pub(crate) fn into_item(self) -> InboxItem {
1230 let subject = self.subject.as_deref().and_then(SubjectKind::parse);
1231 let number = self.number.map(|n| n as u32);
1232 InboxItem {
1233 url: url(self.repo.as_deref(), subject, number, self.run_id.as_deref(), self.link.as_deref()),
1234 id: self.id,
1235 reason: Reason::parse(&self.reason).unwrap_or(Reason::Subscribed),
1236 severity: Severity::parse(&self.severity).unwrap_or(Severity::Info),
1237 title: self.title,
1238 body: self.body,
1239 event: self.event,
1240 repo: self.repo,
1241 workspace: self.workspace,
1242 subject,
1243 number,
1244 actor: self.actor,
1245 count: self.activity.map_or(1, |n| n as u32),
1246 updated_at: self.updated_at.unwrap_or_else(|| self.created_at.clone()),
1247 created_at: self.created_at,
1248 read_at: self.read_at,
1249 done_at: self.done_at,
1250 saved: self.saved != 0.0,
1251 snoozed_until: self.snoozed_until,
1252 }
1253 }
1254}
1255
1256pub(crate) const COLUMNS: &str = "id, reason, severity, title, body, event, workspace, repo_id, repo, subject, number, run_id,
1257 link, actor, activity, created_at, updated_at, read_at, done_at, saved, snoozed_until, bumped";
1258
1259/// The conditions that pick a view's items, after `username = ?1`; `?2` is now.
1260fn view_filter(view: InboxView) -> &'static str {
1261 match view {
1262 InboxView::Inbox => "done_at IS NULL AND (snoozed_until IS NULL OR snoozed_until <= ?2)",
1263 InboxView::Saved => "saved = 1",
1264 InboxView::Done => "done_at IS NOT NULL",
1265 }
1266}
1267
1268/// Whether a list ranks unread warnings first: the inbox itself, unfiltered.
1269fn ranked(a: &ListInboxArgs) -> bool {
1270 a.view == InboxView::Inbox
1271 && a.severity.is_none()
1272 && a.reason.is_none()
1273 && !a.participating
1274 && a.repo_id.is_none()
1275 && !a.unread
1276 && a.since.is_none()
1277 && a.updated_before.is_none()
1278}
1279
1280/// The conditions a list's filters add, and the values they bind, numbered
1281/// after the `bound` values already given.
1282fn filters(a: &ListInboxArgs, bound: usize) -> (Vec<String>, Vec<String>) {
1283 let mut conditions = Vec::new();
1284 let mut values: Vec<String> = Vec::new();
1285 let mut bind = |value: &str, condition: &str| {
1286 values.push(value.to_owned());
1287 conditions.push(condition.replace('?', &format!("?{}", bound + values.len())));
1288 };
1289 if let Some(severity) = a.severity {
1290 bind(severity.as_str(), "severity = ?");
1291 }
1292 if let Some(reason) = a.reason {
1293 bind(reason.as_str(), "reason = ?");
1294 }
1295 if let Some(repo_id) = &a.repo_id {
1296 bind(repo_id, "repo_id = ?");
1297 }
1298 if let Some(since) = &a.since {
1299 bind(since.trim(), "COALESCE(updated_at, created_at) >= ?");
1300 }
1301 if let Some(before) = &a.updated_before {
1302 bind(before.trim(), "COALESCE(updated_at, created_at) < ?");
1303 }
1304 if a.participating {
1305 conditions.push("reason NOT IN ('manual', 'subscribed')".to_owned());
1306 }
1307 if a.unread {
1308 conditions.push("read_at IS NULL".to_owned());
1309 }
1310 (conditions, values)
1311}
1312
1313pub async fn list(db: &D1Database, repos: &Fetcher, a: ListInboxArgs) -> Result<InboxPage> {
1314 let Some(viewer) = &a.viewer else {
1315 return Ok(InboxPage::default());
1316 };
1317 let username = viewer.username.to_lowercase();
1318 let now = rfc3339(now_ms());
1319 let limit = a.limit.unwrap_or(DEFAULT_INBOX_PAGE).clamp(1, MAX_INBOX_PAGE);
1320 let mut values: Vec<JsValue> = vec![username.as_str().into(), now.as_str().into()];
1321 let mut conditions = vec!["username = ?1".to_owned(), view_filter(a.view).to_owned()];
1322 let (filtered, bound) = filters(&a, values.len());
1323 conditions.extend(filtered);
1324 values.extend(bound.iter().map(|value| JsValue::from(value.as_str())));
1325 let ranked = ranked(&a);
1326 // Unread warnings lead the first page, and are left out of the rest.
1327 let leading = "severity = 'warning' AND read_at IS NULL";
1328 let mut rest = conditions.clone();
1329 let mut rest_values = values.clone();
1330 if ranked {
1331 rest.push(format!("NOT ({leading})"));
1332 }
1333 if let Some(before) = &a.before {
1334 rest_values.push(before.as_str().into());
1335 rest.push(format!("bumped < ?{}", rest_values.len()));
1336 }
1337 rest_values.push((limit + 1).into());
1338 let order = if a.view == InboxView::Done { "done_at DESC, bumped DESC" } else { "bumped DESC" };
1339 let mut statements = vec![
1340 db.prepare(format!(
1341 "SELECT {COLUMNS} FROM inbox_items WHERE {} ORDER BY {order} LIMIT ?{}",
1342 rest.join(" AND "),
1343 rest_values.len()
1344 ))
1345 .bind(&rest_values)?,
1346 ];
1347 if ranked && a.before.is_none() {
1348 statements.push(
1349 db.prepare(format!(
1350 "SELECT {COLUMNS} FROM inbox_items WHERE {} AND {leading} ORDER BY bumped DESC LIMIT ?3",
1351 conditions.join(" AND ")
1352 ))
1353 .bind(&[username.as_str().into(), now.as_str().into(), MAX_RANKED.into()])?,
1354 );
1355 }
1356 let results = db.batch(statements).await?;
1357 let mut rows = results.first().map(|result| result.results::<Row>()).transpose()?.unwrap_or_default();
1358 let next = if rows.len() > limit as usize {
1359 rows.truncate(limit as usize);
1360 rows.last().map(|row| row.bumped.clone().unwrap_or_else(|| row.id.clone()))
1361 } else {
1362 None
1363 };
1364 let mut items: Vec<Row> = results.get(1).map(|result| result.results::<Row>()).transpose()?.unwrap_or_default();
1365 items.extend(rows);
1366 let items = readable_only(db, repos, &a.viewer, &username, items).await?;
1367 Ok(InboxPage {
1368 items: items.into_iter().map(Row::into_item).collect(),
1369 next,
1370 })
1371}
1372
1373/// Rows about repositories the viewer can still read; the rest are
1374/// removed from their inbox as they are found.
1375pub(crate) async fn readable_only(db: &D1Database, repos: &Fetcher, viewer: &g1t_contracts::Viewer, username: &str, mut rows: Vec<Row>) -> Result<Vec<Row>> {
1376 let ids: Vec<String> = rows
1377 .iter()
1378 .filter_map(|row| row.repo_id.clone())
1379 .collect::<HashSet<_>>()
1380 .into_iter()
1381 .collect();
1382 if ids.is_empty() {
1383 return Ok(rows);
1384 }
1385 let readable: Vec<Repo> = g1t_kit::call(
1386 repos,
1387 "readable",
1388 &ReadableArgs {
1389 ids: ids.clone(),
1390 viewer: viewer.clone(),
1391 },
1392 )
1393 .await?;
1394 let readable: HashSet<String> = readable.into_iter().map(|repo| repo.id).collect();
1395 let gone: Vec<String> = ids.into_iter().filter(|id| !readable.contains(id)).collect();
1396 if !gone.is_empty() {
1397 forget(db, username, &gone).await?;
1398 rows.retain(|row| row.repo_id.as_ref().is_none_or(|id| !gone.contains(id)));
1399 }
1400 Ok(rows)
1401}
1402
1403/// Removes a person's items about repositories they cannot read.
1404async fn forget(db: &D1Database, username: &str, repo_ids: &[String]) -> Result<()> {
1405 let marks = vec!["?"; repo_ids.len()].join(", ");
1406 let mut values: Vec<JsValue> = vec![username.into()];
1407 values.extend(repo_ids.iter().map(|id| JsValue::from(id.as_str())));
1408 db.batch(vec![
1409 db.prepare(format!(
1410 "DELETE FROM inbox_activity WHERE item_id IN (SELECT id FROM inbox_items WHERE username = ? AND repo_id IN ({marks}))"
1411 ))
1412 .bind(&values)?,
1413 db.prepare(format!("DELETE FROM inbox_items WHERE username = ? AND repo_id IN ({marks})"))
1414 .bind(&values)?,
1415 ])
1416 .await?;
1417 Ok(())
1418}
1419
1420#[derive(Deserialize)]
1421struct CountRow {
1422 severity: String,
1423 n: f64,
1424}
1425
1426pub async fn counts(db: &D1Database, a: InboxCountsArgs) -> Result<InboxCounts> {
1427 let rows = db
1428 .prepare(
1429 "SELECT severity, count(*) AS n FROM inbox_items
1430 WHERE username = ? AND read_at IS NULL AND done_at IS NULL
1431 AND (snoozed_until IS NULL OR snoozed_until <= ?)
1432 GROUP BY severity",
1433 )
1434 .bind(&[a.username.to_lowercase().into(), rfc3339(now_ms()).into()])?
1435 .all()
1436 .await?
1437 .results::<CountRow>()?;
1438 let mut counts = InboxCounts::default();
1439 for row in rows {
1440 let n = row.n as u32;
1441 counts.unread += n;
1442 match Severity::parse(&row.severity) {
1443 Some(Severity::Error) => counts.error += n,
1444 Some(Severity::Warning) => counts.warning += n,
1445 Some(Severity::Success) => counts.success += n,
1446 Some(Severity::Info) | None => counts.info += n,
1447 }
1448 }
1449 Ok(counts)
1450}
1451
1452/// The change a mark makes, as a `SET` clause; `?1` is now.
1453fn mark_change(mark: InboxMark) -> &'static str {
1454 match mark {
1455 InboxMark::Read => "read_at = COALESCE(read_at, ?1)",
1456 InboxMark::Unread => "read_at = NULL",
1457 InboxMark::Done => "done_at = ?1, read_at = COALESCE(read_at, ?1)",
1458 InboxMark::Undone => "done_at = NULL",
1459 InboxMark::Save => "saved = 1",
1460 InboxMark::Unsave => "saved = 0",
1461 InboxMark::Snooze => "snoozed_until = ?2, read_at = COALESCE(read_at, ?1)",
1462 InboxMark::Unsnooze => "snoozed_until = NULL",
1463 }
1464}
1465
1466#[derive(Deserialize)]
1467struct IdRow {
1468 #[allow(dead_code)]
1469 id: String,
1470}
1471
1472/// Changes the person's own items. Returns how many changed.
1473pub async fn mark(db: &D1Database, a: MarkInboxArgs) -> Result<u32> {
1474 let now = rfc3339(now_ms());
1475 let until = match (a.mark, a.until.as_deref().map(str::trim)) {
1476 // Times compare as text (`g1t_contracts::time`); a snooze is for later.
1477 (InboxMark::Snooze, Some(until)) if until.len() == now.len() && until > now.as_str() => until.to_owned(),
1478 (InboxMark::Snooze, _) => return Ok(0),
1479 _ => String::new(),
1480 };
1481 let mut values: Vec<JsValue> = vec![now.as_str().into(), until.as_str().into(), a.username.to_lowercase().into()];
1482 let target = if a.all && a.ids.is_empty() {
1483 let mut target = "done_at IS NULL".to_owned();
1484 let mut bind = |values: &mut Vec<JsValue>, value: &str, condition: &str| {
1485 values.push(value.into());
1486 target.push_str(&format!(" AND {}", condition.replace('?', &format!("?{}", values.len()))));
1487 };
1488 if let Some(severity) = a.severity {
1489 bind(&mut values, severity.as_str(), "severity = ?");
1490 }
1491 if let Some(repo_id) = &a.repo_id {
1492 bind(&mut values, repo_id, "repo_id = ?");
1493 }
1494 if let Some(last_read_at) = a.last_read_at.as_deref().map(str::trim).filter(|at| !at.is_empty()) {
1495 bind(&mut values, last_read_at, "COALESCE(updated_at, created_at) <= ?");
1496 }
1497 target
1498 } else {
1499 let ids: Vec<&String> = a.ids.iter().take(MAX_MARK).collect();
1500 if ids.is_empty() {
1501 return Ok(0);
1502 }
1503 let first = values.len() + 1;
1504 values.extend(ids.iter().map(|id| JsValue::from(id.as_str())));
1505 let marks: Vec<String> = (first..values.len() + 1).map(|at| format!("?{at}")).collect();
1506 format!("id IN ({})", marks.join(", "))
1507 };
1508 let changed = db
1509 .prepare(format!(
1510 "UPDATE inbox_items SET {} WHERE username = ?3 AND {target} RETURNING id",
1511 mark_change(a.mark)
1512 ))
1513 .bind(&values)?
1514 .all()
1515 .await?
1516 .results::<IdRow>()?;
1517 Ok(changed.len() as u32)
1518}
1519
1520#[derive(Deserialize)]
1521struct ActivityRow {
1522 reason: String,
1523 severity: String,
1524 title: String,
1525 body: String,
1526 event: Option<String>,
1527 actor: Option<String>,
1528 created_at: String,
1529}
1530
1531/// One of the viewer's threads, with its history and their subscription.
1532pub async fn thread(db: &D1Database, repos: &Fetcher, work: &Fetcher, a: ThreadArgs) -> Result<Option<InboxThread>> {
1533 let Some(viewer) = &a.viewer else {
1534 return Ok(None);
1535 };
1536 let username = viewer.username.to_lowercase();
1537 let row = db
1538 .prepare(format!("SELECT {COLUMNS} FROM inbox_items WHERE id = ? AND username = ?"))
1539 .bind(&[a.id.as_str().into(), username.as_str().into()])?
1540 .first::<Row>(None)
1541 .await?;
1542 let Some(row) = readable_only(db, repos, &a.viewer, &username, row.into_iter().collect()).await?.pop() else {
1543 return Ok(None);
1544 };
1545 let activity = db
1546 .prepare(
1547 "SELECT reason, severity, title, body, event, actor, created_at FROM inbox_activity
1548 WHERE item_id = ? ORDER BY id DESC LIMIT ?",
1549 )
1550 .bind(&[row.id.as_str().into(), MAX_ACTIVITY.into()])?
1551 .all()
1552 .await?
1553 .results::<ActivityRow>()?
1554 .into_iter()
1555 .map(|row| InboxActivity {
1556 reason: Reason::parse(&row.reason).unwrap_or(Reason::Subscribed),
1557 severity: Severity::parse(&row.severity).unwrap_or(Severity::Info),
1558 title: row.title,
1559 body: row.body,
1560 event: row.event,
1561 actor: row.actor,
1562 created_at: row.created_at,
1563 })
1564 .collect();
1565 let item = row.into_item();
1566 let subscription = match (item.subject, item.number) {
1567 (Some(SubjectKind::Issue | SubjectKind::Pull), Some(_)) => {
1568 subscriptions::subscription(
1569 db,
1570 work,
1571 SubscriptionArgs {
1572 viewer: a.viewer.clone(),
1573 id: Some(item.id.clone()),
1574 ..SubscriptionArgs::default()
1575 },
1576 )
1577 .await?
1578 }
1579 _ => None,
1580 };
1581 Ok(Some(InboxThread {
1582 item,
1583 activity,
1584 subscription,
1585 }))
1586}
1587
1588// --- Keeping up ------------------------------------------------------------------
1589
1590/// Moves rows with renamed workspaces and repositories, and drops those
1591/// of purged repositories, deleted workspaces and purged accounts.
1592pub async fn follow(db: &D1Database, events: &[Event]) -> Result<()> {
1593 let mut statements = Vec::new();
1594 for event in events {
1595 let text = |key: &str| event.data[key].as_str().unwrap_or_default().to_lowercase();
1596 match event.kind.as_str() {
1597 "workspace.renamed" => {
1598 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
1599 continue;
1600 };
1601 let (from, to) = (renamed.from.to_lowercase(), renamed.to.to_lowercase());
1602 if from == to {
1603 continue;
1604 }
1605 statements.push(
1606 db.prepare(
1607 "UPDATE inbox_items SET repo = ?2 || substr(repo, length(?1) + 1), workspace = ?2,
1608 link = CASE WHEN link LIKE '/' || ?1 || '/%' THEN '/' || ?2 || substr(link, length(?1) + 2) ELSE link END
1609 WHERE workspace = ?1",
1610 )
1611 .bind(&[from.as_str().into(), to.as_str().into()])?,
1612 );
1613 statements.push(
1614 db.prepare("UPDATE inbox_watching SET repo = ?2 || substr(repo, length(?1) + 1) WHERE repo LIKE ?1 || '/%'")
1615 .bind(&[from.as_str().into(), to.as_str().into()])?,
1616 );
1617 }
1618 "repo.renamed" => {
1619 let path = format!("{}/{}", text("namespace"), text("to"));
1620 for table in ["inbox_items", "inbox_watching"] {
1621 statements.push(
1622 db.prepare(format!("UPDATE {table} SET repo = ? WHERE repo_id = ?"))
1623 .bind(&[path.as_str().into(), text("repoId").into()])?,
1624 );
1625 }
1626 }
1627 "repo.transferred" => {
1628 let path = format!("{}/{}", text("to"), text("name"));
1629 statements.push(
1630 db.prepare("UPDATE inbox_items SET repo = ?, workspace = ? WHERE repo_id = ?")
1631 .bind(&[path.as_str().into(), text("to").into(), text("repoId").into()])?,
1632 );
1633 statements.push(
1634 db.prepare("UPDATE inbox_watching SET repo = ? WHERE repo_id = ?")
1635 .bind(&[path.as_str().into(), text("repoId").into()])?,
1636 );
1637 }
1638 "repo.purged" => {
1639 for sql in [
1640 "DELETE FROM inbox_activity WHERE item_id IN (SELECT id FROM inbox_items WHERE repo_id = ?)",
1641 "DELETE FROM inbox_items WHERE repo_id = ?",
1642 "DELETE FROM inbox_subscriptions WHERE repo_id = ?",
1643 "DELETE FROM inbox_watching WHERE repo_id = ?",
1644 ] {
1645 statements.push(db.prepare(sql).bind(&[text("repoId").into()])?);
1646 }
1647 }
1648 "workspace.deleted" => {
1649 statements.push(
1650 db.prepare("DELETE FROM inbox_items WHERE workspace = ?")
1651 .bind(&[text("slug").into()])?,
1652 );
1653 statements.push(
1654 db.prepare("DELETE FROM inbox_watching WHERE repo LIKE ? || '/%'")
1655 .bind(&[text("slug").into()])?,
1656 );
1657 }
1658 // An account purged: its inbox, what it watched and its
1659 // settings go with it (identity's account_deletion.rs).
1660 "user.deleted" => {
1661 let username = text("username");
1662 if username.is_empty() {
1663 continue;
1664 }
1665 for sql in [
1666 "DELETE FROM inbox_activity WHERE username = ?",
1667 "DELETE FROM inbox_items WHERE username = ?",
1668 "DELETE FROM inbox_subscriptions WHERE username = ?",
1669 "DELETE FROM inbox_watching WHERE username = ?",
1670 "DELETE FROM inbox_settings WHERE username = ?",
1671 ] {
1672 statements.push(db.prepare(sql).bind(&[username.as_str().into()])?);
1673 }
1674 }
1675 _ => {}
1676 }
1677 }
1678 if !statements.is_empty() {
1679 db.batch(statements).await?;
1680 }
1681 Ok(())
1682}
1683
1684/// Removes items done more than [`DONE_DAYS`] ago, and any not saved older
1685/// than [`MAX_DAYS`], with their history. Returns how many went.
1686pub async fn purge(db: &D1Database, now: u64) -> Result<u32> {
1687 let done = crate::audit::keep_from(now, DONE_DAYS);
1688 let oldest = crate::audit::keep_from(now, MAX_DAYS);
1689 let statements: Vec<D1PreparedStatement> = vec![
1690 db.prepare("DELETE FROM inbox_items WHERE saved = 0 AND done_at < ?")
1691 .bind(&[done.into()])?,
1692 db.prepare("DELETE FROM inbox_items WHERE saved = 0 AND COALESCE(updated_at, created_at) < ?")
1693 .bind(&[oldest.into()])?,
1694 db.prepare("DELETE FROM inbox_activity WHERE NOT EXISTS (SELECT 1 FROM inbox_items WHERE inbox_items.id = inbox_activity.item_id)"),
1695 ];
1696 let results = db.batch(statements).await?;
1697 let mut removed = 0;
1698 for result in results.into_iter().take(2) {
1699 removed += result.meta()?.and_then(|meta| meta.changes).unwrap_or(0) as u32;
1700 }
1701 Ok(removed)
1702}
1703
1704#[cfg(test)]
1705mod tests {
1706 use super::*;
1707 use crate::subscriptions::{State, Subscription, Watcher};
1708 use serde_json::json;
1709
1710 fn event(kind: &str, actor: Option<&str>, data: serde_json::Value) -> Event {
1711 Event {
1712 id: "evt_1".into(),
1713 kind: kind.into(),
1714 source: "work".into(),
1715 time: "2026-10-07T12:00:00.000Z".into(),
1716 repo_id: Some("rep_1".into()),
1717 actor: actor.map(str::to_owned),
1718 data,
1719 }
1720 }
1721
1722 fn person(id: &str, username: &str) -> Principal {
1723 Principal {
1724 id: id.into(),
1725 username: username.into(),
1726 }
1727 }
1728
1729 fn actor(id: &str, username: &str) -> Actor {
1730 Actor {
1731 id: Some(id.into()),
1732 username: Some(username.into()),
1733 }
1734 }
1735
1736 /// A pull request ana opened, assigned to bo, for an issue cy filed and dee is assigned.
1737 fn pull() -> InboxSubject {
1738 InboxSubject {
1739 kind: Some(SubjectKind::Pull),
1740 title: "Add the inbox".into(),
1741 author: person("usr_ana", "ana"),
1742 assignees: vec!["bo".into()],
1743 issue: Some(Box::new(InboxSubject {
1744 kind: Some(SubjectKind::Issue),
1745 title: "An inbox".into(),
1746 author: person("usr_cy", "cy"),
1747 assignees: vec!["dee".into()],
1748 ..InboxSubject::default()
1749 })),
1750 ..InboxSubject::default()
1751 }
1752 }
1753
1754 /// A change g1t made for ana.
1755 fn g1t_pull() -> InboxSubject {
1756 InboxSubject {
1757 author: person(AGENT_ID, "g1t"),
1758 requested_by: Some(person("usr_ana", "ana")),
1759 ..pull()
1760 }
1761 }
1762
1763 fn nobody() -> Audience {
1764 Audience::default()
1765 }
1766
1767 fn told(notices: &[Notice]) -> Vec<(&str, Reason, Severity)> {
1768 notices
1769 .iter()
1770 .map(|notice| (notice.username.as_str(), notice.reason, notice.severity))
1771 .collect()
1772 }
1773
1774 fn comment(author: Principal, body: &str, mentions: &[&str]) -> InboxSubject {
1775 InboxSubject {
1776 comment: Some(InboxComment {
1777 author,
1778 excerpt: body.into(),
1779 mentions: mentions.iter().map(|name| (*name).to_owned()).collect(),
1780 ..InboxComment::default()
1781 }),
1782 ..pull()
1783 }
1784 }
1785
1786 fn subscribed(rows: &[(&str, State, Option<Reason>)]) -> Audience {
1787 Audience {
1788 subscriptions: rows
1789 .iter()
1790 .map(|(name, state, reason)| Subscription {
1791 username: (*name).into(),
1792 state: *state,
1793 reason: *reason,
1794 })
1795 .collect(),
1796 watchers: Vec::new(),
1797 }
1798 }
1799
1800 fn watched(rows: &[(&str, WatchLevel, &[&str])]) -> Audience {
1801 Audience {
1802 subscriptions: Vec::new(),
1803 watchers: rows
1804 .iter()
1805 .map(|(name, level, events)| Watcher {
1806 username: (*name).into(),
1807 level: *level,
1808 events: events.iter().map(|kind| (*kind).to_owned()).collect(),
1809 })
1810 .collect(),
1811 }
1812 }
1813
1814 #[test]
1815 fn only_events_that_tell_someone_are_read() {
1816 let asked = |kind: &str, data| wants(&event(kind, None, data));
1817 assert_eq!(
1818 asked("checks.completed", json!({ "repoId": "rep_1", "number": 4, "status": "failed" })),
1819 Some(Wanted { repo_id: "rep_1".into(), number: Some(4), comment_id: None })
1820 );
1821 assert_eq!(asked("checks.completed", json!({ "repoId": "rep_1", "number": 4, "status": "passed" })), None);
1822 assert_eq!(asked("review.completed", json!({ "repoId": "rep_1", "number": 4 })), None);
1823 assert_eq!(asked("workflow.completed", json!({ "repoId": "rep_1", "conclusion": "success" })), None);
1824 assert_eq!(
1825 asked("workflow.completed", json!({ "repoId": "rep_1", "conclusion": "failure" })),
1826 Some(Wanted { repo_id: "rep_1".into(), number: None, comment_id: None })
1827 );
1828 assert_eq!(
1829 asked("comment.created", json!({ "repoId": "rep_1", "number": 2, "commentId": "cmt_1" })).map(|w| w.comment_id),
1830 Some(Some("cmt_1".into()))
1831 );
1832 assert_eq!(asked("git.push", json!({ "repoId": "rep_1" })), None);
1833 for kind in ["issue.opened", "pull.review_requested", "pull.stalled", "issue.assigned", "pull.closed", "issue.reopened"] {
1834 assert!(asked(kind, json!({ "repoId": "rep_1", "number": 1 })).is_some(), "{kind}");
1835 }
1836 assert_eq!(asked("deployment.failed", json!({ "repoId": "rep_1" })).map(|w| w.number), Some(None));
1837 }
1838
1839 #[test]
1840 fn a_waiting_or_stopped_agent_needs_the_pull_requests_and_the_issues_people_first() {
1841 let asked = event("agent.asked", Some("usr_agent"), json!({ "number": 7 }));
1842 let notices_ = notices(&asked, "acme/rocket", &Actor::default(), Some(&pull()), &nobody());
1843 assert_eq!(
1844 told(&notices_),
1845 vec![
1846 ("ana", Reason::Agent, Severity::Warning),
1847 ("cy", Reason::Agent, Severity::Warning),
1848 ("dee", Reason::Agent, Severity::Warning),
1849 ]
1850 );
1851 assert_eq!(notices_[0].title, "An agent is waiting on acme/rocket#7");
1852 assert_eq!(notices_[0].body, "Add the inbox");
1853 // Stopping says why; even someone who unsubscribed hears of it.
1854 let stalled = event("pull.stalled", None, json!({ "number": 7, "detail": "Its checks could not be run." }));
1855 let audience = subscribed(&[("ana", State::Unsubscribed, None)]);
1856 let notices_ = notices(&stalled, "acme/rocket", &Actor::default(), Some(&g1t_pull()), &audience);
1857 assert_eq!(notices_[0].username, "ana");
1858 assert_eq!(notices_[0].title, "g1t stopped on acme/rocket#7 and needs you");
1859 assert_eq!(notices_[0].body, "Its checks could not be run.");
1860 // Ignoring the thread silences even that.
1861 let audience = subscribed(&[("ana", State::Ignored, None)]);
1862 let notices_ = notices(&stalled, "acme/rocket", &Actor::default(), Some(&g1t_pull()), &audience);
1863 assert!(!notices_.iter().any(|notice| notice.username == "ana"));
1864 }
1865
1866 #[test]
1867 fn whatever_an_agent_waits_on_closes_when_it_goes_on() {
1868 for kind in RESUMES {
1869 assert_eq!(resolves(&event(kind, None, json!({ "repoId": "rep_1", "number": 7 }))), Some("rep_1#7".into()));
1870 }
1871 assert_eq!(resolves(&event("comment.created", None, json!({ "repoId": "rep_1", "number": 7 }))), None);
1872 }
1873
1874 #[test]
1875 fn a_deployments_reviewers_are_told() {
1876 let asked = event(
1877 "deployment.review_requested",
1878 Some("usr_ana"),
1879 json!({
1880 "repoId": "rep_1", "runId": "run_7", "environment": "production", "workflow": "Deploy",
1881 "title": "Ship it", "notify": ["cy", "ana"], "link": "/acme/rocket/actions/runs/run_7"
1882 }),
1883 );
1884 let wanted = wants(&asked).unwrap();
1885 assert_eq!(wanted.number, None);
1886 let thread = thread_of(&asked, &wanted, None);
1887 assert_eq!(thread.key, "rep_1/review/run_7/production");
1888 assert_eq!(thread.link.as_deref(), Some("/acme/rocket/actions/runs/run_7"));
1889 let told = notices(&asked, "acme/rocket", &actor("usr_ana", "ana"), None, &nobody());
1890 let names: Vec<&str> = told.iter().map(|notice| notice.username.as_str()).collect();
1891 // Whoever started the run is told too, when they review it.
1892 assert_eq!(names, ["cy", "ana"]);
1893 assert!(told[0].title.contains("waiting for your review to deploy to production"));
1894 }
1895
1896 #[test]
1897 fn teams_asked_to_review_tell_the_people_they_name() {
1898 let requested = event(
1899 "pull.review_requested",
1900 Some("usr_ana"),
1901 json!({ "number": 7, "reviewers": ["cy"], "teams": [
1902 { "team": "acme/backend", "notified": ["cy", "dee", "ana"], "assigned": ["cy"] }
1903 ] }),
1904 );
1905 let notices_ = notices(&requested, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
1906 // cy was picked and is told as a reviewer; dee through the team;
1907 // never whoever asked.
1908 assert_eq!(
1909 told(&notices_),
1910 vec![("cy", Reason::ReviewRequested, Severity::Warning), ("dee", Reason::ReviewRequested, Severity::Warning)]
1911 );
1912 assert_eq!(notices_[0].title, "ana asked you to review acme/rocket#7");
1913 assert_eq!(notices_[1].title, "ana asked @acme/backend to review acme/rocket#7");
1914 assert_eq!(
1915 subscribes(&requested, Some(&pull())),
1916 vec![("cy".into(), Reason::ReviewRequested), ("dee".into(), Reason::ReviewRequested), ("ana".into(), Reason::ReviewRequested)]
1917 );
1918 // Asked by the CODEOWNERS file.
1919 let owned = event(
1920 "pull.review_requested",
1921 None,
1922 json!({ "number": 7, "reviewers": ["bo"], "codeOwners": true, "teams": [{ "team": "acme/docs", "notified": ["wren"], "assigned": [] }] }),
1923 );
1924 let notices_ = notices(&owned, "acme/rocket", &Actor::default(), Some(&pull()), &nobody());
1925 assert_eq!(notices_[0].title, "acme/rocket#7 changes files you own");
1926 assert_eq!(notices_[1].title, "acme/rocket#7 changes files @acme/docs owns");
1927 }
1928
1929 #[test]
1930 fn a_team_mention_tells_its_people_once_and_a_name_wins() {
1931 let mut on = comment(person("usr_bo", "bo"), "cc @ana @acme/backend", &["ana"]);
1932 if let Some(comment) = on.comment.as_mut() {
1933 comment.team_mentions = vec![TeamMentioned {
1934 team: "acme/backend".into(),
1935 members: vec!["ana".into(), "cy".into(), "bo".into()],
1936 }];
1937 }
1938 let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
1939 let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &nobody());
1940 let mentioned: Vec<(&str, Reason)> = notices_.iter().map(|n| (n.username.as_str(), n.reason)).filter(|(_, r)| matches!(r, Reason::Mention | Reason::TeamMention)).collect();
1941 assert_eq!(mentioned, vec![("ana", Reason::Mention), ("cy", Reason::TeamMention)]);
1942 assert_eq!(
1943 notices_.iter().find(|n| n.username == "cy").unwrap().title,
1944 "bo mentioned @acme/backend on acme/rocket#7"
1945 );
1946 let subscribed = subscribes(&created, Some(&on));
1947 assert!(subscribed.contains(&("cy".into(), Reason::TeamMention)));
1948 assert!(subscribed.contains(&("ana".into(), Reason::Mention)));
1949 }
1950
1951 #[test]
1952 fn a_description_tells_the_people_and_teams_it_mentions() {
1953 let mut on = pull();
1954 on.mentions = vec!["dee".into()];
1955 on.team_mentions = vec![TeamMentioned {
1956 team: "acme/web".into(),
1957 members: vec!["eve".into()],
1958 }];
1959 let opened = event("pull.opened", Some("usr_ana"), json!({ "number": 7 }));
1960 let notices_ = notices(&opened, "acme/rocket", &actor("usr_ana", "ana"), Some(&on), &nobody());
1961 assert!(told(&notices_).contains(&("dee", Reason::Mention, Severity::Info)));
1962 assert!(told(&notices_).contains(&("eve", Reason::TeamMention, Severity::Info)));
1963 assert_eq!(
1964 subscribes(&opened, Some(&on)),
1965 vec![("dee".into(), Reason::Mention), ("eve".into(), Reason::TeamMention)]
1966 );
1967 }
1968
1969 #[test]
1970 fn reviewers_and_assignees_asked_are_told_never_whoever_asked() {
1971 let requested = event("pull.review_requested", Some("usr_ana"), json!({ "number": 7, "reviewers": ["bo", "g1t", "ana"] }));
1972 let notices_ = notices(&requested, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
1973 assert_eq!(told(&notices_), vec![("bo", Reason::ReviewRequested, Severity::Warning)]);
1974 assert_eq!(notices_[0].title, "ana asked you to review acme/rocket#7");
1975 let assigned = event("issue.assigned", Some("usr_ana"), json!({ "number": 3, "added": ["ana", "eve"] }));
1976 let notices_ = notices(&assigned, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
1977 assert_eq!(told(&notices_), vec![("eve", Reason::Assign, Severity::Info)]);
1978 assert_eq!(notices_[0].title, "ana assigned you to acme/rocket#3");
1979 // Both subscribe whoever they name, never g1t.
1980 assert_eq!(subscribes(&requested, Some(&pull())), vec![("bo".into(), Reason::ReviewRequested), ("ana".into(), Reason::ReviewRequested)]);
1981 assert_eq!(subscribes(&assigned, Some(&pull())), vec![("ana".into(), Reason::Assign), ("eve".into(), Reason::Assign)]);
1982 }
1983
1984 #[test]
1985 fn failures_go_to_whoever_answers_for_the_change() {
1986 let failed = event("checks.completed", None, json!({ "number": 7, "status": "failed" }));
1987 assert_eq!(
1988 told(&notices(&failed, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())),
1989 vec![("ana", Reason::CiActivity, Severity::Error)]
1990 );
1991 // g1t's change is the person's who asked for it, never g1t's.
1992 let notices_ = notices(&failed, "acme/rocket", &Actor::default(), Some(&g1t_pull()), &nobody());
1993 assert_eq!(told(&notices_), vec![("ana", Reason::CiActivity, Severity::Error)]);
1994 let errored = event("checks.completed", None, json!({ "number": 7, "status": "errored" }));
1995 assert_eq!(
1996 notices(&errored, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())[0].title,
1997 "Checks could not run on acme/rocket#7"
1998 );
1999 }
2000
2001 #[test]
2002 fn a_workflow_that_fails_tells_its_pull_requests_owner_even_if_they_pushed() {
2003 let failed = event("workflow.completed", Some("usr_ana"), json!({ "pull": 7, "workflow": "CI", "conclusion": "failure" }));
2004 let notices_ = notices(&failed, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
2005 assert_eq!(told(&notices_), vec![("ana", Reason::CiActivity, Severity::Error)]);
2006 assert_eq!(notices_[0].title, "CI failed on acme/rocket#7");
2007 // On a branch: whoever pushed.
2008 let pushed = event(
2009 "workflow.completed",
2010 Some("usr_bo"),
2011 json!({ "workflow": "Deploy", "path": ".g1t/workflows/deploy.yml", "conclusion": "failure", "ref": "refs/heads/main", "number": 12, "sha": "abcdef0123", "runId": "run_9" }),
2012 );
2013 let notices_ = notices(&pushed, "acme/rocket", &actor("usr_bo", "bo"), None, &nobody());
2014 assert_eq!(told(&notices_), vec![("bo", Reason::CiActivity, Severity::Error)]);
2015 assert_eq!(notices_[0].title, "Deploy failed on main in acme/rocket");
2016 assert_eq!(notices_[0].body, "Run 12 at abcdef0");
2017 // Nobody to tell when g1t pushed.
2018 assert!(notices(&pushed, "acme/rocket", &actor("g1t", "g1t"), None, &nobody()).is_empty());
2019 // Every failure of a workflow on a branch is one thread.
2020 let wanted = wants(&pushed).unwrap();
2021 let thread = thread_of(&pushed, &wanted, None);
2022 assert_eq!(thread.key, "rep_1/run/.g1t/workflows/deploy.yml@main");
2023 assert_eq!(thread.run_id.as_deref(), Some("run_9"));
2024 }
2025
2026 #[test]
2027 fn deployments_tell_whoever_answers_for_them_and_watchers() {
2028 let failed = event(
2029 "deployment.failed",
2030 Some("usr_bo"),
2031 json!({ "repoId": "rep_1", "projectId": "prj_1", "project": "rocket", "deploymentId": "dpl_1", "error": "The build failed.", "path": "/acme/rocket/deployments/dpl_1" }),
2032 );
2033 let audience = watched(&[("cy", WatchLevel::Custom, &["deployments"]), ("dee", WatchLevel::Custom, &["issues"]), ("eve", WatchLevel::All, &[])]);
2034 let notices_ = notices(&failed, "acme/rocket", &actor("usr_bo", "bo"), None, &audience);
2035 assert_eq!(
2036 told(&notices_),
2037 vec![
2038 ("bo", Reason::CiActivity, Severity::Error),
2039 ("cy", Reason::Subscribed, Severity::Error),
2040 ("eve", Reason::Subscribed, Severity::Error),
2041 ]
2042 );
2043 assert_eq!(notices_[0].title, "Production of rocket failed to deploy");
2044 assert_eq!(notices_[0].body, "The build failed.");
2045 let thread = thread_of(&failed, &wants(&failed).unwrap(), None);
2046 assert_eq!(thread.key, "rep_1/deploy/prj_1/production");
2047 assert_eq!(url(Some("acme/rocket"), thread.kind, None, None, thread.link.as_deref()), "/acme/rocket/deployments/dpl_1");
2048 // A success is news to the owner only after a failure.
2049 let live = event("deployment.succeeded", Some("usr_bo"), json!({ "repoId": "rep_1", "projectId": "prj_1", "project": "rocket", "commit": "abcdef0123" }));
2050 assert!(notices(&live, "acme/rocket", &actor("usr_bo", "bo"), None, &nobody()).is_empty());
2051 let recovered = event("deployment.succeeded", Some("usr_bo"), json!({ "repoId": "rep_1", "projectId": "prj_1", "project": "rocket", "recovered": true }));
2052 let notices_ = notices(&recovered, "acme/rocket", &actor("usr_bo", "bo"), None, &nobody());
2053 assert_eq!(told(&notices_), vec![("bo", Reason::CiActivity, Severity::Success)]);
2054 assert_eq!(notices_[0].title, "Production of rocket is live again");
2055 // A preview's is the pull request's owner's.
2056 let preview = event("deployment.failed", None, json!({ "repoId": "rep_1", "projectId": "prj_1", "number": 7, "triggeredBy": "g1t" }));
2057 assert_eq!(told(&notices(&preview, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())), vec![("ana", Reason::CiActivity, Severity::Error)]);
2058 assert_eq!(thread_of(&preview, &wants(&preview).unwrap(), Some(&pull())).key, "rep_1/deploy/prj_1/7");
2059 }
2060
2061 #[test]
2062 fn a_mirror_tells_only_the_people_its_settings_name() {
2063 let down = event(
2064 "mirror.unreachable",
2065 None,
2066 json!({
2067 "repoId": "rep_1", "repo": "acme/rocket", "remoteId": "rmt_1", "remote": "github.com/acme/rocket",
2068 "state": "standby", "title": "github.com/acme/rocket is not answering. acme/rocket keeps its copy.",
2069 "detail": "Take over acme/rocket to keep working on g1t until it is back.",
2070 "notify": ["ana"], "link": "/acme/rocket/settings/mirroring"
2071 }),
2072 );
2073 assert_eq!(wants(&down), Some(Wanted { repo_id: "rep_1".into(), number: None, comment_id: None }));
2074 // Watchers of everything hear nothing: only those named.
2075 let audience = watched(&[("eve", WatchLevel::All, &[])]);
2076 let notices_ = notices(&down, "acme/rocket", &Actor::default(), None, &audience);
2077 assert_eq!(told(&notices_), vec![("ana", Reason::StateChange, Severity::Warning)]);
2078 assert!(notices_[0].body.contains("Take over"));
2079 let thread = thread_of(&down, &wants(&down).unwrap(), None);
2080 assert_eq!(thread.key, "rep_1/mirror");
2081 assert_eq!(url(Some("acme/rocket"), thread.kind, None, None, thread.link.as_deref()), "/acme/rocket/settings/mirroring");
2082 // The banner-only default names nobody, so nobody is told.
2083 let quiet = event("mirror.unreachable", None, json!({ "repoId": "rep_1", "title": "x", "link": "/x" }));
2084 assert!(notices(&quiet, "acme/rocket", &Actor::default(), None, &audience).is_empty());
2085 let back = event("mirror.state_changed", None, json!({ "repoId": "rep_1", "state": "standby", "title": "handed back", "notify": ["ana"], "link": "/x" }));
2086 assert_eq!(told(&notices(&back, "acme/rocket", &Actor::default(), None, &audience)), vec![("ana", Reason::StateChange, Severity::Success)]);
2087 }
2088
2089 #[test]
2090 fn security_alerts_tell_the_pusher_the_owners_and_member_watchers() {
2091 let blocked = event(
2092 "secret_scanning_alert.created",
2093 None,
2094 json!({
2095 "repoId": "rep_1", "alertId": "sec_1", "alertType": "secret_scanning", "severity": "critical",
2096 "title": "An AWS access key in config/prod.env", "link": "/acme/rocket/security/secret-scanning/sec_1",
2097 "state": "open", "pusher": "bo", "notify": ["ana"], "members": ["ana", "bo", "cy"]
2098 }),
2099 );
2100 assert_eq!(wants(&blocked), Some(Wanted { repo_id: "rep_1".into(), number: None, comment_id: None }));
2101 // cy watches for security alerts and is a member; eve watches
2102 // everything but cannot see findings; dee watches issues only.
2103 let audience = watched(&[
2104 ("cy", WatchLevel::Custom, &["security"]),
2105 ("dee", WatchLevel::Custom, &["issues"]),
2106 ("eve", WatchLevel::All, &[]),
2107 ]);
2108 let notices_ = notices(&blocked, "acme/rocket", &Actor::default(), None, &audience);
2109 assert_eq!(
2110 told(&notices_),
2111 vec![
2112 ("bo", Reason::SecurityAlert, Severity::Error),
2113 ("ana", Reason::SecurityAlert, Severity::Error),
2114 ("cy", Reason::SecurityAlert, Severity::Error),
2115 ]
2116 );
2117 assert_eq!(notices_[0].title, "A push to acme/rocket was blocked: it adds a secret");
2118 assert_eq!(notices_[0].body, "An AWS access key in config/prod.env");
2119 let thread = thread_of(&blocked, &wants(&blocked).unwrap(), None);
2120 assert_eq!(thread.key, "rep_1/security/sec_1");
2121 assert_eq!(url(Some("acme/rocket"), thread.kind, None, None, thread.link.as_deref()), "/acme/rocket/security/secret-scanning/sec_1");
2122 // A bypass request goes to the reviewers named, never to watchers.
2123 let requested = event(
2124 "secret_scanning.bypass_requested",
2125 Some("usr_bo"),
2126 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" }),
2127 );
2128 let notices_ = notices(&requested, "acme/rocket", &actor("usr_bo", "bo"), None, &audience);
2129 assert_eq!(told(&notices_), vec![("ana", Reason::SecurityAlert, Severity::Warning)]);
2130 assert_eq!(notices_[0].title, "bo asked to bypass push protection in acme/rocket");
2131 // Fixes and dismissals are not news to the inbox.
2132 assert_eq!(wants(&event("code_scanning_alert.fixed", None, json!({ "repoId": "rep_1" }))), None);
2133 }
2134
2135 #[test]
2136 fn nobody_hears_of_what_they_did_themselves() {
2137 let merged = event("pull.merged", Some("usr_ana"), json!({ "number": 7 }));
2138 let by_ana = notices(&merged, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
2139 // The assignee still hears; ana, who merged, does not.
2140 assert_eq!(told(&by_ana), vec![("bo", Reason::StateChange, Severity::Success)]);
2141 let by_bo = notices(&merged, "acme/rocket", &actor("usr_bo", "bo"), Some(&pull()), &nobody());
2142 assert_eq!(told(&by_bo), vec![("ana", Reason::StateChange, Severity::Success)]);
2143 assert_eq!(by_bo[0].title, "bo merged acme/rocket#7");
2144 let by_queue = notices(&merged, "acme/rocket", &Actor::default(), Some(&pull()), &nobody());
2145 assert_eq!(by_queue[0].title, "acme/rocket#7 was merged");
2146 }
2147
2148 #[test]
2149 fn closing_and_reopening_tells_everyone_subscribed_and_watchers() {
2150 let closed = event("issue.closed", Some("usr_bo"), json!({ "number": 3, "reason": "completed" }));
2151 let issue = InboxSubject {
2152 kind: Some(SubjectKind::Issue),
2153 title: "Crash".into(),
2154 author: person("usr_cy", "cy"),
2155 assignees: vec!["bo".into()],
2156 ..InboxSubject::default()
2157 };
2158 let mut audience = subscribed(&[("eve", State::Subscribed, Some(Reason::Comment)), ("fay", State::Unsubscribed, None)]);
2159 audience.watchers = watched(&[("gus", WatchLevel::All, &[]), ("hal", WatchLevel::Custom, &["pulls"])]).watchers;
2160 let notices_ = notices(&closed, "acme/rocket", &actor("usr_bo", "bo"), Some(&issue), &audience);
2161 assert_eq!(
2162 told(&notices_),
2163 vec![
2164 ("cy", Reason::StateChange, Severity::Info),
2165 ("eve", Reason::StateChange, Severity::Info),
2166 ("gus", Reason::Subscribed, Severity::Info),
2167 ]
2168 );
2169 assert_eq!(notices_[0].title, "bo closed acme/rocket#3");
2170 let by_pull = event("issue.closed", None, json!({ "number": 3, "resolvedBy": 9 }));
2171 assert_eq!(notices(&by_pull, "acme/rocket", &Actor::default(), Some(&issue), &nobody())[0].title, "acme/rocket#3 was closed by #9");
2172 let reopened = event("issue.reopened", Some("usr_cy"), json!({ "number": 3 }));
2173 assert_eq!(told(&notices(&reopened, "acme/rocket", &actor("usr_cy", "cy"), Some(&issue), &nobody())), vec![("bo", Reason::StateChange, Severity::Info)]);
2174 }
2175
2176 #[test]
2177 fn opening_tells_who_it_names_and_watchers_of_its_kind() {
2178 let opened = event("pull.opened", Some("usr_ana"), json!({ "number": 7 }));
2179 let subject = InboxSubject { reviewers: vec!["cy".into(), "g1t".into()], ..pull() };
2180 let audience = watched(&[("bo", WatchLevel::All, &[]), ("dee", WatchLevel::Custom, &["issues"]), ("eve", WatchLevel::Custom, &["pulls"]), ("fay", WatchLevel::Participating, &[])]);
2181 let notices_ = notices(&opened, "acme/rocket", &actor("usr_ana", "ana"), Some(&subject), &audience);
2182 assert_eq!(
2183 told(&notices_),
2184 vec![
2185 ("bo", Reason::Assign, Severity::Info),
2186 ("cy", Reason::ReviewRequested, Severity::Warning),
2187 ("eve", Reason::Subscribed, Severity::Info),
2188 ]
2189 );
2190 assert_eq!(notices_[2].title, "ana opened acme/rocket#7");
2191 }
2192
2193 #[test]
2194 fn ignoring_a_repository_silences_it_and_unsubscribing_keeps_only_what_is_asked() {
2195 let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
2196 let on = comment(person("usr_bo", "bo"), "@cy have a look", &["cy"]);
2197 let audience = watched(&[("cy", WatchLevel::Ignore, &[]), ("ana", WatchLevel::All, &[])]);
2198 assert!(notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &audience).iter().all(|notice| notice.username != "cy"));
2199 let audience = subscribed(&[("cy", State::Unsubscribed, None), ("ana", State::Unsubscribed, None)]);
2200 let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &audience);
2201 // cy was mentioned, which is asked of them; ana unsubscribed from the conversation.
2202 assert_eq!(told(&notices_), vec![("cy", Reason::Mention, Severity::Info)]);
2203 }
2204
2205 #[test]
2206 fn g1t_finishing_or_reviewing_tells_the_person_it_worked_for() {
2207 // The agent acts as the person it works for: still an outcome they hear of.
2208 let ready = event("pull.ready", Some("usr_ana"), json!({ "number": 7 }));
2209 let notices_ = notices(&ready, "acme/rocket", &actor("usr_ana", "ana"), Some(&g1t_pull()), &nobody());
2210 assert_eq!(told(&notices_), vec![("ana", Reason::Author, Severity::Success)]);
2211 assert_eq!(notices_[0].title, "g1t finished acme/rocket#7");
2212 // A person's draft marked ready is not news to them.
2213 assert!(notices(&ready, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody()).is_empty());
2214
2215 let approve = event("review.completed", None, json!({ "number": 7, "verdict": "approve" }));
2216 assert_eq!(
2217 told(&notices(&approve, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())),
2218 vec![("ana", Reason::Author, Severity::Success)]
2219 );
2220 let changes = event("review.completed", None, json!({ "number": 7, "verdict": "request_changes" }));
2221 assert_eq!(
2222 told(&notices(&changes, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())),
2223 vec![("ana", Reason::Author, Severity::Info)]
2224 );
2225 }
2226
2227 #[test]
2228 fn comments_tell_those_mentioned_then_everyone_subscribed_never_the_writer() {
2229 let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
2230 let on = comment(person("usr_bo", "bo"), "@cy @bo have a look", &["cy", "bo", "g1t"]);
2231 let audience = subscribed(&[("eve", State::Subscribed, Some(Reason::Comment)), ("fay", State::Subscribed, Some(Reason::Manual))]);
2232 let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &audience);
2233 assert_eq!(
2234 told(&notices_),
2235 vec![
2236 ("cy", Reason::Mention, Severity::Info),
2237 ("ana", Reason::Author, Severity::Info),
2238 ("eve", Reason::Comment, Severity::Info),
2239 ("fay", Reason::Manual, Severity::Info),
2240 ]
2241 );
2242 assert_eq!(notices_[0].title, "bo mentioned you on acme/rocket#7");
2243 assert_eq!(notices_[1].title, "bo commented on acme/rocket#7");
2244 assert_eq!(notices_[0].body, "@cy @bo have a look");
2245 // Mentioned and the owner: told once, as mentioned.
2246 let on = comment(person("usr_bo", "bo"), "@ana", &["ana"]);
2247 assert_eq!(
2248 told(&notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &nobody())),
2249 vec![("ana", Reason::Mention, Severity::Info)]
2250 );
2251 // The owner's own comment tells the assignee, not the owner.
2252 let on = comment(person("usr_ana", "ana"), "thanks", &[]);
2253 assert_eq!(
2254 told(&notices(&created, "acme/rocket", &actor("usr_ana", "ana"), Some(&on), &nobody())),
2255 vec![("bo", Reason::Assign, Severity::Info)]
2256 );
2257 // Something that happened, not something written, tells nobody.
2258 let mut on = comment(person("usr_bo", "bo"), "assigned cy", &[]);
2259 on.comment.as_mut().unwrap().event = true;
2260 assert!(notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &nobody()).is_empty());
2261 // An approval is good news.
2262 let mut on = comment(person("usr_cy", "cy"), "", &[]);
2263 on.comment.as_mut().unwrap().verdict = Some("approve".into());
2264 let notices_ = notices(&created, "acme/rocket", &actor("usr_cy", "cy"), Some(&on), &nobody());
2265 assert_eq!(told(&notices_)[0], ("ana", Reason::Author, Severity::Success));
2266 assert_eq!(notices_[0].body, "Add the inbox");
2267 // Writing and being mentioned subscribe, never g1t.
2268 let on = comment(person("usr_bo", "bo"), "@cy", &["cy"]);
2269 assert_eq!(subscribes(&created, Some(&on)), vec![("bo".into(), Reason::Comment), ("cy".into(), Reason::Mention)]);
2270 let on = comment(person(AGENT_ID, "g1t"), "done", &[]);
2271 assert!(subscribes(&created, Some(&on)).is_empty());
2272 }
2273
2274 #[test]
2275 fn issues_and_pull_requests_are_one_thread_each() {
2276 let created = event("comment.created", Some("usr_bo"), json!({ "repoId": "rep_1", "number": 7, "commentId": "cmt_1" }));
2277 let wanted = wants(&created).unwrap();
2278 let thread = thread_of(&created, &wanted, Some(&pull()));
2279 assert_eq!(thread.key, "rep_1#7");
2280 assert_eq!(thread.kind, Some(SubjectKind::Pull));
2281 let merged = event("pull.merged", None, json!({ "repoId": "rep_1", "number": 7 }));
2282 assert_eq!(thread_of(&merged, &wants(&merged).unwrap(), Some(&pull())).key, thread.key);
2283 }
2284
2285 #[test]
2286 fn items_link_to_what_they_are_about() {
2287 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Pull), Some(7), None, None), "/acme/rocket/pull/7");
2288 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Issue), Some(3), None, None), "/acme/rocket/issues/3");
2289 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Run), None, Some("run_1"), None), "/acme/rocket/actions/runs/run_1");
2290 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Deploy), None, None, Some("/acme/site/deployments/dpl_1")), "/acme/site/deployments/dpl_1");
2291 assert_eq!(url(Some("acme/rocket"), None, None, None, None), "/acme/rocket");
2292 assert_eq!(url(None, None, None, None, None), "/inbox");
2293 }
2294
2295 #[test]
2296 fn long_titles_are_cut_to_a_line() {
2297 let long = "x".repeat(400);
2298 assert_eq!(clip(&long, MAX_TITLE).chars().count(), MAX_TITLE);
2299 assert_eq!(clip(" short ", MAX_TITLE), "short");
2300 }
2301
2302 #[test]
2303 fn marks_change_only_what_they_say() {
2304 assert_eq!(mark_change(InboxMark::Read), "read_at = COALESCE(read_at, ?1)");
2305 assert!(mark_change(InboxMark::Done).contains("done_at = ?1"));
2306 assert_eq!(mark_change(InboxMark::Unsnooze), "snoozed_until = NULL");
2307 assert!(ranked(&ListInboxArgs::default()));
2308 assert!(!ranked(&ListInboxArgs { reason: Some(Reason::Mention), ..ListInboxArgs::default() }));
2309 assert!(!ranked(&ListInboxArgs { participating: true, ..ListInboxArgs::default() }));
2310 }
2311
2312 #[test]
2313 fn filters_bind_in_order_after_what_is_bound() {
2314 let a = ListInboxArgs {
2315 reason: Some(Reason::Mention),
2316 repo_id: Some("rep_1".into()),
2317 participating: true,
2318 unread: true,
2319 ..ListInboxArgs::default()
2320 };
2321 // Two values come first: the username and now.
2322 let (conditions, values) = filters(&a, 2);
2323 assert_eq!(
2324 conditions,
2325 vec!["reason = ?3", "repo_id = ?4", "reason NOT IN ('manual', 'subscribed')", "read_at IS NULL"]
2326 );
2327 assert_eq!(values, vec!["mention", "rep_1"]);
2328 }
2329
2330 #[test]
2331 fn the_bump_keeps_whatever_is_most_urgent_while_unread() {
2332 assert!(BUMP.contains("WHERE NOT EXISTS (SELECT 1 FROM inbox_activity WHERE event_id = ?4 AND username = ?2)"));
2333 assert!(BUMP.contains("ON CONFLICT (username, thread) DO UPDATE"));
2334 assert!(BUMP.contains("activity = inbox_items.activity + 1"));
2335 // Its numbered parameters run from 1 to 19.
2336 for at in 1..=19 {
2337 assert!(BUMP.contains(&format!("?{at}")), "?{at}");
2338 }
2339 assert!(!BUMP.contains("?20"));
2340 }
2341
2342 #[test]
2343 fn token_approvals_tell_the_people_named_about_no_repository() {
2344 let mut asked = event(
2345 "token.approval_requested",
2346 Some("usr_ana"),
2347 json!({ "workspace": "acme", "tokenId": "tok_1", "notify": ["Bo", "cy", "bo", "g1t"], "title": "ana asks", "body": "ci: contents: write", "link": "/acme/-/settings/tokens" }),
2348 );
2349 asked.repo_id = None;
2350 let (thread, told) = workspace_notices(&asked, None);
2351 assert_eq!(thread, "workspace:acme/token/tok_1");
2352 let mut names: Vec<&str> = told.iter().map(|notice| notice.username.as_str()).collect();
2353 names.sort();
2354 assert_eq!(names, ["bo", "cy"], "each once, lowercased, never g1t");
2355 assert!(told.iter().all(|notice| notice.reason == Reason::ReviewRequested && notice.severity == Severity::Warning));
2356 assert!(wants(&asked).is_none(), "about no repository");
2357 let reviewed = event("token.approval_reviewed", None, json!({ "workspace": "acme", "tokenId": "tok_1", "notify": ["ana"], "title": "approved" }));
2358 let (_, told) = workspace_notices(&reviewed, None);
2359 assert_eq!(told[0].reason, Reason::Author);
2360 assert!(WORKSPACE_EVENTS.contains(&"token.approval_reviewed"));
2361 }
2362}