use super::*;
use crate::client::Client;
use crate::stack::Controller;
use std::io::{Read, Write};
use std::net::TcpListener;
use std::sync::atomic::AtomicU16;
fn service(dir: &Path) -> (Monitors, Controller) {
let k = crate::secrets::Keyring::new(age::x25519::Identity::generate(), vec![]);
let secrets = Arc::new(Secrets::new(crate::secrets::LocalDriver::new(
dir,
Arc::new(k),
)));
let client = Client::with_socket("/nonexistent/isb-test/incus.sock");
let store = crate::stack::Store::open(dir).unwrap();
let ctl = Controller::start(
client.clone(),
store,
Duration::from_secs(60),
secrets.clone(),
)
.unwrap();
let apps = Apps::new(dir, client, ctl.clone(), secrets.clone());
let m = Monitors::new(
dir,
apps,
secrets,
Arc::new(|| true),
Some("https://isb.example.com/".into()),
);
(m, ctl)
}
fn flappy() -> (u16, Arc<AtomicU16>) {
let l = TcpListener::bind("127.0.0.1:0").unwrap();
let port = l.local_addr().unwrap().port();
let status = Arc::new(AtomicU16::new(200));
let st = status.clone();
std::thread::spawn(move || {
for s in l.incoming() {
let Ok(mut s) = s else { continue };
let mut b = [0u8; 4096];
let _ = s.read(&mut b);
let code = st.load(Ordering::SeqCst);
let _ = write!(s, "HTTP/1.1 {code} X\r\nContent-Length: 2\r\n\r\nok");
}
});
(port, status)
}
fn kinds(ctl: &Controller) -> Vec<(String, String)> {
ctl.events(0, 1000)
.1
.into_iter()
.filter_map(|e| Some((e.kind?, e.message)))
.filter(|(k, _)| k.starts_with("monitor."))
.collect()
}
#[test]
fn down_once_up_once_with_details() {
let dir = tempfile::tempdir().unwrap();
let (svc, ctl) = service(dir.path());
let org = OrgId::new("acme").unwrap();
let (port, status) = flappy();
let mut m = Monitor::new("shop", Kind::Http);
m.url = Some(format!("http://127.0.0.1:{port}/?token=s3cret"));
m.interval = 30;
svc.create(&org, m.clone()).unwrap();
let run = |n: usize| {
for _ in 0..n {
let o = svc.check(&org, &m);
svc.record(&org, &m, o).unwrap();
}
};
run(2);
assert!(kinds(&ctl).is_empty());
status.store(503, Ordering::SeqCst);
run(5);
let k = kinds(&ctl);
assert_eq!(k.len(), 1, "{k:?}");
assert_eq!(k[0].0, "monitor.down");
assert!(k[0].1.contains("HTTP 503 (expected 200-399)"), "{}", k[0].1);
assert!(!k[0].1.contains("s3cret"));
let d = svc.details(&org, "monitor.down", &k[0].1).unwrap();
assert_eq!(d["status"], 503);
assert_eq!(d["monitor"], "shop");
assert_eq!(d["link"], "https://isb.example.com/orgs/acme/uptime/shop");
assert!(d["latency_ms"].is_u64());
status.store(200, Ordering::SeqCst);
run(4);
let k = kinds(&ctl);
assert_eq!(
k.iter().map(|x| x.0.as_str()).collect::<Vec<_>>(),
["monitor.down", "monitor.up"]
);
let d = svc.details(&org, "monitor.up", &k[1].1).unwrap();
assert!(d["downtime_ms"].is_u64() && d["downtime"].is_string());
let db = svc.db(&org).unwrap();
let db = db.lock().unwrap();
let inc = db.incidents(None, 10, now_ms()).unwrap();
assert_eq!(inc.len(), 1);
assert!(inc[0].ended.is_some());
assert_eq!(db.recent("shop", 100).unwrap().len(), 11);
let e = ctl.events(0, 1000).1;
let e = e
.iter()
.find(|e| e.kind.as_deref() == Some("monitor.up"))
.unwrap();
assert_eq!(
(e.stack.as_str(), e.service.as_str()),
("acme/@monitors", "shop")
);
}
#[test]
fn edits_pause_and_delete() {
let dir = tempfile::tempdir().unwrap();
let (svc, _) = service(dir.path());
let org = OrgId::new("acme").unwrap();
let mut m = Monitor::new("db", Kind::Tcp);
(m.host, m.port) = (Some("db.example.com".into()), Some(5432));
svc.create(&org, m.clone()).unwrap();
assert!(svc.create(&org, m.clone()).is_err(), "duplicate");
let mut p = serde_json::Map::new();
p.insert("interval".into(), json!(120));
p.insert("auto".into(), json!(true));
let n = svc.update(&org, "db", p).unwrap();
assert_eq!(n.interval, 120);
assert!(!n.auto, "auto is not the caller's to set");
let mut p = serde_json::Map::new();
p.insert("interval".into(), Value::Null);
assert_eq!(svc.update(&org, "db", p).unwrap().interval, 60);
let mut p = serde_json::Map::new();
p.insert("keyword".into(), json!("x"));
assert!(
svc.update(&org, "db", p).is_err(),
"a keyword on a tcp monitor"
);
let mut p = serde_json::Map::new();
p.insert("name".into(), json!("db2"));
assert!(svc.update(&org, "db", p).is_err());
assert!(svc.set_paused(&org, "db", true).unwrap().paused);
svc.record(
&org,
&m,
Outcome {
at: 1,
..Default::default()
},
)
.unwrap();
assert!(svc.stored(&org, "db").unwrap().last.is_none());
assert!(!svc.set_paused(&org, "db", false).unwrap().paused);
let mut h = Monitor::new("h", Kind::Http);
h.url = Some("https://a.example.com/".into());
h.headers = vec![super::super::Header {
name: "X-Key".into(),
value: None,
secret: Some("NOPE".into()),
}];
assert!(
svc.create(&org, h)
.unwrap_err()
.to_string()
.contains("no secret NOPE")
);
let mut a = Monitor::new("a", Kind::App);
a.app = Some("ghost".into());
assert!(svc.create(&org, a).is_err());
svc.delete(&org, "db").unwrap();
assert!(svc.get(&org, "db").is_err());
assert!(svc.list(&OrgId::new("beta").unwrap()).unwrap().is_empty());
}
#[test]
fn an_apps_own_monitor_stays_deleted() {
let dir = tempfile::tempdir().unwrap();
let (svc, _) = service(dir.path());
let org = OrgId::default_org();
let mut m = Monitor::new("app-shop", Kind::App);
m.app = Some("shop".into());
m.auto = true;
svc.save(&org, &[m.clone()]).unwrap();
svc.sync_auto(&org).unwrap();
assert!(svc.list(&org).unwrap().is_empty());
svc.save(&org, &[m]).unwrap();
svc.delete(&org, "app-shop").unwrap();
assert_eq!(svc.settings(&org).unwrap().exclude_apps, ["shop"]);
assert!(svc.orgs().contains(&org));
}
#[test]
fn the_scheduler_queues_each_check_once() {
let dir = tempfile::tempdir().unwrap();
let (svc, _) = service(dir.path());
let org = OrgId::new("acme").unwrap();
for n in ["a", "b", "c"] {
let mut m = Monitor::new(n, Kind::Tcp);
(m.host, m.port) = (Some("127.0.0.1".into()), Some(9));
svc.create(&org, m).unwrap();
}
svc.set_paused(&org, "c", true).unwrap();
let (tx, rx) = std::sync::mpsc::sync_channel(QUEUE);
let now = now_ms() + 2000;
svc.tick(now, &tx);
let got: Vec<String> = rx.try_iter().map(|j| j.monitor.name).collect();
assert_eq!(got, ["a", "b"]);
svc.tick(now + 5000, &tx);
assert_eq!(rx.try_iter().count(), 0);
let (tx1, rx1) = std::sync::mpsc::sync_channel(1);
svc.inner
.slots
.lock()
.unwrap()
.values_mut()
.for_each(|s| s.running = false);
svc.tick(now + 6000, &tx1);
assert_eq!(rx1.try_iter().count(), 1);
svc.delete(&org, "a").unwrap();
svc.tick(now + 7000, &tx);
assert!(
!svc.inner
.slots
.lock()
.unwrap()
.keys()
.any(|(_, n)| n == "a")
);
assert!(jitter(0) == 0 && jitter(100).abs() <= 100);
}
#[test]
fn certificates_are_warned_about_once() {
let dir = tempfile::tempdir().unwrap();
let (svc, _) = service(dir.path());
let mut m = Monitor::new("x", Kind::Http);
m.url = Some("https://a.example.com/".into());
let mut s = State::default();
let now = 1_800_000_000_000u64;
let o = |days: u64| Outcome {
at: now,
ok: true,
cert_expires: Some(now / 1000 + days * 86_400 + 60),
url: Some("https://a.example.com/".into()),
..Default::default()
};
assert_eq!(svc.cert_due(&m, &o(30), &mut s), None);
assert_eq!(svc.cert_due(&m, &o(9), &mut s), Some(9));
assert_eq!(
svc.cert_due(&m, &o(9), &mut s),
None,
"once per certificate"
);
assert_eq!(
svc.cert_due(&m, &o(8), &mut s),
Some(8),
"a new certificate"
);
m.cert_expiry_days = 0;
assert_eq!(svc.cert_due(&m, &o(1), &mut s), None);
let msg = cert_message(&m, &o(9), 9);
assert!(msg.contains("expires in 9 days (2027-01-"), "{msg}");
assert_eq!(civil(0), (1970, 1, 1));
assert_eq!(
civil(super::super::probe::days_from_civil(2024, 2, 29)),
(2024, 2, 29)
);
}
#[test]
fn messages() {
let mut m = Monitor::new("web", Kind::Tcp);
(m.host, m.port) = (Some("db".into()), Some(5432));
let o = Outcome {
error: Some("timed out connecting".into()),
..Default::default()
};
let d = down_message(&m, &o, 2, false);
assert_eq!(
d,
"Monitor web is DOWN: db:5432: timed out connecting (2 failed checks in a row)"
);
assert!(down_message(&m, &o, 1, true).contains("flapping"));
let o = Outcome {
latency_ms: Some(12),
url: Some("db:5432".into()),
..Default::default()
};
assert_eq!(
up_message(&m, &o, 252_000),
"Monitor web is UP again after 4m 12s: db:5432 answered in 12 ms"
);
}
fn failed(at: u64) -> Outcome {
Outcome {
at,
ok: false,
error: Some("connection refused".into()),
..Default::default()
}
}
#[test]
fn failures_before_the_first_success_are_pending_not_downtime() {
let dir = tempfile::tempdir().unwrap();
let (svc, ctl) = service(dir.path());
let org = OrgId::new("acme").unwrap();
let (port, status) = flappy();
status.store(503, Ordering::SeqCst);
let mut m = Monitor::new("shop", Kind::Http);
m.url = Some(format!("http://127.0.0.1:{port}/"));
svc.create(&org, m.clone()).unwrap();
let run = |n: usize| {
for _ in 0..n {
let o = svc.check(&org, &m);
svc.record(&org, &m, o).unwrap();
}
};
run(5);
let s = svc.summary(&org, &m).unwrap();
assert_eq!(s["status"], "pending");
assert_eq!(s["never_up"], false);
assert!(s["incident"].is_null());
assert_eq!(s["uptime"]["24h"], Value::Null);
assert!(kinds(&ctl).is_empty(), "{:?}", kinds(&ctl));
{
let db = svc.db(&org).unwrap();
let db = db.lock().unwrap();
assert!(db.incidents(None, 10, now_ms()).unwrap().is_empty());
assert!(db.recent("shop", 10).unwrap().iter().all(|c| c.pending));
}
status.store(200, Ordering::SeqCst);
run(1);
let s = svc.summary(&org, &m).unwrap();
assert_eq!(s["status"], "up");
assert_eq!(s["uptime"]["24h"], 100.0);
assert!(kinds(&ctl).is_empty());
status.store(503, Ordering::SeqCst);
run(2);
let k = kinds(&ctl);
assert_eq!(k.len(), 1, "{k:?}");
assert_eq!(k[0].0, "monitor.down");
}
#[test]
fn a_monitor_that_never_comes_up_says_so_once() {
let dir = tempfile::tempdir().unwrap();
let (svc, ctl) = service(dir.path());
let org = OrgId::new("acme").unwrap();
let mut m = Monitor::new("db", Kind::Tcp);
(m.host, m.port) = (Some("db.example.com".into()), Some(5432));
svc.create(&org, m.clone()).unwrap();
let t0 = now_ms();
svc.record(&org, &m, failed(t0)).unwrap();
svc.record(&org, &m, failed(t0 + 10 * 60_000)).unwrap();
assert!(kinds(&ctl).is_empty());
let late = t0 + super::super::state::NEVER_UP_MS + 1000;
svc.record(&org, &m, failed(late)).unwrap();
svc.record(&org, &m, failed(late + 60_000)).unwrap();
let k = kinds(&ctl);
assert_eq!(k.len(), 1, "{k:?}");
assert_eq!(k[0].0, "monitor.down");
assert!(k[0].1.contains("never came up"), "{}", k[0].1);
let d = svc.details(&org, "monitor.down", &k[0].1).unwrap();
assert_eq!(d["never_up"], true);
let s = svc.summary(&org, &m).unwrap();
assert_eq!(
(s["status"].as_str(), s["never_up"].as_bool()),
(Some("pending"), Some(true))
);
let ok = Outcome {
at: late + 120_000,
ok: true,
latency_ms: Some(5),
..Default::default()
};
svc.record(&org, &m, ok).unwrap();
let k = kinds(&ctl);
assert_eq!(
k.iter().map(|x| x.0.as_str()).collect::<Vec<_>>(),
["monitor.down", "monitor.up"]
);
let db = svc.db(&org).unwrap();
let inc = db.lock().unwrap().incidents(None, 10, now_ms()).unwrap();
assert_eq!(inc.len(), 1);
assert!(inc[0].ended.is_some());
}
#[test]
fn an_app_monitor_waits_for_a_live_deployment() {
let dir = tempfile::tempdir().unwrap();
let (svc, ctl) = service(dir.path());
let org = OrgId::new("acme").unwrap();
let mut m = Monitor::new("app-web", Kind::App);
m.app = Some("web".into());
m.auto = true;
svc.save(&org, &[m.clone()]).unwrap();
let o = svc.check(&org, &m);
assert_eq!(o.error.as_deref(), Some(WAITING_FOR_APP));
svc.record(&org, &m, o).unwrap();
let st = svc.stored(&org, "app-web").unwrap();
assert_eq!(st.state.status, Status::Pending);
assert!(kinds(&ctl).is_empty());
let db = svc.db(&org).unwrap();
let c = db.lock().unwrap().recent("app-web", 5).unwrap();
assert!(c[0].pending);
}
type Domains = BTreeMap<(String, String), Vec<crate::ingress::DomainStatus>>;
struct FakeIngress(Mutex<Domains>);
impl crate::stack::controller::Observer for FakeIngress {
fn rotation(&self, _: &str, _: &str, _: &[std::net::IpAddr]) {}
fn drain(&self, _: &str, _: &str, _: std::net::IpAddr, _: Duration) {}
fn stacks_changed(&self, _: Vec<Arc<crate::stack::StackDef>>) {}
fn domains(&self, stack: &str, service: &str) -> Vec<crate::ingress::DomainStatus> {
let m = self.0.lock().unwrap();
m.get(&(stack.to_string(), service.to_string()))
.cloned()
.unwrap_or_default()
}
}
impl FakeIngress {
fn set(&self, stack: &str, service: &str, d: Vec<crate::ingress::DomainStatus>) {
self.0
.lock()
.unwrap()
.insert((stack.to_string(), service.to_string()), d);
}
}
fn serving(url: &str, upstream: &str) -> crate::ingress::DomainStatus {
crate::ingress::DomainStatus {
host: "wiki.acme.dev".into(),
path: "/".into(),
url: Some(url.into()),
provider: "caddy".into(),
state: "serving".into(),
cert: "none".into(),
upstreams: vec![upstream.into()],
..Default::default()
}
}
fn service_with(
dir: &Path,
stacks: &[(&str, &str)],
allow_private: bool,
) -> (Monitors, Arc<FakeIngress>) {
let store = crate::stack::Store::open(dir).unwrap();
for (name, y) in stacks {
store
.save(&crate::stack::StackDef {
source: None,
domains: Default::default(),
name: (*name).into(),
org: OrgId::new("acme").unwrap(),
file: serde_yaml_ng::from_str(y).unwrap(),
base_dir: "/".into(),
secrets: Default::default(),
force: Default::default(),
images: Default::default(),
deployed_at: 0,
deployed_by: String::new(),
previous: None,
})
.unwrap();
}
let k = crate::secrets::Keyring::new(age::x25519::Identity::generate(), vec![]);
let secrets = Arc::new(Secrets::new(crate::secrets::LocalDriver::new(
dir,
Arc::new(k),
)));
let client = Client::with_socket("/nonexistent/isb-test/incus.sock");
let ing = Arc::new(FakeIngress(Mutex::new(BTreeMap::new())));
let ctl = Controller::start_with(
client.clone(),
store,
Duration::from_secs(3600),
secrets.clone(),
Some(ing.clone()),
)
.unwrap();
let apps = Apps::new(dir, client, ctl, secrets.clone());
let m = Monitors::new(dir, apps, secrets, Arc::new(move || allow_private), None);
(m, ing)
}
fn answering(head: &'static str) -> u16 {
let l = TcpListener::bind("127.0.0.1:0").unwrap();
let port = l.local_addr().unwrap().port();
std::thread::spawn(move || {
for s in l.incoming() {
let Ok(mut s) = s else { continue };
let mut b = [0u8; 4096];
let _ = s.read(&mut b);
let _ = write!(s, "{head}\r\nContent-Length: 2\r\n\r\nok");
}
});
port
}
const WIKI: &str = "services:\n web: {image: x, domains: [{host: wiki.acme.dev, port: 80}]}\n redis: {image: x}\n";
#[test]
fn compose_stack_services_with_a_domain_get_their_own_monitor() {
let dir = tempfile::tempdir().unwrap();
let app_web = "services:\n web: {image: x, labels: {isb.app: web}, domains: [{host: shop.acme.dev, port: 80}]}\n";
let (svc, ing) = service_with(
dir.path(),
&[
("wiki", WIKI),
("isb-tunnel", "services:\n cloudflared: {image: x}\n"),
("shop-production", app_web),
("shop-production-pr-3", app_web),
],
true,
);
let org = OrgId::new("acme").unwrap();
for s in ["acme/shop-production", "acme/shop-production-pr-3"] {
ing.set(
s,
"web",
vec![serving("https://shop.acme.dev/", "10.0.0.6:80")],
);
}
svc.sync_auto(&org).unwrap();
assert!(svc.list(&org).unwrap().is_empty());
ing.set(
"acme/wiki",
"web",
vec![serving("https://wiki.acme.dev/", "10.0.0.5:80")],
);
svc.sync_auto(&org).unwrap();
let all = svc.list(&org).unwrap();
assert_eq!(all.len(), 1, "{all:?}");
let m = &all[0];
assert_eq!(m.name, "stack-wiki-web");
assert_eq!((m.kind, m.auto), (Kind::Service, true));
assert_eq!(m.target(), "service wiki/web");
let o = svc.check(&org, m);
assert_eq!(o.error.as_deref(), Some(WAITING_FOR_SERVICE));
svc.delete(&org, "stack-wiki-web").unwrap();
assert_eq!(svc.settings(&org).unwrap().exclude_services, ["wiki/web"]);
svc.sync_auto(&org).unwrap();
assert!(svc.list(&org).unwrap().is_empty());
let mut h = Monitor::new("w", Kind::Service);
(h.stack, h.service) = (Some("wiki".into()), Some("nope".into()));
assert!(svc.create(&org, h.clone()).is_err());
h.service = Some("web".into());
svc.create(&org, h).unwrap();
}
fn recording(status: &'static str) -> (u16, Arc<Mutex<Vec<String>>>) {
let l = TcpListener::bind("127.0.0.1:0").unwrap();
let port = l.local_addr().unwrap().port();
let seen = Arc::new(Mutex::new(Vec::new()));
let keep = seen.clone();
std::thread::spawn(move || {
for s in l.incoming() {
let Ok(mut s) = s else { continue };
let mut b = [0u8; 4096];
let n = s.read(&mut b).unwrap_or(0);
keep.lock()
.unwrap()
.push(String::from_utf8_lossy(&b[..n]).to_string());
let _ = write!(s, "HTTP/1.1 {status}\r\nContent-Length: 2\r\n\r\nok");
}
});
(port, seen)
}
const ACCESS: &str =
"HTTP/1.1 302 Found\r\nLocation: https://team.cloudflareaccess.com/cdn-cgi/access/login/x";
fn check_once(
stacks: &[(&str, &str)],
allow_private: bool,
d: crate::ingress::DomainStatus,
m: &Monitor,
) -> Outcome {
let org = OrgId::new("acme").unwrap();
let dir = tempfile::tempdir().unwrap();
let (svc, ing) = service_with(dir.path(), stacks, allow_private);
let stack = format!("acme/{}", m.stack.as_deref().unwrap_or_default());
ing.set(&stack, "web", vec![d]);
svc.create(&org, m.clone()).unwrap();
svc.edit_state(&org, &m.name, |s| s.status = Status::Up)
.unwrap();
svc.check(&org, m)
}
fn wiki_monitor() -> Monitor {
let mut m = Monitor::new("stack-wiki-web", Kind::Service);
(m.stack, m.service) = (Some("wiki".into()), Some("web".into()));
m
}
fn hop_names(o: &Outcome) -> Vec<(&str, bool)> {
o.hops.iter().map(|h| (h.hop.as_str(), h.ok)).collect()
}
#[test]
fn a_service_behind_access_is_checked_hop_by_hop() {
let access = answering(ACCESS);
let public = format!("http://127.0.0.1:{access}/");
let (ingress, seen) = recording("200 OK");
let origin = format!("http://127.0.0.1:{ingress}");
let replica = answering("HTTP/1.1 200 OK");
let upstream = format!("127.0.0.1:{replica}");
let m = wiki_monitor();
let mut d = serving(&public, &upstream);
d.origin = Some(origin.clone());
let o = check_once(&[("wiki", WIKI)], true, d.clone(), &m);
assert!(o.ok, "{o:?}");
assert_eq!(hop_names(&o), [("edge", true), ("ingress", true)]);
assert_eq!(o.via, Some(format!("ingress {origin}")));
assert_eq!(o.url.as_deref(), Some("http://wiki.acme.dev/"));
let note = o.note.unwrap_or_default();
assert!(
note.contains("checked hop by hop (Cloudflare edge, ingress); the Access policy is not verified: add CF_ACCESS_CLIENT_ID and CF_ACCESS_CLIENT_SECRET secrets"),
"{note}"
);
let head = seen.lock().unwrap().pop().unwrap_or_default();
assert!(head.starts_with("GET / HTTP/1.1\r\n"), "{head}");
assert!(head.contains("\r\nHost: wiki.acme.dev\r\n"), "{head}");
let (bad, _) = recording("502 Bad Gateway");
let mut d502 = d.clone();
d502.origin = Some(format!("http://127.0.0.1:{bad}"));
let o = check_once(&[("wiki", WIKI)], true, d502, &m);
assert!(!o.ok);
assert_eq!(
o.error.as_deref(),
Some("ingress: HTTP 502 (expected 200-399)")
);
let mut dt = d.clone();
(dt.provider, dt.https) = ("cloudflare-tunnel".into(), true);
let tunnel = ("isb-tunnel", "services:\n cloudflared: {image: x}\n");
let o = check_once(&[("wiki", WIKI), tunnel], true, dt, &m);
assert!(!o.ok);
assert_eq!(
o.error.as_deref(),
Some("tunnel: the org's Cloudflare tunnel is not running")
);
assert_eq!(
hop_names(&o),
[("edge", true), ("tunnel", false), ("ingress", true)]
);
assert!(
o.note
.unwrap_or_default()
.contains("(Cloudflare edge, tunnel, ingress)")
);
let head = seen.lock().unwrap().pop().unwrap_or_default();
assert!(head.contains("\r\nX-Forwarded-Proto: https\r\n"), "{head}");
let o = check_once(&[("wiki", WIKI)], true, serving(&public, &upstream), &m);
assert!(o.ok, "{o:?}");
assert_eq!(hop_names(&o), [("edge", true), ("replica", true)]);
assert_eq!(o.via, Some(format!("internal: upstream {upstream}")));
assert!(o.note.unwrap_or_default().contains("stands in"));
let mut mt = m.clone();
mt.headers = vec![
super::super::Header {
name: "CF-Access-Client-Id".into(),
value: Some("id".into()),
secret: None,
},
super::super::Header {
name: "CF-Access-Client-Secret".into(),
value: Some("s".into()),
secret: None,
},
];
let o = check_once(&[("wiki", WIKI)], true, d, &mt);
assert!(
o.hops.first().is_some_and(|h| h.hop == "edge" && h.ok),
"{o:?}"
);
assert!(
o.note
.unwrap_or_default()
.contains("Access does not allow the org's service token"),
);
}
#[test]
fn a_domain_at_a_private_address_is_checked_through_the_ingress() {
let (ingress, seen) = recording("200 OK");
let replica = answering("HTTP/1.1 200 OK");
let mut d = serving(
&format!("http://127.0.0.1:{ingress}/"),
&format!("127.0.0.1:{replica}"),
);
d.origin = Some(format!("http://127.0.0.1:{ingress}"));
let o = check_once(&[("wiki", WIKI)], false, d, &wiki_monitor());
assert!(o.ok, "{o:?}");
assert_eq!(hop_names(&o), [("ingress", true)]);
let note = o.note.unwrap_or_default();
assert!(
note.contains("private address: checked hop by hop (ingress)"),
"{note}"
);
let head = seen.lock().unwrap().pop().unwrap_or_default();
assert!(head.contains("\r\nHost: wiki.acme.dev\r\n"), "{head}");
}
#[test]
fn the_ingress_is_asked_for_the_public_path_of_a_stripped_route() {
let access = answering(ACCESS);
let ingress = answering_only("/api/status");
let replica = answering_only("/status");
let y = format!(
"services:\n web: {{image: x, domains: [{{host: wiki.acme.dev, path: /api, port: {replica}, strip_prefix: true}}]}}\n"
);
let mut d = serving(
&format!("http://127.0.0.1:{access}/api"),
&format!("127.0.0.1:{replica}"),
);
d.path = "/api".into();
d.origin = Some(format!("http://127.0.0.1:{ingress}"));
let mut m = wiki_monitor();
m.path = Some("/api/status".into());
let o = check_once(&[("wiki", &y)], true, d, &m);
assert!(o.ok, "{o:?}");
assert_eq!(o.url.as_deref(), Some("http://wiki.acme.dev/api/status"));
}
fn answering_only(path: &'static str) -> u16 {
let l = TcpListener::bind("127.0.0.1:0").unwrap();
let port = l.local_addr().unwrap().port();
std::thread::spawn(move || {
for s in l.incoming() {
let Ok(mut s) = s else { continue };
let mut b = [0u8; 4096];
let n = s.read(&mut b).unwrap_or(0);
let req = String::from_utf8_lossy(&b[..n]);
let got = req.split(' ').nth(1).unwrap_or_default();
let status = if got == path {
"200 OK"
} else {
"404 Not Found"
};
let _ = write!(s, "HTTP/1.1 {status}\r\nContent-Length: 2\r\n\r\nok");
}
});
port
}
#[test]
fn a_service_monitor_requests_the_path_of_its_healthcheck() {
let org = OrgId::new("acme").unwrap();
let origin = answering_only("/healthz");
let upstream = format!("127.0.0.1:{origin}");
let mut m = Monitor::new("stack-api-web", Kind::Service);
(m.stack, m.service) = (Some("api".into()), Some("web".into()));
for (path, strip, url, via, says) in [
(
"/",
false,
"/healthz",
"public",
"path of the service's healthcheck",
),
(
"/api",
true,
"/api/healthz",
"public",
"path of the service's healthcheck",
),
(
"/sso",
false,
"/healthz",
"internal",
"is not under the domain's path /sso",
),
] {
let dir = tempfile::tempdir().unwrap();
let y = format!(
"services:\n web:\n image: x\n domains: [{{host: wiki.acme.dev, path: {path}, port: {origin}, strip_prefix: {strip}}}]\n healthcheck: {{test: [CMD, wget, -q, -O, /dev/null, 'http://127.0.0.1:{origin}/healthz']}}\n"
);
let (svc, ing) = service_with(dir.path(), &[("api", &y)], true);
let answer = if strip {
answering_only("/api/healthz")
} else {
origin
};
let mut d = serving(&format!("http://127.0.0.1:{answer}{path}"), &upstream);
d.path = path.into();
ing.set("acme/api", "web", vec![d]);
svc.create(&org, m.clone()).unwrap();
svc.edit_state(&org, &m.name, |s| s.status = Status::Up)
.unwrap();
let o = svc.check(&org, &m);
assert!(o.ok, "{path}: {o:?}");
let checked = o.url.clone().unwrap_or_default();
assert!(checked.ends_with(url), "{path}: {o:?}");
assert!(
o.via.as_deref().unwrap_or_default().starts_with(via),
"{path}: {o:?}"
);
assert!(
o.note.as_deref().unwrap_or_default().contains(says),
"{path}: {o:?}"
);
}
}
#[test]
fn a_stripped_prefix_is_stripped_on_the_replica_too() {
let org = OrgId::new("acme").unwrap();
let access = answering(
"HTTP/1.1 302 Found\r\nLocation: https://team.cloudflareaccess.com/cdn-cgi/access/login/x",
);
let origin = answering_only("/status");
let upstream = format!("127.0.0.1:{origin}");
let dir = tempfile::tempdir().unwrap();
let y = format!(
"services:\n web: {{image: x, domains: [{{host: wiki.acme.dev, path: /api, port: {origin}, strip_prefix: true}}]}}\n"
);
let (svc, ing) = service_with(dir.path(), &[("api", &y)], true);
let mut d = serving(&format!("http://127.0.0.1:{access}/api"), &upstream);
d.path = "/api".into();
ing.set("acme/api", "web", vec![d]);
let mut m = Monitor::new("api-status", Kind::Service);
(m.stack, m.service) = (Some("api".into()), Some("web".into()));
m.path = Some("/api/status".into());
svc.create(&org, m.clone()).unwrap();
svc.edit_state(&org, &m.name, |s| s.status = Status::Up)
.unwrap();
let o = svc.check(&org, &m);
assert!(o.ok, "{o:?}");
let checked = o.url.clone().unwrap_or_default();
assert!(checked.ends_with(&format!(":{origin}/status")), "{o:?}");
assert!(!o.note.unwrap_or_default().contains("healthcheck"));
}