Skip to content

g1t/services/events/src/inbox.rs

1,881 lines83,923 bytesCodeBlame

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

Inbox: the events service tells people what needs them as events arrive1//! The inbox, kept beside the event log. See `g1t_contracts::inbox`.
2//!
Inbox: threads, reasons, subscriptions and watching3//! As each batch arrives from the bus, the events that tell someone are
4//! read against what they name (the work service's `inbox_subject`), who
5//! subscribes to it and who watches its repository (subscriptions.rs), and
6//! each person told gets their thread about it brought to the top, unread,
7//! with a line added to its history. Who is told is worked out in
8//! [`notices`], from the event, its subject and that audience alone:
Inbox: the events service tells people what needs them as events arrive9//!
Inbox: threads, reasons, subscriptions and watching10//! | Event | Who | Reason | Severity |
11//! | --- | --- | --- | --- |
12//! | `agent.asked`, `pull.stalled` | the pull request's owner, and its issue's owner and assignees | agent | warning |
13//! | `pull.review_requested` | the reviewers asked | 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, then everyone subscribed | mention, or why they are subscribed | info (success for an approval) |
23//! | `issue.opened`, `pull.opened` | whoever was assigned or asked to review; watchers | assign, review_requested, subscribed | info |
Inbox: the events service tells people what needs them as events arrive24//!
Inbox: threads, reasons, subscriptions and watching25//! Watchers of a repository at `all` (or `custom`, for the kinds they
26//! chose) hear of every issue and pull request opened, commented on,
27//! closed, reopened or merged, and of deployments. Anyone who ignores the
28//! thread or the repository hears of nothing on it; anyone who
29//! unsubscribed hears only of what is asked of them.
30//!
31//! Nobody is told of what they did themselves, though an outcome they set
32//! off (checks, a workflow, a deployment, g1t's work) is theirs to hear of.
33//! g1t is never told. A failure here is logged and the batch goes on: the
34//! bus never waits on the inbox, so an item can be missed, but nothing else
35//! is held up.
Inbox: the events service tells people what needs them as events arrive36
37use std::collections::{HashMap, HashSet};
38
39use g1t_contracts::credentials::Principal;
40use g1t_contracts::events::{Event, WorkspaceRenamed};
41use g1t_contracts::identity::{AGENT_ID, UsernamesArgs};
42use g1t_contracts::inbox::*;
43use g1t_contracts::repos::{PathByIdArgs, ReadableArgs, Repo, RepoPath};
44use g1t_contracts::time::rfc3339;
45use g1t_contracts::{new_id, system};
46use g1t_kit::now_ms;
47use serde::Deserialize;
48use worker::wasm_bindgen::JsValue;
Inbox: threads, reasons, subscriptions and watching49use worker::{D1Database, D1PreparedStatement, Fetcher, Result};
50
51use crate::subscriptions::{self, Audience};
Inbox: the events service tells people what needs them as events arrive52
53/// Items marked done are kept this long, then removed.
54pub const DONE_DAYS: u32 = 30;
55/// No item is kept longer than this, unless it was saved.
56pub const MAX_DAYS: u32 = 180;
57/// Unread warnings shown ahead of everything else on the first page.
58const MAX_RANKED: u32 = 20;
59const MAX_TITLE: usize = 200;
60const MAX_BODY: usize = 300;
61
62/// What the inbox asks about an event before deciding who is told.
63#[derive(Debug, PartialEq, Eq)]
64pub struct Wanted {
65 pub repo_id: String,
66 /// The issue or pull request, when the event names one.
67 pub number: Option<u32>,
68 pub comment_id: Option<String>,
69}
70
71/// One person to tell, and what.
72#[derive(Debug, PartialEq, Eq)]
73pub struct Notice {
74 pub username: String,
Inbox: threads, reasons, subscriptions and watching75 pub reason: Reason,
Inbox: the events service tells people what needs them as events arrive76 pub severity: Severity,
77 pub title: String,
78 pub body: String,
79}
80
81/// Who did it, as far as is known.
82#[derive(Debug, Default)]
83pub struct Actor {
84 pub id: Option<String>,
85 pub username: Option<String>,
86}
87
88impl Actor {
89 fn is(&self, person: &Principal) -> bool {
90 self.id.as_deref().is_some_and(|id| id == person.id) || self.is_named(&person.username)
91 }
92
93 fn is_named(&self, username: &str) -> bool {
94 self.username.as_deref().is_some_and(|name| name.eq_ignore_ascii_case(username))
95 }
Inbox: threads, reasons, subscriptions and watching96
97 /// A person's name, for a title: never g1t's ids.
98 fn name(&self) -> Option<&str> {
99 self.username.as_deref().filter(|name| !is_g1t(name))
100 }
Inbox: the events service tells people what needs them as events arrive101}
102
Inbox: threads, reasons, subscriptions and watching103pub(crate) fn is_g1t(username: &str) -> bool {
Inbox: the events service tells people what needs them as events arrive104 username.eq_ignore_ascii_case(system::USERNAME) || username.eq_ignore_ascii_case("g1t-agent")
105}
106
Inbox: threads, reasons, subscriptions and watching107pub(crate) fn is_g1t_id(id: &str) -> bool {
Inbox: the events service tells people what needs them as events arrive108 system::is_system_id(id) || id == AGENT_ID
109}
110
Inbox: threads, reasons, subscriptions and watching111/// Events that end what an agent was waiting on a person for: it picked
112/// back up, its head moved, a merge was asked for, or it is over.
113const RESUMES: [&str; 5] = ["pull.resumed", "pull.updated", "pull.merge_requested", "pull.merged", "pull.closed"];
114
Inbox: the events service tells people what needs them as events arrive115/// The issue or pull request an event names, and what to read for it.
116/// None for events the inbox does not tell anyone of.
117pub fn wants(event: &Event) -> Option<Wanted> {
118 let data = &event.data;
119 let text = |key: &str| data[key].as_str().filter(|value| !value.is_empty()).map(str::to_owned);
120 let number = |key: &str| data[key].as_u64().and_then(|n| u32::try_from(n).ok());
121 let repo_id = text("repoId").or_else(|| event.repo_id.clone())?;
122 let on = |number: Option<u32>, comment_id: Option<String>| {
123 Some(Wanted {
124 repo_id: repo_id.clone(),
125 number,
126 comment_id,
127 })
128 };
129 match event.kind.as_str() {
Inbox: threads, reasons, subscriptions and watching130 "agent.asked" | "pull.stalled" | "pull.merged" | "pull.closed" | "pull.ready" | "pull.opened"
131 | "pull.review_requested" | "pull.assigned" | "issue.opened" | "issue.closed" | "issue.reopened"
132 | "issue.assigned" => on(Some(number("number")?), None),
Inbox: the events service tells people what needs them as events arrive133 "checks.completed" => match data["status"].as_str() {
134 Some("failed" | "errored") => on(Some(number("number")?), None),
135 _ => None,
136 },
137 "review.completed" => match data["verdict"].as_str() {
138 Some("approve" | "request_changes") => on(Some(number("number")?), None),
139 _ => None,
140 },
141 "workflow.completed" => match data["conclusion"].as_str() {
142 Some("failure") => on(number("pull"), None),
143 _ => None,
144 },
Inbox: threads, reasons, subscriptions and watching145 "deployment.failed" | "deployment.succeeded" => on(number("number"), None),
Inbox: the events service tells people what needs them as events arrive146 "comment.created" => on(Some(number("number")?), Some(text("commentId")?)),
147 _ => None,
148 }
149}
150
Inbox: threads, reasons, subscriptions and watching151/// Where an event's items go: which thread, what it is, and where it is.
152#[derive(Clone, Debug, PartialEq, Eq)]
153pub struct Thread {
154 pub key: String,
155 pub kind: Option<SubjectKind>,
156 pub number: Option<u32>,
157 pub run_id: Option<String>,
158 /// A path, for a subject with a page of its own.
159 pub link: Option<String>,
160}
161
162/// The thread key of an issue or pull request.
163pub fn numbered_thread(repo_id: &str, number: u32) -> String {
164 format!("{repo_id}#{number}")
165}
166
167/// The thread an event's items go to.
168pub fn thread_of(event: &Event, wanted: &Wanted, subject: Option<&InboxSubject>) -> Thread {
169 let data = &event.data;
170 let text = |key: &str| data[key].as_str().filter(|value| !value.is_empty()).map(str::to_owned);
171 match event.kind.as_str() {
172 // A project's production, or one pull request's preview.
173 "deployment.failed" | "deployment.succeeded" => {
174 let which = wanted.number.map_or_else(|| "production".to_owned(), |number| number.to_string());
175 Thread {
176 key: format!("{}/deploy/{}/{which}", wanted.repo_id, text("projectId").unwrap_or_default()),
177 kind: Some(SubjectKind::Deploy),
178 number: wanted.number,
179 run_id: text("deploymentId"),
180 link: text("path"),
181 }
182 }
183 // A workflow on a branch: its next failure bumps the same thread.
184 "workflow.completed" if subject.is_none() => {
185 let branch = text("ref").unwrap_or_default();
186 let branch = branch.strip_prefix("refs/heads/").unwrap_or(&branch).to_owned();
187 let workflow = text("path").or_else(|| text("workflow")).unwrap_or_default();
188 Thread {
189 key: format!("{}/run/{workflow}@{branch}", wanted.repo_id),
190 kind: Some(SubjectKind::Run),
191 number: None,
192 run_id: text("runId"),
193 link: None,
194 }
195 }
196 _ => match wanted.number {
197 Some(number) => Thread {
198 key: numbered_thread(&wanted.repo_id, number),
199 kind: subject.and_then(|subject| subject.kind),
200 number: Some(number),
201 run_id: None,
202 link: None,
203 },
204 None => Thread {
205 key: format!("event/{}", event.id),
206 kind: None,
207 number: None,
208 run_id: None,
209 link: None,
210 },
211 },
212 }
213}
214
215/// The issue or pull request whose agent threads an event closes: what an
216/// agent was waiting on a person for is over.
217pub fn resolves(event: &Event) -> Option<String> {
218 if !RESUMES.contains(&event.kind.as_str()) {
219 return None;
220 }
221 let repo_id = event.data["repoId"].as_str().map(str::to_owned).or_else(|| event.repo_id.clone())?;
222 let number = event.data["number"].as_u64().and_then(|n| u32::try_from(n).ok())?;
223 Some(numbered_thread(&repo_id, number))
224}
225
226/// Collects who is told, each once with the most specific reason, never
227/// the actor, never g1t, and never anyone ignoring the thread.
Inbox: the events service tells people what needs them as events arrive228struct Told<'a> {
229 actor: &'a Actor,
Inbox: threads, reasons, subscriptions and watching230 audience: &'a Audience,
Inbox: the events service tells people what needs them as events arrive231 notices: Vec<Notice>,
232}
233
234impl Told<'_> {
Inbox: threads, reasons, subscriptions and watching235 /// `asked`: something asked of the person directly, told even when
236 /// they unsubscribed from the thread.
237 fn tell(&mut self, username: &str, reason: Reason, severity: Severity, title: &str, body: &str, asked: bool) {
Inbox: the events service tells people what needs them as events arrive238 let username = username.trim().trim_start_matches('@').to_lowercase();
239 if username.is_empty()
240 || is_g1t(&username)
241 || self.actor.is_named(&username)
Inbox: threads, reasons, subscriptions and watching242 || self.audience.ignores(&username)
243 || (!asked && self.audience.unsubscribed(&username))
Inbox: the events service tells people what needs them as events arrive244 {
245 return;
246 }
Inbox: threads, reasons, subscriptions and watching247 let notice = Notice {
Inbox: the events service tells people what needs them as events arrive248 username,
249 reason,
250 severity,
251 title: clip(title, MAX_TITLE),
252 body: clip(body, MAX_BODY),
Inbox: threads, reasons, subscriptions and watching253 };
254 match self.notices.iter_mut().find(|told| told.username == notice.username) {
255 Some(told) if notice.reason.rank() < told.reason.rank() => *told = notice,
256 Some(_) => {}
257 None => self.notices.push(notice),
258 }
Inbox: the events service tells people what needs them as events arrive259 }
260
261 /// A person known by id as well as name, such as an author.
Inbox: threads, reasons, subscriptions and watching262 fn tell_person(&mut self, person: &Principal, reason: Reason, severity: Severity, title: &str, body: &str, asked: bool) {
Inbox: the events service tells people what needs them as events arrive263 if is_g1t_id(&person.id) || self.actor.is(person) {
264 return;
265 }
Inbox: threads, reasons, subscriptions and watching266 self.tell(&person.username, reason, severity, title, body, asked);
267 }
268
269 /// Everyone subscribed to an issue or pull request: its owner and
270 /// author, its assignees and reviewers, and whoever subscribed by
271 /// commenting, being mentioned or by hand. `reason` overrides why
272 /// each is told, as a state change does.
273 fn tell_subscribed(&mut self, subject: &InboxSubject, reason: Option<Reason>, severity: Severity, title: &str, body: &str) {
274 let why = |own: Reason| reason.unwrap_or(own);
275 self.tell_person(subject.owner(), why(Reason::Author), severity, title, body, false);
276 self.tell_person(&subject.author, why(Reason::Author), severity, title, body, false);
277 for name in &subject.assignees {
278 self.tell(name, why(Reason::Assign), severity, title, body, false);
279 }
280 for name in &subject.reviewers {
281 self.tell(name, why(Reason::ReviewRequested), severity, title, body, false);
282 }
283 let subscribed: Vec<(String, Reason)> = self.audience.subscribed().collect();
284 for (name, own) in subscribed {
285 self.tell(&name, why(own), severity, title, body, false);
286 }
Inbox: the events service tells people what needs them as events arrive287 }
Inbox: threads, reasons, subscriptions and watching288
289 /// Everyone watching the repository for this kind of activity.
290 fn tell_watchers(&mut self, kind: &str, severity: Severity, title: &str, body: &str) {
291 let watching: Vec<String> = self.audience.watching(kind).collect();
292 for name in watching {
293 self.tell(&name, Reason::Subscribed, severity, title, body, false);
294 }
295 }
Inbox: the events service tells people what needs them as events arrive296}
297
298fn clip(text: &str, max: usize) -> String {
299 let text = text.trim();
300 if text.chars().count() <= max {
301 return text.to_owned();
302 }
303 let cut: String = text.chars().take(max - 1).collect();
304 format!("{}…", cut.trim_end())
305}
306
Inbox: threads, reasons, subscriptions and watching307/// The kind of activity a watcher chooses: `issues` or `pulls`.
308fn activity_of(subject: &InboxSubject) -> &'static str {
309 match subject.kind {
310 Some(SubjectKind::Pull) => "pulls",
311 _ => "issues",
312 }
313}
314
315/// The names in an event's list field, such as the reviewers just asked.
316fn names(data: &serde_json::Value, key: &str) -> Vec<String> {
317 data[key]
318 .as_array()
319 .map(|names| names.iter().filter_map(|name| name.as_str().map(str::to_owned)).collect())
320 .unwrap_or_default()
321}
322
323/// Who is told of `event`, in `repo` (`owner/name`), given what it names
324/// and who follows it. `actor` is who caused it; outcomes nobody chose
325/// (checks, workflows, deployments, a review, an agent finishing or
326/// stopping) are told whoever caused them.
327pub fn notices(event: &Event, repo: &str, actor: &Actor, subject: Option<&InboxSubject>, audience: &Audience) -> Vec<Notice> {
Inbox: the events service tells people what needs them as events arrive328 let nobody = Actor::default();
329 let data = &event.data;
330 // The issue or pull request, as titles name it: `acme/rocket#12`.
331 let at = match data["number"].as_u64().or(data["pull"].as_u64()) {
332 Some(number) if subject.is_some() => format!("{repo}#{number}"),
333 _ => repo.to_owned(),
334 };
335 let outcome = matches!(
336 event.kind.as_str(),
Inbox: threads, reasons, subscriptions and watching337 "checks.completed"
338 | "workflow.completed"
339 | "review.completed"
340 | "pull.ready"
341 | "agent.asked"
342 | "pull.stalled"
343 | "deployment.failed"
344 | "deployment.succeeded"
Inbox: the events service tells people what needs them as events arrive345 );
346 let mut told = Told {
347 actor: if outcome { &nobody } else { actor },
Inbox: threads, reasons, subscriptions and watching348 audience,
Inbox: the events service tells people what needs them as events arrive349 notices: Vec::new(),
350 };
Inbox: threads, reasons, subscriptions and watching351 // Who did it, as titles name them.
352 let who = actor.name().unwrap_or("g1t");
Inbox: the events service tells people what needs them as events arrive353
354 match (event.kind.as_str(), subject) {
Inbox: threads, reasons, subscriptions and watching355 ("agent.asked" | "pull.stalled", Some(pull)) => {
356 let (title, body) = if event.kind == "agent.asked" {
357 (format!("An agent is waiting on {at}"), pull.title.clone())
358 } else {
359 let detail = data["detail"].as_str().map(str::trim).filter(|detail| !detail.is_empty());
360 (format!("g1t stopped on {at} and needs you"), detail.unwrap_or(&pull.title).to_owned())
361 };
362 told.tell_person(pull.owner(), Reason::Agent, Severity::Warning, &title, &body, true);
Inbox: the events service tells people what needs them as events arrive363 if let Some(issue) = &pull.issue {
Inbox: threads, reasons, subscriptions and watching364 told.tell_person(issue.owner(), Reason::Agent, Severity::Warning, &title, &body, true);
Inbox: the events service tells people what needs them as events arrive365 for name in &issue.assignees {
Inbox: threads, reasons, subscriptions and watching366 told.tell(name, Reason::Agent, Severity::Warning, &title, &body, true);
Inbox: the events service tells people what needs them as events arrive367 }
368 }
369 }
Inbox: threads, reasons, subscriptions and watching370 ("pull.review_requested", Some(pull)) => {
371 let title = format!("{who} asked you to review {at}");
372 for name in names(data, "reviewers") {
373 told.tell(&name, Reason::ReviewRequested, Severity::Warning, &title, &pull.title, true);
374 }
375 }
376 ("issue.assigned" | "pull.assigned", Some(on)) => {
377 let title = format!("{who} assigned you to {at}");
378 for name in names(data, "added") {
379 told.tell(&name, Reason::Assign, Severity::Info, &title, &on.title, true);
380 }
381 }
382 ("issue.opened" | "pull.opened", Some(on)) => {
383 let title = format!("{who} opened {at}");
384 for name in &on.assignees {
385 told.tell(name, Reason::Assign, Severity::Info, &format!("{who} assigned you to {at}"), &on.title, true);
386 }
387 for name in &on.reviewers {
388 told.tell(name, Reason::ReviewRequested, Severity::Warning, &format!("{who} asked you to review {at}"), &on.title, true);
389 }
390 told.tell_watchers(activity_of(on), Severity::Info, &title, &on.title);
391 }
Inbox: the events service tells people what needs them as events arrive392 ("checks.completed", Some(pull)) => {
393 let title = match data["status"].as_str() {
Inbox: threads, reasons, subscriptions and watching394 Some("errored") => format!("Checks could not run on {at}"),
395 _ => format!("Checks failed on {at}"),
Inbox: the events service tells people what needs them as events arrive396 };
Inbox: threads, reasons, subscriptions and watching397 told.tell_person(pull.owner(), Reason::CiActivity, Severity::Error, &title, &pull.title, true);
Inbox: the events service tells people what needs them as events arrive398 }
399 ("workflow.completed", subject) => {
400 let workflow = data["workflow"].as_str().filter(|name| !name.is_empty()).unwrap_or("A workflow");
401 match subject {
402 Some(pull) => {
Inbox: threads, reasons, subscriptions and watching403 let title = format!("{workflow} failed on {at}");
404 told.tell_person(pull.owner(), Reason::CiActivity, Severity::Error, &title, &pull.title, true);
Inbox: the events service tells people what needs them as events arrive405 }
406 None => {
407 // Not on a pull request: whoever pushed the commit it ran on.
408 let branch = data["ref"].as_str().unwrap_or_default();
409 let branch = branch.strip_prefix("refs/heads/").unwrap_or(branch);
410 let title = format!("{workflow} failed on {branch} in {repo}");
411 let body = format!("Run {} at {}", data["number"], short(data["sha"].as_str().unwrap_or_default()));
412 if let Some(id) = &actor.id
413 && !is_g1t_id(id)
414 && let Some(name) = &actor.username
415 {
Inbox: threads, reasons, subscriptions and watching416 told.tell(name, Reason::CiActivity, Severity::Error, &title, &body, true);
417 }
418 }
419 }
420 }
421 ("deployment.failed" | "deployment.succeeded", subject) => {
422 let failed = event.kind == "deployment.failed";
423 let project = data["project"].as_str().filter(|name| !name.is_empty()).unwrap_or(repo);
424 let what = match subject {
425 Some(_) => format!("The preview of {at}"),
426 None => format!("Production of {project}"),
427 };
428 let (title, severity) = if failed {
429 (format!("{what} failed to deploy"), Severity::Error)
430 } else {
431 (format!("{what} is live"), Severity::Success)
432 };
433 let body = match (failed, data["error"].as_str().map(str::trim).filter(|error| !error.is_empty())) {
434 (true, Some(error)) => error.to_owned(),
435 _ => match subject {
436 Some(pull) => pull.title.clone(),
437 None => format!("Commit {}", short(data["commit"].as_str().unwrap_or_default())),
438 },
439 };
440 // A success is news to whoever answers for it only after a failure.
441 if failed || data["recovered"].as_bool() == Some(true) {
442 let title = if failed { title.clone() } else { format!("{what} is live again") };
443 match subject {
444 Some(pull) => told.tell_person(pull.owner(), Reason::CiActivity, severity, &title, &body, true),
445 None => {
446 let pusher = actor
447 .id
448 .as_deref()
449 .filter(|id| !is_g1t_id(id))
450 .and(actor.username.as_deref())
451 .or_else(|| data["triggeredBy"].as_str());
452 if let Some(name) = pusher {
453 told.tell(name, Reason::CiActivity, severity, &title, &body, true);
454 }
Inbox: the events service tells people what needs them as events arrive455 }
456 }
457 }
Inbox: threads, reasons, subscriptions and watching458 told.tell_watchers("deployments", severity, &title, &body);
Inbox: the events service tells people what needs them as events arrive459 }
460 ("review.completed", Some(pull)) => {
Inbox: threads, reasons, subscriptions and watching461 let (title, severity) = match data["verdict"].as_str() {
462 Some("approve") => (format!("g1t approved {at}"), Severity::Success),
463 _ => (format!("g1t asked for changes on {at}"), Severity::Info),
Inbox: the events service tells people what needs them as events arrive464 };
Inbox: threads, reasons, subscriptions and watching465 told.tell_person(pull.owner(), Reason::Author, severity, &title, &pull.title, true);
Inbox: the events service tells people what needs them as events arrive466 }
467 ("pull.ready", Some(pull)) => {
468 // A change g1t made is ready: the agent's run is over.
469 if let Some(owner) = pull.requested_by.as_ref().filter(|_| is_g1t_id(&pull.author.id) || is_g1t(&pull.author.username)) {
Inbox: threads, reasons, subscriptions and watching470 let title = format!("g1t finished {at}");
471 told.tell_person(owner, Reason::Author, Severity::Success, &title, &pull.title, true);
Inbox: the events service tells people what needs them as events arrive472 }
473 }
Inbox: threads, reasons, subscriptions and watching474 ("pull.merged" | "pull.closed" | "issue.closed" | "issue.reopened", Some(on)) => {
475 let (verb, severity) = match event.kind.as_str() {
476 "pull.merged" => ("merged", Severity::Success),
477 "issue.reopened" => ("reopened", Severity::Info),
478 _ => ("closed", Severity::Info),
479 };
480 let title = match (actor.name(), data["resolvedBy"].as_u64()) {
481 (_, Some(pull)) if event.kind == "issue.closed" => format!("{at} was closed by #{pull}"),
482 (Some(name), _) => format!("{name} {verb} {at}"),
483 (None, _) => format!("{at} was {verb}"),
Inbox: the events service tells people what needs them as events arrive484 };
Inbox: threads, reasons, subscriptions and watching485 told.tell_subscribed(on, Some(Reason::StateChange), severity, &title, &on.title);
486 told.tell_watchers(activity_of(on), severity, &title, &on.title);
Inbox: the events service tells people what needs them as events arrive487 }
488 ("comment.created", Some(on)) => {
489 let Some(comment) = on.comment.as_ref().filter(|comment| !comment.event) else {
490 return Vec::new();
491 };
492 // Whoever wrote it is the actor, whatever the event says.
493 let writer = Actor {
494 id: Some(comment.author.id.clone()),
495 username: Some(comment.author.username.clone()),
496 };
497 told.actor = &writer;
498 let who = &comment.author.username;
499 let body = if comment.excerpt.is_empty() { &on.title } else { &comment.excerpt };
500 for name in &comment.mentions {
Inbox: threads, reasons, subscriptions and watching501 let title = format!("{who} mentioned you on {at}");
502 told.tell(name, Reason::Mention, Severity::Info, &title, body, true);
Inbox: the events service tells people what needs them as events arrive503 }
Inbox: threads, reasons, subscriptions and watching504 let (severity, title) = match comment.verdict.as_deref() {
505 Some("approve") => (Severity::Success, format!("{who} approved {at}")),
506 Some("request_changes") => (Severity::Info, format!("{who} asked for changes on {at}")),
507 _ => (Severity::Info, format!("{who} commented on {at}")),
Inbox: the events service tells people what needs them as events arrive508 };
Inbox: threads, reasons, subscriptions and watching509 // An approval or a request for changes is the owner's to hear
510 // of whatever they chose; the rest, as they subscribed.
511 if comment.verdict.is_some() {
512 told.tell_person(on.owner(), Reason::Author, severity, &title, body, true);
513 }
514 told.tell_subscribed(on, None, severity, &title, body);
515 told.tell_watchers(activity_of(on), severity, &title, body);
Inbox: the events service tells people what needs them as events arrive516 return told.notices;
517 }
518 _ => {}
519 }
520 told.notices
521}
522
Inbox: threads, reasons, subscriptions and watching523/// Who an event subscribes to its issue or pull request without asking,
524/// and why: whoever commented, whoever they mentioned, the people assigned
525/// and the reviewers asked. Never g1t.
526pub fn subscribes(event: &Event, subject: Option<&InboxSubject>) -> Vec<(String, Reason)> {
527 let mut people: Vec<(String, Reason)> = Vec::new();
528 let mut add = |name: &str, reason: Reason| {
529 let name = name.trim().trim_start_matches('@').to_lowercase();
530 if !name.is_empty() && !is_g1t(&name) && !people.iter().any(|(had, _)| *had == name) {
531 people.push((name, reason));
532 }
533 };
534 match event.kind.as_str() {
535 "comment.created" => {
536 if let Some(comment) = subject.and_then(|subject| subject.comment.as_ref()).filter(|comment| !comment.event) {
537 if !is_g1t_id(&comment.author.id) {
538 add(&comment.author.username, Reason::Comment);
539 }
540 for name in &comment.mentions {
541 add(name, Reason::Mention);
542 }
543 }
544 }
545 "issue.assigned" | "pull.assigned" => {
546 for name in names(&event.data, "added") {
547 add(&name, Reason::Assign);
548 }
549 }
550 "pull.review_requested" => {
551 for name in names(&event.data, "reviewers") {
552 add(&name, Reason::ReviewRequested);
553 }
554 }
555 _ => {}
556 }
557 people
558}
559
Inbox: the events service tells people what needs them as events arrive560fn short(sha: &str) -> &str {
561 sha.get(..7).unwrap_or(sha)
562}
563
564/// Where an item is on g1t.sh.
Inbox: threads, reasons, subscriptions and watching565pub fn url(repo: Option<&str>, subject: Option<SubjectKind>, number: Option<u32>, run_id: Option<&str>, link: Option<&str>) -> String {
566 if let Some(link) = link.filter(|link| link.starts_with('/')) {
567 return link.to_owned();
568 }
Inbox: the events service tells people what needs them as events arrive569 let Some(repo) = repo else {
570 return "/inbox".to_owned();
571 };
572 match (subject, number, run_id) {
573 (Some(SubjectKind::Pull), Some(number), _) => format!("/{repo}/pull/{number}"),
574 (Some(SubjectKind::Issue), Some(number), _) => format!("/{repo}/issues/{number}"),
575 (Some(SubjectKind::Run), _, Some(run)) => format!("/{repo}/actions/runs/{run}"),
Inbox: threads, reasons, subscriptions and watching576 (Some(SubjectKind::Deploy), Some(number), _) => format!("/{repo}/pull/{number}"),
Inbox: the events service tells people what needs them as events arrive577 _ => format!("/{repo}"),
578 }
579}
580
581// --- Writing ---------------------------------------------------------------
582
583/// The services the inbox reads from as events arrive.
584pub struct Sources<'a> {
585 pub work: &'a Fetcher,
586 pub repos: &'a Fetcher,
587 pub identity: &'a Fetcher,
588}
589
Inbox: threads, reasons, subscriptions and watching590/// One notice written, and whether it was news (not a redelivery): what
591/// may be emailed.
592struct Written {
593 notice: Notice,
594 repo_id: String,
595 url: String,
596}
597
Inbox: the events service tells people what needs them as events arrive598/// Writes the items a batch from the bus calls for. Never fails the batch:
599/// what cannot be worked out is logged and left.
600pub async fn deliver(db: &D1Database, sources: &Sources<'_>, events: &[Event]) {
Inbox: threads, reasons, subscriptions and watching601 if let Err(error) = resolve(db, events).await {
602 worker::console_error!("inbox: agent threads not closed: {error}");
603 }
Inbox: the events service tells people what needs them as events arrive604 let wanted: Vec<(&Event, Wanted)> = events
605 .iter()
606 .filter_map(|event| wants(event).map(|wanted| (event, wanted)))
607 .collect();
Inbox: threads, reasons, subscriptions and watching608 let created: Vec<&Event> = events.iter().filter(|event| event.kind == "repo.created").collect();
609 if wanted.is_empty() && created.is_empty() {
Inbox: the events service tells people what needs them as events arrive610 return;
611 }
612 // Everyone who caused one, named in one call.
613 let ids: Vec<String> = wanted
614 .iter()
Inbox: threads, reasons, subscriptions and watching615 .map(|(event, _)| *event)
616 .chain(created.iter().copied())
617 .filter_map(|event| event.actor.clone())
Inbox: the events service tells people what needs them as events arrive618 .filter(|id| !is_g1t_id(id))
619 .collect::<HashSet<_>>()
620 .into_iter()
621 .collect();
622 let names: HashMap<String, String> = if ids.is_empty() {
623 HashMap::new()
624 } else {
625 g1t_kit::call(sources.identity, "usernames", &UsernamesArgs { ids })
626 .await
627 .unwrap_or_else(|error| {
628 worker::console_error!("inbox: could not name who acted: {error}");
629 HashMap::new()
630 })
631 };
Inbox: threads, reasons, subscriptions and watching632 for event in created {
633 if let Err(error) = subscriptions::watch_created(db, event, &names).await {
634 worker::console_error!("inbox: {} {} not watched: {error}", event.kind, event.id);
635 }
636 }
Inbox: the events service tells people what needs them as events arrive637 let mut paths: HashMap<String, Option<RepoPath>> = HashMap::new();
Inbox: threads, reasons, subscriptions and watching638 let mut watchers: HashMap<String, Vec<subscriptions::Watcher>> = HashMap::new();
639 let mut written: Vec<Written> = Vec::new();
Inbox: the events service tells people what needs them as events arrive640 for (event, wanted) in wanted {
Inbox: threads, reasons, subscriptions and watching641 match deliver_one(db, sources, &names, &mut paths, &mut watchers, event, wanted).await {
642 Ok(mut news) => written.append(&mut news),
643 Err(error) => worker::console_error!("inbox: {} {} not delivered: {error}", event.kind, event.id),
644 }
645 }
646 if let Err(error) = email(db, sources.identity, written).await {
647 worker::console_error!("inbox: emails not sent: {error}");
648 }
649}
650
651/// Closes what an agent was waiting on a person for once it is over.
652async fn resolve(db: &D1Database, events: &[Event]) -> Result<()> {
653 let now = rfc3339(now_ms());
654 let mut statements = Vec::new();
655 for thread in events.iter().filter_map(resolves) {
656 statements.push(
657 db.prepare(
658 "UPDATE inbox_items SET done_at = ?1, read_at = COALESCE(read_at, ?1)
659 WHERE thread = ?2 AND reason = 'agent' AND done_at IS NULL",
660 )
661 .bind(&[now.as_str().into(), thread.into()])?,
662 );
663 }
664 // Reviews no longer asked for are no longer waiting.
665 for event in events.iter().filter(|event| event.kind == "pull.review_request_removed") {
666 let (Some(repo_id), Some(number)) = (
667 event.data["repoId"].as_str().map(str::to_owned).or_else(|| event.repo_id.clone()),
668 event.data["number"].as_u64(),
669 ) else {
670 continue;
671 };
672 for name in names(&event.data, "reviewers") {
673 statements.push(
674 db.prepare(
675 "UPDATE inbox_items SET done_at = ?1, read_at = COALESCE(read_at, ?1)
676 WHERE thread = ?2 AND username = ?3 AND reason = 'review_requested' AND done_at IS NULL",
677 )
678 .bind(&[
679 now.as_str().into(),
680 format!("{repo_id}#{number}").into(),
681 name.to_lowercase().into(),
682 ])?,
683 );
Inbox: the events service tells people what needs them as events arrive684 }
685 }
Inbox: threads, reasons, subscriptions and watching686 if !statements.is_empty() {
687 db.batch(statements).await?;
688 }
689 Ok(())
Inbox: the events service tells people what needs them as events arrive690}
691
692async fn deliver_one(
693 db: &D1Database,
694 sources: &Sources<'_>,
695 names: &HashMap<String, String>,
696 paths: &mut HashMap<String, Option<RepoPath>>,
Inbox: threads, reasons, subscriptions and watching697 watchers: &mut HashMap<String, Vec<subscriptions::Watcher>>,
Inbox: the events service tells people what needs them as events arrive698 event: &Event,
699 wanted: Wanted,
Inbox: threads, reasons, subscriptions and watching700) -> Result<Vec<Written>> {
Inbox: the events service tells people what needs them as events arrive701 if !paths.contains_key(&wanted.repo_id) {
702 let path: Option<RepoPath> = g1t_kit::call(
703 sources.repos,
704 "path_by_id",
705 &PathByIdArgs {
706 id: wanted.repo_id.clone(),
707 },
708 )
709 .await?;
710 paths.insert(wanted.repo_id.clone(), path);
711 }
712 // A repository that is gone tells nobody.
713 let Some(path) = paths.get(&wanted.repo_id).cloned().flatten() else {
Inbox: threads, reasons, subscriptions and watching714 return Ok(Vec::new());
Inbox: the events service tells people what needs them as events arrive715 };
716 let subject: Option<InboxSubject> = match wanted.number {
717 Some(number) => {
718 let found: Option<InboxSubject> = g1t_kit::call(
719 sources.work,
720 "inbox_subject",
721 &InboxSubjectArgs {
722 repo_id: wanted.repo_id.clone(),
723 number,
724 comment_id: wanted.comment_id.clone(),
725 },
726 )
727 .await?;
728 if found.is_none() {
Inbox: threads, reasons, subscriptions and watching729 return Ok(Vec::new());
Inbox: the events service tells people what needs them as events arrive730 }
731 found
732 }
733 None => None,
734 };
735 let actor = Actor {
736 id: event.actor.clone(),
737 username: match event.actor.as_deref() {
738 Some(id) if is_g1t_id(id) => Some(system::USERNAME.to_owned()),
739 Some(id) => names.get(id).map(|name| name.to_lowercase()),
740 None => None,
741 },
742 };
Inbox: threads, reasons, subscriptions and watching743 let thread = thread_of(event, &wanted, subject.as_ref());
744 if !watchers.contains_key(&wanted.repo_id) {
745 watchers.insert(wanted.repo_id.clone(), subscriptions::watchers(db, &wanted.repo_id).await?);
746 }
747 let audience = Audience {
748 subscriptions: match thread.kind {
749 Some(SubjectKind::Issue | SubjectKind::Pull) => subscriptions::of_thread(db, &thread.key).await?,
750 _ => Vec::new(),
751 },
752 watchers: watchers.get(&wanted.repo_id).cloned().unwrap_or_default(),
753 };
Inbox: the events service tells people what needs them as events arrive754 let repo = format!("{}/{}", path.namespace, path.name).to_lowercase();
Inbox: threads, reasons, subscriptions and watching755 let told = notices(event, &repo, &actor, subject.as_ref(), &audience);
756 let mut statements = subscriptions::auto_subscribe(db, &thread.key, &wanted.repo_id, &subscribes(event, subject.as_ref()), &event.time)?;
Inbox: the events service tells people what needs them as events arrive757 if told.is_empty() {
Inbox: threads, reasons, subscriptions and watching758 if !statements.is_empty() {
759 db.batch(statements).await?;
760 }
761 return Ok(Vec::new());
Inbox: the events service tells people what needs them as events arrive762 }
763 let now = now_ms();
Inbox: threads, reasons, subscriptions and watching764 let link = thread.link.as_deref();
765 let item_url = url(Some(&repo), thread.kind, thread.number, thread.run_id.as_deref(), link);
766 // Three statements a notice: its thread brought up (unless this event
767 // was told before), a line of history, and the history kept short.
768 let first = statements.len();
769 for notice in &told {
770 let username = notice.username.as_str();
771 let values: Vec<JsValue> = vec![
772 new_id("ntf", now).into(),
773 username.into(),
774 thread.key.as_str().into(),
775 event.id.as_str().into(),
776 event.kind.as_str().into(),
777 notice.reason.as_str().into(),
778 notice.severity.as_str().into(),
779 notice.title.as_str().into(),
780 notice.body.as_str().into(),
781 path.namespace.to_lowercase().into(),
782 wanted.repo_id.as_str().into(),
783 repo.as_str().into(),
784 thread.kind.map_or(JsValue::NULL, |kind| kind.as_str().into()),
785 thread.number.map_or(JsValue::NULL, JsValue::from),
786 thread.run_id.as_deref().map_or(JsValue::NULL, JsValue::from),
787 link.map_or(JsValue::NULL, JsValue::from),
788 actor.username.as_deref().map_or(JsValue::NULL, JsValue::from),
789 event.time.as_str().into(),
790 // Same prefix as item ids, so old and new sort by time together.
791 new_id("ntf", now).into(),
792 ];
793 statements.push(db.prepare(BUMP).bind(&values)?);
Inbox: the events service tells people what needs them as events arrive794 statements.push(
795 db.prepare(
Inbox: threads, reasons, subscriptions and watching796 "INSERT OR IGNORE INTO inbox_activity (id, item_id, username, event_id, event, reason, severity, title, body, actor, created_at)
797 SELECT ?1, id, ?2, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11 FROM inbox_items WHERE username = ?2 AND thread = ?3
798 RETURNING username",
Inbox: the events service tells people what needs them as events arrive799 )
800 .bind(&[
801 new_id("ntf", now).into(),
Inbox: threads, reasons, subscriptions and watching802 username.into(),
803 thread.key.as_str().into(),
Inbox: the events service tells people what needs them as events arrive804 event.id.as_str().into(),
Inbox: threads, reasons, subscriptions and watching805 event.kind.as_str().into(),
806 notice.reason.as_str().into(),
Inbox: the events service tells people what needs them as events arrive807 notice.severity.as_str().into(),
808 notice.title.as_str().into(),
809 notice.body.as_str().into(),
810 actor.username.as_deref().map_or(JsValue::NULL, JsValue::from),
811 event.time.as_str().into(),
812 ])?,
813 );
Inbox: threads, reasons, subscriptions and watching814 statements.push(
815 db.prepare(
816 "DELETE FROM inbox_activity
817 WHERE item_id = (SELECT id FROM inbox_items WHERE username = ?1 AND thread = ?2)
818 AND id NOT IN (
819 SELECT a.id FROM inbox_activity a
820 WHERE a.item_id = (SELECT id FROM inbox_items WHERE username = ?1 AND thread = ?2)
821 ORDER BY a.id DESC LIMIT ?3)",
822 )
823 .bind(&[username.into(), thread.key.as_str().into(), MAX_ACTIVITY.into()])?,
824 );
825 }
826 let results = db.batch(statements).await?;
827 #[derive(Deserialize)]
828 struct Told {
829 username: String,
Inbox: the events service tells people what needs them as events arrive830 }
Inbox: threads, reasons, subscriptions and watching831 // A line of history written means the event is news to that person.
832 let mut news = HashSet::new();
833 for (at, _) in told.iter().enumerate() {
834 if let Some(result) = results.get(first + at * 3 + 1)
835 && let Ok(rows) = result.results::<Told>()
836 {
837 news.extend(rows.into_iter().map(|row| row.username));
838 }
839 }
840 Ok(told
841 .into_iter()
842 .filter(|notice| news.contains(&notice.username))
843 .map(|notice| Written {
844 notice,
845 repo_id: wanted.repo_id.clone(),
846 url: item_url.clone(),
847 })
848 .collect())
849}
850
851/// Brings a person's thread up with new activity, or starts it. `?1` id,
852/// `?2` username, `?3` thread, `?4` event id, `?5` event type, `?6` reason,
853/// `?7` severity, `?8` title, `?9` body, `?10` workspace, `?11` repo id,
854/// `?12` repo, `?13` subject, `?14` number, `?15` run id, `?16` link, `?17`
855/// actor, `?18` time, `?19` the time-sortable id lists order by. Nothing
856/// happens when the person was already told of this event. While the
857/// thread is unread its most urgent severity is kept; done or not, it
858/// comes back to the inbox. A snooze stands.
859const BUMP: &str = "INSERT INTO inbox_items (id, username, thread, event_id, event, reason, severity, title, body,
860 workspace, repo_id, repo, subject, number, run_id, link, actor, created_at, updated_at, bumped, activity)
861 SELECT ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?18, ?19, 1
862 WHERE NOT EXISTS (SELECT 1 FROM inbox_activity WHERE event_id = ?4 AND username = ?2)
863 ON CONFLICT (username, thread) DO UPDATE SET
864 event_id = excluded.event_id,
865 event = excluded.event,
866 reason = excluded.reason,
867 severity = CASE
868 WHEN inbox_items.read_at IS NULL AND inbox_items.done_at IS NULL
869 AND (CASE inbox_items.severity WHEN 'warning' THEN 0 WHEN 'error' THEN 1 WHEN 'success' THEN 2 ELSE 3 END)
870 < (CASE excluded.severity WHEN 'warning' THEN 0 WHEN 'error' THEN 1 WHEN 'success' THEN 2 ELSE 3 END)
871 THEN inbox_items.severity ELSE excluded.severity END,
872 title = excluded.title,
873 body = excluded.body,
874 workspace = excluded.workspace,
875 repo = excluded.repo,
876 subject = COALESCE(excluded.subject, inbox_items.subject),
877 run_id = COALESCE(excluded.run_id, inbox_items.run_id),
878 link = COALESCE(excluded.link, inbox_items.link),
879 actor = excluded.actor,
880 updated_at = excluded.updated_at,
881 bumped = excluded.bumped,
882 activity = inbox_items.activity + 1,
883 read_at = NULL,
884 done_at = NULL";
885
886/// Emails what was news to people who asked to be emailed for its reason.
887/// Identity sends each, only to a confirmed address and only if the person
888/// can still read the repository.
889async fn email(db: &D1Database, identity: &Fetcher, written: Vec<Written>) -> Result<()> {
890 if written.is_empty() {
891 return Ok(());
892 }
893 let people: Vec<String> = written
894 .iter()
895 .map(|item| item.notice.username.clone())
896 .collect::<HashSet<_>>()
897 .into_iter()
898 .collect();
899 let settings = subscriptions::settings_of(db, &people).await?;
900 for item in written {
901 let wants = settings.get(&item.notice.username).map_or(&DEFAULT_EMAIL[..], |settings| &settings.email[..]);
902 if !wants.contains(&item.notice.reason) {
903 continue;
904 }
905 let sent: Result<bool> = g1t_kit::call(
906 identity,
907 "notify_by_email",
908 &NotifyByEmailArgs {
909 username: item.notice.username.clone(),
910 repo_id: item.repo_id.clone(),
911 subject: item.notice.title.clone(),
912 intro: item.notice.body.clone(),
913 quote: None,
914 path: item.url.clone(),
915 reason: item.notice.reason,
916 },
917 )
918 .await;
919 if let Err(error) = sent {
920 worker::console_error!("inbox: email to {} not sent: {error}", item.notice.username);
921 }
922 }
Inbox: the events service tells people what needs them as events arrive923 Ok(())
924}
925
926// --- Reading and changing ----------------------------------------------------
927
928#[derive(Deserialize)]
Inbox: threads, reasons, subscriptions and watching929pub(crate) struct Row {
930 pub(crate) id: String,
Inbox: the events service tells people what needs them as events arrive931 reason: String,
932 severity: String,
933 title: String,
934 body: String,
Inbox: threads, reasons, subscriptions and watching935 event: Option<String>,
Inbox: the events service tells people what needs them as events arrive936 workspace: Option<String>,
Inbox: threads, reasons, subscriptions and watching937 pub(crate) repo_id: Option<String>,
Inbox: the events service tells people what needs them as events arrive938 repo: Option<String>,
939 subject: Option<String>,
940 number: Option<f64>,
941 run_id: Option<String>,
Inbox: threads, reasons, subscriptions and watching942 link: Option<String>,
Inbox: the events service tells people what needs them as events arrive943 actor: Option<String>,
Inbox: threads, reasons, subscriptions and watching944 activity: Option<f64>,
Inbox: the events service tells people what needs them as events arrive945 created_at: String,
Inbox: threads, reasons, subscriptions and watching946 updated_at: Option<String>,
Inbox: the events service tells people what needs them as events arrive947 read_at: Option<String>,
948 done_at: Option<String>,
949 saved: f64,
950 snoozed_until: Option<String>,
Inbox: threads, reasons, subscriptions and watching951 bumped: Option<String>,
Inbox: the events service tells people what needs them as events arrive952}
953
954impl Row {
Inbox: threads, reasons, subscriptions and watching955 pub(crate) fn into_item(self) -> InboxItem {
Inbox: the events service tells people what needs them as events arrive956 let subject = self.subject.as_deref().and_then(SubjectKind::parse);
957 let number = self.number.map(|n| n as u32);
958 InboxItem {
Inbox: threads, reasons, subscriptions and watching959 url: url(self.repo.as_deref(), subject, number, self.run_id.as_deref(), self.link.as_deref()),
Inbox: the events service tells people what needs them as events arrive960 id: self.id,
Inbox: threads, reasons, subscriptions and watching961 reason: Reason::parse(&self.reason).unwrap_or(Reason::Subscribed),
Inbox: the events service tells people what needs them as events arrive962 severity: Severity::parse(&self.severity).unwrap_or(Severity::Info),
963 title: self.title,
964 body: self.body,
Inbox: threads, reasons, subscriptions and watching965 event: self.event,
Inbox: the events service tells people what needs them as events arrive966 repo: self.repo,
967 workspace: self.workspace,
968 subject,
969 number,
970 actor: self.actor,
Inbox: threads, reasons, subscriptions and watching971 count: self.activity.map_or(1, |n| n as u32),
972 updated_at: self.updated_at.unwrap_or_else(|| self.created_at.clone()),
Inbox: the events service tells people what needs them as events arrive973 created_at: self.created_at,
974 read_at: self.read_at,
975 done_at: self.done_at,
976 saved: self.saved != 0.0,
977 snoozed_until: self.snoozed_until,
978 }
979 }
980}
981
Inbox: threads, reasons, subscriptions and watching982pub(crate) const COLUMNS: &str = "id, reason, severity, title, body, event, workspace, repo_id, repo, subject, number, run_id,
983 link, actor, activity, created_at, updated_at, read_at, done_at, saved, snoozed_until, bumped";
Inbox: the events service tells people what needs them as events arrive984
985/// The conditions that pick a view's items, after `username = ?1`; `?2` is now.
986fn view_filter(view: InboxView) -> &'static str {
987 match view {
988 InboxView::Inbox => "done_at IS NULL AND (snoozed_until IS NULL OR snoozed_until <= ?2)",
989 InboxView::Saved => "saved = 1",
990 InboxView::Done => "done_at IS NOT NULL",
991 }
992}
993
994/// Whether a list ranks unread warnings first: the inbox itself, unfiltered.
995fn ranked(a: &ListInboxArgs) -> bool {
Inbox: threads, reasons, subscriptions and watching996 a.view == InboxView::Inbox
997 && a.severity.is_none()
998 && a.reason.is_none()
999 && !a.participating
1000 && a.repo_id.is_none()
1001 && !a.unread
1002 && a.since.is_none()
1003 && a.updated_before.is_none()
Inbox: the events service tells people what needs them as events arrive1004}
1005
Inbox: threads, reasons, subscriptions and watching1006/// The conditions a list's filters add, and the values they bind, numbered
1007/// after the `bound` values already given.
1008fn filters(a: &ListInboxArgs, bound: usize) -> (Vec<String>, Vec<String>) {
1009 let mut conditions = Vec::new();
1010 let mut values: Vec<String> = Vec::new();
1011 let mut bind = |value: &str, condition: &str| {
1012 values.push(value.to_owned());
1013 conditions.push(condition.replace('?', &format!("?{}", bound + values.len())));
1014 };
1015 if let Some(severity) = a.severity {
1016 bind(severity.as_str(), "severity = ?");
1017 }
1018 if let Some(reason) = a.reason {
1019 bind(reason.as_str(), "reason = ?");
1020 }
1021 if let Some(repo_id) = &a.repo_id {
1022 bind(repo_id, "repo_id = ?");
1023 }
1024 if let Some(since) = &a.since {
1025 bind(since.trim(), "COALESCE(updated_at, created_at) >= ?");
1026 }
1027 if let Some(before) = &a.updated_before {
1028 bind(before.trim(), "COALESCE(updated_at, created_at) < ?");
1029 }
1030 if a.participating {
1031 conditions.push("reason NOT IN ('manual', 'subscribed')".to_owned());
1032 }
1033 if a.unread {
1034 conditions.push("read_at IS NULL".to_owned());
1035 }
1036 (conditions, values)
1037}
1038
Inbox: the events service tells people what needs them as events arrive1039pub async fn list(db: &D1Database, repos: &Fetcher, a: ListInboxArgs) -> Result<InboxPage> {
1040 let Some(viewer) = &a.viewer else {
1041 return Ok(InboxPage::default());
1042 };
1043 let username = viewer.username.to_lowercase();
1044 let now = rfc3339(now_ms());
1045 let limit = a.limit.unwrap_or(DEFAULT_INBOX_PAGE).clamp(1, MAX_INBOX_PAGE);
1046 let mut values: Vec<JsValue> = vec![username.as_str().into(), now.as_str().into()];
Inbox: threads, reasons, subscriptions and watching1047 let mut conditions = vec!["username = ?1".to_owned(), view_filter(a.view).to_owned()];
1048 let (filtered, bound) = filters(&a, values.len());
1049 conditions.extend(filtered);
1050 values.extend(bound.iter().map(|value| JsValue::from(value.as_str())));
Inbox: the events service tells people what needs them as events arrive1051 let ranked = ranked(&a);
1052 // Unread warnings lead the first page, and are left out of the rest.
1053 let leading = "severity = 'warning' AND read_at IS NULL";
1054 let mut rest = conditions.clone();
Inbox: threads, reasons, subscriptions and watching1055 let mut rest_values = values.clone();
Inbox: the events service tells people what needs them as events arrive1056 if ranked {
1057 rest.push(format!("NOT ({leading})"));
1058 }
1059 if let Some(before) = &a.before {
Inbox: threads, reasons, subscriptions and watching1060 rest_values.push(before.as_str().into());
1061 rest.push(format!("bumped < ?{}", rest_values.len()));
Inbox: the events service tells people what needs them as events arrive1062 }
Inbox: threads, reasons, subscriptions and watching1063 rest_values.push((limit + 1).into());
1064 let order = if a.view == InboxView::Done { "done_at DESC, bumped DESC" } else { "bumped DESC" };
Inbox: the events service tells people what needs them as events arrive1065 let mut statements = vec![
1066 db.prepare(format!(
1067 "SELECT {COLUMNS} FROM inbox_items WHERE {} ORDER BY {order} LIMIT ?{}",
1068 rest.join(" AND "),
Inbox: threads, reasons, subscriptions and watching1069 rest_values.len()
Inbox: the events service tells people what needs them as events arrive1070 ))
Inbox: threads, reasons, subscriptions and watching1071 .bind(&rest_values)?,
Inbox: the events service tells people what needs them as events arrive1072 ];
1073 if ranked && a.before.is_none() {
1074 statements.push(
1075 db.prepare(format!(
Inbox: threads, reasons, subscriptions and watching1076 "SELECT {COLUMNS} FROM inbox_items WHERE {} AND {leading} ORDER BY bumped DESC LIMIT ?3",
1077 conditions.join(" AND ")
Inbox: the events service tells people what needs them as events arrive1078 ))
1079 .bind(&[username.as_str().into(), now.as_str().into(), MAX_RANKED.into()])?,
1080 );
1081 }
1082 let results = db.batch(statements).await?;
1083 let mut rows = results.first().map(|result| result.results::<Row>()).transpose()?.unwrap_or_default();
1084 let next = if rows.len() > limit as usize {
1085 rows.truncate(limit as usize);
Inbox: threads, reasons, subscriptions and watching1086 rows.last().map(|row| row.bumped.clone().unwrap_or_else(|| row.id.clone()))
Inbox: the events service tells people what needs them as events arrive1087 } else {
1088 None
1089 };
1090 let mut items: Vec<Row> = results.get(1).map(|result| result.results::<Row>()).transpose()?.unwrap_or_default();
1091 items.extend(rows);
Inbox: threads, reasons, subscriptions and watching1092 let items = readable_only(db, repos, &a.viewer, &username, items).await?;
1093 Ok(InboxPage {
1094 items: items.into_iter().map(Row::into_item).collect(),
1095 next,
1096 })
1097}
Inbox: the events service tells people what needs them as events arrive1098
Inbox: threads, reasons, subscriptions and watching1099/// Rows about repositories the viewer can still read; the rest are
1100/// removed from their inbox as they are found.
1101pub(crate) async fn readable_only(db: &D1Database, repos: &Fetcher, viewer: &g1t_contracts::Viewer, username: &str, mut rows: Vec<Row>) -> Result<Vec<Row>> {
1102 let ids: Vec<String> = rows
Inbox: the events service tells people what needs them as events arrive1103 .iter()
1104 .filter_map(|row| row.repo_id.clone())
1105 .collect::<HashSet<_>>()
1106 .into_iter()
1107 .collect();
Inbox: threads, reasons, subscriptions and watching1108 if ids.is_empty() {
1109 return Ok(rows);
1110 }
1111 let readable: Vec<Repo> = g1t_kit::call(
1112 repos,
1113 "readable",
1114 &ReadableArgs {
1115 ids: ids.clone(),
1116 viewer: viewer.clone(),
1117 },
1118 )
1119 .await?;
1120 let readable: HashSet<String> = readable.into_iter().map(|repo| repo.id).collect();
1121 let gone: Vec<String> = ids.into_iter().filter(|id| !readable.contains(id)).collect();
1122 if !gone.is_empty() {
1123 forget(db, username, &gone).await?;
1124 rows.retain(|row| row.repo_id.as_ref().is_none_or(|id| !gone.contains(id)));
Inbox: the events service tells people what needs them as events arrive1125 }
Inbox: threads, reasons, subscriptions and watching1126 Ok(rows)
Inbox: the events service tells people what needs them as events arrive1127}
1128
1129/// Removes a person's items about repositories they cannot read.
1130async fn forget(db: &D1Database, username: &str, repo_ids: &[String]) -> Result<()> {
1131 let marks = vec!["?"; repo_ids.len()].join(", ");
1132 let mut values: Vec<JsValue> = vec![username.into()];
1133 values.extend(repo_ids.iter().map(|id| JsValue::from(id.as_str())));
Inbox: threads, reasons, subscriptions and watching1134 db.batch(vec![
1135 db.prepare(format!(
1136 "DELETE FROM inbox_activity WHERE item_id IN (SELECT id FROM inbox_items WHERE username = ? AND repo_id IN ({marks}))"
1137 ))
1138 .bind(&values)?,
1139 db.prepare(format!("DELETE FROM inbox_items WHERE username = ? AND repo_id IN ({marks})"))
1140 .bind(&values)?,
1141 ])
1142 .await?;
Inbox: the events service tells people what needs them as events arrive1143 Ok(())
1144}
1145
1146#[derive(Deserialize)]
1147struct CountRow {
1148 severity: String,
1149 n: f64,
1150}
1151
1152pub async fn counts(db: &D1Database, a: InboxCountsArgs) -> Result<InboxCounts> {
1153 let rows = db
1154 .prepare(
1155 "SELECT severity, count(*) AS n FROM inbox_items
1156 WHERE username = ? AND read_at IS NULL AND done_at IS NULL
1157 AND (snoozed_until IS NULL OR snoozed_until <= ?)
1158 GROUP BY severity",
1159 )
1160 .bind(&[a.username.to_lowercase().into(), rfc3339(now_ms()).into()])?
1161 .all()
1162 .await?
1163 .results::<CountRow>()?;
1164 let mut counts = InboxCounts::default();
1165 for row in rows {
1166 let n = row.n as u32;
1167 counts.unread += n;
1168 match Severity::parse(&row.severity) {
1169 Some(Severity::Error) => counts.error += n,
1170 Some(Severity::Warning) => counts.warning += n,
1171 Some(Severity::Success) => counts.success += n,
1172 Some(Severity::Info) | None => counts.info += n,
1173 }
1174 }
1175 Ok(counts)
1176}
1177
1178/// The change a mark makes, as a `SET` clause; `?1` is now.
1179fn mark_change(mark: InboxMark) -> &'static str {
1180 match mark {
1181 InboxMark::Read => "read_at = COALESCE(read_at, ?1)",
1182 InboxMark::Unread => "read_at = NULL",
1183 InboxMark::Done => "done_at = ?1, read_at = COALESCE(read_at, ?1)",
1184 InboxMark::Undone => "done_at = NULL",
1185 InboxMark::Save => "saved = 1",
1186 InboxMark::Unsave => "saved = 0",
1187 InboxMark::Snooze => "snoozed_until = ?2, read_at = COALESCE(read_at, ?1)",
Inbox: threads, reasons, subscriptions and watching1188 InboxMark::Unsnooze => "snoozed_until = NULL",
Inbox: the events service tells people what needs them as events arrive1189 }
1190}
1191
1192#[derive(Deserialize)]
1193struct IdRow {
1194 #[allow(dead_code)]
1195 id: String,
1196}
1197
1198/// Changes the person's own items. Returns how many changed.
1199pub async fn mark(db: &D1Database, a: MarkInboxArgs) -> Result<u32> {
1200 let now = rfc3339(now_ms());
1201 let until = match (a.mark, a.until.as_deref().map(str::trim)) {
1202 // Times compare as text (`g1t_contracts::time`); a snooze is for later.
1203 (InboxMark::Snooze, Some(until)) if until.len() == now.len() && until > now.as_str() => until.to_owned(),
1204 (InboxMark::Snooze, _) => return Ok(0),
1205 _ => String::new(),
1206 };
1207 let mut values: Vec<JsValue> = vec![now.as_str().into(), until.as_str().into(), a.username.to_lowercase().into()];
1208 let target = if a.all && a.ids.is_empty() {
1209 let mut target = "done_at IS NULL".to_owned();
Inbox: threads, reasons, subscriptions and watching1210 let mut bind = |values: &mut Vec<JsValue>, value: &str, condition: &str| {
1211 values.push(value.into());
1212 target.push_str(&format!(" AND {}", condition.replace('?', &format!("?{}", values.len()))));
1213 };
Inbox: the events service tells people what needs them as events arrive1214 if let Some(severity) = a.severity {
Inbox: threads, reasons, subscriptions and watching1215 bind(&mut values, severity.as_str(), "severity = ?");
1216 }
1217 if let Some(repo_id) = &a.repo_id {
1218 bind(&mut values, repo_id, "repo_id = ?");
1219 }
1220 if let Some(last_read_at) = a.last_read_at.as_deref().map(str::trim).filter(|at| !at.is_empty()) {
1221 bind(&mut values, last_read_at, "COALESCE(updated_at, created_at) <= ?");
Inbox: the events service tells people what needs them as events arrive1222 }
1223 target
1224 } else {
1225 let ids: Vec<&String> = a.ids.iter().take(MAX_MARK).collect();
1226 if ids.is_empty() {
1227 return Ok(0);
1228 }
1229 let first = values.len() + 1;
1230 values.extend(ids.iter().map(|id| JsValue::from(id.as_str())));
1231 let marks: Vec<String> = (first..values.len() + 1).map(|at| format!("?{at}")).collect();
1232 format!("id IN ({})", marks.join(", "))
1233 };
1234 let changed = db
1235 .prepare(format!(
1236 "UPDATE inbox_items SET {} WHERE username = ?3 AND {target} RETURNING id",
1237 mark_change(a.mark)
1238 ))
1239 .bind(&values)?
1240 .all()
1241 .await?
1242 .results::<IdRow>()?;
1243 Ok(changed.len() as u32)
1244}
1245
Inbox: threads, reasons, subscriptions and watching1246#[derive(Deserialize)]
1247struct ActivityRow {
1248 reason: String,
1249 severity: String,
1250 title: String,
1251 body: String,
1252 event: Option<String>,
1253 actor: Option<String>,
1254 created_at: String,
1255}
1256
1257/// One of the viewer's threads, with its history and their subscription.
1258pub async fn thread(db: &D1Database, repos: &Fetcher, work: &Fetcher, a: ThreadArgs) -> Result<Option<InboxThread>> {
1259 let Some(viewer) = &a.viewer else {
1260 return Ok(None);
1261 };
1262 let username = viewer.username.to_lowercase();
1263 let row = db
1264 .prepare(format!("SELECT {COLUMNS} FROM inbox_items WHERE id = ? AND username = ?"))
1265 .bind(&[a.id.as_str().into(), username.as_str().into()])?
1266 .first::<Row>(None)
1267 .await?;
1268 let Some(row) = readable_only(db, repos, &a.viewer, &username, row.into_iter().collect()).await?.pop() else {
1269 return Ok(None);
1270 };
1271 let activity = db
1272 .prepare(
1273 "SELECT reason, severity, title, body, event, actor, created_at FROM inbox_activity
1274 WHERE item_id = ? ORDER BY id DESC LIMIT ?",
1275 )
1276 .bind(&[row.id.as_str().into(), MAX_ACTIVITY.into()])?
1277 .all()
1278 .await?
1279 .results::<ActivityRow>()?
1280 .into_iter()
1281 .map(|row| InboxActivity {
1282 reason: Reason::parse(&row.reason).unwrap_or(Reason::Subscribed),
1283 severity: Severity::parse(&row.severity).unwrap_or(Severity::Info),
1284 title: row.title,
1285 body: row.body,
1286 event: row.event,
1287 actor: row.actor,
1288 created_at: row.created_at,
1289 })
1290 .collect();
1291 let item = row.into_item();
1292 let subscription = match (item.subject, item.number) {
1293 (Some(SubjectKind::Issue | SubjectKind::Pull), Some(_)) => {
1294 subscriptions::subscription(
1295 db,
1296 work,
1297 SubscriptionArgs {
1298 viewer: a.viewer.clone(),
1299 id: Some(item.id.clone()),
1300 ..SubscriptionArgs::default()
1301 },
1302 )
1303 .await?
1304 }
1305 _ => None,
1306 };
1307 Ok(Some(InboxThread {
1308 item,
1309 activity,
1310 subscription,
1311 }))
1312}
1313
Inbox: the events service tells people what needs them as events arrive1314// --- Keeping up ------------------------------------------------------------------
1315
1316/// Moves rows with renamed workspaces and repositories, and drops those
1317/// of purged repositories and deleted workspaces.
1318pub async fn follow(db: &D1Database, events: &[Event]) -> Result<()> {
1319 let mut statements = Vec::new();
1320 for event in events {
1321 let text = |key: &str| event.data[key].as_str().unwrap_or_default().to_lowercase();
1322 match event.kind.as_str() {
1323 "workspace.renamed" => {
1324 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
1325 continue;
1326 };
1327 let (from, to) = (renamed.from.to_lowercase(), renamed.to.to_lowercase());
1328 if from == to {
1329 continue;
1330 }
1331 statements.push(
1332 db.prepare(
Inbox: threads, reasons, subscriptions and watching1333 "UPDATE inbox_items SET repo = ?2 || substr(repo, length(?1) + 1), workspace = ?2,
1334 link = CASE WHEN link LIKE '/' || ?1 || '/%' THEN '/' || ?2 || substr(link, length(?1) + 2) ELSE link END
Inbox: the events service tells people what needs them as events arrive1335 WHERE workspace = ?1",
1336 )
1337 .bind(&[from.as_str().into(), to.as_str().into()])?,
1338 );
Inbox: threads, reasons, subscriptions and watching1339 statements.push(
1340 db.prepare("UPDATE inbox_watching SET repo = ?2 || substr(repo, length(?1) + 1) WHERE repo LIKE ?1 || '/%'")
1341 .bind(&[from.as_str().into(), to.as_str().into()])?,
1342 );
Inbox: the events service tells people what needs them as events arrive1343 }
Inbox: threads, reasons, subscriptions and watching1344 "repo.renamed" => {
1345 let path = format!("{}/{}", text("namespace"), text("to"));
1346 for table in ["inbox_items", "inbox_watching"] {
1347 statements.push(
1348 db.prepare(format!("UPDATE {table} SET repo = ? WHERE repo_id = ?"))
1349 .bind(&[path.as_str().into(), text("repoId").into()])?,
1350 );
1351 }
1352 }
1353 "repo.transferred" => {
1354 let path = format!("{}/{}", text("to"), text("name"));
1355 statements.push(
1356 db.prepare("UPDATE inbox_items SET repo = ?, workspace = ? WHERE repo_id = ?")
1357 .bind(&[path.as_str().into(), text("to").into(), text("repoId").into()])?,
1358 );
1359 statements.push(
1360 db.prepare("UPDATE inbox_watching SET repo = ? WHERE repo_id = ?")
1361 .bind(&[path.as_str().into(), text("repoId").into()])?,
1362 );
1363 }
1364 "repo.purged" => {
1365 for sql in [
1366 "DELETE FROM inbox_activity WHERE item_id IN (SELECT id FROM inbox_items WHERE repo_id = ?)",
1367 "DELETE FROM inbox_items WHERE repo_id = ?",
1368 "DELETE FROM inbox_subscriptions WHERE repo_id = ?",
1369 "DELETE FROM inbox_watching WHERE repo_id = ?",
1370 ] {
1371 statements.push(db.prepare(sql).bind(&[text("repoId").into()])?);
1372 }
1373 }
1374 "workspace.deleted" => {
1375 statements.push(
1376 db.prepare("DELETE FROM inbox_items WHERE workspace = ?")
1377 .bind(&[text("slug").into()])?,
1378 );
1379 statements.push(
1380 db.prepare("DELETE FROM inbox_watching WHERE repo LIKE ? || '/%'")
1381 .bind(&[text("slug").into()])?,
1382 );
1383 }
Inbox: the events service tells people what needs them as events arrive1384 _ => {}
1385 }
1386 }
1387 if !statements.is_empty() {
1388 db.batch(statements).await?;
1389 }
1390 Ok(())
1391}
1392
1393/// Removes items done more than [`DONE_DAYS`] ago, and any not saved older
Inbox: threads, reasons, subscriptions and watching1394/// than [`MAX_DAYS`], with their history. Returns how many went.
Inbox: the events service tells people what needs them as events arrive1395pub async fn purge(db: &D1Database, now: u64) -> Result<u32> {
1396 let done = crate::audit::keep_from(now, DONE_DAYS);
1397 let oldest = crate::audit::keep_from(now, MAX_DAYS);
Inbox: threads, reasons, subscriptions and watching1398 let statements: Vec<D1PreparedStatement> = vec![
1399 db.prepare("DELETE FROM inbox_items WHERE saved = 0 AND done_at < ?")
1400 .bind(&[done.into()])?,
1401 db.prepare("DELETE FROM inbox_items WHERE saved = 0 AND COALESCE(updated_at, created_at) < ?")
1402 .bind(&[oldest.into()])?,
1403 db.prepare("DELETE FROM inbox_activity WHERE NOT EXISTS (SELECT 1 FROM inbox_items WHERE inbox_items.id = inbox_activity.item_id)"),
1404 ];
1405 let results = db.batch(statements).await?;
Inbox: the events service tells people what needs them as events arrive1406 let mut removed = 0;
Inbox: threads, reasons, subscriptions and watching1407 for result in results.into_iter().take(2) {
Inbox: the events service tells people what needs them as events arrive1408 removed += result.meta()?.and_then(|meta| meta.changes).unwrap_or(0) as u32;
1409 }
1410 Ok(removed)
1411}
1412
1413#[cfg(test)]
1414mod tests {
1415 use super::*;
Inbox: threads, reasons, subscriptions and watching1416 use crate::subscriptions::{State, Subscription, Watcher};
Inbox: the events service tells people what needs them as events arrive1417 use serde_json::json;
1418
1419 fn event(kind: &str, actor: Option<&str>, data: serde_json::Value) -> Event {
1420 Event {
1421 id: "evt_1".into(),
1422 kind: kind.into(),
1423 source: "work".into(),
1424 time: "2026-10-07T12:00:00.000Z".into(),
1425 repo_id: Some("rep_1".into()),
1426 actor: actor.map(str::to_owned),
1427 data,
1428 }
1429 }
1430
1431 fn person(id: &str, username: &str) -> Principal {
1432 Principal {
1433 id: id.into(),
1434 username: username.into(),
1435 }
1436 }
1437
1438 fn actor(id: &str, username: &str) -> Actor {
1439 Actor {
1440 id: Some(id.into()),
1441 username: Some(username.into()),
1442 }
1443 }
1444
1445 /// A pull request ana opened, assigned to bo, for an issue cy filed and dee is assigned.
1446 fn pull() -> InboxSubject {
1447 InboxSubject {
1448 kind: Some(SubjectKind::Pull),
1449 title: "Add the inbox".into(),
1450 author: person("usr_ana", "ana"),
1451 assignees: vec!["bo".into()],
1452 issue: Some(Box::new(InboxSubject {
1453 kind: Some(SubjectKind::Issue),
1454 title: "An inbox".into(),
1455 author: person("usr_cy", "cy"),
1456 assignees: vec!["dee".into()],
1457 ..InboxSubject::default()
1458 })),
1459 ..InboxSubject::default()
1460 }
1461 }
1462
1463 /// A change g1t made for ana.
1464 fn g1t_pull() -> InboxSubject {
1465 InboxSubject {
1466 author: person(AGENT_ID, "g1t"),
1467 requested_by: Some(person("usr_ana", "ana")),
1468 ..pull()
1469 }
1470 }
1471
Inbox: threads, reasons, subscriptions and watching1472 fn nobody() -> Audience {
1473 Audience::default()
1474 }
1475
1476 fn told(notices: &[Notice]) -> Vec<(&str, Reason, Severity)> {
Inbox: the events service tells people what needs them as events arrive1477 notices
1478 .iter()
1479 .map(|notice| (notice.username.as_str(), notice.reason, notice.severity))
1480 .collect()
1481 }
1482
1483 fn comment(author: Principal, body: &str, mentions: &[&str]) -> InboxSubject {
1484 InboxSubject {
1485 comment: Some(InboxComment {
1486 author,
1487 excerpt: body.into(),
1488 mentions: mentions.iter().map(|name| (*name).to_owned()).collect(),
1489 ..InboxComment::default()
1490 }),
1491 ..pull()
1492 }
1493 }
1494
Inbox: threads, reasons, subscriptions and watching1495 fn subscribed(rows: &[(&str, State, Option<Reason>)]) -> Audience {
1496 Audience {
1497 subscriptions: rows
1498 .iter()
1499 .map(|(name, state, reason)| Subscription {
1500 username: (*name).into(),
1501 state: *state,
1502 reason: *reason,
1503 })
1504 .collect(),
1505 watchers: Vec::new(),
1506 }
1507 }
1508
1509 fn watched(rows: &[(&str, WatchLevel, &[&str])]) -> Audience {
1510 Audience {
1511 subscriptions: Vec::new(),
1512 watchers: rows
1513 .iter()
1514 .map(|(name, level, events)| Watcher {
1515 username: (*name).into(),
1516 level: *level,
1517 events: events.iter().map(|kind| (*kind).to_owned()).collect(),
1518 })
1519 .collect(),
1520 }
1521 }
1522
Inbox: the events service tells people what needs them as events arrive1523 #[test]
1524 fn only_events_that_tell_someone_are_read() {
1525 let asked = |kind: &str, data| wants(&event(kind, None, data));
1526 assert_eq!(
1527 asked("checks.completed", json!({ "repoId": "rep_1", "number": 4, "status": "failed" })),
1528 Some(Wanted { repo_id: "rep_1".into(), number: Some(4), comment_id: None })
1529 );
1530 assert_eq!(asked("checks.completed", json!({ "repoId": "rep_1", "number": 4, "status": "passed" })), None);
1531 assert_eq!(asked("review.completed", json!({ "repoId": "rep_1", "number": 4 })), None);
1532 assert_eq!(asked("workflow.completed", json!({ "repoId": "rep_1", "conclusion": "success" })), None);
1533 assert_eq!(
1534 asked("workflow.completed", json!({ "repoId": "rep_1", "conclusion": "failure" })),
1535 Some(Wanted { repo_id: "rep_1".into(), number: None, comment_id: None })
1536 );
1537 assert_eq!(
1538 asked("comment.created", json!({ "repoId": "rep_1", "number": 2, "commentId": "cmt_1" })).map(|w| w.comment_id),
1539 Some(Some("cmt_1".into()))
1540 );
1541 assert_eq!(asked("git.push", json!({ "repoId": "rep_1" })), None);
Inbox: threads, reasons, subscriptions and watching1542 for kind in ["issue.opened", "pull.review_requested", "pull.stalled", "issue.assigned", "pull.closed", "issue.reopened"] {
1543 assert!(asked(kind, json!({ "repoId": "rep_1", "number": 1 })).is_some(), "{kind}");
1544 }
1545 assert_eq!(asked("deployment.failed", json!({ "repoId": "rep_1" })).map(|w| w.number), Some(None));
Inbox: the events service tells people what needs them as events arrive1546 }
1547
1548 #[test]
Inbox: threads, reasons, subscriptions and watching1549 fn a_waiting_or_stopped_agent_needs_the_pull_requests_and_the_issues_people_first() {
Inbox: the events service tells people what needs them as events arrive1550 let asked = event("agent.asked", Some("usr_agent"), json!({ "number": 7 }));
Inbox: threads, reasons, subscriptions and watching1551 let notices_ = notices(&asked, "acme/rocket", &Actor::default(), Some(&pull()), &nobody());
Inbox: the events service tells people what needs them as events arrive1552 assert_eq!(
Inbox: threads, reasons, subscriptions and watching1553 told(&notices_),
Inbox: the events service tells people what needs them as events arrive1554 vec![
Inbox: threads, reasons, subscriptions and watching1555 ("ana", Reason::Agent, Severity::Warning),
1556 ("cy", Reason::Agent, Severity::Warning),
1557 ("dee", Reason::Agent, Severity::Warning),
Inbox: the events service tells people what needs them as events arrive1558 ]
1559 );
Inbox: threads, reasons, subscriptions and watching1560 assert_eq!(notices_[0].title, "An agent is waiting on acme/rocket#7");
1561 assert_eq!(notices_[0].body, "Add the inbox");
1562 // Stopping says why; even someone who unsubscribed hears of it.
1563 let stalled = event("pull.stalled", None, json!({ "number": 7, "detail": "Its checks could not be run." }));
1564 let audience = subscribed(&[("ana", State::Unsubscribed, None)]);
1565 let notices_ = notices(&stalled, "acme/rocket", &Actor::default(), Some(&g1t_pull()), &audience);
1566 assert_eq!(notices_[0].username, "ana");
1567 assert_eq!(notices_[0].title, "g1t stopped on acme/rocket#7 and needs you");
1568 assert_eq!(notices_[0].body, "Its checks could not be run.");
1569 // Ignoring the thread silences even that.
1570 let audience = subscribed(&[("ana", State::Ignored, None)]);
1571 let notices_ = notices(&stalled, "acme/rocket", &Actor::default(), Some(&g1t_pull()), &audience);
1572 assert!(!notices_.iter().any(|notice| notice.username == "ana"));
1573 }
1574
1575 #[test]
1576 fn whatever_an_agent_waits_on_closes_when_it_goes_on() {
1577 for kind in RESUMES {
1578 assert_eq!(resolves(&event(kind, None, json!({ "repoId": "rep_1", "number": 7 }))), Some("rep_1#7".into()));
1579 }
1580 assert_eq!(resolves(&event("comment.created", None, json!({ "repoId": "rep_1", "number": 7 }))), None);
Inbox: the events service tells people what needs them as events arrive1581 }
1582
1583 #[test]
Inbox: threads, reasons, subscriptions and watching1584 fn reviewers_and_assignees_asked_are_told_never_whoever_asked() {
1585 let requested = event("pull.review_requested", Some("usr_ana"), json!({ "number": 7, "reviewers": ["bo", "g1t", "ana"] }));
1586 let notices_ = notices(&requested, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
1587 assert_eq!(told(&notices_), vec![("bo", Reason::ReviewRequested, Severity::Warning)]);
1588 assert_eq!(notices_[0].title, "ana asked you to review acme/rocket#7");
1589 let assigned = event("issue.assigned", Some("usr_ana"), json!({ "number": 3, "added": ["ana", "eve"] }));
1590 let notices_ = notices(&assigned, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
1591 assert_eq!(told(&notices_), vec![("eve", Reason::Assign, Severity::Info)]);
1592 assert_eq!(notices_[0].title, "ana assigned you to acme/rocket#3");
1593 // Both subscribe whoever they name, never g1t.
1594 assert_eq!(subscribes(&requested, Some(&pull())), vec![("bo".into(), Reason::ReviewRequested), ("ana".into(), Reason::ReviewRequested)]);
1595 assert_eq!(subscribes(&assigned, Some(&pull())), vec![("ana".into(), Reason::Assign), ("eve".into(), Reason::Assign)]);
1596 }
1597
1598 #[test]
Inbox: the events service tells people what needs them as events arrive1599 fn failures_go_to_whoever_answers_for_the_change() {
1600 let failed = event("checks.completed", None, json!({ "number": 7, "status": "failed" }));
Inbox: threads, reasons, subscriptions and watching1601 assert_eq!(
1602 told(&notices(&failed, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())),
1603 vec![("ana", Reason::CiActivity, Severity::Error)]
1604 );
Inbox: the events service tells people what needs them as events arrive1605 // g1t's change is the person's who asked for it, never g1t's.
Inbox: threads, reasons, subscriptions and watching1606 let notices_ = notices(&failed, "acme/rocket", &Actor::default(), Some(&g1t_pull()), &nobody());
1607 assert_eq!(told(&notices_), vec![("ana", Reason::CiActivity, Severity::Error)]);
Inbox: the events service tells people what needs them as events arrive1608 let errored = event("checks.completed", None, json!({ "number": 7, "status": "errored" }));
Inbox: threads, reasons, subscriptions and watching1609 assert_eq!(
1610 notices(&errored, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())[0].title,
1611 "Checks could not run on acme/rocket#7"
1612 );
Inbox: the events service tells people what needs them as events arrive1613 }
1614
1615 #[test]
1616 fn a_workflow_that_fails_tells_its_pull_requests_owner_even_if_they_pushed() {
1617 let failed = event("workflow.completed", Some("usr_ana"), json!({ "pull": 7, "workflow": "CI", "conclusion": "failure" }));
Inbox: threads, reasons, subscriptions and watching1618 let notices_ = notices(&failed, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
1619 assert_eq!(told(&notices_), vec![("ana", Reason::CiActivity, Severity::Error)]);
Inbox: the events service tells people what needs them as events arrive1620 assert_eq!(notices_[0].title, "CI failed on acme/rocket#7");
1621 // On a branch: whoever pushed.
1622 let pushed = event(
1623 "workflow.completed",
1624 Some("usr_bo"),
Inbox: threads, reasons, subscriptions and watching1625 json!({ "workflow": "Deploy", "path": ".g1t/workflows/deploy.yml", "conclusion": "failure", "ref": "refs/heads/main", "number": 12, "sha": "abcdef0123", "runId": "run_9" }),
Inbox: the events service tells people what needs them as events arrive1626 );
Inbox: threads, reasons, subscriptions and watching1627 let notices_ = notices(&pushed, "acme/rocket", &actor("usr_bo", "bo"), None, &nobody());
1628 assert_eq!(told(&notices_), vec![("bo", Reason::CiActivity, Severity::Error)]);
Inbox: the events service tells people what needs them as events arrive1629 assert_eq!(notices_[0].title, "Deploy failed on main in acme/rocket");
1630 assert_eq!(notices_[0].body, "Run 12 at abcdef0");
1631 // Nobody to tell when g1t pushed.
Inbox: threads, reasons, subscriptions and watching1632 assert!(notices(&pushed, "acme/rocket", &actor("g1t", "g1t"), None, &nobody()).is_empty());
1633 // Every failure of a workflow on a branch is one thread.
1634 let wanted = wants(&pushed).unwrap();
1635 let thread = thread_of(&pushed, &wanted, None);
1636 assert_eq!(thread.key, "rep_1/run/.g1t/workflows/deploy.yml@main");
1637 assert_eq!(thread.run_id.as_deref(), Some("run_9"));
1638 }
1639
1640 #[test]
1641 fn deployments_tell_whoever_answers_for_them_and_watchers() {
1642 let failed = event(
1643 "deployment.failed",
1644 Some("usr_bo"),
1645 json!({ "repoId": "rep_1", "projectId": "prj_1", "project": "rocket", "deploymentId": "dpl_1", "error": "The build failed.", "path": "/acme/rocket/deployments/dpl_1" }),
1646 );
1647 let audience = watched(&[("cy", WatchLevel::Custom, &["deployments"]), ("dee", WatchLevel::Custom, &["issues"]), ("eve", WatchLevel::All, &[])]);
1648 let notices_ = notices(&failed, "acme/rocket", &actor("usr_bo", "bo"), None, &audience);
1649 assert_eq!(
1650 told(&notices_),
1651 vec![
1652 ("bo", Reason::CiActivity, Severity::Error),
1653 ("cy", Reason::Subscribed, Severity::Error),
1654 ("eve", Reason::Subscribed, Severity::Error),
1655 ]
1656 );
1657 assert_eq!(notices_[0].title, "Production of rocket failed to deploy");
1658 assert_eq!(notices_[0].body, "The build failed.");
1659 let thread = thread_of(&failed, &wants(&failed).unwrap(), None);
1660 assert_eq!(thread.key, "rep_1/deploy/prj_1/production");
1661 assert_eq!(url(Some("acme/rocket"), thread.kind, None, None, thread.link.as_deref()), "/acme/rocket/deployments/dpl_1");
1662 // A success is news to the owner only after a failure.
1663 let live = event("deployment.succeeded", Some("usr_bo"), json!({ "repoId": "rep_1", "projectId": "prj_1", "project": "rocket", "commit": "abcdef0123" }));
1664 assert!(notices(&live, "acme/rocket", &actor("usr_bo", "bo"), None, &nobody()).is_empty());
1665 let recovered = event("deployment.succeeded", Some("usr_bo"), json!({ "repoId": "rep_1", "projectId": "prj_1", "project": "rocket", "recovered": true }));
1666 let notices_ = notices(&recovered, "acme/rocket", &actor("usr_bo", "bo"), None, &nobody());
1667 assert_eq!(told(&notices_), vec![("bo", Reason::CiActivity, Severity::Success)]);
1668 assert_eq!(notices_[0].title, "Production of rocket is live again");
1669 // A preview's is the pull request's owner's.
1670 let preview = event("deployment.failed", None, json!({ "repoId": "rep_1", "projectId": "prj_1", "number": 7, "triggeredBy": "g1t" }));
1671 assert_eq!(told(&notices(&preview, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())), vec![("ana", Reason::CiActivity, Severity::Error)]);
1672 assert_eq!(thread_of(&preview, &wants(&preview).unwrap(), Some(&pull())).key, "rep_1/deploy/prj_1/7");
Inbox: the events service tells people what needs them as events arrive1673 }
1674
1675 #[test]
1676 fn nobody_hears_of_what_they_did_themselves() {
1677 let merged = event("pull.merged", Some("usr_ana"), json!({ "number": 7 }));
Inbox: threads, reasons, subscriptions and watching1678 let by_ana = notices(&merged, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
1679 // The assignee still hears; ana, who merged, does not.
1680 assert_eq!(told(&by_ana), vec![("bo", Reason::StateChange, Severity::Success)]);
1681 let by_bo = notices(&merged, "acme/rocket", &actor("usr_bo", "bo"), Some(&pull()), &nobody());
1682 assert_eq!(told(&by_bo), vec![("ana", Reason::StateChange, Severity::Success)]);
Inbox: the events service tells people what needs them as events arrive1683 assert_eq!(by_bo[0].title, "bo merged acme/rocket#7");
Inbox: threads, reasons, subscriptions and watching1684 let by_queue = notices(&merged, "acme/rocket", &Actor::default(), Some(&pull()), &nobody());
Inbox: the events service tells people what needs them as events arrive1685 assert_eq!(by_queue[0].title, "acme/rocket#7 was merged");
1686 }
1687
1688 #[test]
Inbox: threads, reasons, subscriptions and watching1689 fn closing_and_reopening_tells_everyone_subscribed_and_watchers() {
1690 let closed = event("issue.closed", Some("usr_bo"), json!({ "number": 3, "reason": "completed" }));
1691 let issue = InboxSubject {
1692 kind: Some(SubjectKind::Issue),
1693 title: "Crash".into(),
1694 author: person("usr_cy", "cy"),
1695 assignees: vec!["bo".into()],
1696 ..InboxSubject::default()
1697 };
1698 let mut audience = subscribed(&[("eve", State::Subscribed, Some(Reason::Comment)), ("fay", State::Unsubscribed, None)]);
1699 audience.watchers = watched(&[("gus", WatchLevel::All, &[]), ("hal", WatchLevel::Custom, &["pulls"])]).watchers;
1700 let notices_ = notices(&closed, "acme/rocket", &actor("usr_bo", "bo"), Some(&issue), &audience);
1701 assert_eq!(
1702 told(&notices_),
1703 vec![
1704 ("cy", Reason::StateChange, Severity::Info),
1705 ("eve", Reason::StateChange, Severity::Info),
1706 ("gus", Reason::Subscribed, Severity::Info),
1707 ]
1708 );
1709 assert_eq!(notices_[0].title, "bo closed acme/rocket#3");
1710 let by_pull = event("issue.closed", None, json!({ "number": 3, "resolvedBy": 9 }));
1711 assert_eq!(notices(&by_pull, "acme/rocket", &Actor::default(), Some(&issue), &nobody())[0].title, "acme/rocket#3 was closed by #9");
1712 let reopened = event("issue.reopened", Some("usr_cy"), json!({ "number": 3 }));
1713 assert_eq!(told(&notices(&reopened, "acme/rocket", &actor("usr_cy", "cy"), Some(&issue), &nobody())), vec![("bo", Reason::StateChange, Severity::Info)]);
1714 }
1715
1716 #[test]
1717 fn opening_tells_who_it_names_and_watchers_of_its_kind() {
1718 let opened = event("pull.opened", Some("usr_ana"), json!({ "number": 7 }));
1719 let subject = InboxSubject { reviewers: vec!["cy".into(), "g1t".into()], ..pull() };
1720 let audience = watched(&[("bo", WatchLevel::All, &[]), ("dee", WatchLevel::Custom, &["issues"]), ("eve", WatchLevel::Custom, &["pulls"]), ("fay", WatchLevel::Participating, &[])]);
1721 let notices_ = notices(&opened, "acme/rocket", &actor("usr_ana", "ana"), Some(&subject), &audience);
1722 assert_eq!(
1723 told(&notices_),
1724 vec![
1725 ("bo", Reason::Assign, Severity::Info),
1726 ("cy", Reason::ReviewRequested, Severity::Warning),
1727 ("eve", Reason::Subscribed, Severity::Info),
1728 ]
1729 );
1730 assert_eq!(notices_[2].title, "ana opened acme/rocket#7");
1731 }
1732
1733 #[test]
1734 fn ignoring_a_repository_silences_it_and_unsubscribing_keeps_only_what_is_asked() {
1735 let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
1736 let on = comment(person("usr_bo", "bo"), "@cy have a look", &["cy"]);
1737 let audience = watched(&[("cy", WatchLevel::Ignore, &[]), ("ana", WatchLevel::All, &[])]);
1738 assert!(notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &audience).iter().all(|notice| notice.username != "cy"));
1739 let audience = subscribed(&[("cy", State::Unsubscribed, None), ("ana", State::Unsubscribed, None)]);
1740 let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &audience);
1741 // cy was mentioned, which is asked of them; ana unsubscribed from the conversation.
1742 assert_eq!(told(&notices_), vec![("cy", Reason::Mention, Severity::Info)]);
1743 }
1744
1745 #[test]
Inbox: the events service tells people what needs them as events arrive1746 fn g1t_finishing_or_reviewing_tells_the_person_it_worked_for() {
1747 // The agent acts as the person it works for: still an outcome they hear of.
1748 let ready = event("pull.ready", Some("usr_ana"), json!({ "number": 7 }));
Inbox: threads, reasons, subscriptions and watching1749 let notices_ = notices(&ready, "acme/rocket", &actor("usr_ana", "ana"), Some(&g1t_pull()), &nobody());
1750 assert_eq!(told(&notices_), vec![("ana", Reason::Author, Severity::Success)]);
Inbox: the events service tells people what needs them as events arrive1751 assert_eq!(notices_[0].title, "g1t finished acme/rocket#7");
1752 // A person's draft marked ready is not news to them.
Inbox: threads, reasons, subscriptions and watching1753 assert!(notices(&ready, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody()).is_empty());
Inbox: the events service tells people what needs them as events arrive1754
1755 let approve = event("review.completed", None, json!({ "number": 7, "verdict": "approve" }));
Inbox: threads, reasons, subscriptions and watching1756 assert_eq!(
1757 told(&notices(&approve, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())),
1758 vec![("ana", Reason::Author, Severity::Success)]
1759 );
Inbox: the events service tells people what needs them as events arrive1760 let changes = event("review.completed", None, json!({ "number": 7, "verdict": "request_changes" }));
1761 assert_eq!(
Inbox: threads, reasons, subscriptions and watching1762 told(&notices(&changes, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())),
1763 vec![("ana", Reason::Author, Severity::Info)]
Inbox: the events service tells people what needs them as events arrive1764 );
1765 }
1766
1767 #[test]
Inbox: threads, reasons, subscriptions and watching1768 fn comments_tell_those_mentioned_then_everyone_subscribed_never_the_writer() {
Inbox: the events service tells people what needs them as events arrive1769 let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
1770 let on = comment(person("usr_bo", "bo"), "@cy @bo have a look", &["cy", "bo", "g1t"]);
Inbox: threads, reasons, subscriptions and watching1771 let audience = subscribed(&[("eve", State::Subscribed, Some(Reason::Comment)), ("fay", State::Subscribed, Some(Reason::Manual))]);
1772 let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &audience);
Inbox: the events service tells people what needs them as events arrive1773 assert_eq!(
1774 told(&notices_),
Inbox: threads, reasons, subscriptions and watching1775 vec![
1776 ("cy", Reason::Mention, Severity::Info),
1777 ("ana", Reason::Author, Severity::Info),
1778 ("eve", Reason::Comment, Severity::Info),
1779 ("fay", Reason::Manual, Severity::Info),
1780 ]
Inbox: the events service tells people what needs them as events arrive1781 );
1782 assert_eq!(notices_[0].title, "bo mentioned you on acme/rocket#7");
Inbox: threads, reasons, subscriptions and watching1783 assert_eq!(notices_[1].title, "bo commented on acme/rocket#7");
Inbox: the events service tells people what needs them as events arrive1784 assert_eq!(notices_[0].body, "@cy @bo have a look");
1785 // Mentioned and the owner: told once, as mentioned.
1786 let on = comment(person("usr_bo", "bo"), "@ana", &["ana"]);
Inbox: threads, reasons, subscriptions and watching1787 assert_eq!(
1788 told(&notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &nobody())),
1789 vec![("ana", Reason::Mention, Severity::Info)]
1790 );
1791 // The owner's own comment tells the assignee, not the owner.
Inbox: the events service tells people what needs them as events arrive1792 let on = comment(person("usr_ana", "ana"), "thanks", &[]);
Inbox: threads, reasons, subscriptions and watching1793 assert_eq!(
1794 told(&notices(&created, "acme/rocket", &actor("usr_ana", "ana"), Some(&on), &nobody())),
1795 vec![("bo", Reason::Assign, Severity::Info)]
1796 );
Inbox: the events service tells people what needs them as events arrive1797 // Something that happened, not something written, tells nobody.
1798 let mut on = comment(person("usr_bo", "bo"), "assigned cy", &[]);
1799 on.comment.as_mut().unwrap().event = true;
Inbox: threads, reasons, subscriptions and watching1800 assert!(notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &nobody()).is_empty());
Inbox: the events service tells people what needs them as events arrive1801 // An approval is good news.
Inbox: threads, reasons, subscriptions and watching1802 let mut on = comment(person("usr_cy", "cy"), "", &[]);
Inbox: the events service tells people what needs them as events arrive1803 on.comment.as_mut().unwrap().verdict = Some("approve".into());
Inbox: threads, reasons, subscriptions and watching1804 let notices_ = notices(&created, "acme/rocket", &actor("usr_cy", "cy"), Some(&on), &nobody());
1805 assert_eq!(told(&notices_)[0], ("ana", Reason::Author, Severity::Success));
Inbox: the events service tells people what needs them as events arrive1806 assert_eq!(notices_[0].body, "Add the inbox");
Inbox: threads, reasons, subscriptions and watching1807 // Writing and being mentioned subscribe, never g1t.
1808 let on = comment(person("usr_bo", "bo"), "@cy", &["cy"]);
1809 assert_eq!(subscribes(&created, Some(&on)), vec![("bo".into(), Reason::Comment), ("cy".into(), Reason::Mention)]);
1810 let on = comment(person(AGENT_ID, "g1t"), "done", &[]);
1811 assert!(subscribes(&created, Some(&on)).is_empty());
1812 }
1813
1814 #[test]
1815 fn issues_and_pull_requests_are_one_thread_each() {
1816 let created = event("comment.created", Some("usr_bo"), json!({ "repoId": "rep_1", "number": 7, "commentId": "cmt_1" }));
1817 let wanted = wants(&created).unwrap();
1818 let thread = thread_of(&created, &wanted, Some(&pull()));
1819 assert_eq!(thread.key, "rep_1#7");
1820 assert_eq!(thread.kind, Some(SubjectKind::Pull));
1821 let merged = event("pull.merged", None, json!({ "repoId": "rep_1", "number": 7 }));
1822 assert_eq!(thread_of(&merged, &wants(&merged).unwrap(), Some(&pull())).key, thread.key);
Inbox: the events service tells people what needs them as events arrive1823 }
1824
1825 #[test]
1826 fn items_link_to_what_they_are_about() {
Inbox: threads, reasons, subscriptions and watching1827 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Pull), Some(7), None, None), "/acme/rocket/pull/7");
1828 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Issue), Some(3), None, None), "/acme/rocket/issues/3");
1829 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Run), None, Some("run_1"), None), "/acme/rocket/actions/runs/run_1");
1830 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Deploy), None, None, Some("/acme/site/deployments/dpl_1")), "/acme/site/deployments/dpl_1");
1831 assert_eq!(url(Some("acme/rocket"), None, None, None, None), "/acme/rocket");
1832 assert_eq!(url(None, None, None, None, None), "/inbox");
Inbox: the events service tells people what needs them as events arrive1833 }
1834
1835 #[test]
1836 fn long_titles_are_cut_to_a_line() {
1837 let long = "x".repeat(400);
1838 assert_eq!(clip(&long, MAX_TITLE).chars().count(), MAX_TITLE);
1839 assert_eq!(clip(" short ", MAX_TITLE), "short");
1840 }
1841
1842 #[test]
1843 fn marks_change_only_what_they_say() {
1844 assert_eq!(mark_change(InboxMark::Read), "read_at = COALESCE(read_at, ?1)");
1845 assert!(mark_change(InboxMark::Done).contains("done_at = ?1"));
Inbox: threads, reasons, subscriptions and watching1846 assert_eq!(mark_change(InboxMark::Unsnooze), "snoozed_until = NULL");
1847 assert!(ranked(&ListInboxArgs::default()));
1848 assert!(!ranked(&ListInboxArgs { reason: Some(Reason::Mention), ..ListInboxArgs::default() }));
1849 assert!(!ranked(&ListInboxArgs { participating: true, ..ListInboxArgs::default() }));
1850 }
1851
1852 #[test]
1853 fn filters_bind_in_order_after_what_is_bound() {
1854 let a = ListInboxArgs {
1855 reason: Some(Reason::Mention),
1856 repo_id: Some("rep_1".into()),
1857 participating: true,
1858 unread: true,
1859 ..ListInboxArgs::default()
1860 };
1861 // Two values come first: the username and now.
1862 let (conditions, values) = filters(&a, 2);
1863 assert_eq!(
1864 conditions,
1865 vec!["reason = ?3", "repo_id = ?4", "reason NOT IN ('manual', 'subscribed')", "read_at IS NULL"]
1866 );
1867 assert_eq!(values, vec!["mention", "rep_1"]);
1868 }
1869
1870 #[test]
1871 fn the_bump_keeps_whatever_is_most_urgent_while_unread() {
1872 assert!(BUMP.contains("WHERE NOT EXISTS (SELECT 1 FROM inbox_activity WHERE event_id = ?4 AND username = ?2)"));
1873 assert!(BUMP.contains("ON CONFLICT (username, thread) DO UPDATE"));
1874 assert!(BUMP.contains("activity = inbox_items.activity + 1"));
1875 // Its numbered parameters run from 1 to 19.
1876 for at in 1..=19 {
1877 assert!(BUMP.contains(&format!("?{at}")), "?{at}");
1878 }
1879 assert!(!BUMP.contains("?20"));
Inbox: the events service tells people what needs them as events arrive1880 }
1881}