Skip to content

g1t/apps/api/src/notifications.rs

462 lines20,065 bytesCodeBlame

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

API: notifications over REST and MCP, with notifications scopes1//! Notifications: a person's inbox, its threads, their subscriptions to
2//! issues and pull requests, and how they watch repositories. The events
3//! service keeps all of it (`g1t_contracts::inbox`); this is its public
4//! shape, which follows the inbox's own: a thread is one item about one
5//! thing, brought back to the top as things happen to it.
6//!
7//! Every operation here is the person's own: a personal access token or a
8//! session, never a workspace's token or g1t's agents (which act as g1t,
9//! and g1t is never told anything).
10
11use g1t_contracts::inbox::*;
12use g1t_contracts::repos::{GetArgs, Repo, RepoPath};
13use g1t_contracts::time::{parse_rfc3339, rfc3339};
14use g1t_contracts::{FailureCode, Outcome, PrincipalKind, User, Viewer};
15use serde_json::{Value, json};
16use worker::Result;
17
18use crate::operations::{Op, Services};
19
20fn failed(code: FailureCode, message: &str) -> Result<Outcome<Value>> {
21 Ok(Outcome::fail(code, message))
22}
23
24fn ok<T: serde::Serialize>(value: &T) -> Result<Outcome<Value>> {
25 Ok(Outcome::Ok(serde_json::to_value(value)?))
26}
27
28const NO_THREAD: &str = "No such notification thread.";
29const NO_SUBJECT: &str = "Name a thread by id, or an issue or pull request by repo and number.";
30
31fn text(input: &Value, key: &str) -> Option<String> {
32 input[key].as_str().map(str::trim).filter(|value| !value.is_empty()).map(str::to_owned)
33}
34
35/// A yes or no, given as a boolean or as `true`/`false` (a query string).
36fn flag(input: &Value, key: &str) -> Option<bool> {
37 match &input[key] {
38 Value::Bool(value) => Some(*value),
39 Value::String(text) => match text.trim() {
40 "true" | "1" => Some(true),
41 "false" | "0" => Some(false),
42 _ => None,
43 },
44 _ => None,
45 }
46}
47
48fn whole(input: &Value, key: &str) -> Option<u32> {
49 match &input[key] {
50 Value::Number(number) => number.as_u64().and_then(|n| u32::try_from(n).ok()),
51 Value::String(digits) => digits.trim().parse().ok(),
52 _ => None,
53 }
54}
55
56/// A time given as RFC 3339, as the inbox stores times (to the
57/// millisecond, in UTC), or why it is not one.
58pub(crate) fn instant(input: &Value, key: &str) -> std::result::Result<Option<String>, String> {
59 match text(input, key) {
60 None => Ok(None),
61 Some(given) => parse_rfc3339(&given)
62 .map(|ms| Some(rfc3339(ms)))
63 .ok_or_else(|| format!("{key} is a time like 2026-10-07T12:00:00Z, not {given}.")),
64 }
65}
66
67/// The repository `repo` names, if the viewer may read it.
68async fn readable_repo(services: &Services, viewer: &Viewer, input: &Value) -> Result<Option<Outcome<Repo>>> {
69 let Some(path) = crate::operations::repo_path(input) else {
70 return Ok(None);
71 };
72 let found: Outcome<Repo> = g1t_kit::call(
73 &services.repos,
74 "get",
75 &GetArgs {
76 path: RepoPath {
77 namespace: path.namespace,
78 name: path.name,
79 },
80 viewer: viewer.clone(),
81 },
82 )
83 .await?;
84 Ok(Some(found))
85}
86
87/// What `list_notifications` reads from its input.
88pub(crate) fn list_args(viewer: &Viewer, input: &Value) -> std::result::Result<ListInboxArgs, String> {
89 let view = match text(input, "view") {
90 None => InboxView::Inbox,
91 Some(view) => InboxView::parse(&view).ok_or_else(|| format!("view is inbox, saved or done, not {view}."))?,
92 };
93 let reason = match text(input, "reason") {
94 None => None,
95 Some(reason) => Some(Reason::parse(&reason).ok_or_else(|| {
96 format!(
97 "{reason} is not a reason. Give one of {}.",
98 Reason::ALL.map(Reason::as_str).join(", ")
99 )
100 })?),
101 };
102 let severity = match text(input, "severity") {
103 None => None,
104 Some(severity) => Some(
105 Severity::parse(&severity).ok_or_else(|| format!("severity is error, warning, success or info, not {severity}."))?,
106 ),
107 };
108 // As it is elsewhere: what is unread, unless all is asked for. Saved
109 // and done are lists of their own, read or not.
110 let all = flag(input, "all").unwrap_or(false);
111 let unread = flag(input, "unread").unwrap_or(view == InboxView::Inbox && !all);
112 Ok(ListInboxArgs {
113 viewer: viewer.clone(),
114 view,
115 severity,
116 reason,
117 participating: flag(input, "participating").unwrap_or(false),
118 repo_id: None,
119 unread,
120 since: instant(input, "since")?,
121 updated_before: instant(input, "before")?,
122 before: text(input, "cursor"),
123 limit: whole(input, "per_page").map(|n| n.clamp(1, MAX_INBOX_PAGE)),
124 })
125}
126
127/// How a person watches a repository, as the API shows it: its level and
128/// kinds, and the plain answers to whether they get its activity and
129/// whether they ignore it.
130pub(crate) fn watching_json(watching: &Watching, repo: Option<&str>) -> Value {
131 json!({
132 "repo": repo.map(str::to_owned).or_else(|| watching.repo.clone()),
133 "level": watching.level,
134 "events": watching.events,
135 "subscribed": matches!(watching.level, WatchLevel::All | WatchLevel::Custom),
136 "ignored": watching.level == WatchLevel::Ignore,
137 "updated_at": watching.updated_at,
138 })
139}
140
141/// The level `set_repo_subscription` asks for: `level` (with `events`), or
142/// the yes-or-no of `subscribed` and `ignored`.
143pub(crate) fn watch_level(input: &Value) -> std::result::Result<(WatchLevel, Vec<String>), String> {
144 let events: Vec<String> = input["events"]
145 .as_array()
146 .map(|events| events.iter().filter_map(|event| event.as_str().map(str::to_owned)).collect())
147 .unwrap_or_default();
148 if let Some(level) = text(input, "level") {
149 let level = WatchLevel::parse(&level)
150 .ok_or_else(|| format!("level is participating, all, ignore or custom, not {level}."))?;
151 if level == WatchLevel::Custom {
152 let unknown: Vec<&String> = events
153 .iter()
154 .filter(|event| !WATCH_EVENTS.contains(&event.trim().to_lowercase().as_str()))
155 .collect();
156 if let Some(event) = unknown.first() {
157 return Err(format!("{event} is not something to watch. Give some of {}.", WATCH_EVENTS.join(", ")));
158 }
159 if events.is_empty() {
160 return Err(format!("A custom watch needs events: some of {}.", WATCH_EVENTS.join(", ")));
161 }
162 }
163 return Ok((level, events));
164 }
165 Ok(match (flag(input, "ignored"), flag(input, "subscribed")) {
166 (Some(true), _) => (WatchLevel::Ignore, Vec::new()),
167 (_, Some(false)) => (WatchLevel::Participating, Vec::new()),
168 _ => (WatchLevel::All, Vec::new()),
169 })
170}
171
172/// Which issue or pull request a subscription call names: a thread's id,
173/// or a repository and number. Checks the viewer can read the repository.
174async fn subscription_args(services: &Services, viewer: &Viewer, input: &Value) -> Result<Outcome<SubscriptionArgs>> {
175 if let Some(id) = text(input, "id") {
176 return Ok(Outcome::Ok(SubscriptionArgs {
177 viewer: viewer.clone(),
178 id: Some(id),
179 ..SubscriptionArgs::default()
180 }));
181 }
182 let (Some(found), Some(number)) = (readable_repo(services, viewer, input).await?, whole(input, "number")) else {
183 return Ok(Outcome::fail(FailureCode::Invalid, NO_SUBJECT));
184 };
185 Ok(match found {
186 Outcome::Ok(repo) => Outcome::Ok(SubscriptionArgs {
187 viewer: viewer.clone(),
188 id: None,
189 repo_id: Some(repo.id),
190 number: Some(number),
191 }),
192 Outcome::Fail(failure) => Outcome::Fail(failure),
193 })
194}
195
196/// The viewer as a person: notifications are nobody else's.
197fn person(viewer: &Viewer) -> std::result::Result<&User, &'static str> {
198 match viewer {
199 Some(user) if user.kind == PrincipalKind::User => Ok(user),
200 Some(_) => Err("Notifications are a person's own: use a personal access token, not a workspace's or an agent's."),
201 None => Err("This needs a g1t access token."),
202 }
203}
204
205/// One thread, after a change, or not found.
206async fn thread_after(services: &Services, viewer: &Viewer, id: String) -> Result<Outcome<Value>> {
207 let thread: Option<InboxThread> = g1t_kit::call(&services.events, "inbox_thread", &ThreadArgs { viewer: viewer.clone(), id }).await?;
208 match thread {
209 Some(thread) => ok(&thread),
210 None => failed(FailureCode::NotFound, NO_THREAD),
211 }
212}
213
214/// Marks one of the person's threads, then returns it as it is now.
215async fn mark_one(services: &Services, viewer: &Viewer, user: &User, input: &Value, mark: InboxMark, until: Option<String>) -> Result<Outcome<Value>> {
216 let Some(id) = text(input, "id") else {
217 return failed(FailureCode::Invalid, "Give the thread's id.");
218 };
219 // Only a thread the person can still see; one about a repository they
220 // lost is gone.
221 if let Outcome::Fail(failure) = thread_after(services, viewer, id.clone()).await? {
222 return Ok(Outcome::Fail(failure));
223 }
224 let _: u32 = g1t_kit::call(
225 &services.events,
226 "inbox_mark",
227 &MarkInboxArgs {
228 username: user.username.clone(),
229 mark,
230 ids: vec![id.clone()],
231 all: false,
232 severity: None,
233 repo_id: None,
234 last_read_at: None,
235 until,
236 },
237 )
238 .await?;
239 thread_after(services, viewer, id).await
240}
241
242/// Runs one of the notification operations.
243pub async fn run(op: Op, services: &Services, viewer: &Viewer, input: &Value) -> Result<Outcome<Value>> {
244 let user = match person(viewer) {
245 Ok(user) => user,
246 Err(message) => return failed(FailureCode::Forbidden, message),
247 };
248 let username = user.username.to_lowercase();
249 match op {
250 Op::ListNotifications => {
251 let mut args = match list_args(viewer, input) {
252 Ok(args) => args,
253 Err(message) => return failed(FailureCode::Invalid, &message),
254 };
255 if let Some(found) = readable_repo(services, viewer, input).await? {
256 match found {
257 Outcome::Ok(repo) => args.repo_id = Some(repo.id),
258 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
259 }
260 }
261 let page: InboxPage = g1t_kit::call(&services.events, "inbox_list", &args).await?;
262 ok(&page)
263 }
264 Op::MarkNotificationsRead => {
265 let last_read_at = match instant(input, "last_read_at") {
266 Ok(at) => at.unwrap_or_else(|| rfc3339(g1t_kit::now_ms())),
267 Err(message) => return failed(FailureCode::Invalid, &message),
268 };
269 let repo_id = match readable_repo(services, viewer, input).await? {
270 Some(Outcome::Ok(repo)) => Some(repo.id),
271 Some(Outcome::Fail(failure)) => return Ok(Outcome::Fail(failure)),
272 None => None,
273 };
274 let mark = if flag(input, "read") == Some(false) { InboxMark::Unread } else { InboxMark::Read };
275 let marked: u32 = g1t_kit::call(
276 &services.events,
277 "inbox_mark",
278 &MarkInboxArgs {
279 username,
280 mark,
281 ids: Vec::new(),
282 all: true,
283 severity: None,
284 repo_id,
285 last_read_at: Some(last_read_at.clone()),
286 until: None,
287 },
288 )
289 .await?;
290 ok(&json!({ "marked": marked, "last_read_at": last_read_at }))
291 }
292 Op::GetNotificationThread => match text(input, "id") {
293 Some(id) => thread_after(services, viewer, id).await,
294 None => failed(FailureCode::Invalid, "Give the thread's id."),
295 },
296 Op::MarkThreadRead => {
297 let mark = if flag(input, "read") == Some(false) { InboxMark::Unread } else { InboxMark::Read };
298 mark_one(services, viewer, user, input, mark, None).await
299 }
300 Op::MarkThreadDone => {
301 let mark = if flag(input, "done") == Some(false) { InboxMark::Undone } else { InboxMark::Done };
302 mark_one(services, viewer, user, input, mark, None).await
303 }
304 Op::SaveThread => {
305 let mark = if flag(input, "saved") == Some(false) { InboxMark::Unsave } else { InboxMark::Save };
306 mark_one(services, viewer, user, input, mark, None).await
307 }
308 Op::SnoozeThread => {
309 let until = match instant(input, "until") {
310 Ok(until) => until,
311 Err(message) => return failed(FailureCode::Invalid, &message),
312 };
313 match until {
314 Some(until) if until.as_str() <= rfc3339(g1t_kit::now_ms()).as_str() => {
315 failed(FailureCode::Invalid, "until is a time to come; leave it out to bring the thread back now.")
316 }
317 Some(until) => mark_one(services, viewer, user, input, InboxMark::Snooze, Some(until)).await,
318 None => mark_one(services, viewer, user, input, InboxMark::Unsnooze, None).await,
319 }
320 }
321 Op::GetThreadSubscription | Op::SetThreadSubscription | Op::DeleteThreadSubscription => {
322 let on = match subscription_args(services, viewer, input).await? {
323 Outcome::Ok(on) => on,
324 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
325 };
326 let found: Option<ThreadSubscription> = match op {
327 Op::GetThreadSubscription => g1t_kit::call(&services.events, "inbox_subscription", &on).await?,
328 _ => {
329 let (subscribed, ignored) = match op {
330 Op::SetThreadSubscription => (Some(flag(input, "subscribed").unwrap_or(true)), flag(input, "ignored").unwrap_or(false)),
331 _ => (Some(false), false),
332 };
333 g1t_kit::call(&services.events, "inbox_subscribe", &SubscribeArgs { on, subscribed, ignored }).await?
334 }
335 };
336 match found {
337 Some(subscription) => ok(&subscription),
338 None => failed(FailureCode::NotFound, "No such issue or pull request, or it is a thread nobody subscribes to."),
339 }
340 }
341 Op::GetRepoSubscription | Op::SetRepoSubscription | Op::DeleteRepoSubscription => {
342 let repo = match readable_repo(services, viewer, input).await? {
343 Some(Outcome::Ok(repo)) => repo,
344 Some(Outcome::Fail(failure)) => return Ok(Outcome::Fail(failure)),
345 None => return failed(FailureCode::Invalid, "Give the repository as \"owner/name\"."),
346 };
347 let path = format!("{}/{}", repo.namespace, repo.name);
348 let watching: Watching = match op {
349 Op::GetRepoSubscription => {
350 g1t_kit::call(&services.events, "inbox_watching", &WatchingArgs { username, repo_id: repo.id }).await?
351 }
352 _ => {
353 let (level, events) = match op {
354 Op::SetRepoSubscription => match watch_level(input) {
355 Ok((level, events)) => (Some(level), events),
356 Err(message) => return failed(FailureCode::Invalid, &message),
357 },
358 _ => (None, Vec::new()),
359 };
360 g1t_kit::call(
361 &services.events,
362 "inbox_watch",
363 &WatchArgs {
364 username,
365 repo_id: repo.id,
366 repo: Some(path.clone()),
367 level,
368 events,
369 },
370 )
371 .await?
372 }
373 };
374 Ok(Outcome::Ok(watching_json(&watching, Some(&path))))
375 }
376 Op::ListWatchedRepos => {
377 let watched: Vec<Watching> = g1t_kit::call(&services.events, "inbox_watched", &InboxCountsArgs { username }).await?;
378 Ok(Outcome::Ok(Value::Array(watched.iter().map(|watching| watching_json(watching, None)).collect())))
379 }
380 _ => failed(FailureCode::NotFound, "No such endpoint."),
381 }
382}
383
384#[cfg(test)]
385mod tests {
386 use super::*;
387
388 fn viewer() -> Viewer {
389 Some(User {
390 id: "usr_1".into(),
391 username: "ana".into(),
392 ..User::default()
393 })
394 }
395
396 #[test]
397 fn a_list_shows_what_is_unread_unless_all_is_asked_for() {
398 let args = list_args(&viewer(), &json!({})).unwrap();
399 assert!(args.unread);
400 assert_eq!(args.view, InboxView::Inbox);
401 let args = list_args(&viewer(), &json!({ "all": "true" })).unwrap();
402 assert!(!args.unread);
403 // Saved and done show everything in them.
404 assert!(!list_args(&viewer(), &json!({ "view": "done" })).unwrap().unread);
405 let args = list_args(
406 &viewer(),
407 &json!({ "reason": "review_requested", "participating": true, "since": "2026-10-01T00:00:00Z", "before": "2026-10-07T00:00:00Z", "cursor": "ntf_9", "per_page": "500" }),
408 )
409 .unwrap();
410 assert_eq!(args.reason, Some(Reason::ReviewRequested));
411 assert!(args.participating);
412 assert_eq!(args.since.as_deref(), Some("2026-10-01T00:00:00.000Z"));
413 assert_eq!(args.updated_before.as_deref(), Some("2026-10-07T00:00:00.000Z"));
414 assert_eq!(args.before.as_deref(), Some("ntf_9"));
415 assert_eq!(args.limit, Some(MAX_INBOX_PAGE));
416 }
417
418 #[test]
419 fn a_list_names_what_it_cannot_read() {
420 assert!(list_args(&viewer(), &json!({ "reason": "gossip" })).unwrap_err().contains("not a reason"));
421 assert!(list_args(&viewer(), &json!({ "view": "archive" })).unwrap_err().contains("inbox, saved or done"));
422 assert!(list_args(&viewer(), &json!({ "since": "yesterday" })).unwrap_err().contains("since is a time"));
423 }
424
425 #[test]
426 fn watching_is_a_level_or_a_yes_or_no() {
427 assert_eq!(watch_level(&json!({ "level": "all" })).unwrap().0, WatchLevel::All);
428 assert_eq!(watch_level(&json!({ "ignored": true })).unwrap().0, WatchLevel::Ignore);
429 assert_eq!(watch_level(&json!({ "subscribed": false })).unwrap().0, WatchLevel::Participating);
430 assert_eq!(watch_level(&json!({})).unwrap().0, WatchLevel::All);
431 let (level, events) = watch_level(&json!({ "level": "custom", "events": ["pulls", "deployments"] })).unwrap();
432 assert_eq!((level, events.len()), (WatchLevel::Custom, 2));
433 assert!(watch_level(&json!({ "level": "custom" })).unwrap_err().contains("needs events"));
434 assert!(watch_level(&json!({ "level": "custom", "events": ["releases"] })).unwrap_err().contains("releases"));
435 assert!(watch_level(&json!({ "level": "loud" })).is_err());
436 }
437
438 #[test]
439 fn watching_says_plainly_whether_activity_comes() {
440 let custom = Watching {
441 repo_id: "rep_1".into(),
442 repo: None,
443 level: WatchLevel::Custom,
444 events: vec!["pulls".into()],
445 updated_at: None,
446 };
447 let shown = watching_json(&custom, Some("acme/rocket"));
448 assert_eq!(shown["repo"], "acme/rocket");
449 assert_eq!(shown["level"], "custom");
450 assert_eq!(shown["subscribed"], true);
451 assert_eq!(shown["ignored"], false);
452 assert!(shown.get("repo_id").is_none());
453 }
454
455 #[test]
456 fn notifications_are_a_persons_own() {
457 assert!(person(&viewer()).is_ok());
458 let workspace = Some(User { kind: PrincipalKind::Workspace, ..User::default() });
459 assert!(person(&workspace).unwrap_err().contains("person's own"));
460 assert!(person(&None).is_err());
461 }
462}