1use std::path::{Path, PathBuf};
38
39use anyhow::{Context, Result, bail};
40use jiff::Timestamp;
41use serde::{Deserialize, Serialize};
42
43pub const SCHEMA: u32 = 1;
46
47pub const CAP: usize = 200;
50
51#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
53#[serde(rename_all = "lowercase")]
54pub enum Severity {
55 Info,
57 Warn,
59 Error,
61}
62
63#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
65#[serde(tag = "kind", rename_all = "lowercase")]
66pub enum Link {
67 Run {
69 id: String,
71 },
72 Task {
74 id: String,
76 },
77 Url {
79 url: String,
81 },
82}
83
84#[derive(Debug, Clone, Serialize, Deserialize)]
86pub struct Notice {
87 pub id: String,
89 pub key: String,
91 pub severity: Severity,
93 pub message: String,
95 #[serde(default)]
97 pub link: Option<Link>,
98 pub first_at: Timestamp,
100 pub last_at: Timestamp,
102 #[serde(default = "one")]
104 pub count: u32,
105 #[serde(default)]
107 pub read_at: Option<Timestamp>,
108 #[serde(default)]
110 pub dismissed_at: Option<Timestamp>,
111 #[serde(default)]
114 pub subjects: Vec<String>,
115 #[serde(default)]
119 pub covered_by: Option<String>,
120 #[serde(default)]
124 pub since: Option<Timestamp>,
125 #[serde(default = "schema")]
127 pub schema: u32,
128}
129
130fn one() -> u32 {
131 1
132}
133
134fn schema() -> u32 {
135 SCHEMA
136}
137
138impl Notice {
139 pub fn new(severity: Severity, key: &str, message: impl Into<String>) -> Self {
141 let now = Timestamp::now();
142 Self {
143 id: id_of(key),
144 key: key.to_owned(),
145 severity,
146 message: message.into(),
147 link: None,
148 first_at: now,
149 last_at: now,
150 count: 1,
151 read_at: None,
152 dismissed_at: None,
153 subjects: Vec::new(),
154 covered_by: None,
155 since: None,
156 schema: SCHEMA,
157 }
158 }
159
160 pub fn about<I, S>(mut self, subjects: I) -> Self
162 where
163 I: IntoIterator<Item = S>,
164 S: Into<String>,
165 {
166 self.subjects = subjects.into_iter().map(Into::into).collect();
167 self
168 }
169
170 pub fn info(key: &str, message: impl Into<String>) -> Self {
172 Self::new(Severity::Info, key, message)
173 }
174
175 pub fn warn(key: &str, message: impl Into<String>) -> Self {
177 Self::new(Severity::Warn, key, message)
178 }
179
180 pub fn error(key: &str, message: impl Into<String>) -> Self {
182 Self::new(Severity::Error, key, message)
183 }
184
185 pub fn since(mut self, at: Option<Timestamp>) -> Self {
187 self.since = at;
188 self
189 }
190
191 pub fn link(mut self, link: Link) -> Self {
193 self.link = Some(link);
194 self
195 }
196
197 pub fn unread(&self) -> bool {
199 self.read_at.is_none() && self.dismissed_at.is_none()
200 }
201
202 pub fn raise_again(&mut self, again: &Notice, now: Timestamp) -> bool {
209 self.count = self.count.saturating_add(1);
210 self.last_at = now;
211 if again.link.is_some() {
212 self.link = again.link.clone();
213 }
214 if !again.subjects.is_empty() {
215 self.subjects = again.subjects.clone();
216 }
217 let escalated = again.severity > self.severity;
218 let newer_cause = match (self.since, again.since) {
224 (Some(old), Some(new)) => new > old,
225 (None, Some(_)) => self.covered_by.is_some() && again.covered_by.is_none(),
226 _ => false,
227 };
228 if again.since.is_some() {
229 self.since = again.since;
230 }
231 let relit = again.message != self.message || escalated || newer_cause;
232 if relit {
233 self.message = again.message.clone();
234 self.severity = self.severity.max(again.severity);
235 self.read_at = None;
236 self.dismissed_at = None;
237 self.covered_by = None;
238 }
239 if let Some(q) = &again.covered_by
242 && !escalated
243 && self.unread()
244 {
245 self.read_at = Some(now);
246 self.covered_by = Some(q.clone());
247 }
248 relit && self.covered_by.is_none()
249 }
250
251 pub fn mark_read(&mut self, now: Timestamp) {
253 if self.read_at.is_none() {
254 self.read_at = Some(now);
255 }
256 }
257
258 pub fn dismiss(&mut self, now: Timestamp) {
260 self.mark_read(now);
261 if self.dismissed_at.is_none() {
262 self.dismissed_at = Some(now);
263 }
264 }
265
266 fn keep_rank(&self) -> u8 {
268 if self.dismissed_at.is_some() {
269 0
270 } else if self.read_at.is_some() {
271 1
272 } else {
273 2
274 }
275 }
276}
277
278pub fn id_of(key: &str) -> String {
282 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
283 for b in key.bytes() {
284 hash ^= u64::from(b);
285 hash = hash.wrapping_mul(0x0100_0000_01b3);
286 }
287 let slug: String = key
288 .chars()
289 .map(|c| {
290 if c.is_ascii_alphanumeric() {
291 c.to_ascii_lowercase()
292 } else {
293 '-'
294 }
295 })
296 .take(32)
297 .collect();
298 format!("{slug}-{hash:016x}")
299}
300
301fn valid_id(id: &str) -> bool {
304 !id.is_empty()
305 && id.len() <= 64
306 && id
307 .bytes()
308 .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-')
309}
310
311const LOCK_WAIT: std::time::Duration = std::time::Duration::from_secs(5);
313
314const LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
316
317static TMP_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
319
320#[derive(Debug, Clone)]
322pub struct Notices {
323 root: PathBuf,
324}
325
326impl Notices {
327 pub fn open() -> Self {
329 Self::at(crate::run::home().join("notifications"))
330 }
331
332 pub fn at(root: PathBuf) -> Self {
334 Self { root }
335 }
336
337 pub fn root(&self) -> &Path {
339 &self.root
340 }
341
342 fn path_of(&self, id: &str) -> PathBuf {
343 self.root.join(format!("{id}.json"))
344 }
345
346 fn put(&self, n: &Notice) -> Result<()> {
347 std::fs::create_dir_all(&self.root)
348 .with_context(|| format!("create {}", self.root.display()))?;
349 let body = serde_json::to_string_pretty(n).context("serialize notice")?;
350 let path = self.path_of(&n.id);
351 let tmp = path.with_extension(format!(
354 "json.{}.{}.tmp",
355 std::process::id(),
356 TMP_SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
357 ));
358 std::fs::write(&tmp, body).with_context(|| format!("write {}", tmp.display()))?;
359 if let Err(e) = std::fs::rename(&tmp, &path) {
360 let _ = std::fs::remove_file(&tmp);
361 return Err(e).with_context(|| format!("replace {}", path.display()));
362 }
363 Ok(())
364 }
365
366 pub fn get(&self, id: &str) -> Result<Notice> {
368 if !valid_id(id) {
369 bail!("`{id}` is not a notification id");
370 }
371 read_path(&self.path_of(id))
372 }
373
374 fn locked<T>(&self, id: &str, f: impl FnOnce() -> Result<T>) -> Result<T> {
383 std::fs::create_dir_all(&self.root)
384 .with_context(|| format!("create {}", self.root.display()))?;
385 let lock = self.root.join(format!("{id}.lock"));
386 let start = std::time::Instant::now();
387 let mut held = false;
388 while start.elapsed() < LOCK_WAIT {
389 match std::fs::OpenOptions::new()
390 .write(true)
391 .create_new(true)
392 .open(&lock)
393 {
394 Ok(_) => {
395 held = true;
396 break;
397 }
398 Err(_) => {
399 let stale = std::fs::metadata(&lock)
400 .and_then(|m| m.modified())
401 .ok()
402 .and_then(|t| t.elapsed().ok())
403 .is_some_and(|age| age > LOCK_STALE);
404 if stale {
405 let _ = std::fs::remove_file(&lock);
406 } else {
407 std::thread::sleep(std::time::Duration::from_millis(5));
408 }
409 }
410 }
411 }
412 let out = f();
413 if held {
414 let _ = std::fs::remove_file(&lock);
415 }
416 out
417 }
418
419 pub fn raise(&self, incoming: Notice) -> Result<Notice> {
421 self.raise_paged(incoming).map(|(n, _)| n)
422 }
423
424 pub fn raise_paged(&self, incoming: Notice) -> Result<(Notice, bool)> {
426 let out = self.locked(&incoming.id.clone(), || {
427 let now = Timestamp::now();
428 let (stored, page) = match read_path(&self.path_of(&incoming.id)) {
429 Ok(mut existing) => {
430 let page = existing.raise_again(&incoming, now);
431 (existing, page)
432 }
433 Err(_) => {
434 let mut fresh = incoming;
435 if fresh.covered_by.is_some() {
436 fresh.read_at = Some(now);
437 }
438 let page = fresh.covered_by.is_none();
439 (fresh, page)
440 }
441 };
442 self.put(&stored)?;
443 Ok((stored, page))
444 })?;
445 self.prune();
446 Ok(out)
447 }
448
449 pub fn cover(&self, question: &crate::ask::Question) {
452 for n in self.all() {
453 if n.unread() && covers(question, &n) {
454 let _ = self.update(&n.id, |n| {
455 if n.unread() && covers(question, n) {
458 n.mark_read(Timestamp::now());
459 n.covered_by = Some(question.id.clone());
460 }
461 });
462 }
463 }
464 }
465
466 fn update(&self, id: &str, change: impl FnOnce(&mut Notice)) -> Result<Notice> {
468 if !valid_id(id) {
469 bail!("`{id}` is not a notification id");
470 }
471 self.locked(id, || {
472 let mut n = read_path(&self.path_of(id))?;
473 change(&mut n);
474 self.put(&n)?;
475 Ok(n)
476 })
477 }
478
479 fn all(&self) -> Vec<Notice> {
481 let mut all: Vec<Notice> = std::fs::read_dir(&self.root)
482 .into_iter()
483 .flatten()
484 .flatten()
485 .map(|e| e.path())
486 .filter(|p| p.extension().is_some_and(|x| x == "json"))
487 .filter_map(|p| read_path(&p).ok())
488 .collect();
489 all.sort_by(|a, b| b.last_at.cmp(&a.last_at).then_with(|| a.id.cmp(&b.id)));
490 all
491 }
492
493 pub fn list(&self) -> Vec<Notice> {
495 self.all()
496 .into_iter()
497 .filter(|n| n.dismissed_at.is_none())
498 .collect()
499 }
500
501 pub fn count_unread(&self) -> usize {
503 self.all().iter().filter(|n| n.unread()).count()
504 }
505
506 pub fn mark_read(&self, id: &str) -> Result<Notice> {
508 let now = Timestamp::now();
509 self.update(id, |n| n.mark_read(now))
510 }
511
512 pub fn mark_all_read(&self) -> Result<usize> {
514 let now = Timestamp::now();
515 let mut changed = 0;
516 for n in self.all().into_iter().filter(Notice::unread) {
517 let done = self.update(&n.id, |n| n.mark_read(now));
520 if done.is_ok() {
521 changed += 1;
522 }
523 }
524 Ok(changed)
525 }
526
527 pub fn dismiss(&self, id: &str) -> Result<Notice> {
529 let now = Timestamp::now();
530 self.update(id, |n| n.dismiss(now))
531 }
532
533 fn prune(&self) {
535 let mut all = self.all();
536 if all.len() <= CAP {
537 return;
538 }
539 all.sort_by(|a, b| {
541 b.keep_rank()
542 .cmp(&a.keep_rank())
543 .then_with(|| b.last_at.cmp(&a.last_at))
544 });
545 for n in all.split_off(CAP) {
546 let _ = std::fs::remove_file(self.path_of(&n.id));
547 }
548 }
549
550 pub fn revision(&self) -> u64 {
554 use std::hash::{Hash as _, Hasher as _};
555 let mut entries: Vec<(std::ffi::OsString, u128)> = std::fs::read_dir(&self.root)
556 .into_iter()
557 .flatten()
558 .flatten()
559 .filter(|e| e.path().extension().is_some_and(|x| x == "json"))
560 .map(|e| {
561 let at = e
562 .metadata()
563 .and_then(|m| m.modified())
564 .ok()
565 .and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
566 .map_or(0, |d| d.as_nanos());
567 (e.file_name(), at)
568 })
569 .collect();
570 if entries.is_empty() {
571 return 0;
572 }
573 entries.sort();
574 let mut h = std::collections::hash_map::DefaultHasher::new();
575 entries.hash(&mut h);
576 h.finish()
577 }
578}
579
580fn read_path(path: &Path) -> Result<Notice> {
581 let body = std::fs::read_to_string(path).with_context(|| format!("read {}", path.display()))?;
582 let n: Notice =
583 serde_json::from_str(&body).with_context(|| format!("parse {}", path.display()))?;
584 if n.schema > SCHEMA {
585 bail!(
586 "notice {} was written by a newer magi (schema {}, this build speaks up to {SCHEMA})",
587 n.id,
588 n.schema
589 );
590 }
591 Ok(n)
592}
593
594pub fn run_ended(state: &crate::run::RunState) -> Option<Notice> {
600 use crate::run::RunStatus;
601 let failed = matches!(
602 state.status,
603 RunStatus::Blocked | RunStatus::Stalled | RunStatus::Failed
604 );
605 (failed && !state.parked).then(|| {
606 Notice::error(
607 &format!("run:{}", state.id),
608 format!("Run {} ended {}.", state.short(), state.status.as_str()),
609 )
610 .link(Link::Run {
611 id: state.id.clone(),
612 })
613 .about([state.id.clone()])
614 })
615}
616
617pub fn merged_red(run_id: &str, summary: &str) -> Notice {
623 Notice::warn(&format!("merged-red:{run_id}"), summary).link(Link::Run {
624 id: run_id.to_owned(),
625 })
626}
627
628pub fn task_held(task: &crate::queue::Task) -> Option<Notice> {
642 use crate::queue::{HoldSource, TaskStatus};
643 if task.status != TaskStatus::Held || task.hold_source != Some(HoldSource::Machine) {
644 return None;
645 }
646 let why = task
647 .hold_reason
648 .as_deref()
649 .or(task.last_error.as_deref())
650 .unwrap_or("no reason recorded");
651 Some(
652 Notice::warn(
653 &format!("task:{}", task.id),
654 format!("Task {} is held: {why}.", task.short()),
655 )
656 .link(Link::Task {
657 id: task.id.clone(),
658 })
659 .about([task.id.clone()])
660 .since(task.held_at),
661 )
662}
663
664pub fn run_stopped(id: &str, state: &crate::run::RunState) -> Notice {
666 Notice::error(
667 &format!("run:{id}"),
668 format!("Run {} stopped with an error.", state.short()),
669 )
670 .link(Link::Run { id: id.to_owned() })
671 .about([id.to_owned()])
672}
673
674#[derive(Debug, Clone, PartialEq, Eq)]
676pub struct Page {
677 pub run: String,
679 pub summary: String,
681}
682
683static PAGER: std::sync::OnceLock<()> = std::sync::OnceLock::new();
686
687pub fn install_pager() {
689 let _ = PAGER.set(());
690}
691
692fn send_page(page: Page, notify: Option<crate::config::Notify>) {
697 if PAGER.get().is_none() {
698 return;
699 }
700 let spawned = std::thread::Builder::new()
701 .name("notice-pager".into())
702 .spawn(move || {
703 let notify = notify.unwrap_or_else(machine_notify);
704 if notify.command.is_empty() {
705 return;
706 }
707 let rt = match tokio::runtime::Builder::new_current_thread()
708 .enable_all()
709 .build()
710 {
711 Ok(rt) => rt,
712 Err(e) => return tracing::warn!("could not page about a notification: {e:#}"),
713 };
714 if let Err(e) = rt.block_on(crate::ask::notify_text(¬ify, &page.run, &page.summary))
715 {
716 tracing::warn!("could not page about a notification: {e:#}");
717 }
718 });
719 if let Err(e) = spawned {
720 tracing::warn!("could not page about a notification: {e:#}");
721 }
722}
723
724fn machine_notify() -> crate::config::Notify {
727 let layers: Vec<PathBuf> = crate::config::Config::machine_layer()
728 .into_iter()
729 .filter(|p| p.is_file())
730 .collect();
731 match crate::config::Config::load_layers(&layers) {
732 Ok(c) => c.notify,
733 Err(e) => {
734 tracing::warn!("could not read [notify] for a notification: {e:#}");
735 Default::default()
736 }
737 }
738}
739
740pub fn raise(notice: Notice) {
744 if let Some(home) = crate::run::try_home() {
745 raise_in(&home, notice);
746 }
747}
748
749pub fn raise_with(notice: Notice, notify: &crate::config::Notify) {
751 if let Some(home) = crate::run::try_home() {
752 raise_in_with(&home, notice, &|page| send_page(page, Some(notify.clone())));
753 }
754}
755
756pub fn covers(q: &crate::ask::Question, n: &Notice) -> bool {
773 if !q.status.open() || !n.subjects.contains(&q.run) {
774 return false;
775 }
776 if n.since.is_some_and(|since| q.asked_at < since) {
777 return false;
778 }
779 let about_task = n.key.starts_with("task:") || n.key.starts_with("handover:");
780 let about_run = n.key.starts_with("run:");
781 match q.node.as_str() {
782 crate::bump::NOTICE_NODE | crate::land::APPROVAL_NODE => false,
783 crate::conduct::NODE | crate::triage::NODE | crate::triage::DEPS_NODE => about_task,
784 _ => about_run,
785 }
786}
787
788fn covering_question(home: &Path, n: &Notice) -> Option<String> {
789 if n.subjects.is_empty() {
790 return None;
791 }
792 crate::ask::Questions::at(home.join("questions"))
793 .list()
794 .into_iter()
795 .find(|q| covers(q, n))
796 .map(|q| q.id)
797}
798
799pub fn quiet_for(home: &Path, question: &crate::ask::Question) {
801 Notices::at(home.join("notifications")).cover(question);
802}
803
804pub fn raise_in(home: &Path, notice: Notice) {
807 raise_in_with(home, notice, &|page| send_page(page, None));
808}
809
810pub fn raise_in_with(home: &Path, mut notice: Notice, send: &dyn Fn(Page)) {
813 if let Some(q) = covering_question(home, ¬ice) {
814 notice.covered_by = Some(q);
815 }
816 match Notices::at(home.join("notifications")).raise_paged(notice) {
817 Ok((n, true)) => send(Page {
818 run: match &n.link {
819 Some(Link::Run { id }) => id.clone(),
820 _ => String::new(),
821 },
822 summary: n.message.clone(),
823 }),
824 Ok(_) => {}
825 Err(e) => tracing::warn!("could not file a notification: {e:#}"),
826 }
827}
828
829#[cfg(test)]
830mod tests {
831 use super::*;
832
833 fn store() -> (tempfile::TempDir, Notices) {
834 let dir = tempfile::tempdir().unwrap();
835 let s = Notices::at(dir.path().join("notifications"));
836 (dir, s)
837 }
838
839 #[test]
840 fn a_red_merge_notice_is_keyed_on_the_run_and_raised_once() {
841 let (_d, s) = store();
842 let n = merged_red("run-1", "Merged a/b PR #1 with red checks: x (u)");
843 assert_eq!(n.key, "merged-red:run-1");
844 assert_eq!(
845 n.link,
846 Some(Link::Run {
847 id: "run-1".to_owned()
848 })
849 );
850 s.raise(merged_red(
851 "run-1",
852 "Merged a/b PR #1 with red checks: x (u)",
853 ))
854 .unwrap();
855 s.raise(merged_red(
856 "run-1",
857 "Merged a/b PR #1 with red checks: x (u)",
858 ))
859 .unwrap();
860 assert_eq!(s.list().len(), 1);
861 }
862
863 #[test]
864 fn a_repeat_counts_but_stays_read_and_a_change_relights_it() {
865 let now = Timestamp::now();
866 let mut n = Notice::warn("task:1", "held");
867 n.mark_read(now);
868 n.raise_again(&Notice::warn("task:1", "held"), now);
869 assert_eq!(n.count, 2);
870 assert!(n.read_at.is_some(), "identical repeat stays read");
871 n.raise_again(&Notice::warn("task:1", "held differently"), now);
872 assert!(n.unread());
873 assert_eq!(n.message, "held differently");
874 }
875
876 #[test]
877 fn a_rise_in_severity_resurrects_a_dismissed_notice_but_a_repeat_does_not() {
878 let now = Timestamp::now();
879 let mut n = Notice::warn("task:x", "m");
880 n.dismiss(now);
881 n.raise_again(&Notice::warn("task:x", "m"), now);
882 assert!(n.dismissed_at.is_some(), "tombstone holds");
883 n.raise_again(&Notice::error("k", "m"), now);
884 assert!(n.unread());
885 assert_eq!(n.severity, Severity::Error);
886 n.raise_again(&Notice::info("k", "other"), now);
888 assert_eq!(n.severity, Severity::Error);
889 }
890
891 #[test]
892 fn identical_raises_share_one_file() {
893 let (_d, s) = store();
894 for _ in 0..5 {
895 s.raise(Notice::error("run:abc", "blocked")).unwrap();
896 }
897 let all = s.list();
898 assert_eq!(all.len(), 1);
899 assert_eq!(all[0].count, 5);
900 assert_eq!(s.count_unread(), 1);
901 }
902
903 #[test]
904 fn transitions_persist_and_dismissed_leave_the_list() {
905 let (_d, s) = store();
906 let a = s.raise(Notice::info("a", "one")).unwrap();
907 let b = s.raise(Notice::warn("b", "two")).unwrap();
908 assert_eq!(s.count_unread(), 2);
909 s.mark_read(&a.id).unwrap();
910 assert_eq!(s.count_unread(), 1);
911 s.dismiss(&b.id).unwrap();
912 assert_eq!(s.count_unread(), 0);
913 assert_eq!(s.list().len(), 1);
914 s.raise(Notice::warn("b", "two")).unwrap();
915 assert_eq!(
916 s.list().len(),
917 1,
918 "a tombstone survives a same-message raise"
919 );
920 s.raise(Notice::info("c", "three")).unwrap();
921 assert_eq!(s.mark_all_read().unwrap(), 1);
922 assert_eq!(s.count_unread(), 0);
923 }
924
925 #[test]
926 fn the_list_is_newest_first() {
927 let (_d, s) = store();
928 let mut old = Notice::info("old", "old");
929 old.last_at = "2020-01-01T00:00:00Z".parse().unwrap();
930 s.put(&old).unwrap();
931 s.raise(Notice::info("new", "new")).unwrap();
932 let keys: Vec<_> = s.list().into_iter().map(|n| n.key).collect();
933 assert_eq!(keys, ["new", "old"]);
934 }
935
936 #[test]
937 fn the_cap_prunes_dismissed_then_read_then_oldest() {
938 let (_d, s) = store();
939 let keep = s.raise(Notice::error("keep", "unread")).unwrap();
940 let gone = s.raise(Notice::info("gone", "dismissed")).unwrap();
941 s.dismiss(&gone.id).unwrap();
942 for i in 0..CAP - 1 {
943 s.raise(Notice::info(&format!("k{i}"), "x")).unwrap();
944 }
945 let files = std::fs::read_dir(s.root()).unwrap().count();
946 assert_eq!(files, CAP);
947 assert!(
948 s.get(&keep.id).is_ok(),
949 "an unread notice outlives a dismissed one"
950 );
951 assert!(s.get(&gone.id).is_err(), "the dismissed one went first");
952 }
953
954 #[test]
955 fn revision_moves_on_write_and_on_removal() {
956 let (_d, s) = store();
957 assert_eq!(s.revision(), 0);
958 let a = s.raise(Notice::info("a", "m")).unwrap();
959 let r1 = s.revision();
960 assert_ne!(r1, 0);
961 s.raise(Notice::info("b", "m")).unwrap();
962 let r2 = s.revision();
963 assert_ne!(r1, r2);
964 std::fs::remove_file(s.path_of(&a.id)).unwrap();
965 assert_ne!(s.revision(), r2);
966 }
967
968 #[test]
969 fn ids_are_stable_and_untrusted_ids_are_refused() {
970 assert_eq!(id_of("run:1"), id_of("run:1"));
971 assert_ne!(id_of("run:1"), id_of("run-1"));
972 assert!(valid_id(&id_of("release-bump:20260101-abc")));
973 let (_d, s) = store();
974 for bad in ["", "../x", "a/b", "A", "x.json"] {
975 assert!(s.get(bad).is_err(), "{bad}");
976 }
977 }
978
979 #[test]
980 fn a_newer_schema_is_refused_and_an_older_reads() {
981 let (_d, s) = store();
982 let mut n = Notice::info("k", "m");
983 n.schema = SCHEMA + 1;
984 s.put(&n).unwrap();
985 assert!(s.get(&n.id).is_err());
986 let old = r#"{"id":"x-1","key":"x","severity":"warn","message":"m",
987 "first_at":"2020-01-01T00:00:00Z","last_at":"2020-01-01T00:00:00Z"}"#;
988 std::fs::write(s.path_of("x-1"), old).unwrap();
989 assert_eq!(s.get("x-1").unwrap().count, 1);
990 }
991
992 fn state(status: crate::run::RunStatus) -> crate::run::RunState {
993 let mut st = crate::run::RunState::new(
994 std::path::PathBuf::from("/repo"),
995 "main".to_owned(),
996 "abc1234def".to_owned(),
997 "task".to_owned(),
998 crate::config::Config::default(),
999 );
1000 st.status = status;
1001 st
1002 }
1003
1004 #[test]
1005 fn a_run_that_ended_badly_is_news_unless_it_only_parked() {
1006 use crate::run::RunStatus;
1007 for bad in [RunStatus::Blocked, RunStatus::Stalled, RunStatus::Failed] {
1008 let st = state(bad);
1009 let n = run_ended(&st).expect("news");
1010 assert_eq!(n.key, format!("run:{}", st.id));
1011 assert_eq!(n.severity, Severity::Error);
1012 }
1013 assert!(run_ended(&state(RunStatus::Merged)).is_none());
1014 let mut parked = state(RunStatus::Stalled);
1015 parked.parked = true;
1016 assert!(run_ended(&parked).is_none());
1017 }
1018
1019 #[test]
1020 fn concurrent_writers_of_one_key_all_succeed() {
1021 let (_d, s) = store();
1022 let handles: Vec<_> = (0..8)
1023 .map(|_| {
1024 let s = s.clone();
1025 std::thread::spawn(move || {
1026 for _ in 0..20 {
1027 s.raise(Notice::warn("same", "m")).unwrap();
1028 }
1029 })
1030 })
1031 .collect();
1032 for h in handles {
1033 h.join().unwrap();
1034 }
1035 assert_eq!(s.list().len(), 1);
1036 let stray = std::fs::read_dir(s.root())
1037 .unwrap()
1038 .flatten()
1039 .filter(|e| e.path().extension().is_some_and(|x| x == "tmp"))
1040 .count();
1041 assert_eq!(stray, 0);
1042 }
1043
1044 #[test]
1045 fn a_dead_holders_lock_is_broken_and_no_lock_is_left_behind() {
1046 let (_d, s) = store();
1047 let n = s.raise(Notice::info("k", "m")).unwrap();
1048 let lock = s.root().join(format!("{}.lock", n.id));
1049 std::fs::write(&lock, "").unwrap();
1050 let old = std::time::SystemTime::now() - std::time::Duration::from_secs(60);
1051 std::fs::File::options()
1052 .write(true)
1053 .open(&lock)
1054 .unwrap()
1055 .set_modified(old)
1056 .unwrap();
1057 s.mark_read(&n.id).unwrap();
1058 assert!(!lock.exists());
1059 assert!(s.get(&n.id).unwrap().read_at.is_some());
1060 }
1061
1062 #[test]
1063 fn only_a_machine_hold_is_a_task_notice() {
1064 let mut t = crate::queue::Task::new(
1065 "t".to_owned(),
1066 "do it".to_owned(),
1067 std::path::PathBuf::from("/repo"),
1068 crate::queue::Source::Human,
1069 );
1070 assert!(task_held(&t).is_none());
1071 t.hold_manual(None);
1072 assert!(task_held(&t).is_none(), "the operator's own hold");
1073 t.hold_machine(Some("missing blocker".to_owned()));
1074 let n = task_held(&t).expect("machine hold");
1075 assert_eq!(n.key, format!("task:{}", t.id));
1076 assert!(n.message.contains("missing blocker"));
1077 }
1078
1079 fn held_task() -> crate::queue::Task {
1080 let mut t = crate::queue::Task::new(
1081 "t".to_owned(),
1082 "do it".to_owned(),
1083 std::path::PathBuf::from("/repo"),
1084 crate::queue::Source::Human,
1085 );
1086 t.hold_machine(Some("branch b is checked out".to_owned()));
1087 t
1088 }
1089
1090 fn ask(home: &Path, run: &str, node: &str) -> crate::ask::Question {
1091 let mut q = crate::ask::Question::new(
1092 run.to_owned(),
1093 node.to_owned(),
1094 "conductor".to_owned(),
1095 "cannot resume".to_owned(),
1096 String::new(),
1097 Vec::new(),
1098 );
1099 crate::ask::Questions::at(home.join("questions"))
1100 .put(&mut q)
1101 .unwrap();
1102 q
1103 }
1104
1105 fn unread(home: &Path) -> usize {
1106 Notices::at(home.join("notifications")).count_unread()
1107 }
1108
1109 #[test]
1110 fn a_hold_with_an_open_question_for_the_task_pages_once() {
1111 let dir = tempfile::tempdir().unwrap();
1112 let mut t = held_task();
1113 let q = ask(dir.path(), &t.id, "conduct");
1114 crate::queue::Queue::at(dir.path().join("queue"))
1115 .put(&mut t)
1116 .unwrap();
1117 let s = Notices::at(dir.path().join("notifications"));
1118 assert_eq!(s.count_unread(), 0);
1119 let n = s.list().pop().expect("the record is kept");
1120 assert_eq!(n.covered_by.as_deref(), Some(q.id.as_str()));
1121 assert!(n.message.contains("checked out"));
1122 }
1123
1124 #[test]
1125 fn a_refused_handover_without_a_question_pages_once() {
1126 let dir = tempfile::tempdir().unwrap();
1127 let q = crate::queue::Queue::at(dir.path().join("queue"));
1128 let mut t = held_task();
1129 t.hold_for_handover(Some("magi/x/A".to_owned()), "checked out".to_owned());
1130 q.put(&mut t).unwrap();
1131 q.put(&mut t).unwrap();
1132 let s = Notices::at(dir.path().join("notifications"));
1133 assert_eq!(s.list().len(), 1);
1134 assert_eq!(s.list()[0].key, format!("task:{}", t.id));
1135 assert_eq!(s.count_unread(), 1);
1136 }
1137
1138 #[test]
1139 fn an_older_question_does_not_silence_a_new_hold() {
1140 for node in ["deps", "conduct", "triage"] {
1141 let dir = tempfile::tempdir().unwrap();
1142 let mut t = crate::queue::Task::new(
1143 "t".to_owned(),
1144 "do it".to_owned(),
1145 std::path::PathBuf::from("/repo"),
1146 crate::queue::Source::Human,
1147 );
1148 let old = ask(dir.path(), &t.id, node);
1149 std::thread::sleep(std::time::Duration::from_millis(5));
1150 t.hold_machine(Some("new failure".to_owned()));
1151 crate::queue::Queue::at(dir.path().join("queue"))
1152 .put(&mut t)
1153 .unwrap();
1154 let s = Notices::at(dir.path().join("notifications"));
1155 let n = s.list().pop().unwrap();
1156 assert_eq!(n.covered_by, None, "{node} question {}", old.id);
1157 assert_eq!(s.count_unread(), 1);
1158 }
1159 }
1160
1161 #[test]
1162 fn a_legacy_covered_notice_pages_for_a_new_uncovered_hold() {
1163 let (_d, s) = store();
1164 let mut legacy = Notice::warn("task:x", "held");
1165 legacy.covered_by = Some("q1".to_owned());
1166 let id = legacy.id.clone();
1167 s.raise(legacy).unwrap();
1168 assert!(!s.get(&id).unwrap().unread());
1169 s.raise(Notice::warn("task:x", "held").since(Some(Timestamp::now())))
1170 .unwrap();
1171 let n = s.get(&id).unwrap();
1172 assert!(n.unread());
1173 assert_eq!(n.covered_by, None);
1174 }
1175
1176 #[test]
1177 fn a_newer_cause_relights_a_dismissed_notice() {
1178 let (_d, s) = store();
1179 let t0 = Timestamp::now();
1180 let first = Notice::warn("task:x", "held").since(Some(t0));
1181 let id = first.id.clone();
1182 s.raise(first).unwrap();
1183 s.dismiss(&id).unwrap();
1184 s.raise(Notice::warn("task:x", "held").since(Some(t0)))
1185 .unwrap();
1186 assert!(!s.get(&id).unwrap().unread(), "same cause stays dismissed");
1187 let later = t0 + std::time::Duration::from_secs(1);
1188 s.raise(Notice::warn("task:x", "held").since(Some(later)))
1189 .unwrap();
1190 assert!(s.get(&id).unwrap().unread());
1191 }
1192
1193 #[test]
1194 fn a_hold_with_no_question_still_notifies() {
1195 let dir = tempfile::tempdir().unwrap();
1196 let mut t = held_task();
1197 crate::queue::Queue::at(dir.path().join("queue"))
1198 .put(&mut t)
1199 .unwrap();
1200 assert_eq!(unread(dir.path()), 1);
1201 }
1202
1203 #[test]
1204 fn a_question_for_another_task_does_not_suppress() {
1205 let dir = tempfile::tempdir().unwrap();
1206 let mut t = held_task();
1207 ask(dir.path(), "some-other-task", "conduct");
1208 crate::queue::Queue::at(dir.path().join("queue"))
1209 .put(&mut t)
1210 .unwrap();
1211 assert_eq!(unread(dir.path()), 1);
1212 }
1213
1214 #[test]
1215 fn a_question_filed_after_the_hold_quiets_it_once() {
1216 let dir = tempfile::tempdir().unwrap();
1217 let mut t = held_task();
1218 crate::queue::Queue::at(dir.path().join("queue"))
1219 .put(&mut t)
1220 .unwrap();
1221 assert_eq!(unread(dir.path()), 1);
1222 let q = ask(dir.path(), &t.id, "conduct");
1223 assert_eq!(unread(dir.path()), 0);
1224 let s = Notices::at(dir.path().join("notifications"));
1226 let id = s.list()[0].id.clone();
1227 s.update(&id, |n| n.read_at = None).unwrap();
1228 let mut again = crate::ask::Questions::at(dir.path().join("questions"))
1229 .get(&q.id)
1230 .unwrap();
1231 crate::ask::Questions::at(dir.path().join("questions"))
1232 .put(&mut again)
1233 .unwrap();
1234 assert_eq!(unread(dir.path()), 1);
1235 }
1236
1237 #[test]
1238 fn a_blocked_run_with_an_open_question_is_quiet() {
1239 let dir = tempfile::tempdir().unwrap();
1240 ask(dir.path(), "run-1", "implement");
1241 raise_in(
1242 dir.path(),
1243 Notice::error("run:run-1", "Run r ended blocked.").about(["run-1"]),
1244 );
1245 assert_eq!(unread(dir.path()), 0);
1246 raise_in(
1247 dir.path(),
1248 Notice::error("run:run-2", "Run r ended blocked.").about(["run-2"]),
1249 );
1250 assert_eq!(unread(dir.path()), 1);
1251 }
1252
1253 #[test]
1254 fn escalation_makes_a_covered_notice_unread_again() {
1255 let dir = tempfile::tempdir().unwrap();
1256 ask(dir.path(), "t1", "conduct");
1257 raise_in(dir.path(), Notice::warn("task:t1", "held").about(["t1"]));
1258 assert_eq!(unread(dir.path()), 0);
1259 raise_in(dir.path(), Notice::error("task:t1", "held").about(["t1"]));
1260 assert_eq!(unread(dir.path()), 1);
1261 }
1262
1263 #[test]
1264 fn covers_matches_run_exactly_and_ignores_release_questions() {
1265 let dir = tempfile::tempdir().unwrap();
1266 let n = Notice::warn("task:x", "m").about(["t1"]);
1267 let q = ask(dir.path(), "t1", "conduct");
1268 assert!(covers(&q, &n));
1269 assert!(!covers(&q, &Notice::warn("task:x", "m")));
1270 assert!(!covers(&q, &Notice::warn("task:x", "m").about(["t"])));
1271 let release = ask(dir.path(), "t1", crate::bump::NOTICE_NODE);
1272 assert!(!covers(&release, &n));
1273 }
1274
1275 #[test]
1276 fn a_run_question_does_not_swallow_a_task_hold_and_vice_versa() {
1277 let dir = tempfile::tempdir().unwrap();
1278 ask(dir.path(), "t1", "implement");
1279 raise_in(dir.path(), Notice::warn("task:t1", "held").about(["t1"]));
1280 assert_eq!(unread(dir.path()), 1);
1281 ask(dir.path(), "r1", "conduct");
1282 raise_in(dir.path(), Notice::error("run:r1", "ended").about(["r1"]));
1283 assert_eq!(unread(dir.path()), 2);
1284 }
1285
1286 #[test]
1287 fn a_conduct_question_covers_the_hold_however_late_it_was_filed() {
1288 let dir = tempfile::tempdir().unwrap();
1289 let mut q = crate::ask::Question::new(
1290 "t1".to_owned(),
1291 "conduct".to_owned(),
1292 "conductor".to_owned(),
1293 "s".to_owned(),
1294 String::new(),
1295 Vec::new(),
1296 );
1297 q.asked_at = Timestamp::from_second(Timestamp::now().as_second() - 3600).unwrap();
1298 crate::ask::Questions::at(dir.path().join("questions"))
1299 .put(&mut q)
1300 .unwrap();
1301 raise_in(dir.path(), Notice::warn("task:t1", "held").about(["t1"]));
1302 assert_eq!(unread(dir.path()), 0);
1303 }
1304
1305 fn pages(home: &Path, n: Notice) -> Vec<Page> {
1306 let sent = std::cell::RefCell::new(Vec::new());
1307 raise_in_with(home, n, &|p| sent.borrow_mut().push(p));
1308 sent.into_inner()
1309 }
1310
1311 #[test]
1312 fn a_page_fires_on_new_changed_and_escalated_only() {
1313 let dir = tempfile::tempdir().unwrap();
1314 let h = dir.path();
1315 let n = || Notice::warn("run:p1", "stopped").link(Link::Run { id: "p1".into() });
1316 let first = pages(h, n());
1317 assert_eq!(
1318 first,
1319 vec![Page {
1320 run: "p1".into(),
1321 summary: "stopped".into()
1322 }]
1323 );
1324 assert!(pages(h, n()).is_empty(), "identical re-raise is silent");
1325 assert_eq!(pages(h, Notice::warn("run:p1", "stopped again")).len(), 1);
1326 assert_eq!(pages(h, Notice::error("run:p1", "stopped again")).len(), 1);
1327 assert!(pages(h, Notice::error("run:p1", "stopped again")).is_empty());
1328 }
1329
1330 #[test]
1331 fn a_non_run_link_leaves_run_empty() {
1332 let dir = tempfile::tempdir().unwrap();
1333 let n = Notice::warn("task:t1", "held").link(Link::Task { id: "t1".into() });
1334 assert_eq!(pages(dir.path(), n)[0].run, "");
1335 }
1336
1337 #[test]
1338 fn a_dismissed_tombstone_stays_silent_until_it_changes() {
1339 let dir = tempfile::tempdir().unwrap();
1340 let h = dir.path();
1341 pages(h, Notice::warn("disk:r", "low"));
1342 Notices::at(h.join("notifications"))
1343 .dismiss(&id_of("disk:r"))
1344 .unwrap();
1345 assert!(pages(h, Notice::warn("disk:r", "low")).is_empty());
1346 assert_eq!(pages(h, Notice::warn("disk:r", "lower")).len(), 1);
1347 }
1348
1349 #[test]
1350 fn a_newer_cause_pages_again() {
1351 let dir = tempfile::tempdir().unwrap();
1352 let h = dir.path();
1353 let at = |s: i64| Some(Timestamp::from_second(s).unwrap());
1354 let n = |s| Notice::warn("task:t2", "held").since(at(s));
1355 assert_eq!(pages(h, n(100)).len(), 1);
1356 assert!(pages(h, n(100)).is_empty());
1357 assert_eq!(pages(h, n(200)).len(), 1);
1358 }
1359
1360 #[test]
1361 fn a_covered_notice_does_not_page() {
1362 let dir = tempfile::tempdir().unwrap();
1363 let h = dir.path();
1364 let q = ask(h, "r9", "implement");
1365 let n = || Notice::error("run:r9", "ended").about(["r9".to_owned()]);
1366 assert!(pages(h, n()).is_empty());
1367 let stored = Notices::at(h.join("notifications"))
1368 .get(&id_of("run:r9"))
1369 .unwrap();
1370 assert_eq!(stored.covered_by.as_deref(), Some(q.id.as_str()));
1371 assert!(pages(h, n()).is_empty());
1372 }
1373
1374 #[test]
1375 fn a_failing_sender_path_never_propagates() {
1376 let notify = crate::config::Notify {
1378 command: vec!["magi-no-such-program-xyz".into()],
1379 };
1380 let rt = tokio::runtime::Builder::new_current_thread()
1381 .enable_all()
1382 .build()
1383 .unwrap();
1384 let r = rt.block_on(crate::ask::notify_text(¬ify, "", "x"));
1385 assert!(r.is_err());
1386 send_page(
1388 Page {
1389 run: String::new(),
1390 summary: "x".into(),
1391 },
1392 Some(notify),
1393 );
1394 }
1395}