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 task the machine or its attempt budget has held, if it is.
522///
523/// Called from [`crate::queue::Queue::put`], which every task transition goes
524/// through, so a hold made anywhere - the loop, the conductor, triage, a
525/// dependency removed from under a blocked task - is announced. A hold the
526/// operator placed by hand is their own action and is not news. Keyed on the
527/// task, with wording free of anything that varies between retries.
528pub 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
545/// The notice for a run whose graph returned an error before it settled.
546pub 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
554/// The one function producers call. Best-effort by construction: a notice that
555/// cannot be filed is a `tracing::warn`, never a reason to fail the run, the
556/// loop or the request that wanted to mention it.
557pub fn raise(notice: Notice) {
558    if let Some(home) = crate::run::try_home() {
559        raise_in(&home, notice);
560    }
561}
562
563/// [`raise`] into an explicit magi home, for callers that already carry one
564/// (the janitor) and so must not reach for the process-global.
565pub 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        // Severity never falls on a repeat.
605        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}