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 = "schema")]
122 pub schema: u32,
123}
124
125fn one() -> u32 {
126 1
127}
128
129fn schema() -> u32 {
130 SCHEMA
131}
132
133impl Notice {
134 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 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 pub fn info(key: &str, message: impl Into<String>) -> Self {
166 Self::new(Severity::Info, key, message)
167 }
168
169 pub fn warn(key: &str, message: impl Into<String>) -> Self {
171 Self::new(Severity::Warn, key, message)
172 }
173
174 pub fn error(key: &str, message: impl Into<String>) -> Self {
176 Self::new(Severity::Error, key, message)
177 }
178
179 pub fn link(mut self, link: Link) -> Self {
181 self.link = Some(link);
182 self
183 }
184
185 pub fn unread(&self) -> bool {
187 self.read_at.is_none() && self.dismissed_at.is_none()
188 }
189
190 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 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 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 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 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
247pub 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
270fn 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
280const LOCK_WAIT: std::time::Duration = std::time::Duration::from_secs(5);
282
283const LOCK_STALE: std::time::Duration = std::time::Duration::from_secs(10);
285
286static TMP_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
288
289#[derive(Debug, Clone)]
291pub struct Notices {
292 root: PathBuf,
293}
294
295impl Notices {
296 pub fn open() -> Self {
298 Self::at(crate::run::home().join("notifications"))
299 }
300
301 pub fn at(root: PathBuf) -> Self {
303 Self { root }
304 }
305
306 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 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 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 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 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 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 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 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 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 pub fn count_unread(&self) -> usize {
464 self.all().iter().filter(|n| n.unread()).count()
465 }
466
467 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 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 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 pub fn dismiss(&self, id: &str) -> Result<Notice> {
490 let now = Timestamp::now();
491 self.update(id, |n| n.dismiss(now))
492 }
493
494 fn prune(&self) {
496 let mut all = self.all();
497 if all.len() <= CAP {
498 return;
499 }
500 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 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
555pub 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
578pub 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
589pub 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
624pub 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
634pub fn raise(notice: Notice) {
638 if let Some(home) = crate::run::try_home() {
639 raise_in(&home, notice);
640 }
641}
642
643pub 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
677pub fn quiet_for(home: &Path, question: &crate::ask::Question) {
679 Notices::at(home.join("notifications")).cover(question);
680}
681
682pub fn raise_in(home: &Path, mut notice: Notice) {
685 if let Some(q) = covering_question(home, ¬ice) {
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 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 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}