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 task_held(task: &crate::queue::Task) -> Option<Notice> {
535 use crate::queue::{HoldSource, TaskStatus};
536 if task.status != TaskStatus::Held || task.hold_source != Some(HoldSource::Machine) {
537 return None;
538 }
539 let why = task
540 .hold_reason
541 .as_deref()
542 .or(task.last_error.as_deref())
543 .unwrap_or("no reason recorded");
544 Some(
545 Notice::warn(
546 &format!("task:{}", task.id),
547 format!("Task {} is held: {why}.", task.short()),
548 )
549 .link(Link::Task {
550 id: task.id.clone(),
551 }),
552 )
553}
554
555pub fn run_stopped(id: &str, state: &crate::run::RunState) -> Notice {
557 Notice::error(
558 &format!("run:{id}"),
559 format!("Run {} stopped with an error.", state.short()),
560 )
561 .link(Link::Run { id: id.to_owned() })
562}
563
564pub fn raise(notice: Notice) {
568 if let Some(home) = crate::run::try_home() {
569 raise_in(&home, notice);
570 }
571}
572
573pub fn raise_in(home: &Path, notice: Notice) {
576 if let Err(e) = Notices::at(home.join("notifications")).raise(notice) {
577 tracing::warn!("could not file a notification: {e:#}");
578 }
579}
580
581#[cfg(test)]
582mod tests {
583 use super::*;
584
585 fn store() -> (tempfile::TempDir, Notices) {
586 let dir = tempfile::tempdir().unwrap();
587 let s = Notices::at(dir.path().join("notifications"));
588 (dir, s)
589 }
590
591 #[test]
592 fn a_repeat_counts_but_stays_read_and_a_change_relights_it() {
593 let now = Timestamp::now();
594 let mut n = Notice::warn("task:1", "held");
595 n.mark_read(now);
596 n.raise_again(&Notice::warn("task:1", "held"), now);
597 assert_eq!(n.count, 2);
598 assert!(n.read_at.is_some(), "identical repeat stays read");
599 n.raise_again(&Notice::warn("task:1", "held differently"), now);
600 assert!(n.unread());
601 assert_eq!(n.message, "held differently");
602 }
603
604 #[test]
605 fn a_rise_in_severity_resurrects_a_dismissed_notice_but_a_repeat_does_not() {
606 let now = Timestamp::now();
607 let mut n = Notice::warn("k", "m");
608 n.dismiss(now);
609 n.raise_again(&Notice::warn("k", "m"), now);
610 assert!(n.dismissed_at.is_some(), "tombstone holds");
611 n.raise_again(&Notice::error("k", "m"), now);
612 assert!(n.unread());
613 assert_eq!(n.severity, Severity::Error);
614 n.raise_again(&Notice::info("k", "other"), now);
616 assert_eq!(n.severity, Severity::Error);
617 }
618
619 #[test]
620 fn identical_raises_share_one_file() {
621 let (_d, s) = store();
622 for _ in 0..5 {
623 s.raise(Notice::error("run:abc", "blocked")).unwrap();
624 }
625 let all = s.list();
626 assert_eq!(all.len(), 1);
627 assert_eq!(all[0].count, 5);
628 assert_eq!(s.count_unread(), 1);
629 }
630
631 #[test]
632 fn transitions_persist_and_dismissed_leave_the_list() {
633 let (_d, s) = store();
634 let a = s.raise(Notice::info("a", "one")).unwrap();
635 let b = s.raise(Notice::warn("b", "two")).unwrap();
636 assert_eq!(s.count_unread(), 2);
637 s.mark_read(&a.id).unwrap();
638 assert_eq!(s.count_unread(), 1);
639 s.dismiss(&b.id).unwrap();
640 assert_eq!(s.count_unread(), 0);
641 assert_eq!(s.list().len(), 1);
642 s.raise(Notice::warn("b", "two")).unwrap();
643 assert_eq!(
644 s.list().len(),
645 1,
646 "a tombstone survives a same-message raise"
647 );
648 s.raise(Notice::info("c", "three")).unwrap();
649 assert_eq!(s.mark_all_read().unwrap(), 1);
650 assert_eq!(s.count_unread(), 0);
651 }
652
653 #[test]
654 fn the_list_is_newest_first() {
655 let (_d, s) = store();
656 let mut old = Notice::info("old", "old");
657 old.last_at = "2020-01-01T00:00:00Z".parse().unwrap();
658 s.put(&old).unwrap();
659 s.raise(Notice::info("new", "new")).unwrap();
660 let keys: Vec<_> = s.list().into_iter().map(|n| n.key).collect();
661 assert_eq!(keys, ["new", "old"]);
662 }
663
664 #[test]
665 fn the_cap_prunes_dismissed_then_read_then_oldest() {
666 let (_d, s) = store();
667 let keep = s.raise(Notice::error("keep", "unread")).unwrap();
668 let gone = s.raise(Notice::info("gone", "dismissed")).unwrap();
669 s.dismiss(&gone.id).unwrap();
670 for i in 0..CAP - 1 {
671 s.raise(Notice::info(&format!("k{i}"), "x")).unwrap();
672 }
673 let files = std::fs::read_dir(s.root()).unwrap().count();
674 assert_eq!(files, CAP);
675 assert!(
676 s.get(&keep.id).is_ok(),
677 "an unread notice outlives a dismissed one"
678 );
679 assert!(s.get(&gone.id).is_err(), "the dismissed one went first");
680 }
681
682 #[test]
683 fn revision_moves_on_write_and_on_removal() {
684 let (_d, s) = store();
685 assert_eq!(s.revision(), 0);
686 let a = s.raise(Notice::info("a", "m")).unwrap();
687 let r1 = s.revision();
688 assert_ne!(r1, 0);
689 s.raise(Notice::info("b", "m")).unwrap();
690 let r2 = s.revision();
691 assert_ne!(r1, r2);
692 std::fs::remove_file(s.path_of(&a.id)).unwrap();
693 assert_ne!(s.revision(), r2);
694 }
695
696 #[test]
697 fn ids_are_stable_and_untrusted_ids_are_refused() {
698 assert_eq!(id_of("run:1"), id_of("run:1"));
699 assert_ne!(id_of("run:1"), id_of("run-1"));
700 assert!(valid_id(&id_of("release-bump:20260101-abc")));
701 let (_d, s) = store();
702 for bad in ["", "../x", "a/b", "A", "x.json"] {
703 assert!(s.get(bad).is_err(), "{bad}");
704 }
705 }
706
707 #[test]
708 fn a_newer_schema_is_refused_and_an_older_reads() {
709 let (_d, s) = store();
710 let mut n = Notice::info("k", "m");
711 n.schema = SCHEMA + 1;
712 s.put(&n).unwrap();
713 assert!(s.get(&n.id).is_err());
714 let old = r#"{"id":"x-1","key":"x","severity":"warn","message":"m",
715 "first_at":"2020-01-01T00:00:00Z","last_at":"2020-01-01T00:00:00Z"}"#;
716 std::fs::write(s.path_of("x-1"), old).unwrap();
717 assert_eq!(s.get("x-1").unwrap().count, 1);
718 }
719
720 fn state(status: crate::run::RunStatus) -> crate::run::RunState {
721 let mut st = crate::run::RunState::new(
722 std::path::PathBuf::from("/repo"),
723 "main".to_owned(),
724 "abc1234def".to_owned(),
725 "task".to_owned(),
726 crate::config::Config::default(),
727 );
728 st.status = status;
729 st
730 }
731
732 #[test]
733 fn a_run_that_ended_badly_is_news_unless_it_only_parked() {
734 use crate::run::RunStatus;
735 for bad in [RunStatus::Blocked, RunStatus::Stalled, RunStatus::Failed] {
736 let st = state(bad);
737 let n = run_ended(&st).expect("news");
738 assert_eq!(n.key, format!("run:{}", st.id));
739 assert_eq!(n.severity, Severity::Error);
740 }
741 assert!(run_ended(&state(RunStatus::Merged)).is_none());
742 let mut parked = state(RunStatus::Stalled);
743 parked.parked = true;
744 assert!(run_ended(&parked).is_none());
745 }
746
747 #[test]
748 fn concurrent_writers_of_one_key_all_succeed() {
749 let (_d, s) = store();
750 let handles: Vec<_> = (0..8)
751 .map(|_| {
752 let s = s.clone();
753 std::thread::spawn(move || {
754 for _ in 0..20 {
755 s.raise(Notice::warn("same", "m")).unwrap();
756 }
757 })
758 })
759 .collect();
760 for h in handles {
761 h.join().unwrap();
762 }
763 assert_eq!(s.list().len(), 1);
764 let stray = std::fs::read_dir(s.root())
765 .unwrap()
766 .flatten()
767 .filter(|e| e.path().extension().is_some_and(|x| x == "tmp"))
768 .count();
769 assert_eq!(stray, 0);
770 }
771
772 #[test]
773 fn a_dead_holders_lock_is_broken_and_no_lock_is_left_behind() {
774 let (_d, s) = store();
775 let n = s.raise(Notice::info("k", "m")).unwrap();
776 let lock = s.root().join(format!("{}.lock", n.id));
777 std::fs::write(&lock, "").unwrap();
778 let old = std::time::SystemTime::now() - std::time::Duration::from_secs(60);
779 std::fs::File::options()
780 .write(true)
781 .open(&lock)
782 .unwrap()
783 .set_modified(old)
784 .unwrap();
785 s.mark_read(&n.id).unwrap();
786 assert!(!lock.exists());
787 assert!(s.get(&n.id).unwrap().read_at.is_some());
788 }
789
790 #[test]
791 fn only_a_machine_hold_is_a_task_notice() {
792 let mut t = crate::queue::Task::new(
793 "t".to_owned(),
794 "do it".to_owned(),
795 std::path::PathBuf::from("/repo"),
796 crate::queue::Source::Human,
797 );
798 assert!(task_held(&t).is_none());
799 t.hold_manual(None);
800 assert!(task_held(&t).is_none(), "the operator's own hold");
801 t.hold_machine(Some("missing blocker".to_owned()));
802 let n = task_held(&t).expect("machine hold");
803 assert_eq!(n.key, format!("task:{}", t.id));
804 assert!(n.message.contains("missing blocker"));
805 }
806}