Skip to content

Commit

Inbox: the events service tells people what needs them as events arrive

Failed checks and workflows, an agent waiting, g1t finishing or reviewing, merges, comments and mentions each become an item in the inbox of the people they concern, never the person who acted. Work answers what an event names (inbox_subject); the events service keeps the items (migration 0005_inbox) and lists, counts and marks them.

syntaqxcommitted Parentb31c281Browse files
8 files+1543−40/8 viewed
+287−0
1+//! The inbox: what needs a person, or what they follow, as it happens.
2+//!
3+//! The events service keeps it, beside the event log: as events arrive it
4+//! works out who should hear of each (see `services/events/src/inbox.rs`)
5+//! and writes one item per person. Items are kept by username, which never
6+//! changes. Methods, served at `POST /rpc/<method>` on the events service:
7+//!
8+//! - `inbox_list` takes `ListInboxArgs` and returns `InboxPage`. Items about
9+//! a repository the viewer can no longer read are dropped as they are
10+//! found.
11+//! - `inbox_counts` takes `InboxCountsArgs` and returns `InboxCounts`: the
12+//! unread items, by severity. One query, for every page's top bar.
13+//! - `inbox_mark` takes `MarkInboxArgs` and returns how many items changed.
14+//!
15+//! What an event is about (the issue or pull request, its people, the
16+//! comment) comes from the work service's `inbox_subject`, which takes
17+//! `InboxSubjectArgs` and returns `Option<InboxSubject>`.
18+
19+use serde::{Deserialize, Serialize};
20+
21+use crate::Viewer;
22+use crate::credentials::Principal;
23+
24+/// How much an item matters, and how it is shown: a failure, something a
25+/// person must answer (an agent waiting on them), something that went
26+/// well, or something to know.
27+#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
28+#[serde(rename_all = "lowercase")]
29+pub enum Severity {
30+ Error,
31+ Warning,
32+ Success,
33+ Info,
34+}
35+
36+impl Severity {
37+ pub const ALL: [Severity; 4] = [Severity::Error, Severity::Warning, Severity::Success, Severity::Info];
38+
39+ pub fn as_str(self) -> &'static str {
40+ match self {
41+ Severity::Error => "error",
42+ Severity::Warning => "warning",
43+ Severity::Success => "success",
44+ Severity::Info => "info",
45+ }
46+ }
47+
48+ pub fn parse(value: &str) -> Option<Severity> {
49+ Severity::ALL.into_iter().find(|severity| severity.as_str() == value)
50+ }
51+}
52+
53+/// What an item is about.
54+#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
55+#[serde(rename_all = "lowercase")]
56+pub enum SubjectKind {
57+ Issue,
58+ Pull,
59+ /// A workflow run.
60+ Run,
61+}
62+
63+impl SubjectKind {
64+ pub fn as_str(self) -> &'static str {
65+ match self {
66+ SubjectKind::Issue => "issue",
67+ SubjectKind::Pull => "pull",
68+ SubjectKind::Run => "run",
69+ }
70+ }
71+
72+ pub fn parse(value: &str) -> Option<SubjectKind> {
73+ [SubjectKind::Issue, SubjectKind::Pull, SubjectKind::Run]
74+ .into_iter()
75+ .find(|kind| kind.as_str() == value)
76+ }
77+}
78+
79+/// One thing a person was told.
80+#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
81+#[serde(rename_all = "camelCase")]
82+pub struct InboxItem {
83+ pub id: String,
84+ /// Why they were told, such as `checks_failed` or `mentioned`.
85+ pub reason: String,
86+ pub severity: Severity,
87+ /// One line: what happened, and where.
88+ pub title: String,
89+ /// One line: what it happened to, such as the pull request's title.
90+ pub body: String,
91+ /// `owner/name`.
92+ pub repo: Option<String>,
93+ /// The workspace it happened in.
94+ pub workspace: Option<String>,
95+ pub subject: Option<SubjectKind>,
96+ /// The issue or pull request's number.
97+ pub number: Option<u32>,
98+ /// Where it is on g1t.sh: a path such as `/acme/rocket/pull/12`.
99+ pub url: String,
100+ /// Who did it: a username, or `g1t`. Absent when nobody did.
101+ pub actor: Option<String>,
102+ /// RFC 3339.
103+ pub created_at: String,
104+ pub read_at: Option<String>,
105+ pub done_at: Option<String>,
106+ pub saved: bool,
107+ /// While this is in the future the item is out of the list.
108+ pub snoozed_until: Option<String>,
109+}
110+
111+/// Which items a list shows.
112+#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
113+#[serde(rename_all = "lowercase")]
114+pub enum InboxView {
115+ /// Everything not done and not snoozed: the inbox itself.
116+ #[default]
117+ Inbox,
118+ /// Saved items, done or not.
119+ Saved,
120+ /// Items marked done.
121+ Done,
122+}
123+
124+/// `inbox_list`. Newest first, except that unread warnings (an agent
125+/// waiting on the person) come before everything else in the inbox.
126+#[derive(Clone, Debug, Serialize, Deserialize)]
127+#[serde(rename_all = "camelCase")]
128+pub struct ListInboxArgs {
129+ /// Whose inbox: the person signed in. Their memberships decide which
130+ /// repositories they can still read.
131+ pub viewer: Viewer,
132+ #[serde(default)]
133+ pub view: InboxView,
134+ /// Only items of this severity.
135+ #[serde(default)]
136+ pub severity: Option<Severity>,
137+ #[serde(default)]
138+ pub unread: bool,
139+ /// The `next` of the page before.
140+ #[serde(default)]
141+ pub before: Option<String>,
142+ #[serde(default)]
143+ pub limit: Option<u32>,
144+}
145+
146+pub const DEFAULT_INBOX_PAGE: u32 = 30;
147+pub const MAX_INBOX_PAGE: u32 = 100;
148+
149+#[derive(Clone, Debug, Default, Serialize, Deserialize)]
150+#[serde(rename_all = "camelCase")]
151+pub struct InboxPage {
152+ pub items: Vec<InboxItem>,
153+ /// Pass as `before` for the next page; absent on the last.
154+ pub next: Option<String>,
155+}
156+
157+/// `inbox_counts`.
158+#[derive(Clone, Debug, Serialize, Deserialize)]
159+#[serde(rename_all = "camelCase")]
160+pub struct InboxCountsArgs {
161+ pub username: String,
162+}
163+
164+/// Unread items in the inbox view, by severity.
165+#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
166+#[serde(rename_all = "camelCase")]
167+pub struct InboxCounts {
168+ pub unread: u32,
169+ pub error: u32,
170+ pub warning: u32,
171+ pub success: u32,
172+ pub info: u32,
173+}
174+
175+/// What `inbox_mark` does to the items it names.
176+#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
177+#[serde(rename_all = "snake_case")]
178+pub enum InboxMark {
179+ Read,
180+ Unread,
181+ /// Out of the inbox, into Done; read with it.
182+ Done,
183+ /// Back into the inbox.
184+ Undone,
185+ Save,
186+ Unsave,
187+ /// Out of the inbox until `until`.
188+ Snooze,
189+}
190+
191+/// `inbox_mark`: changes the person's own items, by id, or every item in
192+/// their inbox when `ids` is empty and `all` is set (Mark all read).
193+#[derive(Clone, Debug, Serialize, Deserialize)]
194+#[serde(rename_all = "camelCase")]
195+pub struct MarkInboxArgs {
196+ pub username: String,
197+ pub mark: InboxMark,
198+ #[serde(default)]
199+ pub ids: Vec<String>,
200+ #[serde(default)]
201+ pub all: bool,
202+ /// With `all`: only items of this severity.
203+ #[serde(default)]
204+ pub severity: Option<Severity>,
205+ /// For `snooze`: RFC 3339.
206+ #[serde(default)]
207+ pub until: Option<String>,
208+}
209+
210+/// The most ids one `inbox_mark` call changes.
211+pub const MAX_MARK: usize = 100;
212+
213+/// `inbox_subject` on the work service: what an event names, for the
214+/// inbox. A service-to-service read: it checks nobody's access, and what
215+/// it returns is only ever shown to the people it names, or to those who
216+/// can read the repository.
217+#[derive(Clone, Debug, Serialize, Deserialize)]
218+#[serde(rename_all = "camelCase")]
219+pub struct InboxSubjectArgs {
220+ pub repo_id: String,
221+ pub number: u32,
222+ /// The comment the event is about, if any.
223+ #[serde(default)]
224+ pub comment_id: Option<String>,
225+}
226+
227+#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
228+#[serde(rename_all = "camelCase")]
229+pub struct InboxSubject {
230+ /// `issue` or `pull`.
231+ pub kind: Option<SubjectKind>,
232+ pub title: String,
233+ pub author: Principal,
234+ /// For work g1t did: the person it was done for.
235+ #[serde(default)]
236+ pub requested_by: Option<Principal>,
237+ /// Usernames.
238+ #[serde(default)]
239+ pub assignees: Vec<String>,
240+ /// Usernames, and `g1t`. Pull requests only.
241+ #[serde(default)]
242+ pub reviewers: Vec<String>,
243+ /// For a pull request: the issue it is for, with that issue's people.
244+ #[serde(default)]
245+ pub issue: Option<Box<InboxSubject>>,
246+ #[serde(default)]
247+ pub comment: Option<InboxComment>,
248+}
249+
250+impl InboxSubject {
251+ /// Whose it is to answer for: whoever asked g1t for it, or its author.
252+ pub fn owner(&self) -> &Principal {
253+ self.requested_by.as_ref().unwrap_or(&self.author)
254+ }
255+}
256+
257+#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
258+#[serde(rename_all = "camelCase")]
259+pub struct InboxComment {
260+ pub author: Principal,
261+ /// The first line or so, as written.
262+ pub excerpt: String,
263+ /// The people it mentions by name, outside code and quotes. Never `g1t`.
264+ #[serde(default)]
265+ pub mentions: Vec<String>,
266+ /// A review's verdict: `approve` or `request_changes`.
267+ #[serde(default)]
268+ pub verdict: Option<String>,
269+ /// Something that happened (an assignment, a close), not something written.
270+ #[serde(default)]
271+ pub event: bool,
272+}
273+
274+#[cfg(test)]
275+mod tests {
276+ use super::*;
277+
278+ #[test]
279+ fn severities_and_subjects_read_back() {
280+ for severity in Severity::ALL {
281+ assert_eq!(Severity::parse(severity.as_str()), Some(severity));
282+ }
283+ assert_eq!(Severity::parse("fatal"), None);
284+ assert_eq!(SubjectKind::parse("pull"), Some(SubjectKind::Pull));
285+ assert_eq!(serde_json::to_value(InboxMark::Unsave).unwrap(), "unsave");
286+ }
287+}
+1−0
1717 pub mod github;
1818 pub mod guardrails;
1919 pub mod identity;
20+pub mod inbox;
2021 pub mod integrations;
2122 mod ids;
2223 mod names;
+42−0
1+-- The inbox (src/inbox.rs): one row per person told of an event. Written
2+-- as events arrive; a redelivered event finds its rows there already.
3+-- Ids are time-sortable, so ordering by id is ordering by time.
4+CREATE TABLE inbox_items (
5+ id TEXT PRIMARY KEY,
6+ -- Whose: a username, lowercased. Usernames never change.
7+ username TEXT NOT NULL,
8+ -- The event it came from.
9+ event_id TEXT NOT NULL,
10+ -- Why they were told, such as checks_failed or mentioned.
11+ reason TEXT NOT NULL,
12+ -- error | warning | success | info
13+ severity TEXT NOT NULL,
14+ title TEXT NOT NULL,
15+ body TEXT NOT NULL,
16+ -- The workspace's slug, and owner/name. Both follow renames and transfers.
17+ workspace TEXT,
18+ repo_id TEXT,
19+ repo TEXT,
20+ -- issue | pull | run, and its number (an issue or pull request) or id (a run).
21+ subject TEXT,
22+ number INTEGER,
23+ run_id TEXT,
24+ -- Who did it: a username, or g1t.
25+ actor TEXT,
26+ -- RFC 3339 UTC.
27+ created_at TEXT NOT NULL,
28+ read_at TEXT,
29+ done_at TEXT,
30+ saved INTEGER NOT NULL DEFAULT 0,
31+ snoozed_until TEXT,
32+ UNIQUE (event_id, username)
33+);
34+-- A person's inbox, newest first; what is unread in it, for every page's
35+-- count; and what they saved.
36+CREATE INDEX inbox_recent ON inbox_items (username, id) WHERE done_at IS NULL;
37+CREATE INDEX inbox_unread ON inbox_items (username, severity) WHERE read_at IS NULL AND done_at IS NULL;
38+CREATE INDEX inbox_saved ON inbox_items (username, id) WHERE saved = 1;
39+CREATE INDEX inbox_done ON inbox_items (username, done_at) WHERE done_at IS NOT NULL;
40+-- Renames, transfers and purges find a repository's rows by id.
41+CREATE INDEX inbox_repo ON inbox_items (repo_id) WHERE repo_id IS NOT NULL;
42+CREATE INDEX inbox_time ON inbox_items (created_at);
+994−0
1+//! The inbox, kept beside the event log. See `g1t_contracts::inbox`.
2+//!
3+//! As each batch arrives from the bus, the events that need a person are
4+//! read against what they name (the work service's `inbox_subject`) and one
5+//! item is written for each person told. Who is told is worked out in
6+//! [`notices`], from the event and its subject alone:
7+//!
8+//! | Event | Who | Severity |
9+//! | --- | --- | --- |
10+//! | `agent.asked` | the pull request's owner, and its issue's owner and assignees | warning |
11+//! | `checks.completed`, failed or errored | the pull request's owner | error |
12+//! | `workflow.completed`, failed | the pull request's owner, or whoever pushed | error |
13+//! | `review.completed` by g1t | the pull request's owner | success, or info for changes asked |
14+//! | `pull.ready` for a change g1t made | whoever asked g1t for it | success |
15+//! | `pull.merged` | the pull request's owner | success |
16+//! | `comment.created` | everyone mentioned, then the owner | info (success for an approval) |
17+//!
18+//! Nobody is told of what they did themselves, and g1t is never told. A
19+//! failure here is logged and the batch goes on: the bus never waits on
20+//! the inbox, so an item can be missed, but nothing else is held up.
21+
22+use std::collections::{HashMap, HashSet};
23+
24+use g1t_contracts::credentials::Principal;
25+use g1t_contracts::events::{Event, WorkspaceRenamed};
26+use g1t_contracts::identity::{AGENT_ID, UsernamesArgs};
27+use g1t_contracts::inbox::*;
28+use g1t_contracts::repos::{PathByIdArgs, ReadableArgs, Repo, RepoPath};
29+use g1t_contracts::time::rfc3339;
30+use g1t_contracts::{new_id, system};
31+use g1t_kit::now_ms;
32+use serde::Deserialize;
33+use worker::wasm_bindgen::JsValue;
34+use worker::{D1Database, Fetcher, Result};
35+
36+/// Items marked done are kept this long, then removed.
37+pub const DONE_DAYS: u32 = 30;
38+/// No item is kept longer than this, unless it was saved.
39+pub const MAX_DAYS: u32 = 180;
40+/// Unread warnings shown ahead of everything else on the first page.
41+const MAX_RANKED: u32 = 20;
42+const MAX_TITLE: usize = 200;
43+const MAX_BODY: usize = 300;
44+
45+/// What the inbox asks about an event before deciding who is told.
46+#[derive(Debug, PartialEq, Eq)]
47+pub struct Wanted {
48+ pub repo_id: String,
49+ /// The issue or pull request, when the event names one.
50+ pub number: Option<u32>,
51+ pub comment_id: Option<String>,
52+}
53+
54+/// One person to tell, and what.
55+#[derive(Debug, PartialEq, Eq)]
56+pub struct Notice {
57+ pub username: String,
58+ pub reason: &'static str,
59+ pub severity: Severity,
60+ pub title: String,
61+ pub body: String,
62+}
63+
64+/// Who did it, as far as is known.
65+#[derive(Debug, Default)]
66+pub struct Actor {
67+ pub id: Option<String>,
68+ pub username: Option<String>,
69+}
70+
71+impl Actor {
72+ fn is(&self, person: &Principal) -> bool {
73+ self.id.as_deref().is_some_and(|id| id == person.id) || self.is_named(&person.username)
74+ }
75+
76+ fn is_named(&self, username: &str) -> bool {
77+ self.username.as_deref().is_some_and(|name| name.eq_ignore_ascii_case(username))
78+ }
79+}
80+
81+fn is_g1t(username: &str) -> bool {
82+ username.eq_ignore_ascii_case(system::USERNAME) || username.eq_ignore_ascii_case("g1t-agent")
83+}
84+
85+fn is_g1t_id(id: &str) -> bool {
86+ system::is_system_id(id) || id == AGENT_ID
87+}
88+
89+/// The issue or pull request an event names, and what to read for it.
90+/// None for events the inbox does not tell anyone of.
91+pub fn wants(event: &Event) -> Option<Wanted> {
92+ let data = &event.data;
93+ let text = |key: &str| data[key].as_str().filter(|value| !value.is_empty()).map(str::to_owned);
94+ let number = |key: &str| data[key].as_u64().and_then(|n| u32::try_from(n).ok());
95+ let repo_id = text("repoId").or_else(|| event.repo_id.clone())?;
96+ let on = |number: Option<u32>, comment_id: Option<String>| {
97+ Some(Wanted {
98+ repo_id: repo_id.clone(),
99+ number,
100+ comment_id,
101+ })
102+ };
103+ match event.kind.as_str() {
104+ "agent.asked" | "pull.merged" | "pull.ready" => on(Some(number("number")?), None),
105+ "checks.completed" => match data["status"].as_str() {
106+ Some("failed" | "errored") => on(Some(number("number")?), None),
107+ _ => None,
108+ },
109+ "review.completed" => match data["verdict"].as_str() {
110+ Some("approve" | "request_changes") => on(Some(number("number")?), None),
111+ _ => None,
112+ },
113+ "workflow.completed" => match data["conclusion"].as_str() {
114+ Some("failure") => on(number("pull"), None),
115+ _ => None,
116+ },
117+ "comment.created" => on(Some(number("number")?), Some(text("commentId")?)),
118+ _ => None,
119+ }
120+}
121+
122+/// Collects who is told, each once, never the actor and never g1t.
123+struct Told<'a> {
124+ actor: &'a Actor,
125+ notices: Vec<Notice>,
126+}
127+
128+impl Told<'_> {
129+ fn tell(&mut self, username: &str, reason: &'static str, severity: Severity, title: &str, body: &str) {
130+ let username = username.trim().trim_start_matches('@').to_lowercase();
131+ if username.is_empty()
132+ || is_g1t(&username)
133+ || self.actor.is_named(&username)
134+ || self.notices.iter().any(|told| told.username == username)
135+ {
136+ return;
137+ }
138+ self.notices.push(Notice {
139+ username,
140+ reason,
141+ severity,
142+ title: clip(title, MAX_TITLE),
143+ body: clip(body, MAX_BODY),
144+ });
145+ }
146+
147+ /// A person known by id as well as name, such as an author.
148+ fn tell_person(&mut self, person: &Principal, reason: &'static str, severity: Severity, title: &str, body: &str) {
149+ if is_g1t_id(&person.id) || self.actor.is(person) {
150+ return;
151+ }
152+ self.tell(&person.username, reason, severity, title, body);
153+ }
154+}
155+
156+fn clip(text: &str, max: usize) -> String {
157+ let text = text.trim();
158+ if text.chars().count() <= max {
159+ return text.to_owned();
160+ }
161+ let cut: String = text.chars().take(max - 1).collect();
162+ format!("{}…", cut.trim_end())
163+}
164+
165+/// Who is told of `event`, in `repo` (`owner/name`), given what it names.
166+/// `actor` is who caused it; outcomes nobody chose (checks, workflows, a
167+/// review, an agent finishing) are told whoever caused them.
168+pub fn notices(event: &Event, repo: &str, actor: &Actor, subject: Option<&InboxSubject>) -> Vec<Notice> {
169+ let nobody = Actor::default();
170+ let data = &event.data;
171+ // The issue or pull request, as titles name it: `acme/rocket#12`.
172+ let at = match data["number"].as_u64().or(data["pull"].as_u64()) {
173+ Some(number) if subject.is_some() => format!("{repo}#{number}"),
174+ _ => repo.to_owned(),
175+ };
176+ let outcome = matches!(
177+ event.kind.as_str(),
178+ "checks.completed" | "workflow.completed" | "review.completed" | "pull.ready" | "agent.asked"
179+ );
180+ let mut told = Told {
181+ actor: if outcome { &nobody } else { actor },
182+ notices: Vec::new(),
183+ };
184+
185+ match (event.kind.as_str(), subject) {
186+ ("agent.asked", Some(pull)) => {
187+ let title = format!("An agent is waiting on {}", at);
188+ told.tell_person(pull.owner(), "agent_asked", Severity::Warning, &title, &pull.title);
189+ if let Some(issue) = &pull.issue {
190+ told.tell_person(issue.owner(), "agent_asked", Severity::Warning, &title, &pull.title);
191+ for name in &issue.assignees {
192+ told.tell(name, "agent_asked", Severity::Warning, &title, &pull.title);
193+ }
194+ }
195+ }
196+ ("checks.completed", Some(pull)) => {
197+ let title = match data["status"].as_str() {
198+ Some("errored") => format!("Checks could not run on {}", at),
199+ _ => format!("Checks failed on {}", at),
200+ };
201+ told.tell_person(pull.owner(), "checks_failed", Severity::Error, &title, &pull.title);
202+ }
203+ ("workflow.completed", subject) => {
204+ let workflow = data["workflow"].as_str().filter(|name| !name.is_empty()).unwrap_or("A workflow");
205+ match subject {
206+ Some(pull) => {
207+ let title = format!("{workflow} failed on {}", at);
208+ told.tell_person(pull.owner(), "workflow_failed", Severity::Error, &title, &pull.title);
209+ }
210+ None => {
211+ // Not on a pull request: whoever pushed the commit it ran on.
212+ let branch = data["ref"].as_str().unwrap_or_default();
213+ let branch = branch.strip_prefix("refs/heads/").unwrap_or(branch);
214+ let title = format!("{workflow} failed on {branch} in {repo}");
215+ let body = format!("Run {} at {}", data["number"], short(data["sha"].as_str().unwrap_or_default()));
216+ if let Some(id) = &actor.id
217+ && !is_g1t_id(id)
218+ && let Some(name) = &actor.username
219+ {
220+ told.tell(name, "workflow_failed", Severity::Error, &title, &body);
221+ }
222+ }
223+ }
224+ }
225+ ("review.completed", Some(pull)) => {
226+ let (title, severity, reason) = match data["verdict"].as_str() {
227+ Some("approve") => (format!("g1t approved {}", at), Severity::Success, "approved"),
228+ _ => (format!("g1t asked for changes on {}", at), Severity::Info, "changes_requested"),
229+ };
230+ told.tell_person(pull.owner(), reason, severity, &title, &pull.title);
231+ }
232+ ("pull.ready", Some(pull)) => {
233+ // A change g1t made is ready: the agent's run is over.
234+ if let Some(owner) = pull.requested_by.as_ref().filter(|_| is_g1t_id(&pull.author.id) || is_g1t(&pull.author.username)) {
235+ let title = format!("g1t finished {}", at);
236+ told.tell_person(owner, "agent_finished", Severity::Success, &title, &pull.title);
237+ }
238+ }
239+ ("pull.merged", Some(pull)) => {
240+ let title = match &actor.username {
241+ Some(name) if !is_g1t(name) => format!("{name} merged {}", at),
242+ _ => format!("{} was merged", at),
243+ };
244+ told.tell_person(pull.owner(), "merged", Severity::Success, &title, &pull.title);
245+ }
246+ ("comment.created", Some(on)) => {
247+ let Some(comment) = on.comment.as_ref().filter(|comment| !comment.event) else {
248+ return Vec::new();
249+ };
250+ // Whoever wrote it is the actor, whatever the event says.
251+ let writer = Actor {
252+ id: Some(comment.author.id.clone()),
253+ username: Some(comment.author.username.clone()),
254+ };
255+ told.actor = &writer;
256+ let who = &comment.author.username;
257+ let body = if comment.excerpt.is_empty() { &on.title } else { &comment.excerpt };
258+ for name in &comment.mentions {
259+ let title = format!("{who} mentioned you on {}", at);
260+ told.tell(name, "mentioned", Severity::Info, &title, body);
261+ }
262+ let (reason, severity, title) = match comment.verdict.as_deref() {
263+ Some("approve") => ("approved", Severity::Success, format!("{who} approved {}", at)),
264+ Some("request_changes") => ("changes_requested", Severity::Info, format!("{who} asked for changes on {}", at)),
265+ _ => ("commented", Severity::Info, format!("{who} commented on {}", at)),
266+ };
267+ told.tell_person(on.owner(), reason, severity, &title, body);
268+ return told.notices;
269+ }
270+ _ => {}
271+ }
272+ told.notices
273+}
274+
275+fn short(sha: &str) -> &str {
276+ sha.get(..7).unwrap_or(sha)
277+}
278+
279+/// Where an item is on g1t.sh.
280+pub fn url(repo: Option<&str>, subject: Option<SubjectKind>, number: Option<u32>, run_id: Option<&str>) -> String {
281+ let Some(repo) = repo else {
282+ return "/inbox".to_owned();
283+ };
284+ match (subject, number, run_id) {
285+ (Some(SubjectKind::Pull), Some(number), _) => format!("/{repo}/pull/{number}"),
286+ (Some(SubjectKind::Issue), Some(number), _) => format!("/{repo}/issues/{number}"),
287+ (Some(SubjectKind::Run), _, Some(run)) => format!("/{repo}/actions/runs/{run}"),
288+ _ => format!("/{repo}"),
289+ }
290+}
291+
292+// --- Writing ---------------------------------------------------------------
293+
294+/// The services the inbox reads from as events arrive.
295+pub struct Sources<'a> {
296+ pub work: &'a Fetcher,
297+ pub repos: &'a Fetcher,
298+ pub identity: &'a Fetcher,
299+}
300+
301+/// Writes the items a batch from the bus calls for. Never fails the batch:
302+/// what cannot be worked out is logged and left.
303+pub async fn deliver(db: &D1Database, sources: &Sources<'_>, events: &[Event]) {
304+ let wanted: Vec<(&Event, Wanted)> = events
305+ .iter()
306+ .filter_map(|event| wants(event).map(|wanted| (event, wanted)))
307+ .collect();
308+ if wanted.is_empty() {
309+ return;
310+ }
311+ // Everyone who caused one, named in one call.
312+ let ids: Vec<String> = wanted
313+ .iter()
314+ .filter_map(|(event, _)| event.actor.clone())
315+ .filter(|id| !is_g1t_id(id))
316+ .collect::<HashSet<_>>()
317+ .into_iter()
318+ .collect();
319+ let names: HashMap<String, String> = if ids.is_empty() {
320+ HashMap::new()
321+ } else {
322+ g1t_kit::call(sources.identity, "usernames", &UsernamesArgs { ids })
323+ .await
324+ .unwrap_or_else(|error| {
325+ worker::console_error!("inbox: could not name who acted: {error}");
326+ HashMap::new()
327+ })
328+ };
329+ let mut paths: HashMap<String, Option<RepoPath>> = HashMap::new();
330+ for (event, wanted) in wanted {
331+ if let Err(error) = deliver_one(db, sources, &names, &mut paths, event, wanted).await {
332+ worker::console_error!("inbox: {} {} not delivered: {error}", event.kind, event.id);
333+ }
334+ }
335+}
336+
337+async fn deliver_one(
338+ db: &D1Database,
339+ sources: &Sources<'_>,
340+ names: &HashMap<String, String>,
341+ paths: &mut HashMap<String, Option<RepoPath>>,
342+ event: &Event,
343+ wanted: Wanted,
344+) -> Result<()> {
345+ if !paths.contains_key(&wanted.repo_id) {
346+ let path: Option<RepoPath> = g1t_kit::call(
347+ sources.repos,
348+ "path_by_id",
349+ &PathByIdArgs {
350+ id: wanted.repo_id.clone(),
351+ },
352+ )
353+ .await?;
354+ paths.insert(wanted.repo_id.clone(), path);
355+ }
356+ // A repository that is gone tells nobody.
357+ let Some(path) = paths.get(&wanted.repo_id).cloned().flatten() else {
358+ return Ok(());
359+ };
360+ let subject: Option<InboxSubject> = match wanted.number {
361+ Some(number) => {
362+ let found: Option<InboxSubject> = g1t_kit::call(
363+ sources.work,
364+ "inbox_subject",
365+ &InboxSubjectArgs {
366+ repo_id: wanted.repo_id.clone(),
367+ number,
368+ comment_id: wanted.comment_id.clone(),
369+ },
370+ )
371+ .await?;
372+ if found.is_none() {
373+ return Ok(());
374+ }
375+ found
376+ }
377+ None => None,
378+ };
379+ let actor = Actor {
380+ id: event.actor.clone(),
381+ username: match event.actor.as_deref() {
382+ Some(id) if is_g1t_id(id) => Some(system::USERNAME.to_owned()),
383+ Some(id) => names.get(id).map(|name| name.to_lowercase()),
384+ None => None,
385+ },
386+ };
387+ let repo = format!("{}/{}", path.namespace, path.name).to_lowercase();
388+ let told = notices(event, &repo, &actor, subject.as_ref());
389+ if told.is_empty() {
390+ return Ok(());
391+ }
392+ let (kind, number, run_id) = match (&subject, event.kind.as_str()) {
393+ (Some(subject), _) => (subject.kind, wanted.number, None),
394+ (None, "workflow.completed") => (Some(SubjectKind::Run), None, event.data["runId"].as_str()),
395+ (None, _) => (None, None, None),
396+ };
397+ let now = now_ms();
398+ let mut statements = Vec::with_capacity(told.len());
399+ for notice in told {
400+ statements.push(
401+ db.prepare(
402+ // A redelivered event finds its rows there already.
403+ "INSERT OR IGNORE INTO inbox_items (id, username, event_id, reason, severity, title, body,
404+ workspace, repo_id, repo, subject, number, run_id, actor, created_at)
405+ VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
406+ )
407+ .bind(&[
408+ new_id("ntf", now).into(),
409+ notice.username.as_str().into(),
410+ event.id.as_str().into(),
411+ notice.reason.into(),
412+ notice.severity.as_str().into(),
413+ notice.title.as_str().into(),
414+ notice.body.as_str().into(),
415+ path.namespace.to_lowercase().into(),
416+ wanted.repo_id.as_str().into(),
417+ repo.as_str().into(),
418+ kind.map_or(JsValue::NULL, |kind| kind.as_str().into()),
419+ number.map_or(JsValue::NULL, JsValue::from),
420+ run_id.map_or(JsValue::NULL, JsValue::from),
421+ actor.username.as_deref().map_or(JsValue::NULL, JsValue::from),
422+ event.time.as_str().into(),
423+ ])?,
424+ );
425+ }
426+ db.batch(statements).await?;
427+ Ok(())
428+}
429+
430+// --- Reading and changing ----------------------------------------------------
431+
432+#[derive(Deserialize)]
433+struct Row {
434+ id: String,
435+ reason: String,
436+ severity: String,
437+ title: String,
438+ body: String,
439+ workspace: Option<String>,
440+ repo_id: Option<String>,
441+ repo: Option<String>,
442+ subject: Option<String>,
443+ number: Option<f64>,
444+ run_id: Option<String>,
445+ actor: Option<String>,
446+ created_at: String,
447+ read_at: Option<String>,
448+ done_at: Option<String>,
449+ saved: f64,
450+ snoozed_until: Option<String>,
451+}
452+
453+impl Row {
454+ fn into_item(self) -> InboxItem {
455+ let subject = self.subject.as_deref().and_then(SubjectKind::parse);
456+ let number = self.number.map(|n| n as u32);
457+ InboxItem {
458+ url: url(self.repo.as_deref(), subject, number, self.run_id.as_deref()),
459+ id: self.id,
460+ reason: self.reason,
461+ severity: Severity::parse(&self.severity).unwrap_or(Severity::Info),
462+ title: self.title,
463+ body: self.body,
464+ repo: self.repo,
465+ workspace: self.workspace,
466+ subject,
467+ number,
468+ actor: self.actor,
469+ created_at: self.created_at,
470+ read_at: self.read_at,
471+ done_at: self.done_at,
472+ saved: self.saved != 0.0,
473+ snoozed_until: self.snoozed_until,
474+ }
475+ }
476+}
477+
478+const COLUMNS: &str = "id, reason, severity, title, body, workspace, repo_id, repo, subject, number, run_id, actor,
479+ created_at, read_at, done_at, saved, snoozed_until";
480+
481+/// The conditions that pick a view's items, after `username = ?1`; `?2` is now.
482+fn view_filter(view: InboxView) -> &'static str {
483+ match view {
484+ InboxView::Inbox => "done_at IS NULL AND (snoozed_until IS NULL OR snoozed_until <= ?2)",
485+ InboxView::Saved => "saved = 1",
486+ InboxView::Done => "done_at IS NOT NULL",
487+ }
488+}
489+
490+/// Whether a list ranks unread warnings first: the inbox itself, unfiltered.
491+fn ranked(a: &ListInboxArgs) -> bool {
492+ a.view == InboxView::Inbox && a.severity.is_none() && !a.unread
493+}
494+
495+pub async fn list(db: &D1Database, repos: &Fetcher, a: ListInboxArgs) -> Result<InboxPage> {
496+ let Some(viewer) = &a.viewer else {
497+ return Ok(InboxPage::default());
498+ };
499+ let username = viewer.username.to_lowercase();
500+ let now = rfc3339(now_ms());
501+ let limit = a.limit.unwrap_or(DEFAULT_INBOX_PAGE).clamp(1, MAX_INBOX_PAGE);
502+ let mut conditions = vec!["username = ?1".to_owned(), view_filter(a.view).to_owned()];
503+ let mut values: Vec<JsValue> = vec![username.as_str().into(), now.as_str().into()];
504+ if let Some(severity) = a.severity {
505+ values.push(severity.as_str().into());
506+ conditions.push(format!("severity = ?{}", values.len()));
507+ }
508+ if a.unread {
509+ conditions.push("read_at IS NULL".to_owned());
510+ }
511+ let ranked = ranked(&a);
512+ // Unread warnings lead the first page, and are left out of the rest.
513+ let leading = "severity = 'warning' AND read_at IS NULL";
514+ let mut rest = conditions.clone();
515+ if ranked {
516+ rest.push(format!("NOT ({leading})"));
517+ }
518+ if let Some(before) = &a.before {
519+ values.push(before.as_str().into());
520+ rest.push(format!("id < ?{}", values.len()));
521+ }
522+ values.push((limit + 1).into());
523+ let order = if a.view == InboxView::Done { "done_at DESC, id DESC" } else { "id DESC" };
524+ let mut statements = vec![
525+ db.prepare(format!(
526+ "SELECT {COLUMNS} FROM inbox_items WHERE {} ORDER BY {order} LIMIT ?{}",
527+ rest.join(" AND "),
528+ values.len()
529+ ))
530+ .bind(&values)?,
531+ ];
532+ if ranked && a.before.is_none() {
533+ statements.push(
534+ db.prepare(format!(
535+ "SELECT {COLUMNS} FROM inbox_items WHERE {} AND {leading} ORDER BY id DESC LIMIT ?3"
536+ , conditions.join(" AND ")
537+ ))
538+ .bind(&[username.as_str().into(), now.as_str().into(), MAX_RANKED.into()])?,
539+ );
540+ }
541+ let results = db.batch(statements).await?;
542+ let mut rows = results.first().map(|result| result.results::<Row>()).transpose()?.unwrap_or_default();
543+ let next = if rows.len() > limit as usize {
544+ rows.truncate(limit as usize);
545+ rows.last().map(|row| row.id.clone())
546+ } else {
547+ None
548+ };
549+ let mut items: Vec<Row> = results.get(1).map(|result| result.results::<Row>()).transpose()?.unwrap_or_default();
550+ items.extend(rows);
551+
552+ // What is about a repository they can no longer read is dropped.
553+ let ids: Vec<String> = items
554+ .iter()
555+ .filter_map(|row| row.repo_id.clone())
556+ .collect::<HashSet<_>>()
557+ .into_iter()
558+ .collect();
559+ if !ids.is_empty() {
560+ let readable: Vec<Repo> = g1t_kit::call(
561+ repos,
562+ "readable",
563+ &ReadableArgs {
564+ ids: ids.clone(),
565+ viewer: a.viewer.clone(),
566+ },
567+ )
568+ .await?;
569+ let readable: HashSet<String> = readable.into_iter().map(|repo| repo.id).collect();
570+ let gone: Vec<String> = ids.into_iter().filter(|id| !readable.contains(id)).collect();
571+ if !gone.is_empty() {
572+ forget(db, &username, &gone).await?;
573+ items.retain(|row| row.repo_id.as_ref().is_none_or(|id| !gone.contains(id)));
574+ }
575+ }
576+ Ok(InboxPage {
577+ items: items.into_iter().map(Row::into_item).collect(),
578+ next,
579+ })
580+}
581+
582+/// Removes a person's items about repositories they cannot read.
583+async fn forget(db: &D1Database, username: &str, repo_ids: &[String]) -> Result<()> {
584+ let marks = vec!["?"; repo_ids.len()].join(", ");
585+ let mut values: Vec<JsValue> = vec![username.into()];
586+ values.extend(repo_ids.iter().map(|id| JsValue::from(id.as_str())));
587+ db.prepare(format!("DELETE FROM inbox_items WHERE username = ? AND repo_id IN ({marks})"))
588+ .bind(&values)?
589+ .run()
590+ .await?;
591+ Ok(())
592+}
593+
594+#[derive(Deserialize)]
595+struct CountRow {
596+ severity: String,
597+ n: f64,
598+}
599+
600+pub async fn counts(db: &D1Database, a: InboxCountsArgs) -> Result<InboxCounts> {
601+ let rows = db
602+ .prepare(
603+ "SELECT severity, count(*) AS n FROM inbox_items
604+ WHERE username = ? AND read_at IS NULL AND done_at IS NULL
605+ AND (snoozed_until IS NULL OR snoozed_until <= ?)
606+ GROUP BY severity",
607+ )
608+ .bind(&[a.username.to_lowercase().into(), rfc3339(now_ms()).into()])?
609+ .all()
610+ .await?
611+ .results::<CountRow>()?;
612+ let mut counts = InboxCounts::default();
613+ for row in rows {
614+ let n = row.n as u32;
615+ counts.unread += n;
616+ match Severity::parse(&row.severity) {
617+ Some(Severity::Error) => counts.error += n,
618+ Some(Severity::Warning) => counts.warning += n,
619+ Some(Severity::Success) => counts.success += n,
620+ Some(Severity::Info) | None => counts.info += n,
621+ }
622+ }
623+ Ok(counts)
624+}
625+
626+/// The change a mark makes, as a `SET` clause; `?1` is now.
627+fn mark_change(mark: InboxMark) -> &'static str {
628+ match mark {
629+ InboxMark::Read => "read_at = COALESCE(read_at, ?1)",
630+ InboxMark::Unread => "read_at = NULL",
631+ InboxMark::Done => "done_at = ?1, read_at = COALESCE(read_at, ?1)",
632+ InboxMark::Undone => "done_at = NULL",
633+ InboxMark::Save => "saved = 1",
634+ InboxMark::Unsave => "saved = 0",
635+ InboxMark::Snooze => "snoozed_until = ?2, read_at = COALESCE(read_at, ?1)",
636+ }
637+}
638+
639+#[derive(Deserialize)]
640+struct IdRow {
641+ #[allow(dead_code)]
642+ id: String,
643+}
644+
645+/// Changes the person's own items. Returns how many changed.
646+pub async fn mark(db: &D1Database, a: MarkInboxArgs) -> Result<u32> {
647+ let now = rfc3339(now_ms());
648+ let until = match (a.mark, a.until.as_deref().map(str::trim)) {
649+ // Times compare as text (`g1t_contracts::time`); a snooze is for later.
650+ (InboxMark::Snooze, Some(until)) if until.len() == now.len() && until > now.as_str() => until.to_owned(),
651+ (InboxMark::Snooze, _) => return Ok(0),
652+ _ => String::new(),
653+ };
654+ let mut values: Vec<JsValue> = vec![now.as_str().into(), until.as_str().into(), a.username.to_lowercase().into()];
655+ let target = if a.all && a.ids.is_empty() {
656+ let mut target = "done_at IS NULL".to_owned();
657+ if let Some(severity) = a.severity {
658+ values.push(severity.as_str().into());
659+ target.push_str(&format!(" AND severity = ?{}", values.len()));
660+ }
661+ target
662+ } else {
663+ let ids: Vec<&String> = a.ids.iter().take(MAX_MARK).collect();
664+ if ids.is_empty() {
665+ return Ok(0);
666+ }
667+ let first = values.len() + 1;
668+ values.extend(ids.iter().map(|id| JsValue::from(id.as_str())));
669+ let marks: Vec<String> = (first..values.len() + 1).map(|at| format!("?{at}")).collect();
670+ format!("id IN ({})", marks.join(", "))
671+ };
672+ let changed = db
673+ .prepare(format!(
674+ "UPDATE inbox_items SET {} WHERE username = ?3 AND {target} RETURNING id",
675+ mark_change(a.mark)
676+ ))
677+ .bind(&values)?
678+ .all()
679+ .await?
680+ .results::<IdRow>()?;
681+ Ok(changed.len() as u32)
682+}
683+
684+// --- Keeping up ------------------------------------------------------------------
685+
686+/// Moves rows with renamed workspaces and repositories, and drops those
687+/// of purged repositories and deleted workspaces.
688+pub async fn follow(db: &D1Database, events: &[Event]) -> Result<()> {
689+ let mut statements = Vec::new();
690+ for event in events {
691+ let text = |key: &str| event.data[key].as_str().unwrap_or_default().to_lowercase();
692+ match event.kind.as_str() {
693+ "workspace.renamed" => {
694+ let Ok(renamed) = serde_json::from_value::<WorkspaceRenamed>(event.data.clone()) else {
695+ continue;
696+ };
697+ let (from, to) = (renamed.from.to_lowercase(), renamed.to.to_lowercase());
698+ if from == to {
699+ continue;
700+ }
701+ statements.push(
702+ db.prepare(
703+ "UPDATE inbox_items SET repo = ?2 || substr(repo, length(?1) + 1), workspace = ?2
704+ WHERE workspace = ?1",
705+ )
706+ .bind(&[from.as_str().into(), to.as_str().into()])?,
707+ );
708+ }
709+ "repo.renamed" => statements.push(
710+ db.prepare("UPDATE inbox_items SET repo = ? WHERE repo_id = ?").bind(&[
711+ format!("{}/{}", text("namespace"), text("to")).into(),
712+ text("repoId").into(),
713+ ])?,
714+ ),
715+ "repo.transferred" => statements.push(
716+ db.prepare("UPDATE inbox_items SET repo = ?, workspace = ? WHERE repo_id = ?")
717+ .bind(&[
718+ format!("{}/{}", text("to"), text("name")).into(),
719+ text("to").into(),
720+ text("repoId").into(),
721+ ])?,
722+ ),
723+ "repo.purged" => statements.push(
724+ db.prepare("DELETE FROM inbox_items WHERE repo_id = ?")
725+ .bind(&[text("repoId").into()])?,
726+ ),
727+ "workspace.deleted" => statements.push(
728+ db.prepare("DELETE FROM inbox_items WHERE workspace = ?")
729+ .bind(&[text("slug").into()])?,
730+ ),
731+ _ => {}
732+ }
733+ }
734+ if !statements.is_empty() {
735+ db.batch(statements).await?;
736+ }
737+ Ok(())
738+}
739+
740+/// Removes items done more than [`DONE_DAYS`] ago, and any not saved older
741+/// than [`MAX_DAYS`]. Returns how many went.
742+pub async fn purge(db: &D1Database, now: u64) -> Result<u32> {
743+ let done = crate::audit::keep_from(now, DONE_DAYS);
744+ let oldest = crate::audit::keep_from(now, MAX_DAYS);
745+ let results = db
746+ .batch(vec![
747+ db.prepare("DELETE FROM inbox_items WHERE saved = 0 AND done_at < ?")
748+ .bind(&[done.into()])?,
749+ db.prepare("DELETE FROM inbox_items WHERE saved = 0 AND created_at < ?")
750+ .bind(&[oldest.into()])?,
751+ ])
752+ .await?;
753+ let mut removed = 0;
754+ for result in results {
755+ removed += result.meta()?.and_then(|meta| meta.changes).unwrap_or(0) as u32;
756+ }
757+ Ok(removed)
758+}
759+
760+#[cfg(test)]
761+mod tests {
762+ use super::*;
763+ use serde_json::json;
764+
765+ fn event(kind: &str, actor: Option<&str>, data: serde_json::Value) -> Event {
766+ Event {
767+ id: "evt_1".into(),
768+ kind: kind.into(),
769+ source: "work".into(),
770+ time: "2026-10-07T12:00:00.000Z".into(),
771+ repo_id: Some("rep_1".into()),
772+ actor: actor.map(str::to_owned),
773+ data,
774+ }
775+ }
776+
777+ fn person(id: &str, username: &str) -> Principal {
778+ Principal {
779+ id: id.into(),
780+ username: username.into(),
781+ }
782+ }
783+
784+ fn actor(id: &str, username: &str) -> Actor {
785+ Actor {
786+ id: Some(id.into()),
787+ username: Some(username.into()),
788+ }
789+ }
790+
791+ /// A pull request ana opened, assigned to bo, for an issue cy filed and dee is assigned.
792+ fn pull() -> InboxSubject {
793+ InboxSubject {
794+ kind: Some(SubjectKind::Pull),
795+ title: "Add the inbox".into(),
796+ author: person("usr_ana", "ana"),
797+ assignees: vec!["bo".into()],
798+ issue: Some(Box::new(InboxSubject {
799+ kind: Some(SubjectKind::Issue),
800+ title: "An inbox".into(),
801+ author: person("usr_cy", "cy"),
802+ assignees: vec!["dee".into()],
803+ ..InboxSubject::default()
804+ })),
805+ ..InboxSubject::default()
806+ }
807+ }
808+
809+ /// A change g1t made for ana.
810+ fn g1t_pull() -> InboxSubject {
811+ InboxSubject {
812+ author: person(AGENT_ID, "g1t"),
813+ requested_by: Some(person("usr_ana", "ana")),
814+ ..pull()
815+ }
816+ }
817+
818+ fn told(notices: &[Notice]) -> Vec<(&str, &str, Severity)> {
819+ notices
820+ .iter()
821+ .map(|notice| (notice.username.as_str(), notice.reason, notice.severity))
822+ .collect()
823+ }
824+
825+ fn comment(author: Principal, body: &str, mentions: &[&str]) -> InboxSubject {
826+ InboxSubject {
827+ comment: Some(InboxComment {
828+ author,
829+ excerpt: body.into(),
830+ mentions: mentions.iter().map(|name| (*name).to_owned()).collect(),
831+ ..InboxComment::default()
832+ }),
833+ ..pull()
834+ }
835+ }
836+
837+ #[test]
838+ fn only_events_that_tell_someone_are_read() {
839+ let asked = |kind: &str, data| wants(&event(kind, None, data));
840+ assert_eq!(
841+ asked("checks.completed", json!({ "repoId": "rep_1", "number": 4, "status": "failed" })),
842+ Some(Wanted { repo_id: "rep_1".into(), number: Some(4), comment_id: None })
843+ );
844+ assert_eq!(asked("checks.completed", json!({ "repoId": "rep_1", "number": 4, "status": "passed" })), None);
845+ assert_eq!(asked("review.completed", json!({ "repoId": "rep_1", "number": 4 })), None);
846+ assert_eq!(asked("workflow.completed", json!({ "repoId": "rep_1", "conclusion": "success" })), None);
847+ assert_eq!(
848+ asked("workflow.completed", json!({ "repoId": "rep_1", "conclusion": "failure" })),
849+ Some(Wanted { repo_id: "rep_1".into(), number: None, comment_id: None })
850+ );
851+ assert_eq!(
852+ asked("comment.created", json!({ "repoId": "rep_1", "number": 2, "commentId": "cmt_1" })).map(|w| w.comment_id),
853+ Some(Some("cmt_1".into()))
854+ );
855+ assert_eq!(asked("git.push", json!({ "repoId": "rep_1" })), None);
856+ assert_eq!(asked("issue.opened", json!({ "repoId": "rep_1", "number": 1 })), None);
857+ }
858+
859+ #[test]
860+ fn a_waiting_agent_needs_the_pull_requests_and_the_issues_people_first() {
861+ let asked = event("agent.asked", Some("usr_agent"), json!({ "number": 7 }));
862+ let notices = notices(&asked, "acme/rocket", &Actor::default(), Some(&pull()));
863+ assert_eq!(
864+ told(&notices),
865+ vec![
866+ ("ana", "agent_asked", Severity::Warning),
867+ ("cy", "agent_asked", Severity::Warning),
868+ ("dee", "agent_asked", Severity::Warning),
869+ ]
870+ );
871+ assert_eq!(notices[0].title, "An agent is waiting on acme/rocket#7");
872+ assert_eq!(notices[0].body, "Add the inbox");
873+ }
874+
875+ #[test]
876+ fn failures_go_to_whoever_answers_for_the_change() {
877+ let failed = event("checks.completed", None, json!({ "number": 7, "status": "failed" }));
878+ assert_eq!(told(&notices(&failed, "acme/rocket", &Actor::default(), Some(&pull()))), vec![("ana", "checks_failed", Severity::Error)]);
879+ // g1t's change is the person's who asked for it, never g1t's.
880+ let notices_ = notices(&failed, "acme/rocket", &Actor::default(), Some(&g1t_pull()));
881+ assert_eq!(told(&notices_), vec![("ana", "checks_failed", Severity::Error)]);
882+ let errored = event("checks.completed", None, json!({ "number": 7, "status": "errored" }));
883+ assert_eq!(notices(&errored, "acme/rocket", &Actor::default(), Some(&pull()))[0].title, "Checks could not run on acme/rocket#7");
884+ }
885+
886+ #[test]
887+ fn a_workflow_that_fails_tells_its_pull_requests_owner_even_if_they_pushed() {
888+ let failed = event("workflow.completed", Some("usr_ana"), json!({ "pull": 7, "workflow": "CI", "conclusion": "failure" }));
889+ let notices_ = notices(&failed, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull()));
890+ assert_eq!(told(&notices_), vec![("ana", "workflow_failed", Severity::Error)]);
891+ assert_eq!(notices_[0].title, "CI failed on acme/rocket#7");
892+ // On a branch: whoever pushed.
893+ let pushed = event(
894+ "workflow.completed",
895+ Some("usr_bo"),
896+ json!({ "workflow": "Deploy", "conclusion": "failure", "ref": "refs/heads/main", "number": 12, "sha": "abcdef0123" }),
897+ );
898+ let notices_ = notices(&pushed, "acme/rocket", &actor("usr_bo", "bo"), None);
899+ assert_eq!(told(&notices_), vec![("bo", "workflow_failed", Severity::Error)]);
900+ assert_eq!(notices_[0].title, "Deploy failed on main in acme/rocket");
901+ assert_eq!(notices_[0].body, "Run 12 at abcdef0");
902+ // Nobody to tell when g1t pushed.
903+ assert!(notices(&pushed, "acme/rocket", &actor("g1t", "g1t"), None).is_empty());
904+ }
905+
906+ #[test]
907+ fn nobody_hears_of_what_they_did_themselves() {
908+ let merged = event("pull.merged", Some("usr_ana"), json!({ "number": 7 }));
909+ assert!(notices(&merged, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull())).is_empty());
910+ let by_bo = notices(&merged, "acme/rocket", &actor("usr_bo", "bo"), Some(&pull()));
911+ assert_eq!(told(&by_bo), vec![("ana", "merged", Severity::Success)]);
912+ assert_eq!(by_bo[0].title, "bo merged acme/rocket#7");
913+ let by_queue = notices(&merged, "acme/rocket", &Actor::default(), Some(&pull()));
914+ assert_eq!(by_queue[0].title, "acme/rocket#7 was merged");
915+ }
916+
917+ #[test]
918+ fn g1t_finishing_or_reviewing_tells_the_person_it_worked_for() {
919+ // The agent acts as the person it works for: still an outcome they hear of.
920+ let ready = event("pull.ready", Some("usr_ana"), json!({ "number": 7 }));
921+ let notices_ = notices(&ready, "acme/rocket", &actor("usr_ana", "ana"), Some(&g1t_pull()));
922+ assert_eq!(told(&notices_), vec![("ana", "agent_finished", Severity::Success)]);
923+ assert_eq!(notices_[0].title, "g1t finished acme/rocket#7");
924+ // A person's draft marked ready is not news to them.
925+ assert!(notices(&ready, "acme/rocket", &actor("usr_ana", "ana"), Some(&pull())).is_empty());
926+
927+ let approve = event("review.completed", None, json!({ "number": 7, "verdict": "approve" }));
928+ assert_eq!(told(&notices(&approve, "acme/rocket", &Actor::default(), Some(&pull()))), vec![("ana", "approved", Severity::Success)]);
929+ let changes = event("review.completed", None, json!({ "number": 7, "verdict": "request_changes" }));
930+ assert_eq!(
931+ told(&notices(&changes, "acme/rocket", &Actor::default(), Some(&pull()))),
932+ vec![("ana", "changes_requested", Severity::Info)]
933+ );
934+ }
935+
936+ #[test]
937+ fn comments_tell_those_mentioned_then_the_owner_never_the_writer() {
938+ let created = event("comment.created", Some("usr_bo"), json!({ "number": 7, "commentId": "cmt_1" }));
939+ let on = comment(person("usr_bo", "bo"), "@cy @bo have a look", &["cy", "bo", "g1t"]);
940+ let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on));
941+ assert_eq!(
942+ told(&notices_),
943+ vec![("cy", "mentioned", Severity::Info), ("ana", "commented", Severity::Info)]
944+ );
945+ assert_eq!(notices_[0].title, "bo mentioned you on acme/rocket#7");
946+ assert_eq!(notices_[0].body, "@cy @bo have a look");
947+ // Mentioned and the owner: told once, as mentioned.
948+ let on = comment(person("usr_bo", "bo"), "@ana", &["ana"]);
949+ assert_eq!(told(&notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on))), vec![("ana", "mentioned", Severity::Info)]);
950+ // The owner's own comment tells nobody but those they mention.
951+ let on = comment(person("usr_ana", "ana"), "thanks", &[]);
952+ assert!(notices(&created, "acme/rocket", &actor("usr_ana", "ana"), Some(&on)).is_empty());
953+ // Something that happened, not something written, tells nobody.
954+ let mut on = comment(person("usr_bo", "bo"), "assigned cy", &[]);
955+ on.comment.as_mut().unwrap().event = true;
956+ assert!(notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on)).is_empty());
957+ // An approval is good news.
958+ let mut on = comment(person("usr_bo", "bo"), "", &[]);
959+ on.comment.as_mut().unwrap().verdict = Some("approve".into());
960+ let notices_ = notices(&created, "acme/rocket", &actor("usr_bo", "bo"), Some(&on));
961+ assert_eq!(told(&notices_), vec![("ana", "approved", Severity::Success)]);
962+ assert_eq!(notices_[0].body, "Add the inbox");
963+ }
964+
965+ #[test]
966+ fn items_link_to_what_they_are_about() {
967+ assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Pull), Some(7), None), "/acme/rocket/pull/7");
968+ assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Issue), Some(3), None), "/acme/rocket/issues/3");
969+ assert_eq!(url(Some("acme/rocket"), Some(SubjectKind::Run), None, Some("run_1")), "/acme/rocket/actions/runs/run_1");
970+ assert_eq!(url(Some("acme/rocket"), None, None, None), "/acme/rocket");
971+ assert_eq!(url(None, None, None, None), "/inbox");
972+ }
973+
974+ #[test]
975+ fn long_titles_are_cut_to_a_line() {
976+ let long = "x".repeat(400);
977+ assert_eq!(clip(&long, MAX_TITLE).chars().count(), MAX_TITLE);
978+ assert_eq!(clip(" short ", MAX_TITLE), "short");
979+ }
980+
981+ #[test]
982+ fn marks_change_only_what_they_say() {
983+ assert_eq!(mark_change(InboxMark::Read), "read_at = COALESCE(read_at, ?1)");
984+ assert!(mark_change(InboxMark::Done).contains("done_at = ?1"));
985+ assert!(ranked(&ListInboxArgs {
986+ viewer: None,
987+ view: InboxView::Inbox,
988+ severity: None,
989+ unread: false,
990+ before: None,
991+ limit: None,
992+ }));
993+ }
994+}
+31−3
55 //! them to the log and passes each batch on to every subscriber's own
66 //! queue, so a slow or failing subscriber holds up nobody else.
77 //!
8+//! It keeps two more things beside the log: the audit log (audit.rs) and
9+//! each person's inbox (inbox.rs), written as events arrive.
10+//!
811 //! Other services reach it over `POST /rpc/<method>`; see
9−//! `g1t_contracts::events` for the methods and their arguments.
12+//! `g1t_contracts::events`, `audit` and `inbox` for the methods and their
13+//! arguments.
1014
1115 mod audit;
16+mod inbox;
1217
1318 use g1t_contracts::events::{Event, ListArgs, PublishArgs};
1419 use g1t_contracts::new_id;
132137 Ok(rows.into_iter().map(Event::from).collect())
133138 }
134139
135− /// Writes a batch from the bus to the log, then hands it to every
136− /// subscriber.
140+ /// Writes a batch from the bus to the log, hands it to every
141+ /// subscriber, then tells the people it concerns (inbox.rs).
137142 async fn deliver(&self, events: &[Event]) -> Result<()> {
138143 let mut statements = Vec::with_capacity(events.len());
139144 for event in events {
168173 };
169174 send(&js::binding(&self.env, &name)?, events).await?;
170175 }
176+ inbox::follow(&self.db, events).await?;
177+ let (work, repos, identity) = (
178+ self.env.service("WORK")?,
179+ self.env.service("REPOS")?,
180+ self.env.service("IDENTITY")?,
181+ );
182+ let sources = inbox::Sources {
183+ work: &work,
184+ repos: &repos,
185+ identity: &identity,
186+ };
187+ inbox::deliver(&self.db, &sources, events).await;
171188 Ok(())
172189 }
173190 }
187204 "list" => reply(&events.list(args(body)?).await?),
188205 "audit_record" => reply(&audit::record(&events.db, args(body)?).await?),
189206 "audit_list" => reply(&audit::list(&events.db, args(body)?).await?),
207+ "inbox_list" => {
208+ let repos = events.env.service("REPOS")?;
209+ reply(&inbox::list(&events.db, &repos, args(body)?).await?)
210+ }
211+ "inbox_counts" => reply(&inbox::counts(&events.db, args(body)?).await?),
212+ "inbox_mark" => reply(&inbox::mark(&events.db, args(body)?).await?),
190213 _ => Response::error("Unknown method", 404),
191214 }
192215 }
227250 let min_days = days("AUDIT_MIN_DAYS", audit::DEFAULT_MIN_DAYS).min(max_days);
228251 let Ok(db) = env.d1("DB") else { return };
229252 let now = now_ms();
253+ match inbox::purge(&db, now).await {
254+ Ok(removed) if removed > 0 => worker::console_log!("removed {removed} old inbox items"),
255+ Ok(_) => {}
256+ Err(error) => worker::console_error!("could not remove old inbox items: {error}"),
257+ }
230258 match audit::purge(&db, &audit::keep_from(now, max_days), 20).await {
231259 Ok(removed) if removed > 0 => {
232260 worker::console_log!("removed {removed} audit entries older than {max_days} days")
+9−1
4141 },
4242 // Each workspace's audit log keeps as many days as billing says (7 free,
4343 // 90 on the plan, or what staff set), asked for over this binding.
44− "services": [{ "binding": "BILLING", "service": "g1t-billing" }],
44+ // The inbox (src/inbox.rs) reads what each event names from work, where
45+ // its repository is from repos, and who acted from identity; listing
46+ // asks repos which repositories the person can still read.
47+ "services": [
48+ { "binding": "BILLING", "service": "g1t-billing" },
49+ { "binding": "WORK", "service": "g1t-work" },
50+ { "binding": "REPOS", "service": "g1t-repos" },
51+ { "binding": "IDENTITY", "service": "g1t-identity" }
52+ ],
4553 // Once a day, audit entries older than AUDIT_MAX_DAYS are removed for
4654 // everyone (keep it at least billing's AUDIT_MAX_DAYS), then each
4755 // workspace's older than its own days. Only workspaces with entries
+177−0
1+//! What an event names, for the inbox (`g1t_contracts::inbox`): the issue
2+//! or pull request, its people, and the comment, read straight from the
3+//! rows. The events service asks as each event arrives and decides from
4+//! the answer who is told; nobody's access is checked here, so nothing it
5+//! returns is shown to anyone but the people it names, or to those who can
6+//! read the repository.
7+
8+use g1t_contracts::credentials::Principal;
9+use g1t_contracts::inbox::{InboxComment, InboxSubject, InboxSubjectArgs, SubjectKind};
10+use g1t_contracts::is_valid_namespace;
11+use g1t_contracts::work::{Comment, CommentKind, Issue, Pull};
12+use worker::Result;
13+
14+use crate::Work;
15+use crate::mentions::spoken;
16+use crate::rows::CommentRow;
17+
18+/// The most people one comment notifies by name.
19+const MAX_MENTIONS: usize = 20;
20+/// About one line of a comment, shown under its title.
21+const EXCERPT_CHARS: usize = 140;
22+
23+/// The people a comment mentions by name in its own words (not in code or
24+/// a quote), lowercased, each once. Not an email address, a package scope
25+/// or a path, and never a reserved name such as `g1t`.
26+pub(crate) fn people_mentioned(body: &str) -> Vec<String> {
27+ let text = spoken(body);
28+ let chars: Vec<char> = text.chars().collect();
29+ let mut people: Vec<String> = Vec::new();
30+ let mut at = 0;
31+ while at < chars.len() && people.len() < MAX_MENTIONS {
32+ if chars[at] != '@' {
33+ at += 1;
34+ continue;
35+ }
36+ let starts_clean = at == 0 || {
37+ let before = chars[at - 1];
38+ !(before.is_alphanumeric() || "._%+-/\\@`=".contains(before))
39+ };
40+ let name: String = chars[at + 1..]
41+ .iter()
42+ .take_while(|c| c.is_ascii_alphanumeric() || **c == '-')
43+ .collect();
44+ let end = at + 1 + name.chars().count();
45+ // A path or a package such as `@scope/name`, or a domain.
46+ let ends_clean = match chars.get(end) {
47+ None => true,
48+ Some(c) if "_@/\\".contains(*c) => false,
49+ Some('.') => chars.get(end + 1).is_none_or(|c| !c.is_alphanumeric()),
50+ Some(_) => true,
51+ };
52+ let name = name.to_lowercase();
53+ if starts_clean && ends_clean && is_valid_namespace(&name) && !people.contains(&name) {
54+ people.push(name);
55+ }
56+ at = end.max(at + 1);
57+ }
58+ people
59+}
60+
61+/// The first line a comment says something on, cut to about a line.
62+pub(crate) fn excerpt(body: &str) -> String {
63+ let line = body
64+ .lines()
65+ .map(str::trim)
66+ .find(|line| !line.is_empty() && !line.starts_with("```") && !line.starts_with('>'))
67+ .unwrap_or_default();
68+ let line = line.split_whitespace().collect::<Vec<_>>().join(" ");
69+ if line.chars().count() <= EXCERPT_CHARS {
70+ return line;
71+ }
72+ let cut: String = line.chars().take(EXCERPT_CHARS - 1).collect();
73+ format!("{}…", cut.trim_end())
74+}
75+
76+fn principal(id: &str, username: &str) -> Principal {
77+ Principal {
78+ id: id.to_owned(),
79+ username: username.to_lowercase(),
80+ }
81+}
82+
83+fn of_issue(issue: Issue) -> InboxSubject {
84+ InboxSubject {
85+ kind: Some(SubjectKind::Issue),
86+ title: issue.title,
87+ author: principal(&issue.author.id, &issue.author.username),
88+ requested_by: issue.requested_by.map(|user| principal(&user.id, &user.username)),
89+ assignees: issue.assignees,
90+ ..InboxSubject::default()
91+ }
92+}
93+
94+fn of_pull(pull: Pull, issue: Option<Issue>) -> InboxSubject {
95+ InboxSubject {
96+ kind: Some(SubjectKind::Pull),
97+ title: pull.title,
98+ author: principal(&pull.author.id, &pull.author.username),
99+ requested_by: pull.requested_by.map(|user| principal(&user.id, &user.username)),
100+ assignees: pull.assignees,
101+ reviewers: pull.reviewers,
102+ issue: issue.map(|issue| Box::new(of_issue(issue))),
103+ comment: None,
104+ }
105+}
106+
107+fn of_comment(comment: Comment) -> InboxComment {
108+ let event = comment.kind == CommentKind::Event;
109+ InboxComment {
110+ author: principal(&comment.author.id, &comment.author.username),
111+ excerpt: excerpt(&comment.body),
112+ mentions: if event { Vec::new() } else { people_mentioned(&comment.body) },
113+ verdict: comment.verdict.map(|verdict| verdict.as_str().to_owned()),
114+ event,
115+ }
116+}
117+
118+impl Work {
119+ /// The issue or pull request numbered `number`, with its people, and
120+ /// the comment asked about. None when there is no such issue or pull
121+ /// request.
122+ pub(crate) async fn inbox_subject(&self, a: InboxSubjectArgs) -> Result<Option<InboxSubject>> {
123+ let comment = async {
124+ let Some(id) = &a.comment_id else {
125+ return Ok::<_, worker::Error>(None);
126+ };
127+ Ok(self
128+ .db
129+ .prepare("SELECT * FROM comments WHERE id = ? AND repo_id = ?")
130+ .bind(&[id.as_str().into(), a.repo_id.as_str().into()])?
131+ .first::<CommentRow>(None)
132+ .await?
133+ .map(|row| of_comment(Comment::from(row))))
134+ };
135+ let (pull, comment) = futures_util::future::try_join(self.pull(&a.repo_id, a.number), comment).await?;
136+ let mut subject = match pull {
137+ Some(pull) => {
138+ let issue = match pull.issue {
139+ Some(number) => self.issue(&a.repo_id, number).await?,
140+ None => None,
141+ };
142+ of_pull(pull, issue)
143+ }
144+ None => match self.issue(&a.repo_id, a.number).await? {
145+ Some(issue) => of_issue(issue),
146+ None => return Ok(None),
147+ },
148+ };
149+ subject.comment = comment;
150+ Ok(Some(subject))
151+ }
152+}
153+
154+#[cfg(test)]
155+mod tests {
156+ use super::*;
157+
158+ #[test]
159+ fn people_are_mentioned_by_name_in_their_own_words() {
160+ assert_eq!(people_mentioned("@ana can you look? cc @Bob-1."), vec!["ana", "bob-1"]);
161+ // Once each, never g1t, never in code or a quote.
162+ assert_eq!(people_mentioned("@ana @ana @g1t `@carl`\n> @dee said"), vec!["ana"]);
163+ // Not an address, a package, a path or a domain.
164+ assert!(people_mentioned("ops@ana.dev, @scope/pkg, a/@b, @ana.dev").is_empty());
165+ assert_eq!(people_mentioned("(@ana)"), vec!["ana"]);
166+ }
167+
168+ #[test]
169+ fn an_excerpt_is_the_first_line_said() {
170+ assert_eq!(excerpt("\n\n```\ncode\n```"), "code");
171+ assert_eq!(excerpt("> quoted\n Looks good \nmore"), "Looks good");
172+ let long = "word ".repeat(60);
173+ let cut = excerpt(&long);
174+ assert!(cut.ends_with('…'));
175+ assert_eq!(cut.chars().count(), EXCERPT_CHARS);
176+ }
177+}
+2−0
1010 mod compute;
1111 mod confidence;
1212 mod guardrails;
13+mod inbox;
1314 mod lifecycle;
1415 mod memory;
1516 mod mentions;
21342135 "remove_from_queue" => reply(&work.remove_from_queue(args(body)?).await?),
21352136 "message_agent" => reply(&work.message_agent(args(body)?).await?),
21362137 "locate_pull" => reply(&work.locate_pull(args(body)?).await?),
2138+ "inbox_subject" => reply(&work.inbox_subject(args(body)?).await?),
21372139 "answer_message" => reply(&work.answer_message(args(body)?).await?),
21382140 "take_messages" => reply(&work.take_messages(args(body)?).await?),
21392141 "wake_for_messages" => reply(&work.wake_for_messages(args(body)?).await?),