Skip to content
2,576 linesCodeBlameRaw

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

Inbox: the events service tells people what needs them as events arrive1//! The inbox, kept beside the event log. See `g1t_contracts::inbox`.
2//!
Inbox: threads, reasons, subscriptions and watching3//! As each batch arrives from the bus, the events that tell someone are
4//! read against what they name (the work service's `inbox_subject`), who
5//! subscribes to it and who watches its repository (subscriptions.rs), and
6//! each person told gets their thread about it brought to the top, unread,
7//! with a line added to its history. Who is told is worked out in
8//! [`notices`], from the event, its subject and that audience alone:
Inbox: the events service tells people what needs them as events arrive9//!
Inbox: threads, reasons, subscriptions and watching10//! | Event | Who | Reason | Severity |
11//! | --- | --- | --- | --- |
12//! | `agent.asked`, `pull.stalled` | the pull request's owner, and its issue's owner and assignees | agent | warning |
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar13//! | `pull.review_requested` | the reviewers asked, and the people each team asked tells | review_requested | warning |
Inbox: threads, reasons, subscriptions and watching14//! | `issue.assigned`, `pull.assigned` | the people newly assigned | assign | info |
15//! | `checks.completed`, failed or errored | the pull request's owner | ci_activity | error |
16//! | `workflow.completed`, failed | the pull request's owner, or whoever pushed | ci_activity | error |
17//! | `deployment.failed` | the pull request's owner, or whoever pushed; watchers | ci_activity | error |
18//! | `deployment.succeeded` after a failure | the same | ci_activity | success |
19//! | `review.completed` by g1t | the pull request's owner | author | success, or info for changes asked |
20//! | `pull.ready` for a change g1t made | whoever asked g1t for it | author | success |
21//! | `pull.merged`, `pull.closed`, `issue.closed`, `issue.reopened` | everyone subscribed | state_change | success for a merge, else info |
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar22//! | `comment.created` | everyone mentioned, the people of teams mentioned, then everyone subscribed | mention, team_mention, or why they are subscribed | info (success for an approval) |
Cards you act on in chat; agents comment and review as themselves; names shown cleanly; commits on the calendar23//! | `comment.created` by a workspace's agent | the same, shown as "Margo (agent)"; a review also tells the owner, "(advisory)" | the same, or author | the same |
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar24//! | `issue.opened`, `pull.opened` | whoever was assigned, asked to review or mentioned in its description (people and teams); watchers | assign, review_requested, mention, team_mention, subscribed | info |
Inbox: the events service tells people what needs them as events arrive25//!
Inbox: threads, reasons, subscriptions and watching26//! Watchers of a repository at `all` (or `custom`, for the kinds they
27//! chose) hear of every issue and pull request opened, commented on,
28//! closed, reopened or merged, and of deployments. Anyone who ignores the
29//! thread or the repository hears of nothing on it; anyone who
30//! unsubscribed hears only of what is asked of them.
31//!
32//! Nobody is told of what they did themselves, though an outcome they set
33//! off (checks, a workflow, a deployment, g1t's work) is theirs to hear of.
34//! g1t is never told. A failure here is logged and the batch goes on: the
35//! bus never waits on the inbox, so an item can be missed, but nothing else
36//! is held up.
Inbox: the events service tells people what needs them as events arrive37
38use std::collections::{HashMap, HashSet};
39
40use g1t_contracts::credentials::Principal;
41use g1t_contracts::events::{Event, WorkspaceRenamed};
42use g1t_contracts::identity::{AGENT_ID, UsernamesArgs};
43use g1t_contracts::inbox::*;
44use g1t_contracts::repos::{PathByIdArgs, ReadableArgs, Repo, RepoPath};
45use g1t_contracts::time::rfc3339;
46use g1t_contracts::{new_id, system};
47use g1t_kit::now_ms;
48use serde::Deserialize;
49use worker::wasm_bindgen::JsValue;
Inbox: threads, reasons, subscriptions and watching50use worker::{D1Database, D1PreparedStatement, Fetcher, Result};
51
52use crate::subscriptions::{self, Audience};
Inbox: the events service tells people what needs them as events arrive53
54/// Items marked done are kept this long, then removed.
55pub const DONE_DAYS: u32 = 30;
56/// No item is kept longer than this, unless it was saved.
57pub const MAX_DAYS: u32 = 180;
58/// Unread warnings shown ahead of everything else on the first page.
59const MAX_RANKED: u32 = 20;
60const MAX_TITLE: usize = 200;
61const MAX_BODY: usize = 300;
62
63/// What the inbox asks about an event before deciding who is told.
64#[derive(Debug, PartialEq, Eq)]
65pub struct Wanted {
66 pub repo_id: String,
67 /// The issue or pull request, when the event names one.
68 pub number: Option<u32>,
69 pub comment_id: Option<String>,
70}
71
72/// One person to tell, and what.
73#[derive(Debug, PartialEq, Eq)]
74pub struct Notice {
75 pub username: String,
Inbox: threads, reasons, subscriptions and watching76 pub reason: Reason,
Inbox: the events service tells people what needs them as events arrive77 pub severity: Severity,
78 pub title: String,
79 pub body: String,
80}
81
82/// Who did it, as far as is known.
83#[derive(Debug, Default)]
84pub struct Actor {
85 pub id: Option<String>,
86 pub username: Option<String>,
87}
88
89impl Actor {
90 fn is(&self, person: &Principal) -> bool {
91 self.id.as_deref().is_some_and(|id| id == person.id) || self.is_named(&person.username)
92 }
93
94 fn is_named(&self, username: &str) -> bool {
95 self.username.as_deref().is_some_and(|name| name.eq_ignore_ascii_case(username))
96 }
Inbox: threads, reasons, subscriptions and watching97
98 /// A person's name, for a title: never g1t's ids.
99 fn name(&self) -> Option<&str> {
100 self.username.as_deref().filter(|name| !is_g1t(name))
101 }
Inbox: the events service tells people what needs them as events arrive102}
103
Inbox: threads, reasons, subscriptions and watching104pub(crate) fn is_g1t(username: &str) -> bool {
Inbox: the events service tells people what needs them as events arrive105 username.eq_ignore_ascii_case(system::USERNAME) || username.eq_ignore_ascii_case("g1t-agent")
106}
107
Inbox: threads, reasons, subscriptions and watching108pub(crate) fn is_g1t_id(id: &str) -> bool {
Inbox: the events service tells people what needs them as events arrive109 system::is_system_id(id) || id == AGENT_ID
110}
111
Inbox: threads, reasons, subscriptions and watching112/// Events that end what an agent was waiting on a person for: it picked
113/// back up, its head moved, a merge was asked for, or it is over.
114const RESUMES: [&str; 5] = ["pull.resumed", "pull.updated", "pull.merge_requested", "pull.merged", "pull.closed"];
115
Inbox: the events service tells people what needs them as events arrive116/// The issue or pull request an event names, and what to read for it.
117/// None for events the inbox does not tell anyone of.
118pub fn wants(event: &Event) -> Option<Wanted> {
119 let data = &event.data;
120 let text = |key: &str| data[key].as_str().filter(|value| !value.is_empty()).map(str::to_owned);
121 let number = |key: &str| data[key].as_u64().and_then(|n| u32::try_from(n).ok());
122 let repo_id = text("repoId").or_else(|| event.repo_id.clone())?;
123 let on = |number: Option<u32>, comment_id: Option<String>| {
124 Some(Wanted {
125 repo_id: repo_id.clone(),
126 number,
127 comment_id,
128 })
129 };
130 match event.kind.as_str() {
Inbox: threads, reasons, subscriptions and watching131 "agent.asked" | "pull.stalled" | "pull.merged" | "pull.closed" | "pull.ready" | "pull.opened"
132 | "pull.review_requested" | "pull.assigned" | "issue.opened" | "issue.closed" | "issue.reopened"
133 | "issue.assigned" => on(Some(number("number")?), None),
Inbox: the events service tells people what needs them as events arrive134 "checks.completed" => match data["status"].as_str() {
135 Some("failed" | "errored") => on(Some(number("number")?), None),
136 _ => None,
137 },
138 "review.completed" => match data["verdict"].as_str() {
139 Some("approve" | "request_changes") => on(Some(number("number")?), None),
140 _ => None,
141 },
142 "workflow.completed" => match data["conclusion"].as_str() {
143 Some("failure") => on(number("pull"), None),
144 _ => None,
145 },
Inbox: threads, reasons, subscriptions and watching146 "deployment.failed" | "deployment.succeeded" => on(number("number"), None),
Merge branch 'worktree-agent-a3abfcce648e87dca'147 // A workflow run's jobs wait for their environment's reviewers.
148 "deployment.review_requested" => on(None, None),
Inbox: the events service tells people what needs them as events arrive149 "comment.created" => on(Some(number("number")?), Some(text("commentId")?)),
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar150 // Security alerts are threads of their own, not of an issue.
151 kind if SECURITY_EVENTS.contains(&kind) => on(None, None),
Merge branch 'mirroring' into artifacts-mode152 // So is a repository's mirror, told to whom its settings name.
153 kind if MIRROR_EVENTS.contains(&kind) => on(None, None),
Inbox: the events service tells people what needs them as events arrive154 _ => None,
155 }
156}
157
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar158/// The security service's events the inbox tells people of: new alerts,
159/// and push protection bypasses asked for and decided. Fixes and
160/// dismissals are on the Security page and in webhooks.
161pub const SECURITY_EVENTS: [&str; 5] = [
162 "secret_scanning_alert.created",
163 "code_scanning_alert.created",
164 "vulnerability_alert.created",
165 "secret_scanning.bypass_requested",
166 "secret_scanning.bypass_reviewed",
167];
168
Merge branch 'mirroring' into artifacts-mode169/// The integrations service's mirroring events (see
170/// `g1t_contracts::mirrors::MirrorEvent`). People are told only when the
171/// link's settings ask for the inbox, and the event names them.
172pub const MIRROR_EVENTS: [&str; 4] = ["mirror.unreachable", "mirror.reachable", "mirror.state_changed", "mirror.moved_in"];
173
Inbox: threads, reasons, subscriptions and watching174/// Where an event's items go: which thread, what it is, and where it is.
175#[derive(Clone, Debug, PartialEq, Eq)]
176pub struct Thread {
177 pub key: String,
178 pub kind: Option<SubjectKind>,
179 pub number: Option<u32>,
180 pub run_id: Option<String>,
181 /// A path, for a subject with a page of its own.
182 pub link: Option<String>,
183}
184
185/// The thread key of an issue or pull request.
186pub fn numbered_thread(repo_id: &str, number: u32) -> String {
187 format!("{repo_id}#{number}")
188}
189
190/// The thread an event's items go to.
191pub fn thread_of(event: &Event, wanted: &Wanted, subject: Option<&InboxSubject>) -> Thread {
192 let data = &event.data;
193 let text = |key: &str| data[key].as_str().filter(|value| !value.is_empty()).map(str::to_owned);
194 match event.kind.as_str() {
195 // A project's production, or one pull request's preview.
196 "deployment.failed" | "deployment.succeeded" => {
197 let which = wanted.number.map_or_else(|| "production".to_owned(), |number| number.to_string());
198 Thread {
199 key: format!("{}/deploy/{}/{which}", wanted.repo_id, text("projectId").unwrap_or_default()),
200 kind: Some(SubjectKind::Deploy),
201 number: wanted.number,
202 run_id: text("deploymentId"),
203 link: text("path"),
204 }
205 }
Merge branch 'worktree-agent-a3abfcce648e87dca'206 // One run's deployment to one environment, waiting for review.
207 "deployment.review_requested" => Thread {
208 key: format!(
209 "{}/review/{}/{}",
210 wanted.repo_id,
211 text("runId").unwrap_or_default(),
212 text("environment").unwrap_or_default()
213 ),
214 kind: Some(SubjectKind::Run),
215 number: None,
216 run_id: text("runId"),
217 link: text("link"),
218 },
Inbox: threads, reasons, subscriptions and watching219 // A workflow on a branch: its next failure bumps the same thread.
220 "workflow.completed" if subject.is_none() => {
221 let branch = text("ref").unwrap_or_default();
222 let branch = branch.strip_prefix("refs/heads/").unwrap_or(&branch).to_owned();
223 let workflow = text("path").or_else(|| text("workflow")).unwrap_or_default();
224 Thread {
225 key: format!("{}/run/{workflow}@{branch}", wanted.repo_id),
226 kind: Some(SubjectKind::Run),
227 number: None,
228 run_id: text("runId"),
229 link: None,
230 }
231 }
Merge branch 'mirroring' into artifacts-mode232 // A repository's mirror: one thread, its settings page the link.
233 kind if MIRROR_EVENTS.contains(&kind) => Thread {
234 key: format!("{}/mirror", wanted.repo_id),
235 kind: None,
236 number: None,
237 run_id: None,
238 link: text("link"),
239 },
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar240 // One alert, or one bypass request: its page is the link.
241 kind if SECURITY_EVENTS.contains(&kind) => {
242 let which = text("requestId").or_else(|| text("alertId")).unwrap_or_else(|| event.id.clone());
243 Thread {
244 key: format!("{}/security/{which}", wanted.repo_id),
245 kind: None,
246 number: None,
247 run_id: None,
248 link: text("link"),
249 }
250 }
Inbox: threads, reasons, subscriptions and watching251 _ => match wanted.number {
252 Some(number) => Thread {
253 key: numbered_thread(&wanted.repo_id, number),
254 kind: subject.and_then(|subject| subject.kind),
255 number: Some(number),
256 run_id: None,
257 link: None,
258 },
259 None => Thread {
260 key: format!("event/{}", event.id),
261 kind: None,
262 number: None,
263 run_id: None,
264 link: None,
265 },
266 },
267 }
268}
269
270/// The issue or pull request whose agent threads an event closes: what an
271/// agent was waiting on a person for is over.
272pub fn resolves(event: &Event) -> Option<String> {
273 if !RESUMES.contains(&event.kind.as_str()) {
274 return None;
275 }
276 let repo_id = event.data["repoId"].as_str().map(str::to_owned).or_else(|| event.repo_id.clone())?;
277 let number = event.data["number"].as_u64().and_then(|n| u32::try_from(n).ok())?;
278 Some(numbered_thread(&repo_id, number))
279}
280
281/// Collects who is told, each once with the most specific reason, never
282/// the actor, never g1t, and never anyone ignoring the thread.
Inbox: the events service tells people what needs them as events arrive283struct Told<'a> {
284 actor: &'a Actor,
Inbox: threads, reasons, subscriptions and watching285 audience: &'a Audience,
Inbox: the events service tells people what needs them as events arrive286 notices: Vec<Notice>,
287}
288
289impl Told<'_> {
Inbox: threads, reasons, subscriptions and watching290 /// `asked`: something asked of the person directly, told even when
291 /// they unsubscribed from the thread.
292 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 arrive293 let username = username.trim().trim_start_matches('@').to_lowercase();
294 if username.is_empty()
295 || is_g1t(&username)
296 || self.actor.is_named(&username)
Inbox: threads, reasons, subscriptions and watching297 || self.audience.ignores(&username)
298 || (!asked && self.audience.unsubscribed(&username))
Inbox: the events service tells people what needs them as events arrive299 {
300 return;
301 }
Inbox: threads, reasons, subscriptions and watching302 let notice = Notice {
Inbox: the events service tells people what needs them as events arrive303 username,
304 reason,
305 severity,
306 title: clip(title, MAX_TITLE),
307 body: clip(body, MAX_BODY),
Inbox: threads, reasons, subscriptions and watching308 };
309 match self.notices.iter_mut().find(|told| told.username == notice.username) {
310 Some(told) if notice.reason.rank() < told.reason.rank() => *told = notice,
311 Some(_) => {}
312 None => self.notices.push(notice),
313 }
Inbox: the events service tells people what needs them as events arrive314 }
315
316 /// A person known by id as well as name, such as an author.
Inbox: threads, reasons, subscriptions and watching317 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 arrive318 if is_g1t_id(&person.id) || self.actor.is(person) {
319 return;
320 }
Inbox: threads, reasons, subscriptions and watching321 self.tell(&person.username, reason, severity, title, body, asked);
322 }
323
324 /// Everyone subscribed to an issue or pull request: its owner and
325 /// author, its assignees and reviewers, and whoever subscribed by
326 /// commenting, being mentioned or by hand. `reason` overrides why
327 /// each is told, as a state change does.
328 fn tell_subscribed(&mut self, subject: &InboxSubject, reason: Option<Reason>, severity: Severity, title: &str, body: &str) {
329 let why = |own: Reason| reason.unwrap_or(own);
330 self.tell_person(subject.owner(), why(Reason::Author), severity, title, body, false);
331 self.tell_person(&subject.author, why(Reason::Author), severity, title, body, false);
332 for name in &subject.assignees {
333 self.tell(name, why(Reason::Assign), severity, title, body, false);
334 }
335 for name in &subject.reviewers {
336 self.tell(name, why(Reason::ReviewRequested), severity, title, body, false);
337 }
338 let subscribed: Vec<(String, Reason)> = self.audience.subscribed().collect();
339 for (name, own) in subscribed {
340 self.tell(&name, why(own), severity, title, body, false);
341 }
Inbox: the events service tells people what needs them as events arrive342 }
Inbox: threads, reasons, subscriptions and watching343
344 /// Everyone watching the repository for this kind of activity.
345 fn tell_watchers(&mut self, kind: &str, severity: Severity, title: &str, body: &str) {
346 let watching: Vec<String> = self.audience.watching(kind).collect();
347 for name in watching {
348 self.tell(&name, Reason::Subscribed, severity, title, body, false);
349 }
350 }
Inbox: the events service tells people what needs them as events arrive351}
352
353fn clip(text: &str, max: usize) -> String {
354 let text = text.trim();
355 if text.chars().count() <= max {
356 return text.to_owned();
357 }
358 let cut: String = text.chars().take(max - 1).collect();
359 format!("{}…", cut.trim_end())
360}
361
Inbox: threads, reasons, subscriptions and watching362/// The kind of activity a watcher chooses: `issues` or `pulls`.
363fn activity_of(subject: &InboxSubject) -> &'static str {
364 match subject.kind {
365 Some(SubjectKind::Pull) => "pulls",
366 _ => "issues",
367 }
368}
369
370/// The names in an event's list field, such as the reviewers just asked.
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar371/// The teams an event asked to review: each `workspace/slug` with the
372/// people it tells.
373fn teams_asked(data: &serde_json::Value) -> Vec<(String, Vec<String>)> {
374 data["teams"]
375 .as_array()
376 .map(|teams| {
377 teams
378 .iter()
379 .filter_map(|team| Some((team["team"].as_str()?.to_owned(), names(team, "notified"))))
380 .collect()
381 })
382 .unwrap_or_default()
383}
384
Inbox: threads, reasons, subscriptions and watching385fn names(data: &serde_json::Value, key: &str) -> Vec<String> {
386 data[key]
387 .as_array()
388 .map(|names| names.iter().filter_map(|name| name.as_str().map(str::to_owned)).collect())
389 .unwrap_or_default()
390}
391
392/// Who is told of `event`, in `repo` (`owner/name`), given what it names
393/// and who follows it. `actor` is who caused it; outcomes nobody chose
394/// (checks, workflows, deployments, a review, an agent finishing or
395/// stopping) are told whoever caused them.
396pub 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 arrive397 let nobody = Actor::default();
398 let data = &event.data;
399 // The issue or pull request, as titles name it: `acme/rocket#12`.
400 let at = match data["number"].as_u64().or(data["pull"].as_u64()) {
401 Some(number) if subject.is_some() => format!("{repo}#{number}"),
402 _ => repo.to_owned(),
403 };
404 let outcome = matches!(
405 event.kind.as_str(),
Inbox: threads, reasons, subscriptions and watching406 "checks.completed"
407 | "workflow.completed"
408 | "review.completed"
409 | "pull.ready"
410 | "agent.asked"
411 | "pull.stalled"
412 | "deployment.failed"
413 | "deployment.succeeded"
Merge branch 'worktree-agent-a3abfcce648e87dca'414 | "deployment.review_requested"
Merge branch 'mirroring' into artifacts-mode415 ) || SECURITY_EVENTS.contains(&event.kind.as_str())
416 || MIRROR_EVENTS.contains(&event.kind.as_str());
Inbox: the events service tells people what needs them as events arrive417 let mut told = Told {
418 actor: if outcome { &nobody } else { actor },
Inbox: threads, reasons, subscriptions and watching419 audience,
Inbox: the events service tells people what needs them as events arrive420 notices: Vec::new(),
421 };
Inbox: threads, reasons, subscriptions and watching422 // Who did it, as titles name them.
423 let who = actor.name().unwrap_or("g1t");
Inbox: the events service tells people what needs them as events arrive424
425 match (event.kind.as_str(), subject) {
Inbox: threads, reasons, subscriptions and watching426 ("agent.asked" | "pull.stalled", Some(pull)) => {
427 let (title, body) = if event.kind == "agent.asked" {
428 (format!("An agent is waiting on {at}"), pull.title.clone())
429 } else {
430 let detail = data["detail"].as_str().map(str::trim).filter(|detail| !detail.is_empty());
431 (format!("g1t stopped on {at} and needs you"), detail.unwrap_or(&pull.title).to_owned())
432 };
433 told.tell_person(pull.owner(), Reason::Agent, Severity::Warning, &title, &body, true);
Inbox: the events service tells people what needs them as events arrive434 if let Some(issue) = &pull.issue {
Inbox: threads, reasons, subscriptions and watching435 told.tell_person(issue.owner(), Reason::Agent, Severity::Warning, &title, &body, true);
Inbox: the events service tells people what needs them as events arrive436 for name in &issue.assignees {
Inbox: threads, reasons, subscriptions and watching437 told.tell(name, Reason::Agent, Severity::Warning, &title, &body, true);
Inbox: the events service tells people what needs them as events arrive438 }
439 }
440 }
Inbox: threads, reasons, subscriptions and watching441 ("pull.review_requested", Some(pull)) => {
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar442 // Asked because the CODEOWNERS file says they own what changed.
443 let owned = data["codeOwners"].as_bool() == Some(true);
444 let title = if owned {
445 format!("{at} changes files you own")
446 } else {
447 format!("{who} asked you to review {at}")
448 };
Inbox: threads, reasons, subscriptions and watching449 for name in names(data, "reviewers") {
450 told.tell(&name, Reason::ReviewRequested, Severity::Warning, &title, &pull.title, true);
451 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar452 for (team, people) in teams_asked(data) {
453 let title = if owned {
454 format!("{at} changes files @{team} owns")
455 } else {
456 format!("{who} asked @{team} to review {at}")
457 };
458 for name in people {
459 told.tell(&name, Reason::ReviewRequested, Severity::Warning, &title, &pull.title, true);
460 }
461 }
Inbox: threads, reasons, subscriptions and watching462 }
463 ("issue.assigned" | "pull.assigned", Some(on)) => {
464 let title = format!("{who} assigned you to {at}");
465 for name in names(data, "added") {
466 told.tell(&name, Reason::Assign, Severity::Info, &title, &on.title, true);
467 }
468 }
469 ("issue.opened" | "pull.opened", Some(on)) => {
470 let title = format!("{who} opened {at}");
471 for name in &on.assignees {
472 told.tell(name, Reason::Assign, Severity::Info, &format!("{who} assigned you to {at}"), &on.title, true);
473 }
474 for name in &on.reviewers {
475 told.tell(name, Reason::ReviewRequested, Severity::Warning, &format!("{who} asked you to review {at}"), &on.title, true);
476 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar477 for name in &on.mentions {
478 told.tell(name, Reason::Mention, Severity::Info, &format!("{who} mentioned you on {at}"), &on.title, true);
479 }
480 for team in &on.team_mentions {
481 let title = format!("{who} mentioned @{} on {at}", team.team);
482 for name in &team.members {
483 told.tell(name, Reason::TeamMention, Severity::Info, &title, &on.title, true);
484 }
485 }
Inbox: threads, reasons, subscriptions and watching486 told.tell_watchers(activity_of(on), Severity::Info, &title, &on.title);
487 }
Inbox: the events service tells people what needs them as events arrive488 ("checks.completed", Some(pull)) => {
489 let title = match data["status"].as_str() {
Inbox: threads, reasons, subscriptions and watching490 Some("errored") => format!("Checks could not run on {at}"),
491 _ => format!("Checks failed on {at}"),
Inbox: the events service tells people what needs them as events arrive492 };
Inbox: threads, reasons, subscriptions and watching493 told.tell_person(pull.owner(), Reason::CiActivity, Severity::Error, &title, &pull.title, true);
Inbox: the events service tells people what needs them as events arrive494 }
495 ("workflow.completed", subject) => {
496 let workflow = data["workflow"].as_str().filter(|name| !name.is_empty()).unwrap_or("A workflow");
497 match subject {
498 Some(pull) => {
Inbox: threads, reasons, subscriptions and watching499 let title = format!("{workflow} failed on {at}");
500 told.tell_person(pull.owner(), Reason::CiActivity, Severity::Error, &title, &pull.title, true);
Inbox: the events service tells people what needs them as events arrive501 }
502 None => {
503 // Not on a pull request: whoever pushed the commit it ran on.
504 let branch = data["ref"].as_str().unwrap_or_default();
505 let branch = branch.strip_prefix("refs/heads/").unwrap_or(branch);
506 let title = format!("{workflow} failed on {branch} in {repo}");
507 let body = format!("Run {} at {}", data["number"], short(data["sha"].as_str().unwrap_or_default()));
508 if let Some(id) = &actor.id
509 && !is_g1t_id(id)
510 && let Some(name) = &actor.username
511 {
Inbox: threads, reasons, subscriptions and watching512 told.tell(name, Reason::CiActivity, Severity::Error, &title, &body, true);
513 }
514 }
515 }
516 }
517 ("deployment.failed" | "deployment.succeeded", subject) => {
518 let failed = event.kind == "deployment.failed";
519 let project = data["project"].as_str().filter(|name| !name.is_empty()).unwrap_or(repo);
520 let what = match subject {
521 Some(_) => format!("The preview of {at}"),
522 None => format!("Production of {project}"),
523 };
524 let (title, severity) = if failed {
525 (format!("{what} failed to deploy"), Severity::Error)
526 } else {
527 (format!("{what} is live"), Severity::Success)
528 };
529 let body = match (failed, data["error"].as_str().map(str::trim).filter(|error| !error.is_empty())) {
530 (true, Some(error)) => error.to_owned(),
531 _ => match subject {
532 Some(pull) => pull.title.clone(),
533 None => format!("Commit {}", short(data["commit"].as_str().unwrap_or_default())),
534 },
535 };
536 // A success is news to whoever answers for it only after a failure.
537 if failed || data["recovered"].as_bool() == Some(true) {
538 let title = if failed { title.clone() } else { format!("{what} is live again") };
539 match subject {
540 Some(pull) => told.tell_person(pull.owner(), Reason::CiActivity, severity, &title, &body, true),
541 None => {
542 let pusher = actor
543 .id
544 .as_deref()
545 .filter(|id| !is_g1t_id(id))
546 .and(actor.username.as_deref())
547 .or_else(|| data["triggeredBy"].as_str());
548 if let Some(name) = pusher {
549 told.tell(name, Reason::CiActivity, severity, &title, &body, true);
550 }
Inbox: the events service tells people what needs them as events arrive551 }
552 }
553 }
Inbox: threads, reasons, subscriptions and watching554 told.tell_watchers("deployments", severity, &title, &body);
Inbox: the events service tells people what needs them as events arrive555 }
556 ("review.completed", Some(pull)) => {
Inbox: threads, reasons, subscriptions and watching557 let (title, severity) = match data["verdict"].as_str() {
558 Some("approve") => (format!("g1t approved {at}"), Severity::Success),
559 _ => (format!("g1t asked for changes on {at}"), Severity::Info),
Inbox: the events service tells people what needs them as events arrive560 };
Inbox: threads, reasons, subscriptions and watching561 told.tell_person(pull.owner(), Reason::Author, severity, &title, &pull.title, true);
Inbox: the events service tells people what needs them as events arrive562 }
563 ("pull.ready", Some(pull)) => {
564 // A change g1t made is ready: the agent's run is over.
565 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 watching566 let title = format!("g1t finished {at}");
567 told.tell_person(owner, Reason::Author, Severity::Success, &title, &pull.title, true);
Inbox: the events service tells people what needs them as events arrive568 }
569 }
Inbox: threads, reasons, subscriptions and watching570 ("pull.merged" | "pull.closed" | "issue.closed" | "issue.reopened", Some(on)) => {
571 let (verb, severity) = match event.kind.as_str() {
572 "pull.merged" => ("merged", Severity::Success),
573 "issue.reopened" => ("reopened", Severity::Info),
574 _ => ("closed", Severity::Info),
575 };
576 let title = match (actor.name(), data["resolvedBy"].as_u64()) {
577 (_, Some(pull)) if event.kind == "issue.closed" => format!("{at} was closed by #{pull}"),
578 (Some(name), _) => format!("{name} {verb} {at}"),
579 (None, _) => format!("{at} was {verb}"),
Inbox: the events service tells people what needs them as events arrive580 };
Inbox: threads, reasons, subscriptions and watching581 told.tell_subscribed(on, Some(Reason::StateChange), severity, &title, &on.title);
582 told.tell_watchers(activity_of(on), severity, &title, &on.title);
Inbox: the events service tells people what needs them as events arrive583 }
Merge branch 'worktree-agent-a3abfcce648e87dca'584 ("deployment.review_requested", _) => {
585 let environment = data["environment"].as_str().unwrap_or("an environment");
586 let workflow = data["workflow"].as_str().unwrap_or("A workflow");
587 let title = format!("{workflow} is waiting for your review to deploy to {environment} in {repo}");
588 let body = data["title"].as_str().unwrap_or_default();
589 // Those the actions service names: the environment's reviewers.
590 for name in names(data, "notify") {
591 told.tell(&name, Reason::ReviewRequested, Severity::Warning, &title, body, true);
592 }
593 }
Merge branch 'mirroring' into artifacts-mode594 (kind, None) if MIRROR_EVENTS.contains(&kind) => {
595 let severity = match kind {
596 "mirror.unreachable" => Severity::Warning,
597 "mirror.state_changed" if data["state"].as_str() == Some("takeover") => Severity::Warning,
598 "mirror.state_changed" if data["state"].as_str() == Some("standby") => Severity::Success,
599 _ => Severity::Info,
600 };
601 let title = data["title"].as_str().unwrap_or("A mirror changed");
602 let body = data["detail"].as_str().unwrap_or_default();
603 for name in names(data, "notify") {
604 told.tell(&name, Reason::StateChange, severity, title, body, true);
605 }
606 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar607 (kind, None) if SECURITY_EVENTS.contains(&kind) => {
608 let severity = match data["severity"].as_str() {
609 _ if kind == "secret_scanning.bypass_requested" => Severity::Warning,
610 _ if kind == "secret_scanning.bypass_reviewed" => Severity::Info,
611 Some("critical" | "high") => Severity::Error,
612 _ => Severity::Warning,
613 };
614 let what = data["title"].as_str().unwrap_or("A security alert");
615 let title = match kind {
616 "secret_scanning.bypass_requested" => format!("{who} asked to bypass push protection in {repo}"),
617 "secret_scanning.bypass_reviewed" => {
618 let verdict = data["state"].as_str().unwrap_or("reviewed");
619 format!("Your request to bypass push protection in {repo} was {verdict}")
620 }
621 _ if data["pusher"].is_string() => format!("A push to {repo} was blocked: it adds a secret"),
622 _ => format!("New security alert in {repo}"),
623 };
624 // Whoever pushed a blocked secret, and those the service names
625 // (the workspace's owners, or whoever asked for a bypass).
626 if let Some(pusher) = data["pusher"].as_str() {
627 told.tell(pusher, Reason::SecurityAlert, Severity::Error, &title, what, true);
628 }
629 for name in names(data, "notify") {
630 told.tell(&name, Reason::SecurityAlert, severity, &title, what, true);
631 }
632 // Watchers who chose security alerts, if they may see findings:
633 // the service lists who may (`members`), as findings are never
634 // shown to someone who cannot change the code.
635 if !kind.starts_with("secret_scanning.bypass") {
636 let members = names(data, "members");
637 let watching: Vec<String> = told.audience.watching("security").collect();
638 for name in watching.iter().filter(|name| members.iter().any(|member| member.eq_ignore_ascii_case(name))) {
639 told.tell(name, Reason::SecurityAlert, severity, &title, what, false);
640 }
641 }
642 }
Inbox: the events service tells people what needs them as events arrive643 ("comment.created", Some(on)) => {
644 let Some(comment) = on.comment.as_ref().filter(|comment| !comment.event) else {
645 return Vec::new();
646 };
Cards you act on in chat; agents comment and review as themselves; names shown cleanly; commits on the calendar647 // Whoever wrote it is the actor, whatever the event says. One of
648 // the workspace's agents is shown by its name, marked as an
649 // agent; its handle is not a person's, so nobody is left out
650 // for sharing it, and whoever it acted for hears of it as of
651 // anything they set off.
652 let writer = match &comment.agent {
653 Some(agent) => Actor { id: Some(agent.id.clone()), username: None },
654 None => Actor {
655 id: Some(comment.author.id.clone()),
656 username: Some(comment.author.username.clone()),
657 },
Inbox: the events service tells people what needs them as events arrive658 };
659 told.actor = &writer;
Cards you act on in chat; agents comment and review as themselves; names shown cleanly; commits on the calendar660 let who = match &comment.agent {
661 Some(agent) => format!("{} (agent)", agent.display_name),
662 None => comment.author.username.clone(),
663 };
664 let advisory = if comment.advisory { " (advisory)" } else { "" };
Inbox: the events service tells people what needs them as events arrive665 let body = if comment.excerpt.is_empty() { &on.title } else { &comment.excerpt };
666 for name in &comment.mentions {
Inbox: threads, reasons, subscriptions and watching667 let title = format!("{who} mentioned you on {at}");
668 told.tell(name, Reason::Mention, Severity::Info, &title, body, true);
Inbox: the events service tells people what needs them as events arrive669 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar670 for team in &comment.team_mentions {
671 let title = format!("{who} mentioned @{} on {at}", team.team);
672 for name in &team.members {
673 told.tell(name, Reason::TeamMention, Severity::Info, &title, body, true);
674 }
675 }
Inbox: threads, reasons, subscriptions and watching676 let (severity, title) = match comment.verdict.as_deref() {
Cards you act on in chat; agents comment and review as themselves; names shown cleanly; commits on the calendar677 Some("approve") => (Severity::Success, format!("{who} approved {at}{advisory}")),
678 Some("request_changes") => (Severity::Info, format!("{who} asked for changes on {at}{advisory}")),
679 _ if comment.advisory => (Severity::Info, format!("{who} reviewed {at}{advisory}")),
Inbox: threads, reasons, subscriptions and watching680 _ => (Severity::Info, format!("{who} commented on {at}")),
Inbox: the events service tells people what needs them as events arrive681 };
Cards you act on in chat; agents comment and review as themselves; names shown cleanly; commits on the calendar682 // An approval or a request for changes (or an agent's review)
683 // is the owner's to hear of whatever they chose; the rest, as
684 // they subscribed.
685 if comment.verdict.is_some() || comment.advisory {
Inbox: threads, reasons, subscriptions and watching686 told.tell_person(on.owner(), Reason::Author, severity, &title, body, true);
687 }
688 told.tell_subscribed(on, None, severity, &title, body);
689 told.tell_watchers(activity_of(on), severity, &title, body);
Inbox: the events service tells people what needs them as events arrive690 return told.notices;
691 }
692 _ => {}
693 }
694 told.notices
695}
696
Inbox: threads, reasons, subscriptions and watching697/// Who an event subscribes to its issue or pull request without asking,
698/// and why: whoever commented, whoever they mentioned, the people assigned
699/// and the reviewers asked. Never g1t.
700pub fn subscribes(event: &Event, subject: Option<&InboxSubject>) -> Vec<(String, Reason)> {
701 let mut people: Vec<(String, Reason)> = Vec::new();
702 let mut add = |name: &str, reason: Reason| {
703 let name = name.trim().trim_start_matches('@').to_lowercase();
704 if !name.is_empty() && !is_g1t(&name) && !people.iter().any(|(had, _)| *had == name) {
705 people.push((name, reason));
706 }
707 };
708 match event.kind.as_str() {
709 "comment.created" => {
710 if let Some(comment) = subject.and_then(|subject| subject.comment.as_ref()).filter(|comment| !comment.event) {
Cards you act on in chat; agents comment and review as themselves; names shown cleanly; commits on the calendar711 // An agent is no person to subscribe; its handle may be someone's name.
712 if !is_g1t_id(&comment.author.id) && comment.agent.is_none() {
Inbox: threads, reasons, subscriptions and watching713 add(&comment.author.username, Reason::Comment);
714 }
715 for name in &comment.mentions {
716 add(name, Reason::Mention);
717 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar718 for team in &comment.team_mentions {
719 for name in &team.members {
720 add(name, Reason::TeamMention);
721 }
722 }
723 }
724 }
725 "issue.opened" | "pull.opened" => {
726 if let Some(on) = subject {
727 for name in &on.mentions {
728 add(name, Reason::Mention);
729 }
730 for team in &on.team_mentions {
731 for name in &team.members {
732 add(name, Reason::TeamMention);
733 }
734 }
Inbox: threads, reasons, subscriptions and watching735 }
736 }
737 "issue.assigned" | "pull.assigned" => {
738 for name in names(&event.data, "added") {
739 add(&name, Reason::Assign);
740 }
741 }
742 "pull.review_requested" => {
743 for name in names(&event.data, "reviewers") {
744 add(&name, Reason::ReviewRequested);
745 }
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar746 for (_, people) in teams_asked(&event.data) {
747 for name in people {
748 add(&name, Reason::ReviewRequested);
749 }
750 }
Inbox: threads, reasons, subscriptions and watching751 }
752 _ => {}
753 }
754 people
755}
756
Inbox: the events service tells people what needs them as events arrive757fn short(sha: &str) -> &str {
758 sha.get(..7).unwrap_or(sha)
759}
760
761/// Where an item is on g1t.sh.
Inbox: threads, reasons, subscriptions and watching762pub fn url(repo: Option<&str>, subject: Option<SubjectKind>, number: Option<u32>, run_id: Option<&str>, link: Option<&str>) -> String {
763 if let Some(link) = link.filter(|link| link.starts_with('/')) {
764 return link.to_owned();
765 }
Inbox: the events service tells people what needs them as events arrive766 let Some(repo) = repo else {
767 return "/inbox".to_owned();
768 };
769 match (subject, number, run_id) {
770 (Some(SubjectKind::Pull), Some(number), _) => format!("/{repo}/pull/{number}"),
771 (Some(SubjectKind::Issue), Some(number), _) => format!("/{repo}/issues/{number}"),
772 (Some(SubjectKind::Run), _, Some(run)) => format!("/{repo}/actions/runs/{run}"),
Inbox: threads, reasons, subscriptions and watching773 (Some(SubjectKind::Deploy), Some(number), _) => format!("/{repo}/pull/{number}"),
Inbox: the events service tells people what needs them as events arrive774 _ => format!("/{repo}"),
775 }
776}
777
778// --- Writing ---------------------------------------------------------------
779
780/// The services the inbox reads from as events arrive.
781pub struct Sources<'a> {
782 pub work: &'a Fetcher,
783 pub repos: &'a Fetcher,
784 pub identity: &'a Fetcher,
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)785 /// The notify service, told of each new item for live toasts and pushes; None, nobody is.
786 pub notify: Option<&'a Fetcher>,
Inbox: the events service tells people what needs them as events arrive787}
788
Inbox: threads, reasons, subscriptions and watching789/// One notice written, and whether it was news (not a redelivery): what
790/// may be emailed.
791struct Written {
792 notice: Notice,
793 repo_id: String,
794 url: String,
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)795 event_id: String,
Inbox: threads, reasons, subscriptions and watching796}
797
Inbox: the events service tells people what needs them as events arrive798/// Writes the items a batch from the bus calls for. Never fails the batch:
799/// what cannot be worked out is logged and left.
800pub async fn deliver(db: &D1Database, sources: &Sources<'_>, events: &[Event]) {
Inbox: threads, reasons, subscriptions and watching801 if let Err(error) = resolve(db, events).await {
802 worker::console_error!("inbox: agent threads not closed: {error}");
803 }
Fine-grained personal tokens, workspace token rules and approvals in identity804 // A workspace's own notices, about no repository (workspace_notices).
805 for event in events.iter().filter(|event| WORKSPACE_EVENTS.contains(&event.kind.as_str())) {
806 if let Err(error) = deliver_workspace(db, event).await {
807 worker::console_error!("inbox: {} {} not delivered: {error}", event.kind, event.id);
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)808 } else if let Some(notify) = sources.notify {
809 // Notify: the feed tells each person once per event, redelivered or not.
810 let (_, told) = workspace_notices(event, None);
811 let link = event.data["link"].as_str().filter(|link| link.starts_with('/')).unwrap_or("/inbox");
812 let live = told.iter().map(|notice| (notice.username.clone(), live_notification(&event.id, notice, link, &event.time))).collect();
813 tell_live(notify, live).await;
814 for notice in &told {
815 tell_inbox_count(db, notify, &notice.username).await;
816 }
Fine-grained personal tokens, workspace token rules and approvals in identity817 }
818 }
Inbox: the events service tells people what needs them as events arrive819 let wanted: Vec<(&Event, Wanted)> = events
820 .iter()
821 .filter_map(|event| wants(event).map(|wanted| (event, wanted)))
822 .collect();
Inbox: threads, reasons, subscriptions and watching823 let created: Vec<&Event> = events.iter().filter(|event| event.kind == "repo.created").collect();
824 if wanted.is_empty() && created.is_empty() {
Inbox: the events service tells people what needs them as events arrive825 return;
826 }
827 // Everyone who caused one, named in one call.
828 let ids: Vec<String> = wanted
829 .iter()
Inbox: threads, reasons, subscriptions and watching830 .map(|(event, _)| *event)
831 .chain(created.iter().copied())
832 .filter_map(|event| event.actor.clone())
Inbox: the events service tells people what needs them as events arrive833 .filter(|id| !is_g1t_id(id))
834 .collect::<HashSet<_>>()
835 .into_iter()
836 .collect();
837 let names: HashMap<String, String> = if ids.is_empty() {
838 HashMap::new()
839 } else {
840 g1t_kit::call(sources.identity, "usernames", &UsernamesArgs { ids })
841 .await
842 .unwrap_or_else(|error| {
843 worker::console_error!("inbox: could not name who acted: {error}");
844 HashMap::new()
845 })
846 };
Inbox: threads, reasons, subscriptions and watching847 for event in created {
848 if let Err(error) = subscriptions::watch_created(db, event, &names).await {
849 worker::console_error!("inbox: {} {} not watched: {error}", event.kind, event.id);
850 }
851 }
Inbox: the events service tells people what needs them as events arrive852 let mut paths: HashMap<String, Option<RepoPath>> = HashMap::new();
Inbox: threads, reasons, subscriptions and watching853 let mut watchers: HashMap<String, Vec<subscriptions::Watcher>> = HashMap::new();
854 let mut written: Vec<Written> = Vec::new();
Inbox: the events service tells people what needs them as events arrive855 for (event, wanted) in wanted {
Inbox: threads, reasons, subscriptions and watching856 match deliver_one(db, sources, &names, &mut paths, &mut watchers, event, wanted).await {
857 Ok(mut news) => written.append(&mut news),
858 Err(error) => worker::console_error!("inbox: {} {} not delivered: {error}", event.kind, event.id),
859 }
860 }
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)861 // Notify: every item that was news, live in the person's tabs (services/notify).
862 if let Some(notify) = sources.notify {
863 let live = written
864 .iter()
865 .map(|item| (item.notice.username.clone(), live_notification(&item.event_id, &item.notice, &item.url, &rfc3339(now_ms()))))
866 .collect::<Vec<_>>();
867 tell_live(notify, live).await;
868 let people: HashSet<&str> = written.iter().map(|item| item.notice.username.as_str()).collect();
869 for username in people {
870 tell_inbox_count(db, notify, username).await;
871 }
872 }
Inbox: threads, reasons, subscriptions and watching873 if let Err(error) = email(db, sources.identity, written).await {
874 worker::console_error!("inbox: emails not sent: {error}");
875 }
876}
877
Fine-grained personal tokens, workspace token rules and approvals in identity878/// Events about a workspace rather than a repository, each naming who to
879/// tell (`notify`), what to say (`title`, `body`) and where it is (`link`):
880/// identity's personal access token approvals.
881///
882/// | Event | Who | Reason | Severity |
883/// | --- | --- | --- | --- |
884/// | `token.approval_requested` | the workspace's owners | review_requested | warning |
885/// | `token.approval_reviewed` | the token's owner | author | info |
Merge workspace invitations: nobody joins a workspace without saying yes, people are found by username, your own invites can bring someone in, and nobody is left without a workspace (identity 0040)886/// | `workspace_invitation.created` | the person invited | review_requested | warning |
887/// | `workspace_invitation.accepted` | who invited them | author | success |
888/// | `workspace_invitation.declined` | who invited them | author | info |
889///
890/// An invitation's item for the person invited is done once it is
891/// accepted, declined or revoked (`workspace_invitation.revoked`, which
892/// tells nobody).
893pub const WORKSPACE_EVENTS: [&str; 5] = [
894 "token.approval_requested",
895 "token.approval_reviewed",
896 "workspace_invitation.created",
897 "workspace_invitation.accepted",
898 "workspace_invitation.declined",
899];
Fine-grained personal tokens, workspace token rules and approvals in identity900
Merge workspace invitations: nobody joins a workspace without saying yes, people are found by username, your own invites can bring someone in, and nobody is left without a workspace (identity 0040)901/// Events that close what a workspace invitation asked of its person.
902pub const INVITATION_ANSWERED: [&str; 3] =
903 ["workspace_invitation.accepted", "workspace_invitation.declined", "workspace_invitation.revoked"];
904
Fine-grained personal tokens, workspace token rules and approvals in identity905/// The notices a workspace event calls for, and the thread they go to.
906pub fn workspace_notices(event: &Event, actor: Option<&str>) -> (String, Vec<Notice>) {
907 let data = &event.data;
908 let text = |key: &str| data[key].as_str().unwrap_or_default().to_owned();
909 let (reason, severity) = match event.kind.as_str() {
Merge workspace invitations: nobody joins a workspace without saying yes, people are found by username, your own invites can bring someone in, and nobody is left without a workspace (identity 0040)910 "token.approval_requested" | "workspace_invitation.created" => (Reason::ReviewRequested, Severity::Warning),
911 "workspace_invitation.accepted" => (Reason::Author, Severity::Success),
Fine-grained personal tokens, workspace token rules and approvals in identity912 _ => (Reason::Author, Severity::Info),
913 };
Merge workspace invitations: nobody joins a workspace without saying yes, people are found by username, your own invites can bring someone in, and nobody is left without a workspace (identity 0040)914 // An invitation names its own thread; a token approval's is the token's.
915 let thread = match data["thread"].as_str().filter(|thread| !thread.is_empty()) {
916 Some(thread) => thread.to_owned(),
917 None => format!("workspace:{}/token/{}", text("workspace"), text("tokenId")),
918 };
Fine-grained personal tokens, workspace token rules and approvals in identity919 let notices = names(data, "notify")
920 .into_iter()
921 .map(|name| name.to_lowercase())
922 .filter(|name| actor.is_none_or(|actor| !actor.eq_ignore_ascii_case(name)) && !is_g1t(name))
923 .collect::<HashSet<_>>()
924 .into_iter()
925 .map(|username| Notice { username, reason, severity, title: text("title"), body: text("body") })
926 .collect();
927 (thread, notices)
928}
929
930/// Writes a workspace event's notices: a thread each, with no repository.
931async fn deliver_workspace(db: &D1Database, event: &Event) -> Result<()> {
932 let (thread, told) = workspace_notices(event, None);
933 if told.is_empty() {
934 return Ok(());
935 }
936 let now = now_ms();
937 let workspace = event.data["workspace"].as_str().unwrap_or_default().to_lowercase();
938 let link = event.data["link"].as_str().filter(|link| link.starts_with('/'));
939 let mut statements = Vec::new();
940 for notice in &told {
941 let username = notice.username.as_str();
942 statements.push(db.prepare(BUMP).bind(&[
943 new_id("ntf", now).into(),
944 username.into(),
945 thread.as_str().into(),
946 event.id.as_str().into(),
947 event.kind.as_str().into(),
948 notice.reason.as_str().into(),
949 notice.severity.as_str().into(),
950 notice.title.as_str().into(),
951 notice.body.as_str().into(),
952 workspace.as_str().into(),
953 JsValue::NULL,
954 JsValue::NULL,
955 JsValue::NULL,
956 JsValue::NULL,
957 JsValue::NULL,
958 link.map_or(JsValue::NULL, JsValue::from),
959 JsValue::NULL,
960 event.time.as_str().into(),
961 new_id("ntf", now).into(),
962 ])?);
963 statements.push(
964 db.prepare(
965 "INSERT OR IGNORE INTO inbox_activity (id, item_id, username, event_id, event, reason, severity, title, body, actor, created_at)
966 SELECT ?1, id, ?2, ?4, ?5, ?6, ?7, ?8, ?9, NULL, ?10 FROM inbox_items WHERE username = ?2 AND thread = ?3",
967 )
968 .bind(&[
969 new_id("ntf", now).into(),
970 username.into(),
971 thread.as_str().into(),
972 event.id.as_str().into(),
973 event.kind.as_str().into(),
974 notice.reason.as_str().into(),
975 notice.severity.as_str().into(),
976 notice.title.as_str().into(),
977 notice.body.as_str().into(),
978 event.time.as_str().into(),
979 ])?,
980 );
981 }
982 db.batch(statements).await?;
983 Ok(())
984}
985
Inbox: threads, reasons, subscriptions and watching986/// Closes what an agent was waiting on a person for once it is over.
987async fn resolve(db: &D1Database, events: &[Event]) -> Result<()> {
988 let now = rfc3339(now_ms());
989 let mut statements = Vec::new();
990 for thread in events.iter().filter_map(resolves) {
991 statements.push(
992 db.prepare(
993 "UPDATE inbox_items SET done_at = ?1, read_at = COALESCE(read_at, ?1)
994 WHERE thread = ?2 AND reason = 'agent' AND done_at IS NULL",
995 )
996 .bind(&[now.as_str().into(), thread.into()])?,
997 );
998 }
Merge workspace invitations: nobody joins a workspace without saying yes, people are found by username, your own invites can bring someone in, and nobody is left without a workspace (identity 0040)999 // A workspace invitation answered or revoked no longer waits on its person.
1000 for event in events.iter().filter(|event| INVITATION_ANSWERED.contains(&event.kind.as_str())) {
1001 let Some(thread) = event.data["thread"].as_str().filter(|thread| thread.starts_with("invitation:")) else {
1002 continue;
1003 };
1004 statements.push(
1005 db.prepare(
1006 "UPDATE inbox_items SET done_at = ?1, read_at = COALESCE(read_at, ?1)
1007 WHERE thread = ?2 AND reason = 'review_requested' AND done_at IS NULL",
1008 )
1009 .bind(&[now.as_str().into(), thread.into()])?,
1010 );
1011 }
Inbox: threads, reasons, subscriptions and watching1012 // Reviews no longer asked for are no longer waiting.
1013 for event in events.iter().filter(|event| event.kind == "pull.review_request_removed") {
1014 let (Some(repo_id), Some(number)) = (
1015 event.data["repoId"].as_str().map(str::to_owned).or_else(|| event.repo_id.clone()),
1016 event.data["number"].as_u64(),
1017 ) else {
1018 continue;
1019 };
1020 for name in names(&event.data, "reviewers") {
1021 statements.push(
1022 db.prepare(
1023 "UPDATE inbox_items SET done_at = ?1, read_at = COALESCE(read_at, ?1)
1024 WHERE thread = ?2 AND username = ?3 AND reason = 'review_requested' AND done_at IS NULL",
1025 )
1026 .bind(&[
1027 now.as_str().into(),
1028 format!("{repo_id}#{number}").into(),
1029 name.to_lowercase().into(),
1030 ])?,
1031 );
Inbox: the events service tells people what needs them as events arrive1032 }
1033 }
Inbox: threads, reasons, subscriptions and watching1034 if !statements.is_empty() {
1035 db.batch(statements).await?;
1036 }
1037 Ok(())
Inbox: the events service tells people what needs them as events arrive1038}
1039
1040async fn deliver_one(
1041 db: &D1Database,
1042 sources: &Sources<'_>,
1043 names: &HashMap<String, String>,
1044 paths: &mut HashMap<String, Option<RepoPath>>,
Inbox: threads, reasons, subscriptions and watching1045 watchers: &mut HashMap<String, Vec<subscriptions::Watcher>>,
Inbox: the events service tells people what needs them as events arrive1046 event: &Event,
1047 wanted: Wanted,
Inbox: threads, reasons, subscriptions and watching1048) -> Result<Vec<Written>> {
Inbox: the events service tells people what needs them as events arrive1049 if !paths.contains_key(&wanted.repo_id) {
1050 let path: Option<RepoPath> = g1t_kit::call(
1051 sources.repos,
1052 "path_by_id",
1053 &PathByIdArgs {
1054 id: wanted.repo_id.clone(),
1055 },
1056 )
1057 .await?;
1058 paths.insert(wanted.repo_id.clone(), path);
1059 }
1060 // A repository that is gone tells nobody.
1061 let Some(path) = paths.get(&wanted.repo_id).cloned().flatten() else {
Inbox: threads, reasons, subscriptions and watching1062 return Ok(Vec::new());
Inbox: the events service tells people what needs them as events arrive1063 };
1064 let subject: Option<InboxSubject> = match wanted.number {
1065 Some(number) => {
1066 let found: Option<InboxSubject> = g1t_kit::call(
1067 sources.work,
1068 "inbox_subject",
1069 &InboxSubjectArgs {
1070 repo_id: wanted.repo_id.clone(),
1071 number,
1072 comment_id: wanted.comment_id.clone(),
1073 },
1074 )
1075 .await?;
1076 if found.is_none() {
Inbox: threads, reasons, subscriptions and watching1077 return Ok(Vec::new());
Inbox: the events service tells people what needs them as events arrive1078 }
1079 found
1080 }
1081 None => None,
1082 };
1083 let actor = Actor {
1084 id: event.actor.clone(),
1085 username: match event.actor.as_deref() {
1086 Some(id) if is_g1t_id(id) => Some(system::USERNAME.to_owned()),
1087 Some(id) => names.get(id).map(|name| name.to_lowercase()),
1088 None => None,
1089 },
1090 };
Inbox: threads, reasons, subscriptions and watching1091 let thread = thread_of(event, &wanted, subject.as_ref());
1092 if !watchers.contains_key(&wanted.repo_id) {
1093 watchers.insert(wanted.repo_id.clone(), subscriptions::watchers(db, &wanted.repo_id).await?);
1094 }
1095 let audience = Audience {
1096 subscriptions: match thread.kind {
1097 Some(SubjectKind::Issue | SubjectKind::Pull) => subscriptions::of_thread(db, &thread.key).await?,
1098 _ => Vec::new(),
1099 },
1100 watchers: watchers.get(&wanted.repo_id).cloned().unwrap_or_default(),
1101 };
Inbox: the events service tells people what needs them as events arrive1102 let repo = format!("{}/{}", path.namespace, path.name).to_lowercase();
Inbox: threads, reasons, subscriptions and watching1103 let told = notices(event, &repo, &actor, subject.as_ref(), &audience);
1104 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 arrive1105 if told.is_empty() {
Inbox: threads, reasons, subscriptions and watching1106 if !statements.is_empty() {
1107 db.batch(statements).await?;
1108 }
1109 return Ok(Vec::new());
Inbox: the events service tells people what needs them as events arrive1110 }
1111 let now = now_ms();
Inbox: threads, reasons, subscriptions and watching1112 let link = thread.link.as_deref();
1113 let item_url = url(Some(&repo), thread.kind, thread.number, thread.run_id.as_deref(), link);
1114 // Three statements a notice: its thread brought up (unless this event
1115 // was told before), a line of history, and the history kept short.
1116 let first = statements.len();
1117 for notice in &told {
1118 let username = notice.username.as_str();
1119 let values: Vec<JsValue> = vec![
1120 new_id("ntf", now).into(),
1121 username.into(),
1122 thread.key.as_str().into(),
1123 event.id.as_str().into(),
1124 event.kind.as_str().into(),
1125 notice.reason.as_str().into(),
1126 notice.severity.as_str().into(),
1127 notice.title.as_str().into(),
1128 notice.body.as_str().into(),
1129 path.namespace.to_lowercase().into(),
1130 wanted.repo_id.as_str().into(),
1131 repo.as_str().into(),
1132 thread.kind.map_or(JsValue::NULL, |kind| kind.as_str().into()),
1133 thread.number.map_or(JsValue::NULL, JsValue::from),
1134 thread.run_id.as_deref().map_or(JsValue::NULL, JsValue::from),
1135 link.map_or(JsValue::NULL, JsValue::from),
1136 actor.username.as_deref().map_or(JsValue::NULL, JsValue::from),
1137 event.time.as_str().into(),
1138 // Same prefix as item ids, so old and new sort by time together.
1139 new_id("ntf", now).into(),
1140 ];
1141 statements.push(db.prepare(BUMP).bind(&values)?);
Inbox: the events service tells people what needs them as events arrive1142 statements.push(
1143 db.prepare(
Inbox: threads, reasons, subscriptions and watching1144 "INSERT OR IGNORE INTO inbox_activity (id, item_id, username, event_id, event, reason, severity, title, body, actor, created_at)
1145 SELECT ?1, id, ?2, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11 FROM inbox_items WHERE username = ?2 AND thread = ?3
1146 RETURNING username",
Inbox: the events service tells people what needs them as events arrive1147 )
1148 .bind(&[
1149 new_id("ntf", now).into(),
Inbox: threads, reasons, subscriptions and watching1150 username.into(),
1151 thread.key.as_str().into(),
Inbox: the events service tells people what needs them as events arrive1152 event.id.as_str().into(),
Inbox: threads, reasons, subscriptions and watching1153 event.kind.as_str().into(),
1154 notice.reason.as_str().into(),
Inbox: the events service tells people what needs them as events arrive1155 notice.severity.as_str().into(),
1156 notice.title.as_str().into(),
1157 notice.body.as_str().into(),
1158 actor.username.as_deref().map_or(JsValue::NULL, JsValue::from),
1159 event.time.as_str().into(),
1160 ])?,
1161 );
Inbox: threads, reasons, subscriptions and watching1162 statements.push(
1163 db.prepare(
1164 "DELETE FROM inbox_activity
1165 WHERE item_id = (SELECT id FROM inbox_items WHERE username = ?1 AND thread = ?2)
1166 AND id NOT IN (
1167 SELECT a.id FROM inbox_activity a
1168 WHERE a.item_id = (SELECT id FROM inbox_items WHERE username = ?1 AND thread = ?2)
1169 ORDER BY a.id DESC LIMIT ?3)",
1170 )
1171 .bind(&[username.into(), thread.key.as_str().into(), MAX_ACTIVITY.into()])?,
1172 );
1173 }
1174 let results = db.batch(statements).await?;
1175 #[derive(Deserialize)]
1176 struct Told {
1177 username: String,
Inbox: the events service tells people what needs them as events arrive1178 }
Inbox: threads, reasons, subscriptions and watching1179 // A line of history written means the event is news to that person.
1180 let mut news = HashSet::new();
1181 for (at, _) in told.iter().enumerate() {
1182 if let Some(result) = results.get(first + at * 3 + 1)
1183 && let Ok(rows) = result.results::<Told>()
1184 {
1185 news.extend(rows.into_iter().map(|row| row.username));
1186 }
1187 }
1188 Ok(told
1189 .into_iter()
1190 .filter(|notice| news.contains(&notice.username))
1191 .map(|notice| Written {
1192 notice,
1193 repo_id: wanted.repo_id.clone(),
1194 url: item_url.clone(),
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)1195 event_id: event.id.clone(),
Inbox: threads, reasons, subscriptions and watching1196 })
1197 .collect())
1198}
1199
1200/// Brings a person's thread up with new activity, or starts it. `?1` id,
1201/// `?2` username, `?3` thread, `?4` event id, `?5` event type, `?6` reason,
1202/// `?7` severity, `?8` title, `?9` body, `?10` workspace, `?11` repo id,
1203/// `?12` repo, `?13` subject, `?14` number, `?15` run id, `?16` link, `?17`
1204/// actor, `?18` time, `?19` the time-sortable id lists order by. Nothing
1205/// happens when the person was already told of this event. While the
1206/// thread is unread its most urgent severity is kept; done or not, it
1207/// comes back to the inbox. A snooze stands.
1208const BUMP: &str = "INSERT INTO inbox_items (id, username, thread, event_id, event, reason, severity, title, body,
1209 workspace, repo_id, repo, subject, number, run_id, link, actor, created_at, updated_at, bumped, activity)
1210 SELECT ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?18, ?19, 1
1211 WHERE NOT EXISTS (SELECT 1 FROM inbox_activity WHERE event_id = ?4 AND username = ?2)
1212 ON CONFLICT (username, thread) DO UPDATE SET
1213 event_id = excluded.event_id,
1214 event = excluded.event,
1215 reason = excluded.reason,
1216 severity = CASE
1217 WHEN inbox_items.read_at IS NULL AND inbox_items.done_at IS NULL
1218 AND (CASE inbox_items.severity WHEN 'warning' THEN 0 WHEN 'error' THEN 1 WHEN 'success' THEN 2 ELSE 3 END)
1219 < (CASE excluded.severity WHEN 'warning' THEN 0 WHEN 'error' THEN 1 WHEN 'success' THEN 2 ELSE 3 END)
1220 THEN inbox_items.severity ELSE excluded.severity END,
1221 title = excluded.title,
1222 body = excluded.body,
1223 workspace = excluded.workspace,
1224 repo = excluded.repo,
1225 subject = COALESCE(excluded.subject, inbox_items.subject),
1226 run_id = COALESCE(excluded.run_id, inbox_items.run_id),
1227 link = COALESCE(excluded.link, inbox_items.link),
1228 actor = excluded.actor,
1229 updated_at = excluded.updated_at,
1230 bumped = excluded.bumped,
1231 activity = inbox_items.activity + 1,
1232 read_at = NULL,
1233 done_at = NULL";
1234
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)1235/// What the notify service is told of an inbox item (`FeedNotification` in
1236/// @g1t/contracts): its kind from why the person was told, the workspace
1237/// from where it links, and the item's title and first line.
1238pub fn live_notification(event_id: &str, notice: &Notice, url: &str, at: &str) -> serde_json::Value {
1239 let kind = match notice.reason {
1240 Reason::Agent => "agent_waiting",
1241 Reason::ReviewRequested => "approval",
1242 Reason::Mention | Reason::TeamMention => "mention",
1243 _ => "inbox",
1244 };
1245 let path = url.split(['?', '#']).next().unwrap_or_default();
1246 let mut segments = path.trim_start_matches('/').split('/');
1247 let workspace = segments.next().unwrap_or_default().to_lowercase();
1248 let repo = segments.next().map(|name| format!("{workspace}/{name}"));
1249 let name = repo.unwrap_or_else(|| if workspace.is_empty() { "g1t".to_owned() } else { workspace.clone() });
1250 serde_json::json!({
1251 "id": format!("inbox:{event_id}"),
1252 "kind": kind,
1253 "workspace": workspace,
1254 "title": notice.title,
1255 "body": clip(&notice.body, 140),
1256 "href": if url.starts_with('/') { url } else { "/inbox" },
1257 "actor": { "kind": "system", "id": name, "name": name },
1258 "created_at": at,
1259 })
1260}
1261
1262/// What the notify service is told when a person's inbox count changes
1263/// (its `set_inbox`): the count as it now is, never a difference, so a
1264/// call repeated or out of order still ends right.
1265pub fn inbox_count_update(username: &str, unread: u32) -> serde_json::Value {
1266 serde_json::json!({ "username": username.to_lowercase(), "unread": unread })
1267}
1268
1269/// Tells the notify service a person's unread count as it now is, so every
1270/// tab of theirs shows it at once: after items arrive, and after marks made
1271/// anywhere (the site, the API, MCP). Logged and left when it fails.
1272pub async fn tell_inbox_count(db: &D1Database, notify: &Fetcher, username: &str) {
1273 let counted = counts(db, InboxCountsArgs { username: username.to_owned() }).await;
1274 let told: Result<serde_json::Value> = match counted {
1275 Ok(counts) => g1t_kit::call(notify, "set_inbox", &inbox_count_update(username, counts.unread)).await,
1276 Err(error) => Err(error),
1277 };
1278 if let Err(error) = told {
1279 worker::console_error!("inbox: live count not sent: {error}");
1280 }
1281}
1282
1283/// Hands items to the notify service, one call each. Logged and left when it
1284/// cannot be reached: the inbox and email do not wait on it.
1285async fn tell_live(notify: &Fetcher, items: Vec<(String, serde_json::Value)>) {
1286 #[derive(serde::Serialize)]
1287 struct NotifyArgs {
1288 username: String,
1289 notification: serde_json::Value,
1290 }
1291 for (username, notification) in items {
1292 let told: Result<serde_json::Value> = g1t_kit::call(notify, "notify", &NotifyArgs { username, notification }).await;
1293 if let Err(error) = told {
1294 worker::console_error!("inbox: live notification not sent: {error}");
1295 }
1296 }
1297}
1298
Inbox: threads, reasons, subscriptions and watching1299/// Emails what was news to people who asked to be emailed for its reason.
1300/// Identity sends each, only to a confirmed address and only if the person
1301/// can still read the repository.
1302async fn email(db: &D1Database, identity: &Fetcher, written: Vec<Written>) -> Result<()> {
1303 if written.is_empty() {
1304 return Ok(());
1305 }
1306 let people: Vec<String> = written
1307 .iter()
1308 .map(|item| item.notice.username.clone())
1309 .collect::<HashSet<_>>()
1310 .into_iter()
1311 .collect();
1312 let settings = subscriptions::settings_of(db, &people).await?;
1313 for item in written {
1314 let wants = settings.get(&item.notice.username).map_or(&DEFAULT_EMAIL[..], |settings| &settings.email[..]);
1315 if !wants.contains(&item.notice.reason) {
1316 continue;
1317 }
1318 let sent: Result<bool> = g1t_kit::call(
1319 identity,
1320 "notify_by_email",
1321 &NotifyByEmailArgs {
1322 username: item.notice.username.clone(),
1323 repo_id: item.repo_id.clone(),
1324 subject: item.notice.title.clone(),
1325 intro: item.notice.body.clone(),
1326 quote: None,
1327 path: item.url.clone(),
1328 reason: item.notice.reason,
1329 },
1330 )
1331 .await;
1332 if let Err(error) = sent {
1333 worker::console_error!("inbox: email to {} not sent: {error}", item.notice.username);
1334 }
1335 }
Inbox: the events service tells people what needs them as events arrive1336 Ok(())
1337}
1338
1339// --- Reading and changing ----------------------------------------------------
1340
1341#[derive(Deserialize)]
Inbox: threads, reasons, subscriptions and watching1342pub(crate) struct Row {
1343 pub(crate) id: String,
Inbox: the events service tells people what needs them as events arrive1344 reason: String,
1345 severity: String,
1346 title: String,
1347 body: String,
Inbox: threads, reasons, subscriptions and watching1348 event: Option<String>,
Inbox: the events service tells people what needs them as events arrive1349 workspace: Option<String>,
Inbox: threads, reasons, subscriptions and watching1350 pub(crate) repo_id: Option<String>,
Inbox: the events service tells people what needs them as events arrive1351 repo: Option<String>,
1352 subject: Option<String>,
1353 number: Option<f64>,
1354 run_id: Option<String>,
Inbox: threads, reasons, subscriptions and watching1355 link: Option<String>,
Inbox: the events service tells people what needs them as events arrive1356 actor: Option<String>,
Inbox: threads, reasons, subscriptions and watching1357 activity: Option<f64>,
Inbox: the events service tells people what needs them as events arrive1358 created_at: String,
Inbox: threads, reasons, subscriptions and watching1359 updated_at: Option<String>,
Inbox: the events service tells people what needs them as events arrive1360 read_at: Option<String>,
1361 done_at: Option<String>,
1362 saved: f64,
1363 snoozed_until: Option<String>,
Inbox: threads, reasons, subscriptions and watching1364 bumped: Option<String>,
Inbox: the events service tells people what needs them as events arrive1365}
1366
1367impl Row {
Inbox: threads, reasons, subscriptions and watching1368 pub(crate) fn into_item(self) -> InboxItem {
Inbox: the events service tells people what needs them as events arrive1369 let subject = self.subject.as_deref().and_then(SubjectKind::parse);
1370 let number = self.number.map(|n| n as u32);
1371 InboxItem {
Inbox: threads, reasons, subscriptions and watching1372 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 arrive1373 id: self.id,
Inbox: threads, reasons, subscriptions and watching1374 reason: Reason::parse(&self.reason).unwrap_or(Reason::Subscribed),
Inbox: the events service tells people what needs them as events arrive1375 severity: Severity::parse(&self.severity).unwrap_or(Severity::Info),
1376 title: self.title,
1377 body: self.body,
Inbox: threads, reasons, subscriptions and watching1378 event: self.event,
Inbox: the events service tells people what needs them as events arrive1379 repo: self.repo,
1380 workspace: self.workspace,
1381 subject,
1382 number,
1383 actor: self.actor,
Inbox: threads, reasons, subscriptions and watching1384 count: self.activity.map_or(1, |n| n as u32),
1385 updated_at: self.updated_at.unwrap_or_else(|| self.created_at.clone()),
Inbox: the events service tells people what needs them as events arrive1386 created_at: self.created_at,
1387 read_at: self.read_at,
1388 done_at: self.done_at,
1389 saved: self.saved != 0.0,
1390 snoozed_until: self.snoozed_until,
1391 }
1392 }
1393}
1394
Inbox: threads, reasons, subscriptions and watching1395pub(crate) const COLUMNS: &str = "id, reason, severity, title, body, event, workspace, repo_id, repo, subject, number, run_id,
1396 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 arrive1397
1398/// The conditions that pick a view's items, after `username = ?1`; `?2` is now.
1399fn view_filter(view: InboxView) -> &'static str {
1400 match view {
1401 InboxView::Inbox => "done_at IS NULL AND (snoozed_until IS NULL OR snoozed_until <= ?2)",
1402 InboxView::Saved => "saved = 1",
1403 InboxView::Done => "done_at IS NOT NULL",
1404 }
1405}
1406
1407/// Whether a list ranks unread warnings first: the inbox itself, unfiltered.
1408fn ranked(a: &ListInboxArgs) -> bool {
Inbox: threads, reasons, subscriptions and watching1409 a.view == InboxView::Inbox
1410 && a.severity.is_none()
1411 && a.reason.is_none()
1412 && !a.participating
1413 && a.repo_id.is_none()
1414 && !a.unread
1415 && a.since.is_none()
1416 && a.updated_before.is_none()
Inbox: the events service tells people what needs them as events arrive1417}
1418
Inbox: threads, reasons, subscriptions and watching1419/// The conditions a list's filters add, and the values they bind, numbered
1420/// after the `bound` values already given.
1421fn filters(a: &ListInboxArgs, bound: usize) -> (Vec<String>, Vec<String>) {
1422 let mut conditions = Vec::new();
1423 let mut values: Vec<String> = Vec::new();
1424 let mut bind = |value: &str, condition: &str| {
1425 values.push(value.to_owned());
1426 conditions.push(condition.replace('?', &format!("?{}", bound + values.len())));
1427 };
1428 if let Some(severity) = a.severity {
1429 bind(severity.as_str(), "severity = ?");
1430 }
1431 if let Some(reason) = a.reason {
1432 bind(reason.as_str(), "reason = ?");
1433 }
1434 if let Some(repo_id) = &a.repo_id {
1435 bind(repo_id, "repo_id = ?");
1436 }
1437 if let Some(since) = &a.since {
1438 bind(since.trim(), "COALESCE(updated_at, created_at) >= ?");
1439 }
1440 if let Some(before) = &a.updated_before {
1441 bind(before.trim(), "COALESCE(updated_at, created_at) < ?");
1442 }
1443 if a.participating {
1444 conditions.push("reason NOT IN ('manual', 'subscribed')".to_owned());
1445 }
1446 if a.unread {
1447 conditions.push("read_at IS NULL".to_owned());
1448 }
1449 (conditions, values)
1450}
1451
Inbox: the events service tells people what needs them as events arrive1452pub async fn list(db: &D1Database, repos: &Fetcher, a: ListInboxArgs) -> Result<InboxPage> {
1453 let Some(viewer) = &a.viewer else {
1454 return Ok(InboxPage::default());
1455 };
1456 let username = viewer.username.to_lowercase();
1457 let now = rfc3339(now_ms());
1458 let limit = a.limit.unwrap_or(DEFAULT_INBOX_PAGE).clamp(1, MAX_INBOX_PAGE);
1459 let mut values: Vec<JsValue> = vec![username.as_str().into(), now.as_str().into()];
Inbox: threads, reasons, subscriptions and watching1460 let mut conditions = vec!["username = ?1".to_owned(), view_filter(a.view).to_owned()];
1461 let (filtered, bound) = filters(&a, values.len());
1462 conditions.extend(filtered);
1463 values.extend(bound.iter().map(|value| JsValue::from(value.as_str())));
Inbox: the events service tells people what needs them as events arrive1464 let ranked = ranked(&a);
1465 // Unread warnings lead the first page, and are left out of the rest.
1466 let leading = "severity = 'warning' AND read_at IS NULL";
1467 let mut rest = conditions.clone();
Inbox: threads, reasons, subscriptions and watching1468 let mut rest_values = values.clone();
Inbox: the events service tells people what needs them as events arrive1469 if ranked {
1470 rest.push(format!("NOT ({leading})"));
1471 }
1472 if let Some(before) = &a.before {
Inbox: threads, reasons, subscriptions and watching1473 rest_values.push(before.as_str().into());
1474 rest.push(format!("bumped < ?{}", rest_values.len()));
Inbox: the events service tells people what needs them as events arrive1475 }
Inbox: threads, reasons, subscriptions and watching1476 rest_values.push((limit + 1).into());
1477 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 arrive1478 let mut statements = vec![
1479 db.prepare(format!(
1480 "SELECT {COLUMNS} FROM inbox_items WHERE {} ORDER BY {order} LIMIT ?{}",
1481 rest.join(" AND "),
Inbox: threads, reasons, subscriptions and watching1482 rest_values.len()
Inbox: the events service tells people what needs them as events arrive1483 ))
Inbox: threads, reasons, subscriptions and watching1484 .bind(&rest_values)?,
Inbox: the events service tells people what needs them as events arrive1485 ];
1486 if ranked && a.before.is_none() {
1487 statements.push(
1488 db.prepare(format!(
Inbox: threads, reasons, subscriptions and watching1489 "SELECT {COLUMNS} FROM inbox_items WHERE {} AND {leading} ORDER BY bumped DESC LIMIT ?3",
1490 conditions.join(" AND ")
Inbox: the events service tells people what needs them as events arrive1491 ))
1492 .bind(&[username.as_str().into(), now.as_str().into(), MAX_RANKED.into()])?,
1493 );
1494 }
1495 let results = db.batch(statements).await?;
1496 let mut rows = results.first().map(|result| result.results::<Row>()).transpose()?.unwrap_or_default();
1497 let next = if rows.len() > limit as usize {
1498 rows.truncate(limit as usize);
Inbox: threads, reasons, subscriptions and watching1499 rows.last().map(|row| row.bumped.clone().unwrap_or_else(|| row.id.clone()))
Inbox: the events service tells people what needs them as events arrive1500 } else {
1501 None
1502 };
1503 let mut items: Vec<Row> = results.get(1).map(|result| result.results::<Row>()).transpose()?.unwrap_or_default();
1504 items.extend(rows);
Inbox: threads, reasons, subscriptions and watching1505 let items = readable_only(db, repos, &a.viewer, &username, items).await?;
1506 Ok(InboxPage {
1507 items: items.into_iter().map(Row::into_item).collect(),
1508 next,
1509 })
1510}
Inbox: the events service tells people what needs them as events arrive1511
Inbox: threads, reasons, subscriptions and watching1512/// Rows about repositories the viewer can still read; the rest are
1513/// removed from their inbox as they are found.
1514pub(crate) async fn readable_only(db: &D1Database, repos: &Fetcher, viewer: &g1t_contracts::Viewer, username: &str, mut rows: Vec<Row>) -> Result<Vec<Row>> {
1515 let ids: Vec<String> = rows
Inbox: the events service tells people what needs them as events arrive1516 .iter()
1517 .filter_map(|row| row.repo_id.clone())
1518 .collect::<HashSet<_>>()
1519 .into_iter()
1520 .collect();
Inbox: threads, reasons, subscriptions and watching1521 if ids.is_empty() {
1522 return Ok(rows);
1523 }
1524 let readable: Vec<Repo> = g1t_kit::call(
1525 repos,
1526 "readable",
1527 &ReadableArgs {
1528 ids: ids.clone(),
1529 viewer: viewer.clone(),
1530 },
1531 )
1532 .await?;
1533 let readable: HashSet<String> = readable.into_iter().map(|repo| repo.id).collect();
1534 let gone: Vec<String> = ids.into_iter().filter(|id| !readable.contains(id)).collect();
1535 if !gone.is_empty() {
1536 forget(db, username, &gone).await?;
1537 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 arrive1538 }
Inbox: threads, reasons, subscriptions and watching1539 Ok(rows)
Inbox: the events service tells people what needs them as events arrive1540}
1541
1542/// Removes a person's items about repositories they cannot read.
1543async fn forget(db: &D1Database, username: &str, repo_ids: &[String]) -> Result<()> {
1544 let marks = vec!["?"; repo_ids.len()].join(", ");
1545 let mut values: Vec<JsValue> = vec![username.into()];
1546 values.extend(repo_ids.iter().map(|id| JsValue::from(id.as_str())));
Inbox: threads, reasons, subscriptions and watching1547 db.batch(vec![
1548 db.prepare(format!(
1549 "DELETE FROM inbox_activity WHERE item_id IN (SELECT id FROM inbox_items WHERE username = ? AND repo_id IN ({marks}))"
1550 ))
1551 .bind(&values)?,
1552 db.prepare(format!("DELETE FROM inbox_items WHERE username = ? AND repo_id IN ({marks})"))
1553 .bind(&values)?,
1554 ])
1555 .await?;
Inbox: the events service tells people what needs them as events arrive1556 Ok(())
1557}
1558
1559#[derive(Deserialize)]
1560struct CountRow {
1561 severity: String,
1562 n: f64,
1563}
1564
1565pub async fn counts(db: &D1Database, a: InboxCountsArgs) -> Result<InboxCounts> {
1566 let rows = db
1567 .prepare(
1568 "SELECT severity, count(*) AS n FROM inbox_items
1569 WHERE username = ? AND read_at IS NULL AND done_at IS NULL
1570 AND (snoozed_until IS NULL OR snoozed_until <= ?)
1571 GROUP BY severity",
1572 )
1573 .bind(&[a.username.to_lowercase().into(), rfc3339(now_ms()).into()])?
1574 .all()
1575 .await?
1576 .results::<CountRow>()?;
1577 let mut counts = InboxCounts::default();
1578 for row in rows {
1579 let n = row.n as u32;
1580 counts.unread += n;
1581 match Severity::parse(&row.severity) {
1582 Some(Severity::Error) => counts.error += n,
1583 Some(Severity::Warning) => counts.warning += n,
1584 Some(Severity::Success) => counts.success += n,
1585 Some(Severity::Info) | None => counts.info += n,
1586 }
1587 }
1588 Ok(counts)
1589}
1590
1591/// The change a mark makes, as a `SET` clause; `?1` is now.
1592fn mark_change(mark: InboxMark) -> &'static str {
1593 match mark {
1594 InboxMark::Read => "read_at = COALESCE(read_at, ?1)",
1595 InboxMark::Unread => "read_at = NULL",
1596 InboxMark::Done => "done_at = ?1, read_at = COALESCE(read_at, ?1)",
1597 InboxMark::Undone => "done_at = NULL",
1598 InboxMark::Save => "saved = 1",
1599 InboxMark::Unsave => "saved = 0",
1600 InboxMark::Snooze => "snoozed_until = ?2, read_at = COALESCE(read_at, ?1)",
Inbox: threads, reasons, subscriptions and watching1601 InboxMark::Unsnooze => "snoozed_until = NULL",
Inbox: the events service tells people what needs them as events arrive1602 }
1603}
1604
1605#[derive(Deserialize)]
1606struct IdRow {
1607 #[allow(dead_code)]
1608 id: String,
1609}
1610
1611/// Changes the person's own items. Returns how many changed.
1612pub async fn mark(db: &D1Database, a: MarkInboxArgs) -> Result<u32> {
1613 let now = rfc3339(now_ms());
1614 let until = match (a.mark, a.until.as_deref().map(str::trim)) {
1615 // Times compare as text (`g1t_contracts::time`); a snooze is for later.
1616 (InboxMark::Snooze, Some(until)) if until.len() == now.len() && until > now.as_str() => until.to_owned(),
1617 (InboxMark::Snooze, _) => return Ok(0),
1618 _ => String::new(),
1619 };
1620 let mut values: Vec<JsValue> = vec![now.as_str().into(), until.as_str().into(), a.username.to_lowercase().into()];
1621 let target = if a.all && a.ids.is_empty() {
1622 let mut target = "done_at IS NULL".to_owned();
Inbox: threads, reasons, subscriptions and watching1623 let mut bind = |values: &mut Vec<JsValue>, value: &str, condition: &str| {
1624 values.push(value.into());
1625 target.push_str(&format!(" AND {}", condition.replace('?', &format!("?{}", values.len()))));
1626 };
Inbox: the events service tells people what needs them as events arrive1627 if let Some(severity) = a.severity {
Inbox: threads, reasons, subscriptions and watching1628 bind(&mut values, severity.as_str(), "severity = ?");
1629 }
1630 if let Some(repo_id) = &a.repo_id {
1631 bind(&mut values, repo_id, "repo_id = ?");
1632 }
1633 if let Some(last_read_at) = a.last_read_at.as_deref().map(str::trim).filter(|at| !at.is_empty()) {
1634 bind(&mut values, last_read_at, "COALESCE(updated_at, created_at) <= ?");
Inbox: the events service tells people what needs them as events arrive1635 }
1636 target
1637 } else {
1638 let ids: Vec<&String> = a.ids.iter().take(MAX_MARK).collect();
1639 if ids.is_empty() {
1640 return Ok(0);
1641 }
1642 let first = values.len() + 1;
1643 values.extend(ids.iter().map(|id| JsValue::from(id.as_str())));
1644 let marks: Vec<String> = (first..values.len() + 1).map(|at| format!("?{at}")).collect();
1645 format!("id IN ({})", marks.join(", "))
1646 };
1647 let changed = db
1648 .prepare(format!(
1649 "UPDATE inbox_items SET {} WHERE username = ?3 AND {target} RETURNING id",
1650 mark_change(a.mark)
1651 ))
1652 .bind(&values)?
1653 .all()
1654 .await?
1655 .results::<IdRow>()?;
1656 Ok(changed.len() as u32)
1657}
1658
Inbox: threads, reasons, subscriptions and watching1659#[derive(Deserialize)]
1660struct ActivityRow {
1661 reason: String,
1662 severity: String,
1663 title: String,
1664 body: String,
1665 event: Option<String>,
1666 actor: Option<String>,
1667 created_at: String,
1668}
1669
1670/// One of the viewer's threads, with its history and their subscription.
1671pub async fn thread(db: &D1Database, repos: &Fetcher, work: &Fetcher, a: ThreadArgs) -> Result<Option<InboxThread>> {
1672 let Some(viewer) = &a.viewer else {
1673 return Ok(None);
1674 };
1675 let username = viewer.username.to_lowercase();
1676 let row = db
1677 .prepare(format!("SELECT {COLUMNS} FROM inbox_items WHERE id = ? AND username = ?"))
1678 .bind(&[a.id.as_str().into(), username.as_str().into()])?
1679 .first::<Row>(None)
1680 .await?;
1681 let Some(row) = readable_only(db, repos, &a.viewer, &username, row.into_iter().collect()).await?.pop() else {
1682 return Ok(None);
1683 };
1684 let activity = db
1685 .prepare(
1686 "SELECT reason, severity, title, body, event, actor, created_at FROM inbox_activity
1687 WHERE item_id = ? ORDER BY id DESC LIMIT ?",
1688 )
1689 .bind(&[row.id.as_str().into(), MAX_ACTIVITY.into()])?
1690 .all()
1691 .await?
1692 .results::<ActivityRow>()?
1693 .into_iter()
1694 .map(|row| InboxActivity {
1695 reason: Reason::parse(&row.reason).unwrap_or(Reason::Subscribed),
1696 severity: Severity::parse(&row.severity).unwrap_or(Severity::Info),
1697 title: row.title,
1698 body: row.body,
1699 event: row.event,
1700 actor: row.actor,
1701 created_at: row.created_at,
1702 })
1703 .collect();
1704 let item = row.into_item();
1705 let subscription = match (item.subject, item.number) {
1706 (Some(SubjectKind::Issue | SubjectKind::Pull), Some(_)) => {
1707 subscriptions::subscription(
1708 db,
1709 work,
1710 SubscriptionArgs {
1711 viewer: a.viewer.clone(),
1712 id: Some(item.id.clone()),
1713 ..SubscriptionArgs::default()
1714 },
1715 )
1716 .await?
1717 }
1718 _ => None,
1719 };
1720 Ok(Some(InboxThread {
1721 item,
1722 activity,
1723 subscription,
1724 }))
1725}
1726
Inbox: the events service tells people what needs them as events arrive1727// --- Keeping up ------------------------------------------------------------------
1728
1729/// Moves rows with renamed workspaces and repositories, and drops those
Merge account deletion: soft delete for 30 days, staff restore and purge, ghost for what remains (identity 0037)1730/// of purged repositories, deleted workspaces and purged accounts.
Inbox: the events service tells people what needs them as events arrive1731pub async fn follow(db: &D1Database, events: &[Event]) -> Result<()> {
1732 let mut statements = Vec::new();
1733 for event in events {
1734 let text = |key: &str| event.data[key].as_str().unwrap_or_default().to_lowercase();
1735 match event.kind.as_str() {
1736 "workspace.renamed" => {
1737 let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
1738 continue;
1739 };
1740 let (from, to) = (renamed.from.to_lowercase(), renamed.to.to_lowercase());
1741 if from == to {
1742 continue;
1743 }
1744 statements.push(
1745 db.prepare(
Inbox: threads, reasons, subscriptions and watching1746 "UPDATE inbox_items SET repo = ?2 || substr(repo, length(?1) + 1), workspace = ?2,
1747 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 arrive1748 WHERE workspace = ?1",
1749 )
1750 .bind(&[from.as_str().into(), to.as_str().into()])?,
1751 );
Inbox: threads, reasons, subscriptions and watching1752 statements.push(
1753 db.prepare("UPDATE inbox_watching SET repo = ?2 || substr(repo, length(?1) + 1) WHERE repo LIKE ?1 || '/%'")
1754 .bind(&[from.as_str().into(), to.as_str().into()])?,
1755 );
Inbox: the events service tells people what needs them as events arrive1756 }
Inbox: threads, reasons, subscriptions and watching1757 "repo.renamed" => {
1758 let path = format!("{}/{}", text("namespace"), text("to"));
1759 for table in ["inbox_items", "inbox_watching"] {
1760 statements.push(
1761 db.prepare(format!("UPDATE {table} SET repo = ? WHERE repo_id = ?"))
1762 .bind(&[path.as_str().into(), text("repoId").into()])?,
1763 );
1764 }
1765 }
1766 "repo.transferred" => {
1767 let path = format!("{}/{}", text("to"), text("name"));
1768 statements.push(
1769 db.prepare("UPDATE inbox_items SET repo = ?, workspace = ? WHERE repo_id = ?")
1770 .bind(&[path.as_str().into(), text("to").into(), text("repoId").into()])?,
1771 );
1772 statements.push(
1773 db.prepare("UPDATE inbox_watching SET repo = ? WHERE repo_id = ?")
1774 .bind(&[path.as_str().into(), text("repoId").into()])?,
1775 );
1776 }
1777 "repo.purged" => {
1778 for sql in [
1779 "DELETE FROM inbox_activity WHERE item_id IN (SELECT id FROM inbox_items WHERE repo_id = ?)",
1780 "DELETE FROM inbox_items WHERE repo_id = ?",
1781 "DELETE FROM inbox_subscriptions WHERE repo_id = ?",
1782 "DELETE FROM inbox_watching WHERE repo_id = ?",
1783 ] {
1784 statements.push(db.prepare(sql).bind(&[text("repoId").into()])?);
1785 }
1786 }
1787 "workspace.deleted" => {
1788 statements.push(
1789 db.prepare("DELETE FROM inbox_items WHERE workspace = ?")
1790 .bind(&[text("slug").into()])?,
1791 );
1792 statements.push(
1793 db.prepare("DELETE FROM inbox_watching WHERE repo LIKE ? || '/%'")
1794 .bind(&[text("slug").into()])?,
1795 );
1796 }
Merge account deletion: soft delete for 30 days, staff restore and purge, ghost for what remains (identity 0037)1797 // An account purged: its inbox, what it watched and its
1798 // settings go with it (identity's account_deletion.rs).
1799 "user.deleted" => {
1800 let username = text("username");
1801 if username.is_empty() {
1802 continue;
1803 }
1804 for sql in [
1805 "DELETE FROM inbox_activity WHERE username = ?",
1806 "DELETE FROM inbox_items WHERE username = ?",
1807 "DELETE FROM inbox_subscriptions WHERE username = ?",
1808 "DELETE FROM inbox_watching WHERE username = ?",
1809 "DELETE FROM inbox_settings WHERE username = ?",
1810 ] {
1811 statements.push(db.prepare(sql).bind(&[username.as_str().into()])?);
1812 }
1813 }
Inbox: the events service tells people what needs them as events arrive1814 _ => {}
1815 }
1816 }
1817 if !statements.is_empty() {
1818 db.batch(statements).await?;
1819 }
1820 Ok(())
1821}
1822
1823/// Removes items done more than [`DONE_DAYS`] ago, and any not saved older
Inbox: threads, reasons, subscriptions and watching1824/// than [`MAX_DAYS`], with their history. Returns how many went.
Inbox: the events service tells people what needs them as events arrive1825pub async fn purge(db: &D1Database, now: u64) -> Result<u32> {
1826 let done = crate::audit::keep_from(now, DONE_DAYS);
1827 let oldest = crate::audit::keep_from(now, MAX_DAYS);
Inbox: threads, reasons, subscriptions and watching1828 let statements: Vec<D1PreparedStatement> = vec![
1829 db.prepare("DELETE FROM inbox_items WHERE saved = 0 AND done_at < ?")
1830 .bind(&[done.into()])?,
1831 db.prepare("DELETE FROM inbox_items WHERE saved = 0 AND COALESCE(updated_at, created_at) < ?")
1832 .bind(&[oldest.into()])?,
1833 db.prepare("DELETE FROM inbox_activity WHERE NOT EXISTS (SELECT 1 FROM inbox_items WHERE inbox_items.id = inbox_activity.item_id)"),
1834 ];
1835 let results = db.batch(statements).await?;
Inbox: the events service tells people what needs them as events arrive1836 let mut removed = 0;
Inbox: threads, reasons, subscriptions and watching1837 for result in results.into_iter().take(2) {
Inbox: the events service tells people what needs them as events arrive1838 removed += result.meta()?.and_then(|meta| meta.changes).unwrap_or(0) as u32;
1839 }
1840 Ok(removed)
1841}
1842
1843#[cfg(test)]
1844mod tests {
1845 use super::*;
Inbox: threads, reasons, subscriptions and watching1846 use crate::subscriptions::{State, Subscription, Watcher};
Inbox: the events service tells people what needs them as events arrive1847 use serde_json::json;
1848
1849 fn event(kind: &str, actor: Option<&str>, data: serde_json::Value) -> Event {
1850 Event {
1851 id: "evt_1".into(),
1852 kind: kind.into(),
1853 source: "work".into(),
1854 time: "2026-10-07T12:00:00.000Z".into(),
1855 repo_id: Some("rep_1".into()),
1856 actor: actor.map(str::to_owned),
1857 data,
1858 }
1859 }
1860
1861 fn person(id: &str, username: &str) -> Principal {
1862 Principal {
1863 id: id.into(),
1864 username: username.into(),
1865 }
1866 }
1867
1868 fn actor(id: &str, username: &str) -> Actor {
1869 Actor {
1870 id: Some(id.into()),
1871 username: Some(username.into()),
1872 }
1873 }
1874
1875 /// A pull request ana opened, assigned to bo, for an issue cy filed and dee is assigned.
1876 fn pull() -> InboxSubject {
1877 InboxSubject {
1878 kind: Some(SubjectKind::Pull),
1879 title: "Add the inbox".into(),
1880 author: person("usr_ana", "ana"),
1881 assignees: vec!["bo".into()],
1882 issue: Some(Box::new(InboxSubject {
1883 kind: Some(SubjectKind::Issue),
1884 title: "An inbox".into(),
1885 author: person("usr_cy", "cy"),
1886 assignees: vec!["dee".into()],
1887 ..InboxSubject::default()
1888 })),
1889 ..InboxSubject::default()
1890 }
1891 }
1892
1893 /// A change g1t made for ana.
1894 fn g1t_pull() -> InboxSubject {
1895 InboxSubject {
1896 author: person(AGENT_ID, "g1t"),
1897 requested_by: Some(person("usr_ana", "ana")),
1898 ..pull()
1899 }
1900 }
1901
Inbox: threads, reasons, subscriptions and watching1902 fn nobody() -> Audience {
1903 Audience::default()
1904 }
1905
1906 fn told(notices: &[Notice]) -> Vec<(&str, Reason, Severity)> {
Inbox: the events service tells people what needs them as events arrive1907 notices
1908 .iter()
1909 .map(|notice| (notice.username.as_str(), notice.reason, notice.severity))
1910 .collect()
1911 }
1912
1913 fn comment(author: Principal, body: &str, mentions: &[&str]) -> InboxSubject {
1914 InboxSubject {
1915 comment: Some(InboxComment {
1916 author,
1917 excerpt: body.into(),
1918 mentions: mentions.iter().map(|name| (*name).to_owned()).collect(),
1919 ..InboxComment::default()
1920 }),
1921 ..pull()
1922 }
1923 }
1924
Inbox: threads, reasons, subscriptions and watching1925 fn subscribed(rows: &[(&str, State, Option<Reason>)]) -> Audience {
1926 Audience {
1927 subscriptions: rows
1928 .iter()
1929 .map(|(name, state, reason)| Subscription {
1930 username: (*name).into(),
1931 state: *state,
1932 reason: *reason,
1933 })
1934 .collect(),
1935 watchers: Vec::new(),
1936 }
1937 }
1938
1939 fn watched(rows: &[(&str, WatchLevel, &[&str])]) -> Audience {
1940 Audience {
1941 subscriptions: Vec::new(),
1942 watchers: rows
1943 .iter()
1944 .map(|(name, level, events)| Watcher {
1945 username: (*name).into(),
1946 level: *level,
1947 events: events.iter().map(|kind| (*kind).to_owned()).collect(),
1948 })
1949 .collect(),
1950 }
1951 }
1952
Inbox: the events service tells people what needs them as events arrive1953 #[test]
1954 fn only_events_that_tell_someone_are_read() {
1955 let asked = |kind: &str, data| wants(&event(kind, None, data));
1956 assert_eq!(
1957 asked("checks.completed", json!({ "repoId": "rep_1", "number": 4, "status": "failed" })),
1958 Some(Wanted { repo_id: "rep_1".into(), number: Some(4), comment_id: None })
1959 );
1960 assert_eq!(asked("checks.completed", json!({ "repoId": "rep_1", "number": 4, "status": "passed" })), None);
1961 assert_eq!(asked("review.completed", json!({ "repoId": "rep_1", "number": 4 })), None);
1962 assert_eq!(asked("workflow.completed", json!({ "repoId": "rep_1", "conclusion": "success" })), None);
1963 assert_eq!(
1964 asked("workflow.completed", json!({ "repoId": "rep_1", "conclusion": "failure" })),
1965 Some(Wanted { repo_id: "rep_1".into(), number: None, comment_id: None })
1966 );
1967 assert_eq!(
1968 asked("comment.created", json!({ "repoId": "rep_1", "number": 2, "commentId": "cmt_1" })).map(|w| w.comment_id),
1969 Some(Some("cmt_1".into()))
1970 );
1971 assert_eq!(asked("git.push", json!({ "repoId": "rep_1" })), None);
Inbox: threads, reasons, subscriptions and watching1972 for kind in ["issue.opened", "pull.review_requested", "pull.stalled", "issue.assigned", "pull.closed", "issue.reopened"] {
1973 assert!(asked(kind, json!({ "repoId": "rep_1", "number": 1 })).is_some(), "{kind}");
1974 }
1975 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 arrive1976 }
1977
1978 #[test]
Inbox: threads, reasons, subscriptions and watching1979 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 arrive1980 let asked = event("agent.asked", Some("usr_agent"), json!({ "number": 7 }));
Inbox: threads, reasons, subscriptions and watching1981 let notices_ = notices(&asked, "acme/rocket", &Actor::default(), Some(&pull()), &nobody());
Inbox: the events service tells people what needs them as events arrive1982 assert_eq!(
Inbox: threads, reasons, subscriptions and watching1983 told(&notices_),
Inbox: the events service tells people what needs them as events arrive1984 vec![
Inbox: threads, reasons, subscriptions and watching1985 ("ana", Reason::Agent, Severity::Warning),
1986 ("cy", Reason::Agent, Severity::Warning),
1987 ("dee", Reason::Agent, Severity::Warning),
Inbox: the events service tells people what needs them as events arrive1988 ]
1989 );
Inbox: threads, reasons, subscriptions and watching1990 assert_eq!(notices_[0].title, "An agent is waiting on acme/rocket#7");
1991 assert_eq!(notices_[0].body, "Add the inbox");
1992 // Stopping says why; even someone who unsubscribed hears of it.
1993 let stalled = event("pull.stalled", None, json!({ "number": 7, "detail": "Its checks could not be run." }));
1994 let audience = subscribed(&[("ana", State::Unsubscribed, None)]);
1995 let notices_ = notices(&stalled, "acme/rocket", &Actor::default(), Some(&g1t_pull()), &audience);
1996 assert_eq!(notices_[0].username, "ana");
1997 assert_eq!(notices_[0].title, "g1t stopped on acme/rocket#7 and needs you");
1998 assert_eq!(notices_[0].body, "Its checks could not be run.");
1999 // Ignoring the thread silences even that.
2000 let audience = subscribed(&[("ana", State::Ignored, None)]);
2001 let notices_ = notices(&stalled, "acme/rocket", &Actor::default(), Some(&g1t_pull()), &audience);
2002 assert!(!notices_.iter().any(|notice| notice.username == "ana"));
2003 }
2004
2005 #[test]
2006 fn whatever_an_agent_waits_on_closes_when_it_goes_on() {
2007 for kind in RESUMES {
2008 assert_eq!(resolves(&event(kind, None, json!({ "repoId": "rep_1", "number": 7 }))), Some("rep_1#7".into()));
2009 }
2010 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 arrive2011 }
2012
2013 #[test]
Merge branch 'worktree-agent-a3abfcce648e87dca'2014 fn a_deployments_reviewers_are_told() {
2015 let asked = event(
2016 "deployment.review_requested",
2017 Some("usr_ana"),
2018 json!({
2019 "repoId": "rep_1", "runId": "run_7", "environment": "production", "workflow": "Deploy",
2020 "title": "Ship it", "notify": ["cy", "ana"], "link": "/acme/rocket/actions/runs/run_7"
2021 }),
2022 );
2023 let wanted = wants(&asked).unwrap();
2024 assert_eq!(wanted.number, None);
2025 let thread = thread_of(&asked, &wanted, None);
2026 assert_eq!(thread.key, "rep_1/review/run_7/production");
2027 assert_eq!(thread.link.as_deref(), Some("/acme/rocket/actions/runs/run_7"));
2028 let told = notices(&asked, "acme/rocket", &actor("usr_ana", "ana"), None, &nobody());
2029 let names: Vec<&str> = told.iter().map(|notice| notice.username.as_str()).collect();
2030 // Whoever started the run is told too, when they review it.
2031 assert_eq!(names, ["cy", "ana"]);
2032 assert!(told[0].title.contains("waiting for your review to deploy to production"));
2033 }
2034
2035 #[test]
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar2036 fn teams_asked_to_review_tell_the_people_they_name() {
2037 let requested = event(
2038 "pull.review_requested",
2039 Some("usr_ana"),
2040 json!({ "number": 7, "reviewers": ["cy"], "teams": [
2041 { "team": "acme/backend", "notified": ["cy", "dee", "ana"], "assigned": ["cy"] }
2042 ] }),
2043 );
2044 let notices_ = notices(&requested, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
2045 // cy was picked and is told as a reviewer; dee through the team;
2046 // never whoever asked.
2047 assert_eq!(
2048 told(&notices_),
2049 vec![("cy", Reason::ReviewRequested, Severity::Warning), ("dee", Reason::ReviewRequested, Severity::Warning)]
2050 );
2051 assert_eq!(notices_[0].title, "ana asked you to review acme/rocket#7");
2052 assert_eq!(notices_[1].title, "ana asked @acme/backend to review acme/rocket#7");
2053 assert_eq!(
2054 subscribes(&requested, Some(&pull())),
2055 vec![("cy".into(), Reason::ReviewRequested), ("dee".into(), Reason::ReviewRequested), ("ana".into(), Reason::ReviewRequested)]
2056 );
2057 // Asked by the CODEOWNERS file.
2058 let owned = event(
2059 "pull.review_requested",
2060 None,
2061 json!({ "number": 7, "reviewers": ["bo"], "codeOwners": true, "teams": [{ "team": "acme/docs", "notified": ["wren"], "assigned": [] }] }),
2062 );
2063 let notices_ = notices(&owned, "acme/rocket", &Actor::default(), Some(&pull()), &nobody());
2064 assert_eq!(notices_[0].title, "acme/rocket#7 changes files you own");
2065 assert_eq!(notices_[1].title, "acme/rocket#7 changes files @acme/docs owns");
2066 }
2067
2068 #[test]
2069 fn a_team_mention_tells_its_people_once_and_a_name_wins() {
2070 let mut on = comment(person("usr_bo", "bo"), "cc @ana @acme/backend", &["ana"]);
2071 if let Some(comment) = on.comment.as_mut() {
2072 comment.team_mentions = vec![TeamMentioned {
2073 team: "acme/backend".into(),
2074 members: vec!["ana".into(), "cy".into(), "bo".into()],
2075 }];
2076 }
2077 let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
2078 let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &nobody());
2079 let mentioned: Vec<(&str, Reason)> = notices_.iter().map(|n| (n.username.as_str(), n.reason)).filter(|(_, r)| matches!(r, Reason::Mention | Reason::TeamMention)).collect();
2080 assert_eq!(mentioned, vec![("ana", Reason::Mention), ("cy", Reason::TeamMention)]);
2081 assert_eq!(
2082 notices_.iter().find(|n| n.username == "cy").unwrap().title,
2083 "bo mentioned @acme/backend on acme/rocket#7"
2084 );
2085 let subscribed = subscribes(&created, Some(&on));
2086 assert!(subscribed.contains(&("cy".into(), Reason::TeamMention)));
2087 assert!(subscribed.contains(&("ana".into(), Reason::Mention)));
2088 }
2089
2090 #[test]
2091 fn a_description_tells_the_people_and_teams_it_mentions() {
2092 let mut on = pull();
2093 on.mentions = vec!["dee".into()];
2094 on.team_mentions = vec![TeamMentioned {
2095 team: "acme/web".into(),
2096 members: vec!["eve".into()],
2097 }];
2098 let opened = event("pull.opened", Some("usr_ana"), json!({ "number": 7 }));
2099 let notices_ = notices(&opened, "acme/rocket", &actor("usr_ana", "ana"), Some(&on), &nobody());
2100 assert!(told(&notices_).contains(&("dee", Reason::Mention, Severity::Info)));
2101 assert!(told(&notices_).contains(&("eve", Reason::TeamMention, Severity::Info)));
2102 assert_eq!(
2103 subscribes(&opened, Some(&on)),
2104 vec![("dee".into(), Reason::Mention), ("eve".into(), Reason::TeamMention)]
2105 );
2106 }
2107
2108 #[test]
Inbox: threads, reasons, subscriptions and watching2109 fn reviewers_and_assignees_asked_are_told_never_whoever_asked() {
2110 let requested = event("pull.review_requested", Some("usr_ana"), json!({ "number": 7, "reviewers": ["bo", "g1t", "ana"] }));
2111 let notices_ = notices(&requested, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
2112 assert_eq!(told(&notices_), vec![("bo", Reason::ReviewRequested, Severity::Warning)]);
2113 assert_eq!(notices_[0].title, "ana asked you to review acme/rocket#7");
2114 let assigned = event("issue.assigned", Some("usr_ana"), json!({ "number": 3, "added": ["ana", "eve"] }));
2115 let notices_ = notices(&assigned, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
2116 assert_eq!(told(&notices_), vec![("eve", Reason::Assign, Severity::Info)]);
2117 assert_eq!(notices_[0].title, "ana assigned you to acme/rocket#3");
2118 // Both subscribe whoever they name, never g1t.
2119 assert_eq!(subscribes(&requested, Some(&pull())), vec![("bo".into(), Reason::ReviewRequested), ("ana".into(), Reason::ReviewRequested)]);
2120 assert_eq!(subscribes(&assigned, Some(&pull())), vec![("ana".into(), Reason::Assign), ("eve".into(), Reason::Assign)]);
2121 }
2122
2123 #[test]
Inbox: the events service tells people what needs them as events arrive2124 fn failures_go_to_whoever_answers_for_the_change() {
2125 let failed = event("checks.completed", None, json!({ "number": 7, "status": "failed" }));
Inbox: threads, reasons, subscriptions and watching2126 assert_eq!(
2127 told(&notices(&failed, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())),
2128 vec![("ana", Reason::CiActivity, Severity::Error)]
2129 );
Inbox: the events service tells people what needs them as events arrive2130 // g1t's change is the person's who asked for it, never g1t's.
Inbox: threads, reasons, subscriptions and watching2131 let notices_ = notices(&failed, "acme/rocket", &Actor::default(), Some(&g1t_pull()), &nobody());
2132 assert_eq!(told(&notices_), vec![("ana", Reason::CiActivity, Severity::Error)]);
Inbox: the events service tells people what needs them as events arrive2133 let errored = event("checks.completed", None, json!({ "number": 7, "status": "errored" }));
Inbox: threads, reasons, subscriptions and watching2134 assert_eq!(
2135 notices(&errored, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())[0].title,
2136 "Checks could not run on acme/rocket#7"
2137 );
Inbox: the events service tells people what needs them as events arrive2138 }
2139
2140 #[test]
2141 fn a_workflow_that_fails_tells_its_pull_requests_owner_even_if_they_pushed() {
2142 let failed = event("workflow.completed", Some("usr_ana"), json!({ "pull": 7, "workflow": "CI", "conclusion": "failure" }));
Inbox: threads, reasons, subscriptions and watching2143 let notices_ = notices(&failed, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
2144 assert_eq!(told(&notices_), vec![("ana", Reason::CiActivity, Severity::Error)]);
Inbox: the events service tells people what needs them as events arrive2145 assert_eq!(notices_[0].title, "CI failed on acme/rocket#7");
2146 // On a branch: whoever pushed.
2147 let pushed = event(
2148 "workflow.completed",
2149 Some("usr_bo"),
Inbox: threads, reasons, subscriptions and watching2150 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 arrive2151 );
Inbox: threads, reasons, subscriptions and watching2152 let notices_ = notices(&pushed, "acme/rocket", &actor("usr_bo", "bo"), None, &nobody());
2153 assert_eq!(told(&notices_), vec![("bo", Reason::CiActivity, Severity::Error)]);
Inbox: the events service tells people what needs them as events arrive2154 assert_eq!(notices_[0].title, "Deploy failed on main in acme/rocket");
2155 assert_eq!(notices_[0].body, "Run 12 at abcdef0");
2156 // Nobody to tell when g1t pushed.
Inbox: threads, reasons, subscriptions and watching2157 assert!(notices(&pushed, "acme/rocket", &actor("g1t", "g1t"), None, &nobody()).is_empty());
2158 // Every failure of a workflow on a branch is one thread.
2159 let wanted = wants(&pushed).unwrap();
2160 let thread = thread_of(&pushed, &wanted, None);
2161 assert_eq!(thread.key, "rep_1/run/.g1t/workflows/deploy.yml@main");
2162 assert_eq!(thread.run_id.as_deref(), Some("run_9"));
2163 }
2164
2165 #[test]
2166 fn deployments_tell_whoever_answers_for_them_and_watchers() {
2167 let failed = event(
2168 "deployment.failed",
2169 Some("usr_bo"),
2170 json!({ "repoId": "rep_1", "projectId": "prj_1", "project": "rocket", "deploymentId": "dpl_1", "error": "The build failed.", "path": "/acme/rocket/deployments/dpl_1" }),
2171 );
2172 let audience = watched(&[("cy", WatchLevel::Custom, &["deployments"]), ("dee", WatchLevel::Custom, &["issues"]), ("eve", WatchLevel::All, &[])]);
2173 let notices_ = notices(&failed, "acme/rocket", &actor("usr_bo", "bo"), None, &audience);
2174 assert_eq!(
2175 told(&notices_),
2176 vec![
2177 ("bo", Reason::CiActivity, Severity::Error),
2178 ("cy", Reason::Subscribed, Severity::Error),
2179 ("eve", Reason::Subscribed, Severity::Error),
2180 ]
2181 );
2182 assert_eq!(notices_[0].title, "Production of rocket failed to deploy");
2183 assert_eq!(notices_[0].body, "The build failed.");
2184 let thread = thread_of(&failed, &wants(&failed).unwrap(), None);
2185 assert_eq!(thread.key, "rep_1/deploy/prj_1/production");
2186 assert_eq!(url(Some("acme/rocket"), thread.kind, None, None, thread.link.as_deref()), "/acme/rocket/deployments/dpl_1");
2187 // A success is news to the owner only after a failure.
2188 let live = event("deployment.succeeded", Some("usr_bo"), json!({ "repoId": "rep_1", "projectId": "prj_1", "project": "rocket", "commit": "abcdef0123" }));
2189 assert!(notices(&live, "acme/rocket", &actor("usr_bo", "bo"), None, &nobody()).is_empty());
2190 let recovered = event("deployment.succeeded", Some("usr_bo"), json!({ "repoId": "rep_1", "projectId": "prj_1", "project": "rocket", "recovered": true }));
2191 let notices_ = notices(&recovered, "acme/rocket", &actor("usr_bo", "bo"), None, &nobody());
2192 assert_eq!(told(&notices_), vec![("bo", Reason::CiActivity, Severity::Success)]);
2193 assert_eq!(notices_[0].title, "Production of rocket is live again");
2194 // A preview's is the pull request's owner's.
2195 let preview = event("deployment.failed", None, json!({ "repoId": "rep_1", "projectId": "prj_1", "number": 7, "triggeredBy": "g1t" }));
2196 assert_eq!(told(&notices(&preview, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())), vec![("ana", Reason::CiActivity, Severity::Error)]);
2197 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 arrive2198 }
2199
2200 #[test]
Merge branch 'mirroring' into artifacts-mode2201 fn a_mirror_tells_only_the_people_its_settings_name() {
2202 let down = event(
2203 "mirror.unreachable",
2204 None,
2205 json!({
2206 "repoId": "rep_1", "repo": "acme/rocket", "remoteId": "rmt_1", "remote": "github.com/acme/rocket",
2207 "state": "standby", "title": "github.com/acme/rocket is not answering. acme/rocket keeps its copy.",
2208 "detail": "Take over acme/rocket to keep working on g1t until it is back.",
2209 "notify": ["ana"], "link": "/acme/rocket/settings/mirroring"
2210 }),
2211 );
2212 assert_eq!(wants(&down), Some(Wanted { repo_id: "rep_1".into(), number: None, comment_id: None }));
2213 // Watchers of everything hear nothing: only those named.
2214 let audience = watched(&[("eve", WatchLevel::All, &[])]);
2215 let notices_ = notices(&down, "acme/rocket", &Actor::default(), None, &audience);
2216 assert_eq!(told(&notices_), vec![("ana", Reason::StateChange, Severity::Warning)]);
2217 assert!(notices_[0].body.contains("Take over"));
2218 let thread = thread_of(&down, &wants(&down).unwrap(), None);
2219 assert_eq!(thread.key, "rep_1/mirror");
2220 assert_eq!(url(Some("acme/rocket"), thread.kind, None, None, thread.link.as_deref()), "/acme/rocket/settings/mirroring");
2221 // The banner-only default names nobody, so nobody is told.
2222 let quiet = event("mirror.unreachable", None, json!({ "repoId": "rep_1", "title": "x", "link": "/x" }));
2223 assert!(notices(&quiet, "acme/rocket", &Actor::default(), None, &audience).is_empty());
2224 let back = event("mirror.state_changed", None, json!({ "repoId": "rep_1", "state": "standby", "title": "handed back", "notify": ["ana"], "link": "/x" }));
2225 assert_eq!(told(&notices(&back, "acme/rocket", &Actor::default(), None, &audience)), vec![("ana", Reason::StateChange, Severity::Success)]);
2226 }
2227
2228 #[test]
Teams and CODEOWNERS, labels and milestones, dependency updates, the security suite, and a clearer top bar2229 fn security_alerts_tell_the_pusher_the_owners_and_member_watchers() {
2230 let blocked = event(
2231 "secret_scanning_alert.created",
2232 None,
2233 json!({
2234 "repoId": "rep_1", "alertId": "sec_1", "alertType": "secret_scanning", "severity": "critical",
2235 "title": "An AWS access key in config/prod.env", "link": "/acme/rocket/security/secret-scanning/sec_1",
2236 "state": "open", "pusher": "bo", "notify": ["ana"], "members": ["ana", "bo", "cy"]
2237 }),
2238 );
2239 assert_eq!(wants(&blocked), Some(Wanted { repo_id: "rep_1".into(), number: None, comment_id: None }));
2240 // cy watches for security alerts and is a member; eve watches
2241 // everything but cannot see findings; dee watches issues only.
2242 let audience = watched(&[
2243 ("cy", WatchLevel::Custom, &["security"]),
2244 ("dee", WatchLevel::Custom, &["issues"]),
2245 ("eve", WatchLevel::All, &[]),
2246 ]);
2247 let notices_ = notices(&blocked, "acme/rocket", &Actor::default(), None, &audience);
2248 assert_eq!(
2249 told(&notices_),
2250 vec![
2251 ("bo", Reason::SecurityAlert, Severity::Error),
2252 ("ana", Reason::SecurityAlert, Severity::Error),
2253 ("cy", Reason::SecurityAlert, Severity::Error),
2254 ]
2255 );
2256 assert_eq!(notices_[0].title, "A push to acme/rocket was blocked: it adds a secret");
2257 assert_eq!(notices_[0].body, "An AWS access key in config/prod.env");
2258 let thread = thread_of(&blocked, &wants(&blocked).unwrap(), None);
2259 assert_eq!(thread.key, "rep_1/security/sec_1");
2260 assert_eq!(url(Some("acme/rocket"), thread.kind, None, None, thread.link.as_deref()), "/acme/rocket/security/secret-scanning/sec_1");
2261 // A bypass request goes to the reviewers named, never to watchers.
2262 let requested = event(
2263 "secret_scanning.bypass_requested",
2264 Some("usr_bo"),
2265 json!({ "repoId": "rep_1", "alertId": "sec_1", "requestId": "byp_1", "title": "An AWS access key in a.env", "notify": ["ana"], "link": "/acme/-/security/bypass-requests" }),
2266 );
2267 let notices_ = notices(&requested, "acme/rocket", &actor("usr_bo", "bo"), None, &audience);
2268 assert_eq!(told(&notices_), vec![("ana", Reason::SecurityAlert, Severity::Warning)]);
2269 assert_eq!(notices_[0].title, "bo asked to bypass push protection in acme/rocket");
2270 // Fixes and dismissals are not news to the inbox.
2271 assert_eq!(wants(&event("code_scanning_alert.fixed", None, json!({ "repoId": "rep_1" }))), None);
2272 }
2273
2274 #[test]
Inbox: the events service tells people what needs them as events arrive2275 fn nobody_hears_of_what_they_did_themselves() {
2276 let merged = event("pull.merged", Some("usr_ana"), json!({ "number": 7 }));
Inbox: threads, reasons, subscriptions and watching2277 let by_ana = notices(&merged, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()), &nobody());
2278 // The assignee still hears; ana, who merged, does not.
2279 assert_eq!(told(&by_ana), vec![("bo", Reason::StateChange, Severity::Success)]);
2280 let by_bo = notices(&merged, "acme/rocket", &actor("usr_bo", "bo"), Some(&pull()), &nobody());
2281 assert_eq!(told(&by_bo), vec![("ana", Reason::StateChange, Severity::Success)]);
Inbox: the events service tells people what needs them as events arrive2282 assert_eq!(by_bo[0].title, "bo merged acme/rocket#7");
Inbox: threads, reasons, subscriptions and watching2283 let by_queue = notices(&merged, "acme/rocket", &Actor::default(), Some(&pull()), &nobody());
Inbox: the events service tells people what needs them as events arrive2284 assert_eq!(by_queue[0].title, "acme/rocket#7 was merged");
2285 }
2286
2287 #[test]
Inbox: threads, reasons, subscriptions and watching2288 fn closing_and_reopening_tells_everyone_subscribed_and_watchers() {
2289 let closed = event("issue.closed", Some("usr_bo"), json!({ "number": 3, "reason": "completed" }));
2290 let issue = InboxSubject {
2291 kind: Some(SubjectKind::Issue),
2292 title: "Crash".into(),
2293 author: person("usr_cy", "cy"),
2294 assignees: vec!["bo".into()],
2295 ..InboxSubject::default()
2296 };
2297 let mut audience = subscribed(&[("eve", State::Subscribed, Some(Reason::Comment)), ("fay", State::Unsubscribed, None)]);
2298 audience.watchers = watched(&[("gus", WatchLevel::All, &[]), ("hal", WatchLevel::Custom, &["pulls"])]).watchers;
2299 let notices_ = notices(&closed, "acme/rocket", &actor("usr_bo", "bo"), Some(&issue), &audience);
2300 assert_eq!(
2301 told(&notices_),
2302 vec![
2303 ("cy", Reason::StateChange, Severity::Info),
2304 ("eve", Reason::StateChange, Severity::Info),
2305 ("gus", Reason::Subscribed, Severity::Info),
2306 ]
2307 );
2308 assert_eq!(notices_[0].title, "bo closed acme/rocket#3");
2309 let by_pull = event("issue.closed", None, json!({ "number": 3, "resolvedBy": 9 }));
2310 assert_eq!(notices(&by_pull, "acme/rocket", &Actor::default(), Some(&issue), &nobody())[0].title, "acme/rocket#3 was closed by #9");
2311 let reopened = event("issue.reopened", Some("usr_cy"), json!({ "number": 3 }));
2312 assert_eq!(told(&notices(&reopened, "acme/rocket", &actor("usr_cy", "cy"), Some(&issue), &nobody())), vec![("bo", Reason::StateChange, Severity::Info)]);
2313 }
2314
2315 #[test]
2316 fn opening_tells_who_it_names_and_watchers_of_its_kind() {
2317 let opened = event("pull.opened", Some("usr_ana"), json!({ "number": 7 }));
2318 let subject = InboxSubject { reviewers: vec!["cy".into(), "g1t".into()], ..pull() };
2319 let audience = watched(&[("bo", WatchLevel::All, &[]), ("dee", WatchLevel::Custom, &["issues"]), ("eve", WatchLevel::Custom, &["pulls"]), ("fay", WatchLevel::Participating, &[])]);
2320 let notices_ = notices(&opened, "acme/rocket", &actor("usr_ana", "ana"), Some(&subject), &audience);
2321 assert_eq!(
2322 told(&notices_),
2323 vec![
2324 ("bo", Reason::Assign, Severity::Info),
2325 ("cy", Reason::ReviewRequested, Severity::Warning),
2326 ("eve", Reason::Subscribed, Severity::Info),
2327 ]
2328 );
2329 assert_eq!(notices_[2].title, "ana opened acme/rocket#7");
2330 }
2331
2332 #[test]
2333 fn ignoring_a_repository_silences_it_and_unsubscribing_keeps_only_what_is_asked() {
2334 let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
2335 let on = comment(person("usr_bo", "bo"), "@cy have a look", &["cy"]);
2336 let audience = watched(&[("cy", WatchLevel::Ignore, &[]), ("ana", WatchLevel::All, &[])]);
2337 assert!(notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &audience).iter().all(|notice| notice.username != "cy"));
2338 let audience = subscribed(&[("cy", State::Unsubscribed, None), ("ana", State::Unsubscribed, None)]);
2339 let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &audience);
2340 // cy was mentioned, which is asked of them; ana unsubscribed from the conversation.
2341 assert_eq!(told(&notices_), vec![("cy", Reason::Mention, Severity::Info)]);
2342 }
2343
2344 #[test]
Inbox: the events service tells people what needs them as events arrive2345 fn g1t_finishing_or_reviewing_tells_the_person_it_worked_for() {
2346 // The agent acts as the person it works for: still an outcome they hear of.
2347 let ready = event("pull.ready", Some("usr_ana"), json!({ "number": 7 }));
Inbox: threads, reasons, subscriptions and watching2348 let notices_ = notices(&ready, "acme/rocket", &actor("usr_ana", "ana"), Some(&g1t_pull()), &nobody());
2349 assert_eq!(told(&notices_), vec![("ana", Reason::Author, Severity::Success)]);
Inbox: the events service tells people what needs them as events arrive2350 assert_eq!(notices_[0].title, "g1t finished acme/rocket#7");
2351 // A person's draft marked ready is not news to them.
Inbox: threads, reasons, subscriptions and watching2352 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 arrive2353
2354 let approve = event("review.completed", None, json!({ "number": 7, "verdict": "approve" }));
Inbox: threads, reasons, subscriptions and watching2355 assert_eq!(
2356 told(&notices(&approve, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())),
2357 vec![("ana", Reason::Author, Severity::Success)]
2358 );
Inbox: the events service tells people what needs them as events arrive2359 let changes = event("review.completed", None, json!({ "number": 7, "verdict": "request_changes" }));
2360 assert_eq!(
Inbox: threads, reasons, subscriptions and watching2361 told(&notices(&changes, "acme/rocket", &Actor::default(), Some(&pull()), &nobody())),
2362 vec![("ana", Reason::Author, Severity::Info)]
Inbox: the events service tells people what needs them as events arrive2363 );
2364 }
2365
2366 #[test]
Cards you act on in chat; agents comment and review as themselves; names shown cleanly; commits on the calendar2367 fn an_agent_s_review_is_shown_by_its_name_and_marked_advisory() {
2368 // Margo (@margo), acting for bo, asks for changes on ana's pull request.
2369 let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
2370 let mut on = comment(person("agt_1", "ana"), "Trim the name.", &[]);
2371 if let Some(said) = on.comment.as_mut() {
2372 said.agent = Some(g1t_contracts::work::AgentRef {
2373 id: "agt_1".into(),
2374 handle: "ana".into(),
2375 display_name: "Margo".into(),
2376 avatar_seed: "margo".into(),
2377 });
2378 said.acting_for = Some(person("usr_bo", "bo"));
2379 said.verdict = Some("request_changes".into());
2380 said.advisory = true;
2381 }
2382 let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &nobody());
2383 // Its handle sharing ana's name leaves nobody out: ana owns it and hears.
2384 assert!(told(&notices_).contains(&("ana", Reason::Author, Severity::Info)));
2385 assert_eq!(notices_[0].title, "Margo (agent) asked for changes on acme/rocket#7 (advisory)");
2386 // An agent subscribes nobody by its handle.
2387 assert!(subscribes(&created, Some(&on)).iter().all(|(name, _)| name != "ana"));
2388 }
2389
2390 #[test]
Inbox: threads, reasons, subscriptions and watching2391 fn comments_tell_those_mentioned_then_everyone_subscribed_never_the_writer() {
Inbox: the events service tells people what needs them as events arrive2392 let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
2393 let on = comment(person("usr_bo", "bo"), "@cy @bo have a look", &["cy", "bo", "g1t"]);
Inbox: threads, reasons, subscriptions and watching2394 let audience = subscribed(&[("eve", State::Subscribed, Some(Reason::Comment)), ("fay", State::Subscribed, Some(Reason::Manual))]);
2395 let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &audience);
Inbox: the events service tells people what needs them as events arrive2396 assert_eq!(
2397 told(&notices_),
Inbox: threads, reasons, subscriptions and watching2398 vec![
2399 ("cy", Reason::Mention, Severity::Info),
2400 ("ana", Reason::Author, Severity::Info),
2401 ("eve", Reason::Comment, Severity::Info),
2402 ("fay", Reason::Manual, Severity::Info),
2403 ]
Inbox: the events service tells people what needs them as events arrive2404 );
2405 assert_eq!(notices_[0].title, "bo mentioned you on acme/rocket#7");
Inbox: threads, reasons, subscriptions and watching2406 assert_eq!(notices_[1].title, "bo commented on acme/rocket#7");
Inbox: the events service tells people what needs them as events arrive2407 assert_eq!(notices_[0].body, "@cy @bo have a look");
2408 // Mentioned and the owner: told once, as mentioned.
2409 let on = comment(person("usr_bo", "bo"), "@ana", &["ana"]);
Inbox: threads, reasons, subscriptions and watching2410 assert_eq!(
2411 told(&notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on), &nobody())),
2412 vec![("ana", Reason::Mention, Severity::Info)]
2413 );
2414 // The owner's own comment tells the assignee, not the owner.
Inbox: the events service tells people what needs them as events arrive2415 let on = comment(person("usr_ana", "ana"), "thanks", &[]);
Inbox: threads, reasons, subscriptions and watching2416 assert_eq!(
2417 told(&notices(&created, "acme/rocket", &actor("usr_ana", "ana"), Some(&on), &nobody())),
2418 vec![("bo", Reason::Assign, Severity::Info)]
2419 );
Inbox: the events service tells people what needs them as events arrive2420 // Something that happened, not something written, tells nobody.
2421 let mut on = comment(person("usr_bo", "bo"), "assigned cy", &[]);
2422 on.comment.as_mut().unwrap().event = true;
Inbox: threads, reasons, subscriptions and watching2423 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 arrive2424 // An approval is good news.
Inbox: threads, reasons, subscriptions and watching2425 let mut on = comment(person("usr_cy", "cy"), "", &[]);
Inbox: the events service tells people what needs them as events arrive2426 on.comment.as_mut().unwrap().verdict = Some("approve".into());
Inbox: threads, reasons, subscriptions and watching2427 let notices_ = notices(&created, "acme/rocket", &actor("usr_cy", "cy"), Some(&on), &nobody());
2428 assert_eq!(told(&notices_)[0], ("ana", Reason::Author, Severity::Success));
Inbox: the events service tells people what needs them as events arrive2429 assert_eq!(notices_[0].body, "Add the inbox");
Inbox: threads, reasons, subscriptions and watching2430 // Writing and being mentioned subscribe, never g1t.
2431 let on = comment(person("usr_bo", "bo"), "@cy", &["cy"]);
2432 assert_eq!(subscribes(&created, Some(&on)), vec![("bo".into(), Reason::Comment), ("cy".into(), Reason::Mention)]);
2433 let on = comment(person(AGENT_ID, "g1t"), "done", &[]);
2434 assert!(subscribes(&created, Some(&on)).is_empty());
2435 }
2436
2437 #[test]
2438 fn issues_and_pull_requests_are_one_thread_each() {
2439 let created = event("comment.created", Some("usr_bo"), json!({ "repoId": "rep_1", "number": 7, "commentId": "cmt_1" }));
2440 let wanted = wants(&created).unwrap();
2441 let thread = thread_of(&created, &wanted, Some(&pull()));
2442 assert_eq!(thread.key, "rep_1#7");
2443 assert_eq!(thread.kind, Some(SubjectKind::Pull));
2444 let merged = event("pull.merged", None, json!({ "repoId": "rep_1", "number": 7 }));
2445 assert_eq!(thread_of(&merged, &wants(&merged).unwrap(), Some(&pull())).key, thread.key);
Inbox: the events service tells people what needs them as events arrive2446 }
2447
2448 #[test]
2449 fn items_link_to_what_they_are_about() {
Inbox: threads, reasons, subscriptions and watching2450 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Pull), Some(7), None, None), "/acme/rocket/pull/7");
2451 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Issue), Some(3), None, None), "/acme/rocket/issues/3");
2452 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Run), None, Some("run_1"), None), "/acme/rocket/actions/runs/run_1");
2453 assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Deploy), None, None, Some("/acme/site/deployments/dpl_1")), "/acme/site/deployments/dpl_1");
2454 assert_eq!(url(Some("acme/rocket"), None, None, None, None), "/acme/rocket");
2455 assert_eq!(url(None, None, None, None, None), "/inbox");
Inbox: the events service tells people what needs them as events arrive2456 }
2457
2458 #[test]
2459 fn long_titles_are_cut_to_a_line() {
2460 let long = "x".repeat(400);
2461 assert_eq!(clip(&long, MAX_TITLE).chars().count(), MAX_TITLE);
2462 assert_eq!(clip(" short ", MAX_TITLE), "short");
2463 }
2464
2465 #[test]
2466 fn marks_change_only_what_they_say() {
2467 assert_eq!(mark_change(InboxMark::Read), "read_at = COALESCE(read_at, ?1)");
2468 assert!(mark_change(InboxMark::Done).contains("done_at = ?1"));
Inbox: threads, reasons, subscriptions and watching2469 assert_eq!(mark_change(InboxMark::Unsnooze), "snoozed_until = NULL");
2470 assert!(ranked(&ListInboxArgs::default()));
2471 assert!(!ranked(&ListInboxArgs { reason: Some(Reason::Mention), ..ListInboxArgs::default() }));
2472 assert!(!ranked(&ListInboxArgs { participating: true, ..ListInboxArgs::default() }));
2473 }
2474
2475 #[test]
2476 fn filters_bind_in_order_after_what_is_bound() {
2477 let a = ListInboxArgs {
2478 reason: Some(Reason::Mention),
2479 repo_id: Some("rep_1".into()),
2480 participating: true,
2481 unread: true,
2482 ..ListInboxArgs::default()
2483 };
2484 // Two values come first: the username and now.
2485 let (conditions, values) = filters(&a, 2);
2486 assert_eq!(
2487 conditions,
2488 vec!["reason = ?3", "repo_id = ?4", "reason NOT IN ('manual', 'subscribed')", "read_at IS NULL"]
2489 );
2490 assert_eq!(values, vec!["mention", "rep_1"]);
2491 }
2492
2493 #[test]
2494 fn the_bump_keeps_whatever_is_most_urgent_while_unread() {
2495 assert!(BUMP.contains("WHERE NOT EXISTS (SELECT 1 FROM inbox_activity WHERE event_id = ?4 AND username = ?2)"));
2496 assert!(BUMP.contains("ON CONFLICT (username, thread) DO UPDATE"));
2497 assert!(BUMP.contains("activity = inbox_items.activity + 1"));
2498 // Its numbered parameters run from 1 to 19.
2499 for at in 1..=19 {
2500 assert!(BUMP.contains(&format!("?{at}")), "?{at}");
2501 }
2502 assert!(!BUMP.contains("?20"));
Inbox: the events service tells people what needs them as events arrive2503 }
Fine-grained personal tokens, workspace token rules and approvals in identity2504
2505 #[test]
Merge the workspace shell: navigation and phone shell, g1t as orchestrator, agents in roles with audience-checked reads, reactions and custom emoji, live notifications and browser push, the homepage tour (agents 0002, chat 0002)2506 fn live_notifications_say_what_and_where() {
2507 let notice = Notice {
2508 username: "ana".into(),
2509 reason: Reason::Agent,
2510 severity: Severity::Warning,
2511 title: "g1t needs you on acme/api#4".into(),
2512 body: "Which database?".into(),
2513 };
2514 let live = live_notification("evt_1", &notice, "/Acme/api/pull/4#c1", "2026-10-08T00:00:00Z");
2515 assert_eq!(live["id"], "inbox:evt_1");
2516 assert_eq!(live["kind"], "agent_waiting");
2517 assert_eq!(live["workspace"], "acme");
2518 assert_eq!(live["href"], "/Acme/api/pull/4#c1");
2519 assert_eq!(live["actor"]["name"], "acme/api");
2520 let mention = Notice { reason: Reason::TeamMention, ..notice };
2521 assert_eq!(live_notification("evt_2", &mention, "https://elsewhere", "t")["href"], "/inbox");
2522 assert_eq!(live_notification("evt_2", &mention, "/x", "t")["kind"], "mention");
2523 }
2524
2525 #[test]
2526 fn inbox_counts_are_told_as_they_are_by_username() {
2527 assert_eq!(inbox_count_update("Ana", 3), serde_json::json!({ "username": "ana", "unread": 3 }));
2528 assert_eq!(inbox_count_update("bo", 0)["unread"], 0);
2529 }
2530
2531 #[test]
Fine-grained personal tokens, workspace token rules and approvals in identity2532 fn token_approvals_tell_the_people_named_about_no_repository() {
2533 let mut asked = event(
2534 "token.approval_requested",
2535 Some("usr_ana"),
2536 json!({ "workspace": "acme", "tokenId": "tok_1", "notify": ["Bo", "cy", "bo", "g1t"], "title": "ana asks", "body": "ci: contents: write", "link": "/acme/-/settings/tokens" }),
2537 );
2538 asked.repo_id = None;
2539 let (thread, told) = workspace_notices(&asked, None);
2540 assert_eq!(thread, "workspace:acme/token/tok_1");
2541 let mut names: Vec<&str> = told.iter().map(|notice| notice.username.as_str()).collect();
2542 names.sort();
2543 assert_eq!(names, ["bo", "cy"], "each once, lowercased, never g1t");
2544 assert!(told.iter().all(|notice| notice.reason == Reason::ReviewRequested && notice.severity == Severity::Warning));
2545 assert!(wants(&asked).is_none(), "about no repository");
2546 let reviewed = event("token.approval_reviewed", None, json!({ "workspace": "acme", "tokenId": "tok_1", "notify": ["ana"], "title": "approved" }));
2547 let (_, told) = workspace_notices(&reviewed, None);
2548 assert_eq!(told[0].reason, Reason::Author);
2549 assert!(WORKSPACE_EVENTS.contains(&"token.approval_reviewed"));
2550 }
Merge workspace invitations: nobody joins a workspace without saying yes, people are found by username, your own invites can bring someone in, and nobody is left without a workspace (identity 0040)2551
2552 #[test]
2553 fn a_workspace_invitation_asks_its_person_and_tells_whoever_sent_it_the_answer() {
2554 let mut sent = event(
2555 "workspace_invitation.created",
2556 Some("usr_owner"),
2557 json!({ "workspace": "flagon-io", "invitationId": "inv_1", "thread": "invitation:inv_1", "notify": ["daweazl"], "title": "@syntaqx invited you to join Flagon, Inc.", "link": "/invitations" }),
2558 );
2559 sent.repo_id = None;
2560 let (thread, told) = workspace_notices(&sent, None);
2561 assert_eq!(thread, "invitation:inv_1");
2562 assert_eq!(told.len(), 1);
2563 assert_eq!((told[0].username.as_str(), told[0].reason, told[0].severity), ("daweazl", Reason::ReviewRequested, Severity::Warning));
2564 let accepted = event("workspace_invitation.accepted", None, json!({ "workspace": "flagon-io", "thread": "invitation:inv_1", "notify": ["syntaqx"] }));
2565 let (thread, told) = workspace_notices(&accepted, None);
2566 assert_eq!(thread, "invitation:inv_1");
2567 assert_eq!((told[0].reason, told[0].severity), (Reason::Author, Severity::Success));
2568 let declined = event("workspace_invitation.declined", None, json!({ "workspace": "flagon-io", "thread": "invitation:inv_1", "notify": ["syntaqx"] }));
2569 assert_eq!(workspace_notices(&declined, None).1[0].severity, Severity::Info);
2570 // Answered or revoked, the person's item is done; a revocation tells nobody.
2571 for kind in ["workspace_invitation.accepted", "workspace_invitation.declined", "workspace_invitation.revoked"] {
2572 assert!(INVITATION_ANSWERED.contains(&kind));
2573 }
2574 assert!(!WORKSPACE_EVENTS.contains(&"workspace_invitation.revoked"));
2575 }
Inbox: the events service tells people what needs them as events arrive2576}

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