Skip to content

g1t/services/events/src/inbox.rs

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