use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr};
use std::sync::Arc;
use std::time::Duration;
use serde::{Deserialize, Serialize};
use super::health_path;
use super::probe::{self, HttpAnswer, HttpProbe};
use super::{ACCESS_ID_SECRET, ACCESS_SECRET_SECRET, Kind, Monitor, parse_status};
use crate::app::Apps;
use crate::org::OrgId;
use crate::secrets::Secrets;
use crate::stack::Controller;
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct Outcome {
pub at: u64,
pub ok: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub latency_ms: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub status: Option<u16>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub url: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub via: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cert_expires: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub note: Option<String>,
}
impl Outcome {
fn fail(at: u64, error: String) -> Outcome {
Outcome {
at,
error: Some(error),
..Default::default()
}
}
}
pub(crate) struct Ctx<'a> {
pub apps: &'a Apps,
pub ctl: &'a Controller,
pub secrets: &'a Secrets,
pub allow_private: bool,
pub tls: Arc<rustls::ClientConfig>,
}
struct AppView {
public: Option<String>,
internal: Result<(SocketAddr, String), String>,
host: Option<String>,
path: String,
note: Option<String>,
}
fn loopback_for(ip: IpAddr) -> IpAddr {
match ip {
IpAddr::V4(v) if v.is_unspecified() => IpAddr::V4(Ipv4Addr::LOCALHOST),
IpAddr::V6(v) if v.is_unspecified() => IpAddr::V6(Ipv6Addr::LOCALHOST),
ip => ip,
}
}
fn with_path(url: &str, path: &str) -> String {
let after = url.find("://").map(|i| i + 3).unwrap_or(0);
let end = url[after..]
.find('/')
.map(|i| after + i)
.unwrap_or(url.len());
format!("{}{path}", &url[..end])
}
struct Followed {
stack: String,
service: String,
what: String,
noun: &'static str,
port: Option<Result<u16, String>>,
}
fn followed(apps: &Apps, org: &OrgId, m: &Monitor) -> Result<Followed, String> {
if m.kind == Kind::Service {
let (st, sv) = (
m.stack.clone().unwrap_or_default(),
m.service.clone().unwrap_or_default(),
);
return Ok(Followed {
stack: crate::stack::qualified(org, &st),
what: format!("service {st}/{sv}"),
service: sv,
noun: "service",
port: None,
});
}
let name = m.app.as_deref().unwrap_or_default();
let app = apps
.get(org, name)
.map_err(|_| format!("app {name} does not exist"))?;
let stack = app.spec.stack().map_err(|e| e.to_string())?;
Ok(Followed {
stack: crate::stack::qualified(org, &stack),
service: name.to_string(),
what: format!("app {name}"),
noun: "app",
port: Some(
app.spec
.port
.ok_or_else(|| format!("app {name} has no port to check")),
),
})
}
pub(crate) fn is_live(apps: &Apps, org: &OrgId, m: &Monitor) -> bool {
let Ok(f) = followed(apps, org, m) else {
return false;
};
apps.controller()
.status(&f.stack)
.ok()
.and_then(|st| st.services.into_iter().find(|s| s.service == f.service))
.is_some_and(|s| s.healthy > 0 && s.instances.iter().any(|i| i.in_rotation))
}
fn view_paths(
ctx: &Ctx,
m: &Monitor,
f: &Followed,
pick: Option<&crate::ingress::DomainStatus>,
) -> health_path::Paths {
let def = ctx.ctl.definition(&f.stack).ok();
let spec = def.as_ref().and_then(|d| d.file.services.get(&f.service));
let route = pick
.zip(spec)
.and_then(|(d, s)| health_path::spec_of(&s.domains, &d.host, &d.path));
let domain = pick.map_or("/", |d| d.path.as_str());
let strip = route.is_some_and(|r| r.strip_prefix);
let port = route.and_then(|r| r.port).or_else(|| {
pick.and_then(|d| d.upstreams.first())
.and_then(|u| u.parse::<SocketAddr>().ok())
.map(|a| a.port())
.or_else(|| f.port.as_ref().and_then(|p| p.as_ref().ok().copied()))
});
let health = match (&m.path, spec.and_then(|s| s.healthcheck.as_ref()), port) {
(None, Some(h), Some(port)) => health_path::from_healthcheck(&h.test, port),
_ => None,
};
health_path::choose(m.path.as_deref(), domain, strip, health.as_deref(), f.noun)
}
fn app_view(ctx: &Ctx, org: &OrgId, m: &Monitor) -> Result<(AppView, &'static str), String> {
let f = followed(ctx.apps, org, m)?;
let what = &f.what;
let st = ctx
.ctl
.status(&f.stack)
.map_err(|_| format!("{what} is not deployed"))?;
let svc = st
.services
.into_iter()
.find(|s| s.service == f.service)
.ok_or_else(|| format!("{what} is not deployed"))?;
let served: Vec<_> = svc
.domains
.iter()
.filter(|d| d.url.is_some() && matches!(d.state.as_str(), "serving" | "no-replicas"))
.collect();
let pick = match &m.domain {
Some(h) => Some(
*served
.iter()
.find(|d| d.host == *h)
.ok_or_else(|| format!("{what} does not serve {h}"))?,
),
None => served.first().copied(),
};
let paths = view_paths(ctx, m, &f, pick);
let public = pick
.and_then(|d| d.url.as_deref())
.zip(paths.public.as_deref())
.map(|(u, p)| with_path(u, p));
let replica = |ip: IpAddr| {
svc.instances
.iter()
.find(|i| i.in_rotation && i.ip.as_deref() == Some(ip.to_string().as_str()))
.map(|i| i.name.clone())
};
let internal = match (&f.port, svc.ports.first()) {
(Some(_), Some(p)) => p
.listen
.parse::<SocketAddr>()
.map(|a| {
(
SocketAddr::new(loopback_for(a.ip()), a.port()),
format!("published port {}", p.listen),
)
})
.map_err(|_| format!("published port {} is not an address", p.listen)),
(Some(port), None) => {
let r = svc
.instances
.iter()
.find(|i| i.in_rotation)
.and_then(|i| Some((i.name.clone(), i.ip.as_deref()?.parse::<IpAddr>().ok()?)));
match (port, r) {
(Err(e), _) => Err(e.clone()),
(Ok(_), None) => Err(format!("{what} has no replica in rotation")),
(Ok(p), Some((inst, ip))) => {
Ok((SocketAddr::new(ip, *p), format!("replica {inst}")))
}
}
}
(None, _) => pick
.or_else(|| served.first().copied())
.and_then(|d| d.upstreams.first())
.and_then(|u| u.parse::<SocketAddr>().ok())
.map(|a| match replica(a.ip()) {
Some(inst) => (a, format!("replica {inst}")),
None => (a, format!("upstream {a}")),
})
.ok_or_else(|| format!("{what} has no replica in rotation")),
};
Ok((
AppView {
public,
internal,
host: pick.map(|d| d.host.clone()),
path: paths.internal,
note: paths.note,
},
f.noun,
))
}
fn headers(ctx: &Ctx, org: &OrgId, m: &Monitor) -> Result<Vec<(String, String)>, String> {
let read = |s: &str| -> Result<String, String> {
let (v, _) = ctx.secrets.get(org, s).map_err(|e| {
if e.is_not_found() {
format!("secret {s} is gone")
} else {
format!("secret {s}: {e}")
}
})?;
String::from_utf8(v)
.map(|v| v.trim().to_string())
.map_err(|_| format!("secret {s} is not UTF-8 text"))
};
let mut out = Vec::new();
for h in &m.headers {
let v = match (&h.value, &h.secret) {
(Some(v), _) => v.clone(),
(None, Some(s)) => read(s)?,
(None, None) => continue,
};
out.push((h.name.clone(), v));
}
let has_access = out
.iter()
.any(|(k, _)| k.to_ascii_lowercase().starts_with("cf-access-client"));
if m.follows() && !has_access {
if let (Ok(id), Ok(secret)) = (read(ACCESS_ID_SECRET), read(ACCESS_SECRET_SECRET)) {
out.push(("CF-Access-Client-Id".into(), id));
out.push(("CF-Access-Client-Secret".into(), secret));
}
}
Ok(out)
}
pub fn access_redirect(a: &HttpAnswer) -> bool {
(300..400).contains(&a.status)
&& a.location.as_deref().is_some_and(|l| {
crate::net::parse_url(l).is_ok_and(|t| t.host.ends_with(".cloudflareaccess.com"))
})
}
pub fn judge(m: &Monitor, a: &HttpAnswer) -> Result<(), String> {
let ranges = parse_status(&m.expected_status).map_err(|e| e.to_string())?;
if !ranges
.iter()
.any(|(lo, hi)| (*lo..=*hi).contains(&a.status))
{
let mut e = format!("HTTP {} (expected {})", a.status, m.expected_status);
if access_redirect(a) {
e.push_str(": redirected to Cloudflare Access sign-in");
}
return Err(e);
}
if access_redirect(a) {
return Err(format!(
"HTTP {}: redirected to Cloudflare Access sign-in; give the monitor a service token (headers CF-Access-Client-Id and CF-Access-Client-Secret from secrets)",
a.status
));
}
let body = String::from_utf8_lossy(&a.body);
if let Some(k) = &m.keyword {
if !body.contains(k.as_str()) {
return Err(format!(
"HTTP {}: the body does not contain {k:?}",
a.status
));
}
}
if let Some(k) = &m.keyword_absent {
if body.contains(k.as_str()) {
return Err(format!("HTTP {}: the body contains {k:?}", a.status));
}
}
Ok(())
}
fn http_outcome(m: &Monitor, p: &HttpProbe, at: u64, via: &str) -> Outcome {
let mut o = match probe::http(p) {
Ok(a) => {
let r = judge(m, &a);
Outcome {
at,
ok: r.is_ok(),
latency_ms: Some(a.latency.as_millis() as u64),
status: Some(a.status),
error: r.err(),
url: Some(a.final_url.clone()),
cert_expires: a.cert_expires,
..Default::default()
}
}
Err(e) => Outcome {
url: Some(probe::display_url(&p.url)),
..Outcome::fail(at, e)
},
};
o.via = Some(via.into());
o
}
fn probe_for(ctx: &Ctx, m: &Monitor, url: String, hs: Vec<(String, String)>) -> HttpProbe {
HttpProbe {
url,
method: m.method.clone(),
headers: hs,
timeout: Duration::from_secs(m.timeout),
follow_redirects: m.follow_redirects,
allow_private: ctx.allow_private,
connect_to: None,
tls: ctx.tls.clone(),
}
}
pub(crate) fn run(ctx: &Ctx, org: &OrgId, m: &Monitor, at: u64) -> Outcome {
match m.kind {
Kind::Tcp => {
let (h, p) = (m.host.as_deref().unwrap_or_default(), m.port.unwrap_or(0));
let r = probe::tcp(h, p, Duration::from_secs(m.timeout), ctx.allow_private);
Outcome {
ok: r.is_ok(),
latency_ms: r.as_ref().ok().map(|d| d.as_millis() as u64),
error: r.err(),
url: Some(format!("{h}:{p}")),
..Outcome::fail(at, String::new())
}
}
Kind::Http => match headers(ctx, org, m) {
Ok(hs) => http_outcome(
m,
&probe_for(ctx, m, m.url.clone().unwrap_or_default(), hs),
at,
"public",
),
Err(e) => Outcome::fail(at, e),
},
Kind::App | Kind::Service => run_app(ctx, org, m, at),
}
}
fn run_app(ctx: &Ctx, org: &OrgId, m: &Monitor, at: u64) -> Outcome {
let (view, noun) = match app_view(ctx, org, m) {
Ok(v) => v,
Err(e) => return Outcome::fail(at, e),
};
let hs = match headers(ctx, org, m) {
Ok(h) => h,
Err(e) => return Outcome::fail(at, e),
};
let mut note = view.note.clone();
if let Some(url) = &view.public {
let o = http_outcome(m, &probe_for(ctx, m, url.clone(), hs.clone()), at, "public");
let refused = o
.error
.as_deref()
.is_some_and(|e| e.starts_with("refusing "));
let access = o
.error
.as_deref()
.is_some_and(|e| e.contains("Cloudflare Access"));
if !refused && !access {
return Outcome { note, ..o };
}
let why = if access {
format!(
"the domain is behind Cloudflare Access: checked the {noun}'s own endpoint instead (add CF_ACCESS_CLIENT_ID and CF_ACCESS_CLIENT_SECRET secrets to check the public URL)"
)
} else {
format!(
"the domain resolves to a private address: checked the {noun}'s own endpoint instead (a platform admin can allow private targets)"
)
};
note = Some(match note {
Some(n) => format!("{n}; {why}"),
None => why,
});
}
let (addr, what) = match view.internal {
Ok(x) => x,
Err(e) => {
return Outcome {
note,
..Outcome::fail(at, e)
};
}
};
let host = view.host.unwrap_or_else(|| addr.ip().to_string());
let host = if host.contains(':') {
format!("[{host}]")
} else {
host
};
let mut p = probe_for(
ctx,
m,
format!("http://{host}:{}{}", addr.port(), view.path),
hs,
);
p.connect_to = Some(addr);
p.follow_redirects = false;
let mut o = http_outcome(m, &p, at, &format!("internal: {what}"));
o.note = note;
o
}
#[cfg(test)]
mod tests {
use super::*;
fn ans(status: u16, location: Option<&str>, body: &str) -> HttpAnswer {
HttpAnswer {
status,
location: location.map(String::from),
body: body.as_bytes().to_vec(),
..Default::default()
}
}
#[test]
fn judging_answers() {
let mut m = Monitor::new("x", Kind::Http);
assert!(judge(&m, &ans(200, None, "")).is_ok());
assert!(judge(&m, &ans(301, Some("/a"), "")).is_ok());
let e = judge(&m, &ans(503, None, "")).unwrap_err();
assert_eq!(e, "HTTP 503 (expected 200-399)");
let access = ans(
302,
Some("https://team.cloudflareaccess.com/cdn-cgi/access/login/x"),
"",
);
assert!(
judge(&m, &access)
.unwrap_err()
.contains("Cloudflare Access")
);
m.expected_status = "200".into();
assert!(
judge(&m, &access)
.unwrap_err()
.contains("Cloudflare Access")
);
assert!(judge(&m, &ans(204, None, "")).is_err());
m.keyword = Some("Hostname".into());
assert!(judge(&m, &ans(200, None, "Hostname: web-1")).is_ok());
assert!(
judge(&m, &ans(200, None, "nope"))
.unwrap_err()
.contains("does not contain")
);
m.keyword = None;
m.keyword_absent = Some("error".into());
assert!(
judge(&m, &ans(200, None, "an error page"))
.unwrap_err()
.contains("contains")
);
}
#[test]
fn paths_and_addresses() {
assert_eq!(
with_path("https://a.example.com/x?y", "/healthz"),
"https://a.example.com/healthz"
);
assert_eq!(with_path("http://a:8080", "/"), "http://a:8080/");
assert_eq!(
loopback_for("0.0.0.0".parse().unwrap()),
IpAddr::V4(Ipv4Addr::LOCALHOST)
);
assert_eq!(
loopback_for("10.1.2.3".parse().unwrap()).to_string(),
"10.1.2.3"
);
}
}