Skip to main content

magi/
notices.rs

1//! Notices: what needs the operator's attention but is not a question.
2//!
3//! A merged run whose release bump failed, a task held after three attempts, a
4//! disk gate that refused to start work - each used to end in a log line or in
5//! a question-shaped record nobody could answer. A [`Notice`] is the missing
6//! shape: something to *know*, with a severity, a timestamp and an optional
7//! link, that the operator marks read or dismisses. It is deliberately not a
8//! [`crate::ask::Question`]: nothing waits on it and there is nothing to answer.
9//!
10//! # Shape
11//!
12//! The same split as [`crate::queue`] and [`crate::ask`]. [`Notice`] is data
13//! plus *pure* transitions; [`Notices`] owns every filesystem call and is
14//! constructed with its root, so a test drives a real store in a temp
15//! directory and nothing here is process-global.
16//!
17//! One notice is one JSON file, written atomically, because `magi serve`,
18//! `magi web` and the CLI all write here from different processes and a rename
19//! is the only cross-process write that needs no coordination. The file name is
20//! derived from the notice's *key* by a stable function, which is what makes
21//! deduplication lock-free: two processes raising the same problem write the
22//! same file. A raise racing another may lose one increment of `count`; that is
23//! accepted.
24//!
25//! # Deduplication, and what dismissing means
26//!
27//! A key names a kind and a subject (`release-bump:<run>`, `task:<id>`), never
28//! the prose. Raising a key that exists bumps `count` and `last_at`. It returns
29//! the notice to *unread* only when the message changed or the severity rose:
30//! a retry loop re-raising the identical problem must not relight the bell.
31//! Dismissal is a tombstone (`dismissed_at`) rather than a delete, so the same
32//! message re-raised does not resurrect what the operator already waved away.
33//! The cost is that a genuine recurrence with identical wording stays hidden
34//! until the cap prunes the tombstone; producers should therefore keep
35//! messages stable and put anything that varies elsewhere.
36
37use std::path::{Path, PathBuf};
38
39use anyhow::{Context, Result, bail};
40use jiff::Timestamp;
41use serde::{Deserialize, Serialize};
42
43/// On-disk format. A file is refused only when its own schema is *greater*
44/// than this one; every field added later must carry `#[serde(default)]`.
45pub const SCHEMA: u32 = 1;
46
47/// The most notices kept. Older ones are pruned on every write: dismissed
48/// first, then read, then unread, oldest first within each.
49pub const CAP: usize = 200;
50
51/// How serious a notice is. Ordered so a rise is a plain comparison.
52#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
53#[serde(rename_all = "lowercase")]
54pub enum Severity {
55    /// Worth knowing, nothing to do.
56    Info,
57    /// Something needs a look soon.
58    Warn,
59    /// Something failed and stays failed until someone acts.
60    Error,
61}
62
63/// Where a notice points, when it points anywhere.
64#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
65#[serde(tag = "kind", rename_all = "lowercase")]
66pub enum Link {
67    /// A run, by id.
68    Run {
69        /// The run id.
70        id: String,
71    },
72    /// A queued task, by id.
73    Task {
74        /// The task id.
75        id: String,
76    },
77    /// An external page, such as a pull request.
78    Url {
79        /// The address.
80        url: String,
81    },
82}
83
84/// One thing the operator should know.
85#[derive(Debug, Clone, Serialize, Deserialize)]
86pub struct Notice {
87    /// Stable file-name-safe id derived from [`Notice::key`].
88    pub id: String,
89    /// Deduplication key: kind and subject, never prose.
90    pub key: String,
91    /// How serious.
92    pub severity: Severity,
93    /// What happened.
94    pub message: String,
95    /// Where it points, when it points anywhere.
96    #[serde(default)]
97    pub link: Option<Link>,
98    /// First raised.
99    pub first_at: Timestamp,
100    /// Last raised.
101    pub last_at: Timestamp,
102    /// How many times this key was raised.
103    #[serde(default = "one")]
104    pub count: u32,
105    /// When it was marked read.
106    #[serde(default)]
107    pub read_at: Option<Timestamp>,
108    /// When it was dismissed; the tombstone.
109    #[serde(default)]
110    pub dismissed_at: Option<Timestamp>,
111    /// On-disk format version.
112    #[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    /// A notice of `severity` for `key`.
126    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    /// An [`Severity::Info`] notice.
144    pub fn info(key: &str, message: impl Into<String>) -> Self {
145        Self::new(Severity::Info, key, message)
146    }
147
148    /// A [`Severity::Warn`] notice.
149    pub fn warn(key: &str, message: impl Into<String>) -> Self {
150        Self::new(Severity::Warn, key, message)
151    }
152
153    /// An [`Severity::Error`] notice.
154    pub fn error(key: &str, message: impl Into<String>) -> Self {
155        Self::new(Severity::Error, key, message)
156    }
157
158    /// Attach a link.
159    pub fn link(mut self, link: Link) -> Self {
160        self.link = Some(link);
161        self
162    }
163
164    /// Neither read nor dismissed.
165    pub fn unread(&self) -> bool {
166        self.read_at.is_none() && self.dismissed_at.is_none()
167    }
168
169    /// Fold a repeat of this notice into it. Always counts; only a changed
170    /// message or a higher severity makes it unread and visible again.
171    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    /// Mark read, keeping the first read time.
186    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    /// Tombstone: read and hidden from the list.
193    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    /// Pruning order: dismissed go first, then read, then unread.
201    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
212/// A stable, file-name-safe id for a key: a readable slug plus FNV-1a of the
213/// whole key. Not `DefaultHasher`, whose output may change between builds and
214/// would strand every existing file's dedupe.
215pub 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
235/// Whether `id` could have come from [`id_of`]. Checked before any path is
236/// built from a request.
237fn 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
245/// How long a caller waits for a notice's lock before proceeding without it.
246const LOCK_WAIT: std::time::Duration = std::time::Duration::from_secs(5);
247
248/// Age after which a lock file is taken to belong to a dead process.
249const LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
250
251/// Distinguishes temp files written by threads of one process.
252static TMP_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
253
254/// A notice store on disk.
255#[derive(Debug, Clone)]
256pub struct Notices {
257    root: PathBuf,
258}
259
260impl Notices {
261    /// The operator's notices, `<home>/notifications`.
262    pub fn open() -> Self {
263        Self::at(crate::run::home().join("notifications"))
264    }
265
266    /// A store at an explicit root; tests use this.
267    pub fn at(root: PathBuf) -> Self {
268        Self { root }
269    }
270
271    /// Directory holding the notice files.
272    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        // Unique per write: two processes raising or marking the same notice
286        // must not share a temp file, or the second rename finds it gone.
287        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    /// Load one notice by id.
301    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    /// Run a read-modify-write of one notice under its lock file.
309    ///
310    /// `magi serve`, `magi web` and the CLI all update the same files, and a
311    /// unique temp file only makes each write atomic - it does not stop one
312    /// process saving a stale read over another's newer raise. The lock is an
313    /// exclusive-create file beside the notice; a holder that died is broken
314    /// after [`LOCK_STALE`], and a caller that cannot get it in
315    /// [`LOCK_WAIT`] proceeds anyway rather than lose the notice.
316    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    /// File a notice, folding it into an existing one with the same key.
354    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    /// Load, change and save one notice under its lock.
372    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    /// Every notice on disk, unreadable files skipped, newest first.
385    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    /// What the bell lists: everything not dismissed, newest first.
399    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    /// How many notices are unread.
407    pub fn count_unread(&self) -> usize {
408        self.all().iter().filter(|n| n.unread()).count()
409    }
410
411    /// Mark read, keeping the first read time.
412    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    /// Mark every unread notice read; returns how many changed.
418    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            // Re-read under the lock: a raise since the scan may have made it
423            // a different, still-unread notice, which stays as it is.
424            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    /// Tombstone: read and hidden from the list.
433    pub fn dismiss(&self, id: &str) -> Result<Notice> {
434        let now = Timestamp::now();
435        self.update(id, |n| n.dismiss(now))
436    }
437
438    /// Keep at most [`CAP`] files. Best-effort.
439    fn prune(&self) {
440        let mut all = self.all();
441        if all.len() <= CAP {
442            return;
443        }
444        // Most worth keeping first; drop the tail.
445        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    /// Change token for the stream. Hashes every file's name and mtime, so a
456    /// prune or delete moves it as surely as a write does (a max-mtime would
457    /// not move when the newest file is not the one removed).
458    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
499/// The notice for a run that ended Blocked, Stalled or Failed, if it did.
500///
501/// A parked run was only asked to stop, so it is not news. A run that left a
502/// pull request behind still is: the PR is waiting on someone. Keyed on the
503/// run with a message free of anything that varies between attempts.
504pub 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
521/// The notice for a pull request that merged while checks were red.
522///
523/// Keyed on the run. `summary` names the repository, pull request and failing
524/// checks: with no notify command configured this is the only place the
525/// operator learns them. A run merges once, so the wording cannot churn.
526pub 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
532/// The notice for a task the machine or its attempt budget has held, if it is.
533///
534/// Called from [`crate::queue::Queue::put`], which every task transition goes
535/// through, so a hold made anywhere - the loop, the conductor, triage, a
536/// dependency removed from under a blocked task - is announced. A hold the
537/// operator placed by hand is their own action and is not news. Keyed on the
538/// task, with wording free of anything that varies between retries.
539///
540/// Falls back to [`crate::queue::Task::last_error`] when `hold_reason` is
541/// empty, the same fallback `crate::triage`'s question detail uses: a record
542/// written before `Task::fail`/`Task::handed_off` started copying `why` into
543/// `hold_reason` too would otherwise still read "no reason recorded" even
544/// though the run said exactly why it stopped.
545pub 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
566/// The notice for a run whose graph returned an error before it settled.
567pub 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
575/// The one function producers call. Best-effort by construction: a notice that
576/// cannot be filed is a `tracing::warn`, never a reason to fail the run, the
577/// loop or the request that wanted to mention it.
578pub fn raise(notice: Notice) {
579    if let Some(home) = crate::run::try_home() {
580        raise_in(&home, notice);
581    }
582}
583
584/// [`raise`] into an explicit magi home, for callers that already carry one
585/// (the janitor) and so must not reach for the process-global.
586pub 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        // Severity never falls on a repeat.
650        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}