use std::path::{Path, PathBuf};
use anyhow::{Context, Result, bail};
use jiff::Timestamp;
use serde::{Deserialize, Serialize};
pub const SCHEMA: u32 = 1;
pub const CAP: usize = 200;
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum Severity {
Info,
Warn,
Error,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "lowercase")]
pub enum Link {
Run {
id: String,
},
Task {
id: String,
},
Url {
url: String,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Notice {
pub id: String,
pub key: String,
pub severity: Severity,
pub message: String,
#[serde(default)]
pub link: Option<Link>,
pub first_at: Timestamp,
pub last_at: Timestamp,
#[serde(default = "one")]
pub count: u32,
#[serde(default)]
pub read_at: Option<Timestamp>,
#[serde(default)]
pub dismissed_at: Option<Timestamp>,
#[serde(default = "schema")]
pub schema: u32,
}
fn one() -> u32 {
1
}
fn schema() -> u32 {
SCHEMA
}
impl Notice {
pub fn new(severity: Severity, key: &str, message: impl Into<String>) -> Self {
let now = Timestamp::now();
Self {
id: id_of(key),
key: key.to_owned(),
severity,
message: message.into(),
link: None,
first_at: now,
last_at: now,
count: 1,
read_at: None,
dismissed_at: None,
schema: SCHEMA,
}
}
pub fn info(key: &str, message: impl Into<String>) -> Self {
Self::new(Severity::Info, key, message)
}
pub fn warn(key: &str, message: impl Into<String>) -> Self {
Self::new(Severity::Warn, key, message)
}
pub fn error(key: &str, message: impl Into<String>) -> Self {
Self::new(Severity::Error, key, message)
}
pub fn link(mut self, link: Link) -> Self {
self.link = Some(link);
self
}
pub fn unread(&self) -> bool {
self.read_at.is_none() && self.dismissed_at.is_none()
}
pub fn raise_again(&mut self, again: &Notice, now: Timestamp) {
self.count = self.count.saturating_add(1);
self.last_at = now;
if again.link.is_some() {
self.link = again.link.clone();
}
if again.message != self.message || again.severity > self.severity {
self.message = again.message.clone();
self.severity = self.severity.max(again.severity);
self.read_at = None;
self.dismissed_at = None;
}
}
pub fn mark_read(&mut self, now: Timestamp) {
if self.read_at.is_none() {
self.read_at = Some(now);
}
}
pub fn dismiss(&mut self, now: Timestamp) {
self.mark_read(now);
if self.dismissed_at.is_none() {
self.dismissed_at = Some(now);
}
}
fn keep_rank(&self) -> u8 {
if self.dismissed_at.is_some() {
0
} else if self.read_at.is_some() {
1
} else {
2
}
}
}
pub fn id_of(key: &str) -> String {
let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
for b in key.bytes() {
hash ^= u64::from(b);
hash = hash.wrapping_mul(0x0100_0000_01b3);
}
let slug: String = key
.chars()
.map(|c| {
if c.is_ascii_alphanumeric() {
c.to_ascii_lowercase()
} else {
'-'
}
})
.take(32)
.collect();
format!("{slug}-{hash:016x}")
}
fn valid_id(id: &str) -> bool {
!id.is_empty()
&& id.len() <= 64
&& id
.bytes()
.all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-')
}
const LOCK_WAIT: std::time::Duration = std::time::Duration::from_secs(5);
const LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
static TMP_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
#[derive(Debug, Clone)]
pub struct Notices {
root: PathBuf,
}
impl Notices {
pub fn open() -> Self {
Self::at(crate::run::home().join("notifications"))
}
pub fn at(root: PathBuf) -> Self {
Self { root }
}
pub fn root(&self) -> &Path {
&self.root
}
fn path_of(&self, id: &str) -> PathBuf {
self.root.join(format!("{id}.json"))
}
fn put(&self, n: &Notice) -> Result<()> {
std::fs::create_dir_all(&self.root)
.with_context(|| format!("create {}", self.root.display()))?;
let body = serde_json::to_string_pretty(n).context("serialize notice")?;
let path = self.path_of(&n.id);
let tmp = path.with_extension(format!(
"json.{}.{}.tmp",
std::process::id(),
TMP_SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
));
std::fs::write(&tmp, body).with_context(|| format!("write {}", tmp.display()))?;
if let Err(e) = std::fs::rename(&tmp, &path) {
let _ = std::fs::remove_file(&tmp);
return Err(e).with_context(|| format!("replace {}", path.display()));
}
Ok(())
}
pub fn get(&self, id: &str) -> Result<Notice> {
if !valid_id(id) {
bail!("`{id}` is not a notification id");
}
read_path(&self.path_of(id))
}
fn locked<T>(&self, id: &str, f: impl FnOnce() -> Result<T>) -> Result<T> {
std::fs::create_dir_all(&self.root)
.with_context(|| format!("create {}", self.root.display()))?;
let lock = self.root.join(format!("{id}.lock"));
let start = std::time::Instant::now();
let mut held = false;
while start.elapsed() < LOCK_WAIT {
match std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&lock)
{
Ok(_) => {
held = true;
break;
}
Err(_) => {
let stale = std::fs::metadata(&lock)
.and_then(|m| m.modified())
.ok()
.and_then(|t| t.elapsed().ok())
.is_some_and(|age| age > LOCK_STALE);
if stale {
let _ = std::fs::remove_file(&lock);
} else {
std::thread::sleep(std::time::Duration::from_millis(5));
}
}
}
}
let out = f();
if held {
let _ = std::fs::remove_file(&lock);
}
out
}
pub fn raise(&self, incoming: Notice) -> Result<Notice> {
let stored = self.locked(&incoming.id.clone(), || {
let now = Timestamp::now();
let stored = match read_path(&self.path_of(&incoming.id)) {
Ok(mut existing) => {
existing.raise_again(&incoming, now);
existing
}
Err(_) => incoming,
};
self.put(&stored)?;
Ok(stored)
})?;
self.prune();
Ok(stored)
}
fn update(&self, id: &str, change: impl FnOnce(&mut Notice)) -> Result<Notice> {
if !valid_id(id) {
bail!("`{id}` is not a notification id");
}
self.locked(id, || {
let mut n = read_path(&self.path_of(id))?;
change(&mut n);
self.put(&n)?;
Ok(n)
})
}
fn all(&self) -> Vec<Notice> {
let mut all: Vec<Notice> = std::fs::read_dir(&self.root)
.into_iter()
.flatten()
.flatten()
.map(|e| e.path())
.filter(|p| p.extension().is_some_and(|x| x == "json"))
.filter_map(|p| read_path(&p).ok())
.collect();
all.sort_by(|a, b| b.last_at.cmp(&a.last_at).then_with(|| a.id.cmp(&b.id)));
all
}
pub fn list(&self) -> Vec<Notice> {
self.all()
.into_iter()
.filter(|n| n.dismissed_at.is_none())
.collect()
}
pub fn count_unread(&self) -> usize {
self.all().iter().filter(|n| n.unread()).count()
}
pub fn mark_read(&self, id: &str) -> Result<Notice> {
let now = Timestamp::now();
self.update(id, |n| n.mark_read(now))
}
pub fn mark_all_read(&self) -> Result<usize> {
let now = Timestamp::now();
let mut changed = 0;
for n in self.all().into_iter().filter(Notice::unread) {
let done = self.update(&n.id, |n| n.mark_read(now));
if done.is_ok() {
changed += 1;
}
}
Ok(changed)
}
pub fn dismiss(&self, id: &str) -> Result<Notice> {
let now = Timestamp::now();
self.update(id, |n| n.dismiss(now))
}
fn prune(&self) {
let mut all = self.all();
if all.len() <= CAP {
return;
}
all.sort_by(|a, b| {
b.keep_rank()
.cmp(&a.keep_rank())
.then_with(|| b.last_at.cmp(&a.last_at))
});
for n in all.split_off(CAP) {
let _ = std::fs::remove_file(self.path_of(&n.id));
}
}
pub fn revision(&self) -> u64 {
use std::hash::{Hash as _, Hasher as _};
let mut entries: Vec<(std::ffi::OsString, u128)> = std::fs::read_dir(&self.root)
.into_iter()
.flatten()
.flatten()
.filter(|e| e.path().extension().is_some_and(|x| x == "json"))
.map(|e| {
let at = e
.metadata()
.and_then(|m| m.modified())
.ok()
.and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
.map_or(0, |d| d.as_nanos());
(e.file_name(), at)
})
.collect();
if entries.is_empty() {
return 0;
}
entries.sort();
let mut h = std::collections::hash_map::DefaultHasher::new();
entries.hash(&mut h);
h.finish()
}
}
fn read_path(path: &Path) -> Result<Notice> {
let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
let n: Notice =
serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
if n.schema > SCHEMA {
bail!(
"notice {} was written by a newer magi (schema {}, this build speaks up to {SCHEMA})",
n.id,
n.schema
);
}
Ok(n)
}
pub fn run_ended(state: &crate::run::RunState) -> Option<Notice> {
use crate::run::RunStatus;
let failed = matches!(
state.status,
RunStatus::Blocked | RunStatus::Stalled | RunStatus::Failed
);
(failed && !state.parked).then(|| {
Notice::error(
&format!("run:{}", state.id),
format!("Run {} ended {}.", state.short(), state.status.as_str()),
)
.link(Link::Run {
id: state.id.clone(),
})
})
}
pub fn task_held(task: &crate::queue::Task) -> Option<Notice> {
use crate::queue::{HoldSource, TaskStatus};
if task.status != TaskStatus::Held || task.hold_source != Some(HoldSource::Machine) {
return None;
}
let why = task.hold_reason.as_deref().unwrap_or("no reason recorded");
Some(
Notice::warn(
&format!("task:{}", task.id),
format!("Task {} is held: {why}.", task.short()),
)
.link(Link::Task {
id: task.id.clone(),
}),
)
}
pub fn run_stopped(id: &str, state: &crate::run::RunState) -> Notice {
Notice::error(
&format!("run:{id}"),
format!("Run {} stopped with an error.", state.short()),
)
.link(Link::Run { id: id.to_owned() })
}
pub fn raise(notice: Notice) {
if let Some(home) = crate::run::try_home() {
raise_in(&home, notice);
}
}
pub fn raise_in(home: &Path, notice: Notice) {
if let Err(e) = Notices::at(home.join("notifications")).raise(notice) {
tracing::warn!("could not file a notification: {e:#}");
}
}
#[cfg(test)]
mod tests {
use super::*;
fn store() -> (tempfile::TempDir, Notices) {
let dir = tempfile::tempdir().unwrap();
let s = Notices::at(dir.path().join("notifications"));
(dir, s)
}
#[test]
fn a_repeat_counts_but_stays_read_and_a_change_relights_it() {
let now = Timestamp::now();
let mut n = Notice::warn("task:1", "held");
n.mark_read(now);
n.raise_again(&Notice::warn("task:1", "held"), now);
assert_eq!(n.count, 2);
assert!(n.read_at.is_some(), "identical repeat stays read");
n.raise_again(&Notice::warn("task:1", "held differently"), now);
assert!(n.unread());
assert_eq!(n.message, "held differently");
}
#[test]
fn a_rise_in_severity_resurrects_a_dismissed_notice_but_a_repeat_does_not() {
let now = Timestamp::now();
let mut n = Notice::warn("k", "m");
n.dismiss(now);
n.raise_again(&Notice::warn("k", "m"), now);
assert!(n.dismissed_at.is_some(), "tombstone holds");
n.raise_again(&Notice::error("k", "m"), now);
assert!(n.unread());
assert_eq!(n.severity, Severity::Error);
n.raise_again(&Notice::info("k", "other"), now);
assert_eq!(n.severity, Severity::Error);
}
#[test]
fn identical_raises_share_one_file() {
let (_d, s) = store();
for _ in 0..5 {
s.raise(Notice::error("run:abc", "blocked")).unwrap();
}
let all = s.list();
assert_eq!(all.len(), 1);
assert_eq!(all[0].count, 5);
assert_eq!(s.count_unread(), 1);
}
#[test]
fn transitions_persist_and_dismissed_leave_the_list() {
let (_d, s) = store();
let a = s.raise(Notice::info("a", "one")).unwrap();
let b = s.raise(Notice::warn("b", "two")).unwrap();
assert_eq!(s.count_unread(), 2);
s.mark_read(&a.id).unwrap();
assert_eq!(s.count_unread(), 1);
s.dismiss(&b.id).unwrap();
assert_eq!(s.count_unread(), 0);
assert_eq!(s.list().len(), 1);
s.raise(Notice::warn("b", "two")).unwrap();
assert_eq!(
s.list().len(),
1,
"a tombstone survives a same-message raise"
);
s.raise(Notice::info("c", "three")).unwrap();
assert_eq!(s.mark_all_read().unwrap(), 1);
assert_eq!(s.count_unread(), 0);
}
#[test]
fn the_list_is_newest_first() {
let (_d, s) = store();
let mut old = Notice::info("old", "old");
old.last_at = "2020-01-01T00:00:00Z".parse().unwrap();
s.put(&old).unwrap();
s.raise(Notice::info("new", "new")).unwrap();
let keys: Vec<_> = s.list().into_iter().map(|n| n.key).collect();
assert_eq!(keys, ["new", "old"]);
}
#[test]
fn the_cap_prunes_dismissed_then_read_then_oldest() {
let (_d, s) = store();
let keep = s.raise(Notice::error("keep", "unread")).unwrap();
let gone = s.raise(Notice::info("gone", "dismissed")).unwrap();
s.dismiss(&gone.id).unwrap();
for i in 0..CAP - 1 {
s.raise(Notice::info(&format!("k{i}"), "x")).unwrap();
}
let files = std::fs::read_dir(s.root()).unwrap().count();
assert_eq!(files, CAP);
assert!(
s.get(&keep.id).is_ok(),
"an unread notice outlives a dismissed one"
);
assert!(s.get(&gone.id).is_err(), "the dismissed one went first");
}
#[test]
fn revision_moves_on_write_and_on_removal() {
let (_d, s) = store();
assert_eq!(s.revision(), 0);
let a = s.raise(Notice::info("a", "m")).unwrap();
let r1 = s.revision();
assert_ne!(r1, 0);
s.raise(Notice::info("b", "m")).unwrap();
let r2 = s.revision();
assert_ne!(r1, r2);
std::fs::remove_file(s.path_of(&a.id)).unwrap();
assert_ne!(s.revision(), r2);
}
#[test]
fn ids_are_stable_and_untrusted_ids_are_refused() {
assert_eq!(id_of("run:1"), id_of("run:1"));
assert_ne!(id_of("run:1"), id_of("run-1"));
assert!(valid_id(&id_of("release-bump:20260101-abc")));
let (_d, s) = store();
for bad in ["", "../x", "a/b", "A", "x.json"] {
assert!(s.get(bad).is_err(), "{bad}");
}
}
#[test]
fn a_newer_schema_is_refused_and_an_older_reads() {
let (_d, s) = store();
let mut n = Notice::info("k", "m");
n.schema = SCHEMA + 1;
s.put(&n).unwrap();
assert!(s.get(&n.id).is_err());
let old = r#"{"id":"x-1","key":"x","severity":"warn","message":"m",
"first_at":"2020-01-01T00:00:00Z","last_at":"2020-01-01T00:00:00Z"}"#;
std::fs::write(s.path_of("x-1"), old).unwrap();
assert_eq!(s.get("x-1").unwrap().count, 1);
}
fn state(status: crate::run::RunStatus) -> crate::run::RunState {
let mut st = crate::run::RunState::new(
std::path::PathBuf::from("/repo"),
"main".to_owned(),
"abc1234def".to_owned(),
"task".to_owned(),
crate::config::Config::default(),
);
st.status = status;
st
}
#[test]
fn a_run_that_ended_badly_is_news_unless_it_only_parked() {
use crate::run::RunStatus;
for bad in [RunStatus::Blocked, RunStatus::Stalled, RunStatus::Failed] {
let st = state(bad);
let n = run_ended(&st).expect("news");
assert_eq!(n.key, format!("run:{}", st.id));
assert_eq!(n.severity, Severity::Error);
}
assert!(run_ended(&state(RunStatus::Merged)).is_none());
let mut parked = state(RunStatus::Stalled);
parked.parked = true;
assert!(run_ended(&parked).is_none());
}
#[test]
fn concurrent_writers_of_one_key_all_succeed() {
let (_d, s) = store();
let handles: Vec<_> = (0..8)
.map(|_| {
let s = s.clone();
std::thread::spawn(move || {
for _ in 0..20 {
s.raise(Notice::warn("same", "m")).unwrap();
}
})
})
.collect();
for h in handles {
h.join().unwrap();
}
assert_eq!(s.list().len(), 1);
let stray = std::fs::read_dir(s.root())
.unwrap()
.flatten()
.filter(|e| e.path().extension().is_some_and(|x| x == "tmp"))
.count();
assert_eq!(stray, 0);
}
#[test]
fn a_dead_holders_lock_is_broken_and_no_lock_is_left_behind() {
let (_d, s) = store();
let n = s.raise(Notice::info("k", "m")).unwrap();
let lock = s.root().join(format!("{}.lock", n.id));
std::fs::write(&lock, "").unwrap();
let old = std::time::SystemTime::now() - std::time::Duration::from_secs(60);
std::fs::File::options()
.write(true)
.open(&lock)
.unwrap()
.set_modified(old)
.unwrap();
s.mark_read(&n.id).unwrap();
assert!(!lock.exists());
assert!(s.get(&n.id).unwrap().read_at.is_some());
}
#[test]
fn only_a_machine_hold_is_a_task_notice() {
let mut t = crate::queue::Task::new(
"t".to_owned(),
"do it".to_owned(),
std::path::PathBuf::from("/repo"),
crate::queue::Source::Human,
);
assert!(task_held(&t).is_none());
t.hold_manual(None);
assert!(task_held(&t).is_none(), "the operator's own hold");
t.hold_machine(Some("missing blocker".to_owned()));
let n = task_held(&t).expect("machine hold");
assert_eq!(n.key, format!("task:{}", t.id));
assert!(n.message.contains("missing blocker"));
}
}