1use std::collections::{BTreeMap, VecDeque};
12use std::path::{Path, PathBuf};
13use std::sync::atomic::{AtomicBool, Ordering};
14use std::sync::mpsc::{Receiver, SyncSender, TrySendError};
15use std::sync::{Arc, Mutex};
16use std::time::Duration;
17
18use serde::{Deserialize, Serialize};
19use serde_json::{Value, json};
20
21use super::state::{NEVER_UP_MS, Notify, State, Status, Thresholds};
22use super::store::{Check, Db};
23use super::target::{self, Ctx, Outcome};
24use super::{Kind, MAX_PER_ORG, Monitor, Settings, auto, dir, human, read_json, write_json};
25use crate::app::Apps;
26use crate::error::{Error, Result};
27use crate::org::OrgId;
28use crate::secrets::Secrets;
29use crate::stack::controller::now_ms;
30
31pub const WORKERS: usize = 8;
33pub const QUEUE: usize = 256;
35const AUTO_SYNC_MS: u64 = 60_000;
37const ROLLUP_MS: u64 = 300_000;
39const DETAILS_KEPT: usize = 256;
41
42pub const WAITING_FOR_APP: &str = "waiting for the app's first live deployment";
44pub const WAITING_FOR_SERVICE: &str = "waiting for the service's first replica in rotation";
46
47pub type AllowPrivate = Arc<dyn Fn() -> bool + Send + Sync>;
49
50#[derive(Debug, Clone, Default, Serialize, Deserialize)]
52#[serde(default)]
53pub struct Stored {
54 pub state: State,
55 pub last: Option<Outcome>,
56}
57
58#[derive(Debug, Clone, Copy, Default)]
59struct Slot {
60 next: u64,
61 running: bool,
62}
63
64struct Job {
65 org: OrgId,
66 monitor: Monitor,
67}
68
69type DetailKey = (OrgId, String, String);
70
71struct Inner {
72 state: PathBuf,
73 apps: Apps,
74 secrets: Arc<Secrets>,
75 allow_private: AllowPrivate,
76 public_url: Option<String>,
77 edit: Mutex<()>,
79 dbs: Mutex<BTreeMap<OrgId, Arc<Mutex<Db>>>>,
80 slots: Mutex<BTreeMap<(OrgId, String), Slot>>,
81 details: Mutex<VecDeque<(DetailKey, Value)>>,
82 stop: Arc<AtomicBool>,
83 tls: Arc<rustls::ClientConfig>,
84}
85
86#[derive(Clone)]
88pub struct Monitors {
89 inner: Arc<Inner>,
90}
91
92pub(super) fn fnv(s: &str) -> u64 {
93 s.bytes().fold(0xcbf2_9ce4_8422_2325, |h, b| {
94 (h ^ u64::from(b)).wrapping_mul(0x0100_0000_01b3)
95 })
96}
97
98fn jitter(max_ms: u64) -> i64 {
100 if max_ms == 0 {
101 return 0;
102 }
103 let n = std::time::SystemTime::now()
104 .duration_since(std::time::UNIX_EPOCH)
105 .map(|d| d.subsec_nanos())
106 .unwrap_or(0);
107 let r = fnv(&n.to_string()) % (2 * max_ms + 1);
108 r as i64 - max_ms as i64
109}
110
111impl Monitors {
112 pub fn new(
113 state: &Path,
114 apps: Apps,
115 secrets: Arc<Secrets>,
116 allow_private: AllowPrivate,
117 public_url: Option<String>,
118 ) -> Monitors {
119 Monitors {
120 inner: Arc::new(Inner {
121 state: state.to_path_buf(),
122 apps,
123 secrets,
124 allow_private,
125 public_url: public_url.map(|u| u.trim_end_matches('/').to_string()),
126 edit: Mutex::new(()),
127 dbs: Mutex::new(BTreeMap::new()),
128 slots: Mutex::new(BTreeMap::new()),
129 details: Mutex::new(VecDeque::new()),
130 stop: Arc::new(AtomicBool::new(false)),
131 tls: crate::net::default_tls(),
132 }),
133 }
134 }
135
136 pub fn with_tls(mut self, tls: Arc<rustls::ClientConfig>) -> Monitors {
138 if let Some(i) = Arc::get_mut(&mut self.inner) {
139 i.tls = tls;
140 }
141 self
142 }
143
144 fn path(&self, org: &OrgId, file: &str) -> PathBuf {
145 dir(&self.inner.state, org).join(file)
146 }
147
148 pub(crate) fn db(&self, org: &OrgId) -> Result<Arc<Mutex<Db>>> {
149 let mut dbs = self.inner.dbs.lock().unwrap();
150 if let Some(d) = dbs.get(org) {
151 return Ok(d.clone());
152 }
153 let d = Arc::new(Mutex::new(Db::open(&self.path(org, "monitors.db"))?));
154 dbs.insert(org.clone(), d.clone());
155 Ok(d)
156 }
157
158 pub fn orgs(&self) -> Vec<OrgId> {
160 let mut out = vec![OrgId::default_org()];
161 if let Ok(rd) = std::fs::read_dir(self.inner.state.join("orgs")) {
162 for e in rd.flatten() {
163 if let Some(o) = e.file_name().to_str().and_then(|n| OrgId::new(n).ok()) {
164 if !out.contains(&o) {
165 out.push(o);
166 }
167 }
168 }
169 }
170 out
171 }
172
173 pub fn list(&self, org: &OrgId) -> Result<Vec<Monitor>> {
176 read_json(&self.path(org, "monitors.json"))
177 }
178
179 fn save(&self, org: &OrgId, all: &[Monitor]) -> Result<()> {
180 write_json(&self.path(org, "monitors.json"), &all)
181 }
182
183 pub fn get(&self, org: &OrgId, name: &str) -> Result<Monitor> {
184 self.list(org)?
185 .into_iter()
186 .find(|m| m.name == name)
187 .ok_or_else(|| Error::NotFound(format!("monitor {name}")))
188 }
189
190 pub fn settings(&self, org: &OrgId) -> Result<Settings> {
191 read_json(&self.path(org, "settings.json"))
192 }
193
194 pub fn set_settings(&self, org: &OrgId, s: &Settings) -> Result<()> {
195 for a in &s.exclude_apps {
196 crate::app::validate_app_name(a)?;
197 }
198 for x in &s.exclude_services {
199 let ok = x.split_once('/').is_some_and(|(st, sv)| {
200 crate::stack::validate_stack_name(st).is_ok()
201 && super::validate_service_name(sv).is_ok()
202 });
203 if !ok {
204 return Err(Error::invalid(format!(
205 "exclude_services: {x:?} is not <stack>/<service>"
206 )));
207 }
208 }
209 let _g = self.inner.edit.lock().unwrap();
210 write_json(&self.path(org, "settings.json"), s)
211 }
212
213 fn check_refs(&self, org: &OrgId, m: &Monitor) -> Result<()> {
216 if m.kind == Kind::App {
217 let a = m.app.as_deref().unwrap_or_default();
218 self.inner.apps.get(org, a)?;
219 }
220 if m.kind == Kind::Service {
221 let (st, sv) = (
222 m.stack.as_deref().unwrap_or_default(),
223 m.service.as_deref().unwrap_or_default(),
224 );
225 let d = self
226 .inner
227 .apps
228 .controller()
229 .definition(&crate::stack::qualified(org, st))
230 .map_err(|_| Error::NotFound(format!("stack {st} in org {org}")))?;
231 if !d.file.services.contains_key(sv) {
232 return Err(Error::NotFound(format!("service {sv} in stack {st}")));
233 }
234 }
235 for s in m.secrets() {
236 self.inner.secrets.inspect(org, s).map_err(|e| {
237 if e.is_not_found() {
238 Error::invalid(format!("no secret {s} in org {org}: create it first"))
239 } else {
240 e
241 }
242 })?;
243 }
244 Ok(())
245 }
246
247 pub fn create(&self, org: &OrgId, mut m: Monitor) -> Result<Monitor> {
248 m.validate()?;
249 self.check_refs(org, &m)?;
250 let _g = self.inner.edit.lock().unwrap();
251 let mut all = self.list(org)?;
252 if all.iter().any(|x| x.name == m.name) {
253 return Err(Error::invalid(format!("monitor {} exists", m.name)));
254 }
255 if all.len() >= MAX_PER_ORG {
256 return Err(Error::invalid(format!(
257 "at most {MAX_PER_ORG} monitors per org"
258 )));
259 }
260 let now = now_ms() / 1000;
261 (m.created_at, m.updated_at) = (now, now);
262 all.push(m.clone());
263 self.save(org, &all)?;
264 self.due_soon(org, &m.name);
265 Ok(m)
266 }
267
268 pub fn update(
271 &self,
272 org: &OrgId,
273 name: &str,
274 patch: serde_json::Map<String, Value>,
275 ) -> Result<Monitor> {
276 if patch.get("name").is_some_and(|n| n.as_str() != Some(name)) {
277 return Err(Error::invalid(
278 "a monitor's name cannot change; create another",
279 ));
280 }
281 let _g = self.inner.edit.lock().unwrap();
282 let mut all = self.list(org)?;
283 let cur = all
284 .iter_mut()
285 .find(|m| m.name == name)
286 .ok_or_else(|| Error::NotFound(format!("monitor {name}")))?;
287 let mut v = serde_json::to_value(&*cur)?;
288 let o = v.as_object_mut().expect("a monitor is an object");
289 for (k, x) in patch {
290 if matches!(k.as_str(), "created_at" | "updated_at" | "auto") {
291 continue;
292 }
293 if x.is_null() {
294 o.remove(&k);
295 } else {
296 o.insert(k, x);
297 }
298 }
299 let mut n: Monitor =
300 serde_json::from_value(v).map_err(|e| Error::invalid(format!("bad monitor: {e}")))?;
301 n.validate()?;
302 self.check_refs(org, &n)?;
303 n.updated_at = now_ms() / 1000;
304 *cur = n.clone();
305 self.save(org, &all)?;
306 drop(_g);
307 self.edit_state(org, name, State::reset_counts)?;
308 self.due_soon(org, name);
309 Ok(n)
310 }
311
312 pub fn delete(&self, org: &OrgId, name: &str) -> Result<Monitor> {
315 let _g = self.inner.edit.lock().unwrap();
316 let mut all = self.list(org)?;
317 let i = all
318 .iter()
319 .position(|m| m.name == name)
320 .ok_or_else(|| Error::NotFound(format!("monitor {name}")))?;
321 let m = all.remove(i);
322 self.save(org, &all)?;
323 if m.auto {
324 let mut s = self.settings(org)?;
325 let (list, x) = match m.kind {
326 Kind::Service => (
327 &mut s.exclude_services,
328 auto::exclusion(
329 m.stack.as_deref().unwrap_or_default(),
330 m.service.as_deref().unwrap_or_default(),
331 ),
332 ),
333 _ => (&mut s.exclude_apps, m.app.clone().unwrap_or_default()),
334 };
335 if !list.contains(&x) {
336 list.push(x);
337 write_json(&self.path(org, "settings.json"), &s)?;
338 }
339 }
340 drop(_g);
341 self.db(org)?.lock().unwrap().forget(name)?;
342 self.inner
343 .slots
344 .lock()
345 .unwrap()
346 .remove(&(org.clone(), name.to_string()));
347 Ok(m)
348 }
349
350 pub fn set_paused(&self, org: &OrgId, name: &str, paused: bool) -> Result<Monitor> {
351 let mut p = serde_json::Map::new();
352 p.insert("paused".into(), json!(paused));
353 self.update(org, name, p)
354 }
355
356 fn edit_state(&self, org: &OrgId, name: &str, f: impl FnOnce(&mut State)) -> Result<()> {
357 let db = self.db(org)?;
358 let db = db.lock().unwrap();
359 let mut s: Stored = db.load_state(name)?.unwrap_or_default();
360 f(&mut s.state);
361 db.save_state(name, &s)
362 }
363
364 pub(crate) fn stored(&self, org: &OrgId, name: &str) -> Result<Stored> {
365 Ok(self
366 .db(org)?
367 .lock()
368 .unwrap()
369 .load_state(name)?
370 .unwrap_or_default())
371 }
372
373 fn due_soon(&self, org: &OrgId, name: &str) {
374 let mut slots = self.inner.slots.lock().unwrap();
375 let s = slots.entry((org.clone(), name.to_string())).or_default();
376 s.next = now_ms() + 1000;
377 }
378
379 pub fn sync_auto(&self, org: &OrgId) -> Result<()> {
385 let settings = self.settings(org)?;
386 let mut cands = auto::apps(&self.inner.apps, org)?;
387 cands.extend(auto::stack_services(&self.inner.apps, org));
388 let _g = self.inner.edit.lock().unwrap();
389 let mut all = self.list(org)?;
390 let (gone, added) = auto::reconcile(&mut all, &settings, &cands, now_ms() / 1000);
391 if !gone.is_empty() || !added.is_empty() {
392 self.save(org, &all)?;
393 }
394 drop(_g);
395 for n in &gone {
396 self.db(org)?.lock().unwrap().forget(n)?;
397 }
398 for n in &added {
399 eprintln!("isb serve: monitor: {org}/{n}: watching its domain");
400 self.due_soon(org, n);
401 }
402 Ok(())
403 }
404
405 pub fn start(&self) {
409 let (tx, rx) = std::sync::mpsc::sync_channel::<Job>(QUEUE);
410 let rx = Arc::new(Mutex::new(rx));
411 for i in 0..WORKERS {
412 let (me, rx) = (self.clone(), rx.clone());
413 let _ = std::thread::Builder::new()
414 .name(format!("isb-monitor-{i}"))
415 .spawn(move || me.worker(&rx));
416 }
417 let me = self.clone();
418 let _ = std::thread::Builder::new()
419 .name("isb-monitor".into())
420 .spawn(move || me.scheduler(&tx));
421 }
422
423 pub fn shutdown(&self) {
424 self.inner.stop.store(true, Ordering::SeqCst);
425 }
426
427 pub fn stopper(&self) -> Arc<AtomicBool> {
429 self.inner.stop.clone()
430 }
431
432 fn scheduler(&self, tx: &SyncSender<Job>) {
433 let (mut synced, mut rolled) = (0u64, now_ms());
434 while !self.inner.stop.load(Ordering::SeqCst) {
435 let now = now_ms();
436 if now.saturating_sub(synced) >= AUTO_SYNC_MS {
437 for o in self.orgs() {
438 if let Err(e) = self.sync_auto(&o) {
439 eprintln!("isb serve: monitor: {o}: own monitors: {e}");
440 }
441 }
442 synced = now;
443 }
444 if now.saturating_sub(rolled) >= ROLLUP_MS {
445 self.rollup_all(now);
446 rolled = now;
447 }
448 self.tick(now, tx);
449 std::thread::sleep(Duration::from_secs(1));
450 }
451 }
452
453 fn rollup_all(&self, now: u64) {
454 for o in self.orgs() {
455 if !self.path(&o, "monitors.db").exists() {
456 continue;
457 }
458 if let Err(e) = self.db(&o).and_then(|d| d.lock().unwrap().rollup(now)) {
459 eprintln!("isb serve: monitor: {o}: history rollup: {e}");
460 }
461 }
462 }
463
464 fn tick(&self, now: u64, tx: &SyncSender<Job>) {
466 let mut seen = Vec::new();
467 for org in self.orgs() {
468 for m in self.list(&org).unwrap_or_default() {
469 let key = (org.clone(), m.name.clone());
470 seen.push(key.clone());
471 if m.paused {
472 continue;
473 }
474 let mut slots = self.inner.slots.lock().unwrap();
475 let s = slots.entry(key.clone()).or_insert_with(|| Slot {
476 next: now + fnv(&format!("{}/{}", org, m.name)) % (m.interval.min(60) * 1000),
478 running: false,
479 });
480 if s.running || now < s.next {
481 continue;
482 }
483 match tx.try_send(Job {
484 org: org.clone(),
485 monitor: m,
486 }) {
487 Ok(()) => s.running = true,
488 Err(TrySendError::Full(_)) => return,
489 Err(TrySendError::Disconnected(_)) => return,
490 }
491 }
492 }
493 self.inner
494 .slots
495 .lock()
496 .unwrap()
497 .retain(|k, _| seen.contains(k));
498 }
499
500 fn worker(&self, rx: &Mutex<Receiver<Job>>) {
501 loop {
502 let job = match rx.lock().unwrap().recv_timeout(Duration::from_secs(1)) {
503 Ok(j) => j,
504 Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
505 if self.inner.stop.load(Ordering::SeqCst) {
506 return;
507 }
508 continue;
509 }
510 Err(_) => return,
511 };
512 let started = now_ms();
513 let o = self.check(&job.org, &job.monitor);
514 if let Err(e) = self.record(&job.org, &job.monitor, o) {
515 eprintln!("isb serve: monitor: {}/{}: {e}", job.org, job.monitor.name);
516 }
517 let iv = job.monitor.interval * 1000;
518 let mut slots = self.inner.slots.lock().unwrap();
519 if let Some(s) = slots.get_mut(&(job.org, job.monitor.name)) {
520 s.running = false;
521 s.next = (started + iv).saturating_add_signed(jitter((iv / 20).min(2000)));
522 }
523 }
524 }
525
526 pub fn check(&self, org: &OrgId, m: &Monitor) -> Outcome {
530 if m.follows()
531 && self
532 .stored(org, &m.name)
533 .map(|s| s.state.status == Status::Pending)
534 .unwrap_or(true)
535 && !target::is_live(&self.inner.apps, org, m)
536 {
537 let why = if m.kind == Kind::App {
538 WAITING_FOR_APP
539 } else {
540 WAITING_FOR_SERVICE
541 };
542 return Outcome {
543 at: now_ms(),
544 ok: false,
545 error: Some(why.into()),
546 ..Default::default()
547 };
548 }
549 let ctx = Ctx {
550 apps: &self.inner.apps,
551 ctl: self.inner.apps.controller(),
552 secrets: &self.inner.secrets,
553 allow_private: (self.inner.allow_private)(),
554 tls: self.inner.tls.clone(),
555 };
556 target::run(&ctx, org, m, now_ms())
557 }
558
559 pub fn record(&self, org: &OrgId, m: &Monitor, o: Outcome) -> Result<()> {
561 match self.get(org, &m.name) {
563 Ok(cur) if !cur.paused => {}
564 _ => return Ok(()),
565 }
566 let db = self.db(org)?;
567 let db = db.lock().unwrap();
568 let mut s: Stored = db.load_state(&m.name)?.unwrap_or_default();
569 let t = Thresholds {
570 failures: m.failure_threshold,
571 recoveries: m.recovery_threshold,
572 };
573 let step = s.state.observe(o.ok, o.at, t);
574 db.insert(
575 &m.name,
576 &Check {
577 at: o.at,
578 ok: o.ok,
579 latency_ms: o.latency_ms,
580 status: o.status,
581 error: o.error.clone(),
582 pending: step.pending,
583 },
584 )?;
585 if let Notify::NeverUp { since } = step.notify {
586 let why = format!("never came up: {}", o.error.as_deref().unwrap_or("failed"));
587 db.open_incident(&m.name, since, Some(&why))?;
588 }
589 match step.changed {
590 Some(Status::Down) => {
591 let since = s.state.failing_since.unwrap_or(o.at);
592 db.open_incident(&m.name, since, o.error.as_deref())?;
593 }
594 Some(Status::Up) => db.close_incident(&m.name, o.at)?,
595 _ => {}
596 }
597 let cert = self.cert_due(m, &o, &mut s.state);
598 s.last = Some(o.clone());
599 db.save_state(&m.name, &s)?;
600 drop(db);
601 match step.notify {
602 Notify::Down { since, flapping } => {
603 self.emit_down(org, m, &o, since, flapping, s.state.fails)
604 }
605 Notify::NeverUp { since } => self.emit_never_up(org, m, &o, since, s.state.fails),
606 Notify::Up {
607 down_since,
608 downtime_ms,
609 } => self.emit_up(org, m, &o, down_since, downtime_ms),
610 Notify::None => {}
611 }
612 if let Some(days) = cert {
613 self.emit_cert(org, m, &o, days);
614 }
615 Ok(())
616 }
617
618 fn cert_due(&self, m: &Monitor, o: &Outcome, s: &mut State) -> Option<i64> {
621 let exp = o.cert_expires?;
622 if m.cert_expiry_days == 0 || s.cert_warned == Some(exp) {
623 return None;
624 }
625 let days = (exp as i64 - (o.at / 1000) as i64).div_euclid(86_400);
626 (days <= i64::from(m.cert_expiry_days)).then(|| {
627 s.cert_warned = Some(exp);
628 days
629 })
630 }
631
632 fn subject(&self, org: &OrgId, m: &Monitor) -> (String, String) {
638 if let (Kind::Service, Some(st), Some(sv)) = (m.kind, &m.stack, &m.service) {
639 return (crate::stack::qualified(org, st), sv.clone());
640 }
641 if let Some(a) = m
642 .app
643 .as_deref()
644 .and_then(|a| self.inner.apps.get(org, a).ok())
645 {
646 if let Ok(st) = a.spec.stack() {
647 return (crate::stack::qualified(org, &st), a.spec.name);
648 }
649 }
650 (crate::stack::qualified(org, "@monitors"), m.name.clone())
651 }
652
653 pub fn link(&self, org: &OrgId, name: &str) -> Option<String> {
655 let base = self.inner.public_url.as_deref()?;
656 Some(format!("{base}/orgs/{org}/uptime/{name}"))
657 }
658
659 fn base_details(&self, org: &OrgId, m: &Monitor, o: &Outcome) -> Value {
660 json!({
661 "monitor": m.name,
662 "type": m.kind,
663 "app": m.app,
664 "stack": m.stack,
665 "service": m.service,
666 "url": o.url.clone().unwrap_or_else(|| m.target()),
667 "status": o.status,
668 "latency_ms": o.latency_ms,
669 "error": o.error,
670 "via": o.via,
671 "note": o.note,
672 "checked_at": o.at,
673 "link": self.link(org, &m.name),
674 })
675 }
676
677 fn emit(
678 &self,
679 org: &OrgId,
680 m: &Monitor,
681 kind: &str,
682 level: &str,
683 message: String,
684 details: Value,
685 ) {
686 let (stack, service) = self.subject(org, m);
687 {
688 let mut d = self.inner.details.lock().unwrap();
689 if d.len() >= DETAILS_KEPT {
690 d.pop_front();
691 }
692 d.push_back(((org.clone(), kind.to_string(), message.clone()), details));
693 }
694 self.inner
695 .apps
696 .controller()
697 .event(kind, level, &stack, &service, message);
698 }
699
700 fn emit_down(
701 &self,
702 org: &OrgId,
703 m: &Monitor,
704 o: &Outcome,
705 since: u64,
706 flapping: bool,
707 fails: u32,
708 ) {
709 let mut d = self.base_details(org, m, o);
710 d["down_since"] = json!(since);
711 d["failures"] = json!(fails);
712 d["flapping"] = json!(flapping);
713 self.emit(
714 org,
715 m,
716 "monitor.down",
717 "error",
718 down_message(m, o, fails, flapping),
719 d,
720 );
721 }
722
723 fn emit_never_up(&self, org: &OrgId, m: &Monitor, o: &Outcome, since: u64, fails: u32) {
726 let mut d = self.base_details(org, m, o);
727 d["down_since"] = json!(since);
728 d["failures"] = json!(fails);
729 d["flapping"] = json!(false);
730 d["never_up"] = json!(true);
731 self.emit(
732 org,
733 m,
734 "monitor.down",
735 "error",
736 never_up_message(m, o, fails),
737 d,
738 );
739 }
740
741 fn emit_up(&self, org: &OrgId, m: &Monitor, o: &Outcome, down_since: u64, downtime_ms: u64) {
742 let mut d = self.base_details(org, m, o);
743 d["down_since"] = json!(down_since);
744 d["downtime_ms"] = json!(downtime_ms);
745 d["downtime"] = json!(human(downtime_ms));
746 self.emit(
747 org,
748 m,
749 "monitor.up",
750 "info",
751 up_message(m, o, downtime_ms),
752 d,
753 );
754 }
755
756 fn emit_cert(&self, org: &OrgId, m: &Monitor, o: &Outcome, days: i64) {
757 let mut d = self.base_details(org, m, o);
758 d["cert_expires_at"] = json!(o.cert_expires);
759 d["cert_days_left"] = json!(days);
760 self.emit(
761 org,
762 m,
763 "monitor.cert_expiring",
764 "warn",
765 cert_message(m, o, days),
766 d,
767 );
768 }
769
770 pub fn details(&self, org: &OrgId, kind: &str, message: &str) -> Option<Value> {
772 let d = self.inner.details.lock().unwrap();
773 d.iter()
774 .rev()
775 .find(|((o, k, msg), _)| o == org && k == kind && msg == message)
776 .map(|(_, v)| v.clone())
777 }
778}
779
780fn what(m: &Monitor, o: &Outcome) -> String {
781 o.url.clone().unwrap_or_else(|| m.target())
782}
783
784pub fn down_message(m: &Monitor, o: &Outcome, fails: u32, flapping: bool) -> String {
785 let why = o.error.as_deref().unwrap_or("failed");
786 let mut s = format!(
787 "Monitor {} is DOWN: {}: {why} ({fails} failed check{} in a row)",
788 m.name,
789 what(m, o),
790 if fails == 1 { "" } else { "s" }
791 );
792 if flapping {
793 s.push_str("; it is flapping, so further changes are held until it is stable for 30 min");
794 }
795 s
796}
797
798pub fn never_up_message(m: &Monitor, o: &Outcome, fails: u32) -> String {
799 let why = o.error.as_deref().unwrap_or("failed");
800 format!(
801 "Monitor {} never came up: {}: {why} (no successful check in {} min, {fails} failed check{})",
802 m.name,
803 what(m, o),
804 NEVER_UP_MS / 60_000,
805 if fails == 1 { "" } else { "s" }
806 )
807}
808
809pub fn up_message(m: &Monitor, o: &Outcome, downtime_ms: u64) -> String {
810 let answer = match (o.status, o.latency_ms) {
811 (Some(st), Some(l)) => format!("answered HTTP {st} in {l} ms"),
812 (None, Some(l)) => format!("answered in {l} ms"),
813 _ => "answered".into(),
814 };
815 format!(
816 "Monitor {} is UP again after {}: {} {answer}",
817 m.name,
818 human(downtime_ms),
819 what(m, o)
820 )
821}
822
823pub fn cert_message(m: &Monitor, o: &Outcome, days: i64) -> String {
824 let on = o
825 .cert_expires
826 .map(|e| {
827 let d = (e / 86_400) as i64;
828 let (y, mo, da) = civil(d);
829 format!(" ({y:04}-{mo:02}-{da:02})")
830 })
831 .unwrap_or_default();
832 let when = if days < 0 {
833 "has expired".to_string()
834 } else {
835 format!("expires in {days} day{}", if days == 1 { "" } else { "s" })
836 };
837 format!(
838 "Monitor {}: the TLS certificate of {} {when}{on}",
839 m.name,
840 what(m, o)
841 )
842}
843
844fn civil(z: i64) -> (i64, i64, i64) {
846 let z = z + 719_468;
847 let era = z.div_euclid(146_097);
848 let doe = z - era * 146_097;
849 let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
850 let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
851 let mp = (5 * doy + 2) / 153;
852 let d = doy - (153 * mp + 2) / 5 + 1;
853 let m = if mp < 10 { mp + 3 } else { mp - 9 };
854 (yoe + era * 400 + i64::from(m <= 2), m, d)
855}
856
857#[cfg(test)]
858#[path = "service_tests.rs"]
859mod tests;