1use std::path::{Path, PathBuf};
38
39use anyhow::{Context, Result, bail};
40use jiff::Timestamp;
41use serde::{Deserialize, Serialize};
42
43pub const SCHEMA: u32 = 1;
46
47pub const CAP: usize = 200;
50
51#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
53#[serde(rename_all = "lowercase")]
54pub enum Severity {
55 Info,
57 Warn,
59 Error,
61}
62
63#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
65#[serde(tag = "kind", rename_all = "lowercase")]
66pub enum Link {
67 Run {
69 id: String,
71 },
72 Task {
74 id: String,
76 },
77 Url {
79 url: String,
81 },
82}
83
84#[derive(Debug, Clone, Serialize, Deserialize)]
86pub struct Notice {
87 pub id: String,
89 pub key: String,
91 pub severity: Severity,
93 pub message: String,
95 #[serde(default)]
97 pub link: Option<Link>,
98 pub first_at: Timestamp,
100 pub last_at: Timestamp,
102 #[serde(default = "one")]
104 pub count: u32,
105 #[serde(default)]
107 pub read_at: Option<Timestamp>,
108 #[serde(default)]
110 pub dismissed_at: Option<Timestamp>,
111 #[serde(default = "schema")]
113 pub schema: u32,
114}
115
116fn one() -> u32 {
117 1
118}
119
120fn schema() -> u32 {
121 SCHEMA
122}
123
124impl Notice {
125 pub fn new(severity: Severity, key: &str, message: impl Into<String>) -> Self {
127 let now = Timestamp::now();
128 Self {
129 id: id_of(key),
130 key: key.to_owned(),
131 severity,
132 message: message.into(),
133 link: None,
134 first_at: now,
135 last_at: now,
136 count: 1,
137 read_at: None,
138 dismissed_at: None,
139 schema: SCHEMA,
140 }
141 }
142
143 pub fn info(key: &str, message: impl Into<String>) -> Self {
145 Self::new(Severity::Info, key, message)
146 }
147
148 pub fn warn(key: &str, message: impl Into<String>) -> Self {
150 Self::new(Severity::Warn, key, message)
151 }
152
153 pub fn error(key: &str, message: impl Into<String>) -> Self {
155 Self::new(Severity::Error, key, message)
156 }
157
158 pub fn link(mut self, link: Link) -> Self {
160 self.link = Some(link);
161 self
162 }
163
164 pub fn unread(&self) -> bool {
166 self.read_at.is_none() && self.dismissed_at.is_none()
167 }
168
169 pub fn raise_again(&mut self, again: &Notice, now: Timestamp) {
172 self.count = self.count.saturating_add(1);
173 self.last_at = now;
174 if again.link.is_some() {
175 self.link = again.link.clone();
176 }
177 if again.message != self.message || again.severity > self.severity {
178 self.message = again.message.clone();
179 self.severity = self.severity.max(again.severity);
180 self.read_at = None;
181 self.dismissed_at = None;
182 }
183 }
184
185 pub fn mark_read(&mut self, now: Timestamp) {
187 if self.read_at.is_none() {
188 self.read_at = Some(now);
189 }
190 }
191
192 pub fn dismiss(&mut self, now: Timestamp) {
194 self.mark_read(now);
195 if self.dismissed_at.is_none() {
196 self.dismissed_at = Some(now);
197 }
198 }
199
200 fn keep_rank(&self) -> u8 {
202 if self.dismissed_at.is_some() {
203 0
204 } else if self.read_at.is_some() {
205 1
206 } else {
207 2
208 }
209 }
210}
211
212pub fn id_of(key: &str) -> String {
216 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
217 for b in key.bytes() {
218 hash ^= u64::from(b);
219 hash = hash.wrapping_mul(0x0100_0000_01b3);
220 }
221 let slug: String = key
222 .chars()
223 .map(|c| {
224 if c.is_ascii_alphanumeric() {
225 c.to_ascii_lowercase()
226 } else {
227 '-'
228 }
229 })
230 .take(32)
231 .collect();
232 format!("{slug}-{hash:016x}")
233}
234
235fn valid_id(id: &str) -> bool {
238 !id.is_empty()
239 && id.len() <= 64
240 && id
241 .bytes()
242 .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-')
243}
244
245const LOCK_WAIT: std::time::Duration = std::time::Duration::from_secs(5);
247
248const LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
250
251static TMP_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
253
254#[derive(Debug, Clone)]
256pub struct Notices {
257 root: PathBuf,
258}
259
260impl Notices {
261 pub fn open() -> Self {
263 Self::at(crate::run::home().join("notifications"))
264 }
265
266 pub fn at(root: PathBuf) -> Self {
268 Self { root }
269 }
270
271 pub fn root(&self) -> &Path {
273 &self.root
274 }
275
276 fn path_of(&self, id: &str) -> PathBuf {
277 self.root.join(format!("{id}.json"))
278 }
279
280 fn put(&self, n: &Notice) -> Result<()> {
281 std::fs::create_dir_all(&self.root)
282 .with_context(|| format!("create {}", self.root.display()))?;
283 let body = serde_json::to_string_pretty(n).context("serialize notice")?;
284 let path = self.path_of(&n.id);
285 let tmp = path.with_extension(format!(
288 "json.{}.{}.tmp",
289 std::process::id(),
290 TMP_SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
291 ));
292 std::fs::write(&tmp, body).with_context(|| format!("write {}", tmp.display()))?;
293 if let Err(e) = std::fs::rename(&tmp, &path) {
294 let _ = std::fs::remove_file(&tmp);
295 return Err(e).with_context(|| format!("replace {}", path.display()));
296 }
297 Ok(())
298 }
299
300 pub fn get(&self, id: &str) -> Result<Notice> {
302 if !valid_id(id) {
303 bail!("`{id}` is not a notification id");
304 }
305 read_path(&self.path_of(id))
306 }
307
308 fn locked<T>(&self, id: &str, f: impl FnOnce() -> Result<T>) -> Result<T> {
317 std::fs::create_dir_all(&self.root)
318 .with_context(|| format!("create {}", self.root.display()))?;
319 let lock = self.root.join(format!("{id}.lock"));
320 let start = std::time::Instant::now();
321 let mut held = false;
322 while start.elapsed() < LOCK_WAIT {
323 match std::fs::OpenOptions::new()
324 .write(true)
325 .create_new(true)
326 .open(&lock)
327 {
328 Ok(_) => {
329 held = true;
330 break;
331 }
332 Err(_) => {
333 let stale = std::fs::metadata(&lock)
334 .and_then(|m| m.modified())
335 .ok()
336 .and_then(|t| t.elapsed().ok())
337 .is_some_and(|age| age > LOCK_STALE);
338 if stale {
339 let _ = std::fs::remove_file(&lock);
340 } else {
341 std::thread::sleep(std::time::Duration::from_millis(5));
342 }
343 }
344 }
345 }
346 let out = f();
347 if held {
348 let _ = std::fs::remove_file(&lock);
349 }
350 out
351 }
352
353 pub fn raise(&self, incoming: Notice) -> Result<Notice> {
355 let stored = self.locked(&incoming.id.clone(), || {
356 let now = Timestamp::now();
357 let stored = match read_path(&self.path_of(&incoming.id)) {
358 Ok(mut existing) => {
359 existing.raise_again(&incoming, now);
360 existing
361 }
362 Err(_) => incoming,
363 };
364 self.put(&stored)?;
365 Ok(stored)
366 })?;
367 self.prune();
368 Ok(stored)
369 }
370
371 fn update(&self, id: &str, change: impl FnOnce(&mut Notice)) -> Result<Notice> {
373 if !valid_id(id) {
374 bail!("`{id}` is not a notification id");
375 }
376 self.locked(id, || {
377 let mut n = read_path(&self.path_of(id))?;
378 change(&mut n);
379 self.put(&n)?;
380 Ok(n)
381 })
382 }
383
384 fn all(&self) -> Vec<Notice> {
386 let mut all: Vec<Notice> = std::fs::read_dir(&self.root)
387 .into_iter()
388 .flatten()
389 .flatten()
390 .map(|e| e.path())
391 .filter(|p| p.extension().is_some_and(|x| x == "json"))
392 .filter_map(|p| read_path(&p).ok())
393 .collect();
394 all.sort_by(|a, b| b.last_at.cmp(&a.last_at).then_with(|| a.id.cmp(&b.id)));
395 all
396 }
397
398 pub fn list(&self) -> Vec<Notice> {
400 self.all()
401 .into_iter()
402 .filter(|n| n.dismissed_at.is_none())
403 .collect()
404 }
405
406 pub fn count_unread(&self) -> usize {
408 self.all().iter().filter(|n| n.unread()).count()
409 }
410
411 pub fn mark_read(&self, id: &str) -> Result<Notice> {
413 let now = Timestamp::now();
414 self.update(id, |n| n.mark_read(now))
415 }
416
417 pub fn mark_all_read(&self) -> Result<usize> {
419 let now = Timestamp::now();
420 let mut changed = 0;
421 for n in self.all().into_iter().filter(Notice::unread) {
422 let done = self.update(&n.id, |n| n.mark_read(now));
425 if done.is_ok() {
426 changed += 1;
427 }
428 }
429 Ok(changed)
430 }
431
432 pub fn dismiss(&self, id: &str) -> Result<Notice> {
434 let now = Timestamp::now();
435 self.update(id, |n| n.dismiss(now))
436 }
437
438 fn prune(&self) {
440 let mut all = self.all();
441 if all.len() <= CAP {
442 return;
443 }
444 all.sort_by(|a, b| {
446 b.keep_rank()
447 .cmp(&a.keep_rank())
448 .then_with(|| b.last_at.cmp(&a.last_at))
449 });
450 for n in all.split_off(CAP) {
451 let _ = std::fs::remove_file(self.path_of(&n.id));
452 }
453 }
454
455 pub fn revision(&self) -> u64 {
459 use std::hash::{Hash as _, Hasher as _};
460 let mut entries: Vec<(std::ffi::OsString, u128)> = std::fs::read_dir(&self.root)
461 .into_iter()
462 .flatten()
463 .flatten()
464 .filter(|e| e.path().extension().is_some_and(|x| x == "json"))
465 .map(|e| {
466 let at = e
467 .metadata()
468 .and_then(|m| m.modified())
469 .ok()
470 .and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
471 .map_or(0, |d| d.as_nanos());
472 (e.file_name(), at)
473 })
474 .collect();
475 if entries.is_empty() {
476 return 0;
477 }
478 entries.sort();
479 let mut h = std::collections::hash_map::DefaultHasher::new();
480 entries.hash(&mut h);
481 h.finish()
482 }
483}
484
485fn read_path(path: &Path) -> Result<Notice> {
486 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
487 let n: Notice =
488 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
489 if n.schema > SCHEMA {
490 bail!(
491 "notice {} was written by a newer magi (schema {}, this build speaks up to {SCHEMA})",
492 n.id,
493 n.schema
494 );
495 }
496 Ok(n)
497}
498
499pub fn run_ended(state: &crate::run::RunState) -> Option<Notice> {
505 use crate::run::RunStatus;
506 let failed = matches!(
507 state.status,
508 RunStatus::Blocked | RunStatus::Stalled | RunStatus::Failed
509 );
510 (failed && !state.parked).then(|| {
511 Notice::error(
512 &format!("run:{}", state.id),
513 format!("Run {} ended {}.", state.short(), state.status.as_str()),
514 )
515 .link(Link::Run {
516 id: state.id.clone(),
517 })
518 })
519}
520
521pub fn merged_red(run_id: &str, summary: &str) -> Notice {
527 Notice::warn(&format!("merged-red:{run_id}"), summary).link(Link::Run {
528 id: run_id.to_owned(),
529 })
530}
531
532pub fn task_held(task: &crate::queue::Task) -> Option<Notice> {
546 use crate::queue::{HoldSource, TaskStatus};
547 if task.status != TaskStatus::Held || task.hold_source != Some(HoldSource::Machine) {
548 return None;
549 }
550 let why = task
551 .hold_reason
552 .as_deref()
553 .or(task.last_error.as_deref())
554 .unwrap_or("no reason recorded");
555 Some(
556 Notice::warn(
557 &format!("task:{}", task.id),
558 format!("Task {} is held: {why}.", task.short()),
559 )
560 .link(Link::Task {
561 id: task.id.clone(),
562 }),
563 )
564}
565
566pub fn run_stopped(id: &str, state: &crate::run::RunState) -> Notice {
568 Notice::error(
569 &format!("run:{id}"),
570 format!("Run {} stopped with an error.", state.short()),
571 )
572 .link(Link::Run { id: id.to_owned() })
573}
574
575pub fn raise(notice: Notice) {
579 if let Some(home) = crate::run::try_home() {
580 raise_in(&home, notice);
581 }
582}
583
584pub fn raise_in(home: &Path, notice: Notice) {
587 if let Err(e) = Notices::at(home.join("notifications")).raise(notice) {
588 tracing::warn!("could not file a notification: {e:#}");
589 }
590}
591
592#[cfg(test)]
593mod tests {
594 use super::*;
595
596 fn store() -> (tempfile::TempDir, Notices) {
597 let dir = tempfile::tempdir().unwrap();
598 let s = Notices::at(dir.path().join("notifications"));
599 (dir, s)
600 }
601
602 #[test]
603 fn a_red_merge_notice_is_keyed_on_the_run_and_raised_once() {
604 let (_d, s) = store();
605 let n = merged_red("run-1", "Merged a/b PR #1 with red checks: x (u)");
606 assert_eq!(n.key, "merged-red:run-1");
607 assert_eq!(
608 n.link,
609 Some(Link::Run {
610 id: "run-1".to_owned()
611 })
612 );
613 s.raise(merged_red(
614 "run-1",
615 "Merged a/b PR #1 with red checks: x (u)",
616 ))
617 .unwrap();
618 s.raise(merged_red(
619 "run-1",
620 "Merged a/b PR #1 with red checks: x (u)",
621 ))
622 .unwrap();
623 assert_eq!(s.list().len(), 1);
624 }
625
626 #[test]
627 fn a_repeat_counts_but_stays_read_and_a_change_relights_it() {
628 let now = Timestamp::now();
629 let mut n = Notice::warn("task:1", "held");
630 n.mark_read(now);
631 n.raise_again(&Notice::warn("task:1", "held"), now);
632 assert_eq!(n.count, 2);
633 assert!(n.read_at.is_some(), "identical repeat stays read");
634 n.raise_again(&Notice::warn("task:1", "held differently"), now);
635 assert!(n.unread());
636 assert_eq!(n.message, "held differently");
637 }
638
639 #[test]
640 fn a_rise_in_severity_resurrects_a_dismissed_notice_but_a_repeat_does_not() {
641 let now = Timestamp::now();
642 let mut n = Notice::warn("k", "m");
643 n.dismiss(now);
644 n.raise_again(&Notice::warn("k", "m"), now);
645 assert!(n.dismissed_at.is_some(), "tombstone holds");
646 n.raise_again(&Notice::error("k", "m"), now);
647 assert!(n.unread());
648 assert_eq!(n.severity, Severity::Error);
649 n.raise_again(&Notice::info("k", "other"), now);
651 assert_eq!(n.severity, Severity::Error);
652 }
653
654 #[test]
655 fn identical_raises_share_one_file() {
656 let (_d, s) = store();
657 for _ in 0..5 {
658 s.raise(Notice::error("run:abc", "blocked")).unwrap();
659 }
660 let all = s.list();
661 assert_eq!(all.len(), 1);
662 assert_eq!(all[0].count, 5);
663 assert_eq!(s.count_unread(), 1);
664 }
665
666 #[test]
667 fn transitions_persist_and_dismissed_leave_the_list() {
668 let (_d, s) = store();
669 let a = s.raise(Notice::info("a", "one")).unwrap();
670 let b = s.raise(Notice::warn("b", "two")).unwrap();
671 assert_eq!(s.count_unread(), 2);
672 s.mark_read(&a.id).unwrap();
673 assert_eq!(s.count_unread(), 1);
674 s.dismiss(&b.id).unwrap();
675 assert_eq!(s.count_unread(), 0);
676 assert_eq!(s.list().len(), 1);
677 s.raise(Notice::warn("b", "two")).unwrap();
678 assert_eq!(
679 s.list().len(),
680 1,
681 "a tombstone survives a same-message raise"
682 );
683 s.raise(Notice::info("c", "three")).unwrap();
684 assert_eq!(s.mark_all_read().unwrap(), 1);
685 assert_eq!(s.count_unread(), 0);
686 }
687
688 #[test]
689 fn the_list_is_newest_first() {
690 let (_d, s) = store();
691 let mut old = Notice::info("old", "old");
692 old.last_at = "2020-01-01T00:00:00Z".parse().unwrap();
693 s.put(&old).unwrap();
694 s.raise(Notice::info("new", "new")).unwrap();
695 let keys: Vec<_> = s.list().into_iter().map(|n| n.key).collect();
696 assert_eq!(keys, ["new", "old"]);
697 }
698
699 #[test]
700 fn the_cap_prunes_dismissed_then_read_then_oldest() {
701 let (_d, s) = store();
702 let keep = s.raise(Notice::error("keep", "unread")).unwrap();
703 let gone = s.raise(Notice::info("gone", "dismissed")).unwrap();
704 s.dismiss(&gone.id).unwrap();
705 for i in 0..CAP - 1 {
706 s.raise(Notice::info(&format!("k{i}"), "x")).unwrap();
707 }
708 let files = std::fs::read_dir(s.root()).unwrap().count();
709 assert_eq!(files, CAP);
710 assert!(
711 s.get(&keep.id).is_ok(),
712 "an unread notice outlives a dismissed one"
713 );
714 assert!(s.get(&gone.id).is_err(), "the dismissed one went first");
715 }
716
717 #[test]
718 fn revision_moves_on_write_and_on_removal() {
719 let (_d, s) = store();
720 assert_eq!(s.revision(), 0);
721 let a = s.raise(Notice::info("a", "m")).unwrap();
722 let r1 = s.revision();
723 assert_ne!(r1, 0);
724 s.raise(Notice::info("b", "m")).unwrap();
725 let r2 = s.revision();
726 assert_ne!(r1, r2);
727 std::fs::remove_file(s.path_of(&a.id)).unwrap();
728 assert_ne!(s.revision(), r2);
729 }
730
731 #[test]
732 fn ids_are_stable_and_untrusted_ids_are_refused() {
733 assert_eq!(id_of("run:1"), id_of("run:1"));
734 assert_ne!(id_of("run:1"), id_of("run-1"));
735 assert!(valid_id(&id_of("release-bump:20260101-abc")));
736 let (_d, s) = store();
737 for bad in ["", "../x", "a/b", "A", "x.json"] {
738 assert!(s.get(bad).is_err(), "{bad}");
739 }
740 }
741
742 #[test]
743 fn a_newer_schema_is_refused_and_an_older_reads() {
744 let (_d, s) = store();
745 let mut n = Notice::info("k", "m");
746 n.schema = SCHEMA + 1;
747 s.put(&n).unwrap();
748 assert!(s.get(&n.id).is_err());
749 let old = r#"{"id":"x-1","key":"x","severity":"warn","message":"m",
750 "first_at":"2020-01-01T00:00:00Z","last_at":"2020-01-01T00:00:00Z"}"#;
751 std::fs::write(s.path_of("x-1"), old).unwrap();
752 assert_eq!(s.get("x-1").unwrap().count, 1);
753 }
754
755 fn state(status: crate::run::RunStatus) -> crate::run::RunState {
756 let mut st = crate::run::RunState::new(
757 std::path::PathBuf::from("/repo"),
758 "main".to_owned(),
759 "abc1234def".to_owned(),
760 "task".to_owned(),
761 crate::config::Config::default(),
762 );
763 st.status = status;
764 st
765 }
766
767 #[test]
768 fn a_run_that_ended_badly_is_news_unless_it_only_parked() {
769 use crate::run::RunStatus;
770 for bad in [RunStatus::Blocked, RunStatus::Stalled, RunStatus::Failed] {
771 let st = state(bad);
772 let n = run_ended(&st).expect("news");
773 assert_eq!(n.key, format!("run:{}", st.id));
774 assert_eq!(n.severity, Severity::Error);
775 }
776 assert!(run_ended(&state(RunStatus::Merged)).is_none());
777 let mut parked = state(RunStatus::Stalled);
778 parked.parked = true;
779 assert!(run_ended(&parked).is_none());
780 }
781
782 #[test]
783 fn concurrent_writers_of_one_key_all_succeed() {
784 let (_d, s) = store();
785 let handles: Vec<_> = (0..8)
786 .map(|_| {
787 let s = s.clone();
788 std::thread::spawn(move || {
789 for _ in 0..20 {
790 s.raise(Notice::warn("same", "m")).unwrap();
791 }
792 })
793 })
794 .collect();
795 for h in handles {
796 h.join().unwrap();
797 }
798 assert_eq!(s.list().len(), 1);
799 let stray = std::fs::read_dir(s.root())
800 .unwrap()
801 .flatten()
802 .filter(|e| e.path().extension().is_some_and(|x| x == "tmp"))
803 .count();
804 assert_eq!(stray, 0);
805 }
806
807 #[test]
808 fn a_dead_holders_lock_is_broken_and_no_lock_is_left_behind() {
809 let (_d, s) = store();
810 let n = s.raise(Notice::info("k", "m")).unwrap();
811 let lock = s.root().join(format!("{}.lock", n.id));
812 std::fs::write(&lock, "").unwrap();
813 let old = std::time::SystemTime::now() - std::time::Duration::from_secs(60);
814 std::fs::File::options()
815 .write(true)
816 .open(&lock)
817 .unwrap()
818 .set_modified(old)
819 .unwrap();
820 s.mark_read(&n.id).unwrap();
821 assert!(!lock.exists());
822 assert!(s.get(&n.id).unwrap().read_at.is_some());
823 }
824
825 #[test]
826 fn only_a_machine_hold_is_a_task_notice() {
827 let mut t = crate::queue::Task::new(
828 "t".to_owned(),
829 "do it".to_owned(),
830 std::path::PathBuf::from("/repo"),
831 crate::queue::Source::Human,
832 );
833 assert!(task_held(&t).is_none());
834 t.hold_manual(None);
835 assert!(task_held(&t).is_none(), "the operator's own hold");
836 t.hold_machine(Some("missing blocker".to_owned()));
837 let n = task_held(&t).expect("machine hold");
838 assert_eq!(n.key, format!("task:{}", t.id));
839 assert!(n.message.contains("missing blocker"));
840 }
841}