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