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    /// The task and run ids this notice is about. A question whose `run`
112    /// names one of them already pages the operator for the same cause.
113    #[serde(default)]
114    pub subjects: Vec<String>,
115    /// The question that carries this notice as context, when it was filed
116    /// already read for that reason. Never a tombstone: a recurrence with a
117    /// changed message or higher severity is unread again.
118    #[serde(default)]
119    pub covered_by: Option<String>,
120    /// On-disk format version.
121    #[serde(default = "schema")]
122    pub schema: u32,
123}
124
125fn one() -> u32 {
126    1
127}
128
129fn schema() -> u32 {
130    SCHEMA
131}
132
133impl Notice {
134    /// A notice of `severity` for `key`.
135    pub fn new(severity: Severity, key: &str, message: impl Into<String>) -> Self {
136        let now = Timestamp::now();
137        Self {
138            id: id_of(key),
139            key: key.to_owned(),
140            severity,
141            message: message.into(),
142            link: None,
143            first_at: now,
144            last_at: now,
145            count: 1,
146            read_at: None,
147            dismissed_at: None,
148            subjects: Vec::new(),
149            covered_by: None,
150            schema: SCHEMA,
151        }
152    }
153
154    /// Name the task / run ids this notice is about.
155    pub fn about<I, S>(mut self, subjects: I) -> Self
156    where
157        I: IntoIterator<Item = S>,
158        S: Into<String>,
159    {
160        self.subjects = subjects.into_iter().map(Into::into).collect();
161        self
162    }
163
164    /// An [`Severity::Info`] notice.
165    pub fn info(key: &str, message: impl Into<String>) -> Self {
166        Self::new(Severity::Info, key, message)
167    }
168
169    /// A [`Severity::Warn`] notice.
170    pub fn warn(key: &str, message: impl Into<String>) -> Self {
171        Self::new(Severity::Warn, key, message)
172    }
173
174    /// An [`Severity::Error`] notice.
175    pub fn error(key: &str, message: impl Into<String>) -> Self {
176        Self::new(Severity::Error, key, message)
177    }
178
179    /// Attach a link.
180    pub fn link(mut self, link: Link) -> Self {
181        self.link = Some(link);
182        self
183    }
184
185    /// Neither read nor dismissed.
186    pub fn unread(&self) -> bool {
187        self.read_at.is_none() && self.dismissed_at.is_none()
188    }
189
190    /// Fold a repeat of this notice into it. Always counts; only a changed
191    /// message or a higher severity makes it unread and visible again.
192    pub fn raise_again(&mut self, again: &Notice, now: Timestamp) {
193        self.count = self.count.saturating_add(1);
194        self.last_at = now;
195        if again.link.is_some() {
196            self.link = again.link.clone();
197        }
198        if !again.subjects.is_empty() {
199            self.subjects = again.subjects.clone();
200        }
201        let escalated = again.severity > self.severity;
202        if again.message != self.message || escalated {
203            self.message = again.message.clone();
204            self.severity = self.severity.max(again.severity);
205            self.read_at = None;
206            self.dismissed_at = None;
207            self.covered_by = None;
208        }
209        // A question already pages for this cause: keep the record, skip the
210        // second page. A higher severity is news the question did not carry.
211        if let Some(q) = &again.covered_by
212            && !escalated
213            && self.unread()
214        {
215            self.read_at = Some(now);
216            self.covered_by = Some(q.clone());
217        }
218    }
219
220    /// Mark read, keeping the first read time.
221    pub fn mark_read(&mut self, now: Timestamp) {
222        if self.read_at.is_none() {
223            self.read_at = Some(now);
224        }
225    }
226
227    /// Tombstone: read and hidden from the list.
228    pub fn dismiss(&mut self, now: Timestamp) {
229        self.mark_read(now);
230        if self.dismissed_at.is_none() {
231            self.dismissed_at = Some(now);
232        }
233    }
234
235    /// Pruning order: dismissed go first, then read, then unread.
236    fn keep_rank(&self) -> u8 {
237        if self.dismissed_at.is_some() {
238            0
239        } else if self.read_at.is_some() {
240            1
241        } else {
242            2
243        }
244    }
245}
246
247/// A stable, file-name-safe id for a key: a readable slug plus FNV-1a of the
248/// whole key. Not `DefaultHasher`, whose output may change between builds and
249/// would strand every existing file's dedupe.
250pub fn id_of(key: &str) -> String {
251    let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
252    for b in key.bytes() {
253        hash ^= u64::from(b);
254        hash = hash.wrapping_mul(0x0100_0000_01b3);
255    }
256    let slug: String = key
257        .chars()
258        .map(|c| {
259            if c.is_ascii_alphanumeric() {
260                c.to_ascii_lowercase()
261            } else {
262                '-'
263            }
264        })
265        .take(32)
266        .collect();
267    format!("{slug}-{hash:016x}")
268}
269
270/// Whether `id` could have come from [`id_of`]. Checked before any path is
271/// built from a request.
272fn valid_id(id: &str) -> bool {
273    !id.is_empty()
274        && id.len() <= 64
275        && id
276            .bytes()
277            .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-')
278}
279
280/// How long a caller waits for a notice's lock before proceeding without it.
281const LOCK_WAIT: std::time::Duration = std::time::Duration::from_secs(5);
282
283/// Age after which a lock file is taken to belong to a dead process.
284const LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
285
286/// Distinguishes temp files written by threads of one process.
287static TMP_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
288
289/// A notice store on disk.
290#[derive(Debug, Clone)]
291pub struct Notices {
292    root: PathBuf,
293}
294
295impl Notices {
296    /// The operator's notices, `<home>/notifications`.
297    pub fn open() -> Self {
298        Self::at(crate::run::home().join("notifications"))
299    }
300
301    /// A store at an explicit root; tests use this.
302    pub fn at(root: PathBuf) -> Self {
303        Self { root }
304    }
305
306    /// Directory holding the notice files.
307    pub fn root(&self) -> &Path {
308        &self.root
309    }
310
311    fn path_of(&self, id: &str) -> PathBuf {
312        self.root.join(format!("{id}.json"))
313    }
314
315    fn put(&self, n: &Notice) -> Result<()> {
316        std::fs::create_dir_all(&self.root)
317            .with_context(|| format!("create {}", self.root.display()))?;
318        let body = serde_json::to_string_pretty(n).context("serialize notice")?;
319        let path = self.path_of(&n.id);
320        // Unique per write: two processes raising or marking the same notice
321        // must not share a temp file, or the second rename finds it gone.
322        let tmp = path.with_extension(format!(
323            "json.{}.{}.tmp",
324            std::process::id(),
325            TMP_SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
326        ));
327        std::fs::write(&tmp, body).with_context(|| format!("write {}", tmp.display()))?;
328        if let Err(e) = std::fs::rename(&tmp, &path) {
329            let _ = std::fs::remove_file(&tmp);
330            return Err(e).with_context(|| format!("replace {}", path.display()));
331        }
332        Ok(())
333    }
334
335    /// Load one notice by id.
336    pub fn get(&self, id: &str) -> Result<Notice> {
337        if !valid_id(id) {
338            bail!("`{id}` is not a notification id");
339        }
340        read_path(&self.path_of(id))
341    }
342
343    /// Run a read-modify-write of one notice under its lock file.
344    ///
345    /// `magi serve`, `magi web` and the CLI all update the same files, and a
346    /// unique temp file only makes each write atomic - it does not stop one
347    /// process saving a stale read over another's newer raise. The lock is an
348    /// exclusive-create file beside the notice; a holder that died is broken
349    /// after [`LOCK_STALE`], and a caller that cannot get it in
350    /// [`LOCK_WAIT`] proceeds anyway rather than lose the notice.
351    fn locked<T>(&self, id: &str, f: impl FnOnce() -> Result<T>) -> Result<T> {
352        std::fs::create_dir_all(&self.root)
353            .with_context(|| format!("create {}", self.root.display()))?;
354        let lock = self.root.join(format!("{id}.lock"));
355        let start = std::time::Instant::now();
356        let mut held = false;
357        while start.elapsed() < LOCK_WAIT {
358            match std::fs::OpenOptions::new()
359                .write(true)
360                .create_new(true)
361                .open(&lock)
362            {
363                Ok(_) => {
364                    held = true;
365                    break;
366                }
367                Err(_) => {
368                    let stale = std::fs::metadata(&lock)
369                        .and_then(|m| m.modified())
370                        .ok()
371                        .and_then(|t| t.elapsed().ok())
372                        .is_some_and(|age| age > LOCK_STALE);
373                    if stale {
374                        let _ = std::fs::remove_file(&lock);
375                    } else {
376                        std::thread::sleep(std::time::Duration::from_millis(5));
377                    }
378                }
379            }
380        }
381        let out = f();
382        if held {
383            let _ = std::fs::remove_file(&lock);
384        }
385        out
386    }
387
388    /// File a notice, folding it into an existing one with the same key.
389    pub fn raise(&self, incoming: Notice) -> Result<Notice> {
390        let stored = self.locked(&incoming.id.clone(), || {
391            let now = Timestamp::now();
392            let stored = match read_path(&self.path_of(&incoming.id)) {
393                Ok(mut existing) => {
394                    existing.raise_again(&incoming, now);
395                    existing
396                }
397                Err(_) => {
398                    let mut fresh = incoming;
399                    if fresh.covered_by.is_some() {
400                        fresh.read_at = Some(now);
401                    }
402                    fresh
403                }
404            };
405            self.put(&stored)?;
406            Ok(stored)
407        })?;
408        self.prune();
409        Ok(stored)
410    }
411
412    /// Mark every unread notice about `subject` read, recording `question` as
413    /// the card that carries it. Best-effort per notice.
414    pub fn cover(&self, question: &crate::ask::Question) {
415        for n in self.all() {
416            if n.unread() && covers(question, &n) {
417                let _ = self.update(&n.id, |n| {
418                    if n.unread() {
419                        n.mark_read(Timestamp::now());
420                        n.covered_by = Some(question.id.clone());
421                    }
422                });
423            }
424        }
425    }
426
427    /// Load, change and save one notice under its lock.
428    fn update(&self, id: &str, change: impl FnOnce(&mut Notice)) -> Result<Notice> {
429        if !valid_id(id) {
430            bail!("`{id}` is not a notification id");
431        }
432        self.locked(id, || {
433            let mut n = read_path(&self.path_of(id))?;
434            change(&mut n);
435            self.put(&n)?;
436            Ok(n)
437        })
438    }
439
440    /// Every notice on disk, unreadable files skipped, newest first.
441    fn all(&self) -> Vec<Notice> {
442        let mut all: Vec<Notice> = std::fs::read_dir(&self.root)
443            .into_iter()
444            .flatten()
445            .flatten()
446            .map(|e| e.path())
447            .filter(|p| p.extension().is_some_and(|x| x == "json"))
448            .filter_map(|p| read_path(&p).ok())
449            .collect();
450        all.sort_by(|a, b| b.last_at.cmp(&a.last_at).then_with(|| a.id.cmp(&b.id)));
451        all
452    }
453
454    /// What the bell lists: everything not dismissed, newest first.
455    pub fn list(&self) -> Vec<Notice> {
456        self.all()
457            .into_iter()
458            .filter(|n| n.dismissed_at.is_none())
459            .collect()
460    }
461
462    /// How many notices are unread.
463    pub fn count_unread(&self) -> usize {
464        self.all().iter().filter(|n| n.unread()).count()
465    }
466
467    /// Mark read, keeping the first read time.
468    pub fn mark_read(&self, id: &str) -> Result<Notice> {
469        let now = Timestamp::now();
470        self.update(id, |n| n.mark_read(now))
471    }
472
473    /// Mark every unread notice read; returns how many changed.
474    pub fn mark_all_read(&self) -> Result<usize> {
475        let now = Timestamp::now();
476        let mut changed = 0;
477        for n in self.all().into_iter().filter(Notice::unread) {
478            // Re-read under the lock: a raise since the scan may have made it
479            // a different, still-unread notice, which stays as it is.
480            let done = self.update(&n.id, |n| n.mark_read(now));
481            if done.is_ok() {
482                changed += 1;
483            }
484        }
485        Ok(changed)
486    }
487
488    /// Tombstone: read and hidden from the list.
489    pub fn dismiss(&self, id: &str) -> Result<Notice> {
490        let now = Timestamp::now();
491        self.update(id, |n| n.dismiss(now))
492    }
493
494    /// Keep at most [`CAP`] files. Best-effort.
495    fn prune(&self) {
496        let mut all = self.all();
497        if all.len() <= CAP {
498            return;
499        }
500        // Most worth keeping first; drop the tail.
501        all.sort_by(|a, b| {
502            b.keep_rank()
503                .cmp(&a.keep_rank())
504                .then_with(|| b.last_at.cmp(&a.last_at))
505        });
506        for n in all.split_off(CAP) {
507            let _ = std::fs::remove_file(self.path_of(&n.id));
508        }
509    }
510
511    /// Change token for the stream. Hashes every file's name and mtime, so a
512    /// prune or delete moves it as surely as a write does (a max-mtime would
513    /// not move when the newest file is not the one removed).
514    pub fn revision(&self) -> u64 {
515        use std::hash::{Hash as _, Hasher as _};
516        let mut entries: Vec<(std::ffi::OsString, u128)> = std::fs::read_dir(&self.root)
517            .into_iter()
518            .flatten()
519            .flatten()
520            .filter(|e| e.path().extension().is_some_and(|x| x == "json"))
521            .map(|e| {
522                let at = e
523                    .metadata()
524                    .and_then(|m| m.modified())
525                    .ok()
526                    .and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
527                    .map_or(0, |d| d.as_nanos());
528                (e.file_name(), at)
529            })
530            .collect();
531        if entries.is_empty() {
532            return 0;
533        }
534        entries.sort();
535        let mut h = std::collections::hash_map::DefaultHasher::new();
536        entries.hash(&mut h);
537        h.finish()
538    }
539}
540
541fn read_path(path: &Path) -> Result<Notice> {
542    let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
543    let n: Notice =
544        serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
545    if n.schema > SCHEMA {
546        bail!(
547            "notice {} was written by a newer magi (schema {}, this build speaks up to {SCHEMA})",
548            n.id,
549            n.schema
550        );
551    }
552    Ok(n)
553}
554
555/// The notice for a run that ended Blocked, Stalled or Failed, if it did.
556///
557/// A parked run was only asked to stop, so it is not news. A run that left a
558/// pull request behind still is: the PR is waiting on someone. Keyed on the
559/// run with a message free of anything that varies between attempts.
560pub fn run_ended(state: &crate::run::RunState) -> Option<Notice> {
561    use crate::run::RunStatus;
562    let failed = matches!(
563        state.status,
564        RunStatus::Blocked | RunStatus::Stalled | RunStatus::Failed
565    );
566    (failed && !state.parked).then(|| {
567        Notice::error(
568            &format!("run:{}", state.id),
569            format!("Run {} ended {}.", state.short(), state.status.as_str()),
570        )
571        .link(Link::Run {
572            id: state.id.clone(),
573        })
574        .about([state.id.clone()])
575    })
576}
577
578/// The notice for a pull request that merged while checks were red.
579///
580/// Keyed on the run. `summary` names the repository, pull request and failing
581/// checks: with no notify command configured this is the only place the
582/// operator learns them. A run merges once, so the wording cannot churn.
583pub fn merged_red(run_id: &str, summary: &str) -> Notice {
584    Notice::warn(&format!("merged-red:{run_id}"), summary).link(Link::Run {
585        id: run_id.to_owned(),
586    })
587}
588
589/// The notice for a task the machine or its attempt budget has held, if it is.
590///
591/// Called from [`crate::queue::Queue::put`], which every task transition goes
592/// through, so a hold made anywhere - the loop, the conductor, triage, a
593/// dependency removed from under a blocked task - is announced. A hold the
594/// operator placed by hand is their own action and is not news. Keyed on the
595/// task, with wording free of anything that varies between retries.
596///
597/// Falls back to [`crate::queue::Task::last_error`] when `hold_reason` is
598/// empty, the same fallback `crate::triage`'s question detail uses: a record
599/// written before `Task::fail`/`Task::handed_off` started copying `why` into
600/// `hold_reason` too would otherwise still read "no reason recorded" even
601/// though the run said exactly why it stopped.
602pub fn task_held(task: &crate::queue::Task) -> Option<Notice> {
603    use crate::queue::{HoldSource, TaskStatus};
604    if task.status != TaskStatus::Held || task.hold_source != Some(HoldSource::Machine) {
605        return None;
606    }
607    let why = task
608        .hold_reason
609        .as_deref()
610        .or(task.last_error.as_deref())
611        .unwrap_or("no reason recorded");
612    Some(
613        Notice::warn(
614            &format!("task:{}", task.id),
615            format!("Task {} is held: {why}.", task.short()),
616        )
617        .link(Link::Task {
618            id: task.id.clone(),
619        })
620        .about([task.id.clone()]),
621    )
622}
623
624/// The notice for a run whose graph returned an error before it settled.
625pub fn run_stopped(id: &str, state: &crate::run::RunState) -> Notice {
626    Notice::error(
627        &format!("run:{id}"),
628        format!("Run {} stopped with an error.", state.short()),
629    )
630    .link(Link::Run { id: id.to_owned() })
631    .about([id.to_owned()])
632}
633
634/// The one function producers call. Best-effort by construction: a notice that
635/// cannot be filed is a `tracing::warn`, never a reason to fail the run, the
636/// loop or the request that wanted to mention it.
637pub fn raise(notice: Notice) {
638    if let Some(home) = crate::run::try_home() {
639        raise_in(&home, notice);
640    }
641}
642
643/// Does open question `q` already page the operator for what `n` reports?
644///
645/// The cause is decided by what each side is about, not by when it happened:
646/// a conductor / triage question exists *because* its task is held, so it
647/// covers that task's hold and handover notices whenever it was filed; a
648/// question raised from inside a run (any other seat node) covers that run's
649/// ended / stopped notices and never a task hold. Exact match on
650/// `Question::run` (a task id for the former, a run id for the latter). The
651/// land approval and release questions are to-dos, not duplicates, and cover
652/// nothing.
653pub fn covers(q: &crate::ask::Question, n: &Notice) -> bool {
654    if !q.status.open() || !n.subjects.contains(&q.run) {
655        return false;
656    }
657    let about_task = n.key.starts_with("task:") || n.key.starts_with("handover:");
658    let about_run = n.key.starts_with("run:");
659    match q.node.as_str() {
660        crate::bump::NOTICE_NODE | crate::land::APPROVAL_NODE => false,
661        crate::conduct::NODE | crate::triage::NODE | crate::triage::DEPS_NODE => about_task,
662        _ => about_run,
663    }
664}
665
666fn covering_question(home: &Path, n: &Notice) -> Option<String> {
667    if n.subjects.is_empty() {
668        return None;
669    }
670    crate::ask::Questions::at(home.join("questions"))
671        .list()
672        .into_iter()
673        .find(|q| covers(q, n))
674        .map(|q| q.id)
675}
676
677/// A question was just filed for `run`: quiet the unread notices about it.
678pub fn quiet_for(home: &Path, question: &crate::ask::Question) {
679    Notices::at(home.join("notifications")).cover(question);
680}
681
682/// [`raise`] into an explicit magi home, for callers that already carry one
683/// (the janitor) and so must not reach for the process-global.
684pub fn raise_in(home: &Path, mut notice: Notice) {
685    if let Some(q) = covering_question(home, &notice) {
686        notice.covered_by = Some(q);
687    }
688    if let Err(e) = Notices::at(home.join("notifications")).raise(notice) {
689        tracing::warn!("could not file a notification: {e:#}");
690    }
691}
692
693#[cfg(test)]
694mod tests {
695    use super::*;
696
697    fn store() -> (tempfile::TempDir, Notices) {
698        let dir = tempfile::tempdir().unwrap();
699        let s = Notices::at(dir.path().join("notifications"));
700        (dir, s)
701    }
702
703    #[test]
704    fn a_red_merge_notice_is_keyed_on_the_run_and_raised_once() {
705        let (_d, s) = store();
706        let n = merged_red("run-1", "Merged a/b PR #1 with red checks: x (u)");
707        assert_eq!(n.key, "merged-red:run-1");
708        assert_eq!(
709            n.link,
710            Some(Link::Run {
711                id: "run-1".to_owned()
712            })
713        );
714        s.raise(merged_red(
715            "run-1",
716            "Merged a/b PR #1 with red checks: x (u)",
717        ))
718        .unwrap();
719        s.raise(merged_red(
720            "run-1",
721            "Merged a/b PR #1 with red checks: x (u)",
722        ))
723        .unwrap();
724        assert_eq!(s.list().len(), 1);
725    }
726
727    #[test]
728    fn a_repeat_counts_but_stays_read_and_a_change_relights_it() {
729        let now = Timestamp::now();
730        let mut n = Notice::warn("task:1", "held");
731        n.mark_read(now);
732        n.raise_again(&Notice::warn("task:1", "held"), now);
733        assert_eq!(n.count, 2);
734        assert!(n.read_at.is_some(), "identical repeat stays read");
735        n.raise_again(&Notice::warn("task:1", "held differently"), now);
736        assert!(n.unread());
737        assert_eq!(n.message, "held differently");
738    }
739
740    #[test]
741    fn a_rise_in_severity_resurrects_a_dismissed_notice_but_a_repeat_does_not() {
742        let now = Timestamp::now();
743        let mut n = Notice::warn("task:x", "m");
744        n.dismiss(now);
745        n.raise_again(&Notice::warn("task:x", "m"), now);
746        assert!(n.dismissed_at.is_some(), "tombstone holds");
747        n.raise_again(&Notice::error("k", "m"), now);
748        assert!(n.unread());
749        assert_eq!(n.severity, Severity::Error);
750        // Severity never falls on a repeat.
751        n.raise_again(&Notice::info("k", "other"), now);
752        assert_eq!(n.severity, Severity::Error);
753    }
754
755    #[test]
756    fn identical_raises_share_one_file() {
757        let (_d, s) = store();
758        for _ in 0..5 {
759            s.raise(Notice::error("run:abc", "blocked")).unwrap();
760        }
761        let all = s.list();
762        assert_eq!(all.len(), 1);
763        assert_eq!(all[0].count, 5);
764        assert_eq!(s.count_unread(), 1);
765    }
766
767    #[test]
768    fn transitions_persist_and_dismissed_leave_the_list() {
769        let (_d, s) = store();
770        let a = s.raise(Notice::info("a", "one")).unwrap();
771        let b = s.raise(Notice::warn("b", "two")).unwrap();
772        assert_eq!(s.count_unread(), 2);
773        s.mark_read(&a.id).unwrap();
774        assert_eq!(s.count_unread(), 1);
775        s.dismiss(&b.id).unwrap();
776        assert_eq!(s.count_unread(), 0);
777        assert_eq!(s.list().len(), 1);
778        s.raise(Notice::warn("b", "two")).unwrap();
779        assert_eq!(
780            s.list().len(),
781            1,
782            "a tombstone survives a same-message raise"
783        );
784        s.raise(Notice::info("c", "three")).unwrap();
785        assert_eq!(s.mark_all_read().unwrap(), 1);
786        assert_eq!(s.count_unread(), 0);
787    }
788
789    #[test]
790    fn the_list_is_newest_first() {
791        let (_d, s) = store();
792        let mut old = Notice::info("old", "old");
793        old.last_at = "2020-01-01T00:00:00Z".parse().unwrap();
794        s.put(&old).unwrap();
795        s.raise(Notice::info("new", "new")).unwrap();
796        let keys: Vec<_> = s.list().into_iter().map(|n| n.key).collect();
797        assert_eq!(keys, ["new", "old"]);
798    }
799
800    #[test]
801    fn the_cap_prunes_dismissed_then_read_then_oldest() {
802        let (_d, s) = store();
803        let keep = s.raise(Notice::error("keep", "unread")).unwrap();
804        let gone = s.raise(Notice::info("gone", "dismissed")).unwrap();
805        s.dismiss(&gone.id).unwrap();
806        for i in 0..CAP - 1 {
807            s.raise(Notice::info(&format!("k{i}"), "x")).unwrap();
808        }
809        let files = std::fs::read_dir(s.root()).unwrap().count();
810        assert_eq!(files, CAP);
811        assert!(
812            s.get(&keep.id).is_ok(),
813            "an unread notice outlives a dismissed one"
814        );
815        assert!(s.get(&gone.id).is_err(), "the dismissed one went first");
816    }
817
818    #[test]
819    fn revision_moves_on_write_and_on_removal() {
820        let (_d, s) = store();
821        assert_eq!(s.revision(), 0);
822        let a = s.raise(Notice::info("a", "m")).unwrap();
823        let r1 = s.revision();
824        assert_ne!(r1, 0);
825        s.raise(Notice::info("b", "m")).unwrap();
826        let r2 = s.revision();
827        assert_ne!(r1, r2);
828        std::fs::remove_file(s.path_of(&a.id)).unwrap();
829        assert_ne!(s.revision(), r2);
830    }
831
832    #[test]
833    fn ids_are_stable_and_untrusted_ids_are_refused() {
834        assert_eq!(id_of("run:1"), id_of("run:1"));
835        assert_ne!(id_of("run:1"), id_of("run-1"));
836        assert!(valid_id(&id_of("release-bump:20260101-abc")));
837        let (_d, s) = store();
838        for bad in ["", "../x", "a/b", "A", "x.json"] {
839            assert!(s.get(bad).is_err(), "{bad}");
840        }
841    }
842
843    #[test]
844    fn a_newer_schema_is_refused_and_an_older_reads() {
845        let (_d, s) = store();
846        let mut n = Notice::info("k", "m");
847        n.schema = SCHEMA + 1;
848        s.put(&n).unwrap();
849        assert!(s.get(&n.id).is_err());
850        let old = r#"{"id":"x-1","key":"x","severity":"warn","message":"m",
851            "first_at":"2020-01-01T00:00:00Z","last_at":"2020-01-01T00:00:00Z"}"#;
852        std::fs::write(s.path_of("x-1"), old).unwrap();
853        assert_eq!(s.get("x-1").unwrap().count, 1);
854    }
855
856    fn state(status: crate::run::RunStatus) -> crate::run::RunState {
857        let mut st = crate::run::RunState::new(
858            std::path::PathBuf::from("/repo"),
859            "main".to_owned(),
860            "abc1234def".to_owned(),
861            "task".to_owned(),
862            crate::config::Config::default(),
863        );
864        st.status = status;
865        st
866    }
867
868    #[test]
869    fn a_run_that_ended_badly_is_news_unless_it_only_parked() {
870        use crate::run::RunStatus;
871        for bad in [RunStatus::Blocked, RunStatus::Stalled, RunStatus::Failed] {
872            let st = state(bad);
873            let n = run_ended(&st).expect("news");
874            assert_eq!(n.key, format!("run:{}", st.id));
875            assert_eq!(n.severity, Severity::Error);
876        }
877        assert!(run_ended(&state(RunStatus::Merged)).is_none());
878        let mut parked = state(RunStatus::Stalled);
879        parked.parked = true;
880        assert!(run_ended(&parked).is_none());
881    }
882
883    #[test]
884    fn concurrent_writers_of_one_key_all_succeed() {
885        let (_d, s) = store();
886        let handles: Vec<_> = (0..8)
887            .map(|_| {
888                let s = s.clone();
889                std::thread::spawn(move || {
890                    for _ in 0..20 {
891                        s.raise(Notice::warn("same", "m")).unwrap();
892                    }
893                })
894            })
895            .collect();
896        for h in handles {
897            h.join().unwrap();
898        }
899        assert_eq!(s.list().len(), 1);
900        let stray = std::fs::read_dir(s.root())
901            .unwrap()
902            .flatten()
903            .filter(|e| e.path().extension().is_some_and(|x| x == "tmp"))
904            .count();
905        assert_eq!(stray, 0);
906    }
907
908    #[test]
909    fn a_dead_holders_lock_is_broken_and_no_lock_is_left_behind() {
910        let (_d, s) = store();
911        let n = s.raise(Notice::info("k", "m")).unwrap();
912        let lock = s.root().join(format!("{}.lock", n.id));
913        std::fs::write(&lock, "").unwrap();
914        let old = std::time::SystemTime::now() - std::time::Duration::from_secs(60);
915        std::fs::File::options()
916            .write(true)
917            .open(&lock)
918            .unwrap()
919            .set_modified(old)
920            .unwrap();
921        s.mark_read(&n.id).unwrap();
922        assert!(!lock.exists());
923        assert!(s.get(&n.id).unwrap().read_at.is_some());
924    }
925
926    #[test]
927    fn only_a_machine_hold_is_a_task_notice() {
928        let mut t = crate::queue::Task::new(
929            "t".to_owned(),
930            "do it".to_owned(),
931            std::path::PathBuf::from("/repo"),
932            crate::queue::Source::Human,
933        );
934        assert!(task_held(&t).is_none());
935        t.hold_manual(None);
936        assert!(task_held(&t).is_none(), "the operator's own hold");
937        t.hold_machine(Some("missing blocker".to_owned()));
938        let n = task_held(&t).expect("machine hold");
939        assert_eq!(n.key, format!("task:{}", t.id));
940        assert!(n.message.contains("missing blocker"));
941    }
942
943    fn held_task() -> crate::queue::Task {
944        let mut t = crate::queue::Task::new(
945            "t".to_owned(),
946            "do it".to_owned(),
947            std::path::PathBuf::from("/repo"),
948            crate::queue::Source::Human,
949        );
950        t.hold_machine(Some("branch b is checked out".to_owned()));
951        t
952    }
953
954    fn ask(home: &Path, run: &str, node: &str) -> crate::ask::Question {
955        let mut q = crate::ask::Question::new(
956            run.to_owned(),
957            node.to_owned(),
958            "conductor".to_owned(),
959            "cannot resume".to_owned(),
960            String::new(),
961            Vec::new(),
962        );
963        crate::ask::Questions::at(home.join("questions"))
964            .put(&mut q)
965            .unwrap();
966        q
967    }
968
969    fn unread(home: &Path) -> usize {
970        Notices::at(home.join("notifications")).count_unread()
971    }
972
973    #[test]
974    fn a_hold_with_an_open_question_for_the_task_pages_once() {
975        let dir = tempfile::tempdir().unwrap();
976        let mut t = held_task();
977        let q = ask(dir.path(), &t.id, "conduct");
978        crate::queue::Queue::at(dir.path().join("queue"))
979            .put(&mut t)
980            .unwrap();
981        let s = Notices::at(dir.path().join("notifications"));
982        assert_eq!(s.count_unread(), 0);
983        let n = s.list().pop().expect("the record is kept");
984        assert_eq!(n.covered_by.as_deref(), Some(q.id.as_str()));
985        assert!(n.message.contains("checked out"));
986    }
987
988    #[test]
989    fn a_hold_with_no_question_still_notifies() {
990        let dir = tempfile::tempdir().unwrap();
991        let mut t = held_task();
992        crate::queue::Queue::at(dir.path().join("queue"))
993            .put(&mut t)
994            .unwrap();
995        assert_eq!(unread(dir.path()), 1);
996    }
997
998    #[test]
999    fn a_question_for_another_task_does_not_suppress() {
1000        let dir = tempfile::tempdir().unwrap();
1001        let mut t = held_task();
1002        ask(dir.path(), "some-other-task", "conduct");
1003        crate::queue::Queue::at(dir.path().join("queue"))
1004            .put(&mut t)
1005            .unwrap();
1006        assert_eq!(unread(dir.path()), 1);
1007    }
1008
1009    #[test]
1010    fn a_question_filed_after_the_hold_quiets_it_once() {
1011        let dir = tempfile::tempdir().unwrap();
1012        let mut t = held_task();
1013        crate::queue::Queue::at(dir.path().join("queue"))
1014            .put(&mut t)
1015            .unwrap();
1016        assert_eq!(unread(dir.path()), 1);
1017        let q = ask(dir.path(), &t.id, "conduct");
1018        assert_eq!(unread(dir.path()), 0);
1019        // A later update of the same question does not run the hook again.
1020        let s = Notices::at(dir.path().join("notifications"));
1021        let id = s.list()[0].id.clone();
1022        s.update(&id, |n| n.read_at = None).unwrap();
1023        let mut again = crate::ask::Questions::at(dir.path().join("questions"))
1024            .get(&q.id)
1025            .unwrap();
1026        crate::ask::Questions::at(dir.path().join("questions"))
1027            .put(&mut again)
1028            .unwrap();
1029        assert_eq!(unread(dir.path()), 1);
1030    }
1031
1032    #[test]
1033    fn a_blocked_run_with_an_open_question_is_quiet() {
1034        let dir = tempfile::tempdir().unwrap();
1035        ask(dir.path(), "run-1", "implement");
1036        raise_in(
1037            dir.path(),
1038            Notice::error("run:run-1", "Run r ended blocked.").about(["run-1"]),
1039        );
1040        assert_eq!(unread(dir.path()), 0);
1041        raise_in(
1042            dir.path(),
1043            Notice::error("run:run-2", "Run r ended blocked.").about(["run-2"]),
1044        );
1045        assert_eq!(unread(dir.path()), 1);
1046    }
1047
1048    #[test]
1049    fn escalation_makes_a_covered_notice_unread_again() {
1050        let dir = tempfile::tempdir().unwrap();
1051        ask(dir.path(), "t1", "conduct");
1052        raise_in(dir.path(), Notice::warn("task:t1", "held").about(["t1"]));
1053        assert_eq!(unread(dir.path()), 0);
1054        raise_in(dir.path(), Notice::error("task:t1", "held").about(["t1"]));
1055        assert_eq!(unread(dir.path()), 1);
1056    }
1057
1058    #[test]
1059    fn covers_matches_run_exactly_and_ignores_release_questions() {
1060        let dir = tempfile::tempdir().unwrap();
1061        let n = Notice::warn("task:x", "m").about(["t1"]);
1062        let q = ask(dir.path(), "t1", "conduct");
1063        assert!(covers(&q, &n));
1064        assert!(!covers(&q, &Notice::warn("task:x", "m")));
1065        assert!(!covers(&q, &Notice::warn("task:x", "m").about(["t"])));
1066        let release = ask(dir.path(), "t1", crate::bump::NOTICE_NODE);
1067        assert!(!covers(&release, &n));
1068    }
1069
1070    #[test]
1071    fn a_run_question_does_not_swallow_a_task_hold_and_vice_versa() {
1072        let dir = tempfile::tempdir().unwrap();
1073        ask(dir.path(), "t1", "implement");
1074        raise_in(dir.path(), Notice::warn("task:t1", "held").about(["t1"]));
1075        assert_eq!(unread(dir.path()), 1);
1076        ask(dir.path(), "r1", "conduct");
1077        raise_in(dir.path(), Notice::error("run:r1", "ended").about(["r1"]));
1078        assert_eq!(unread(dir.path()), 2);
1079    }
1080
1081    #[test]
1082    fn a_conduct_question_covers_the_hold_however_late_it_was_filed() {
1083        let dir = tempfile::tempdir().unwrap();
1084        let mut q = crate::ask::Question::new(
1085            "t1".to_owned(),
1086            "conduct".to_owned(),
1087            "conductor".to_owned(),
1088            "s".to_owned(),
1089            String::new(),
1090            Vec::new(),
1091        );
1092        q.asked_at = Timestamp::from_second(Timestamp::now().as_second() - 3600).unwrap();
1093        crate::ask::Questions::at(dir.path().join("questions"))
1094            .put(&mut q)
1095            .unwrap();
1096        raise_in(dir.path(), Notice::warn("task:t1", "held").about(["t1"]));
1097        assert_eq!(unread(dir.path()), 0);
1098    }
1099}