use anyhow::{bail, Context, Result};
use serde_json::{json, Value};
use tracing::{debug, info};
use crate::config::IngressProvider;
use crate::{CloudflareClient, MirrorConfig};
const ZONE_FIELD: &str = "zone";
const PORT_FIELD: &str = "port";
const MACHINE_FIELD: &str = "machine";
const UPSTREAM_HOST_FIELD: &str = "upstream_host";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IngressRule {
pub hostname: String,
pub port: u16,
pub slot: String,
pub provider_id: Option<String>,
pub machine: Option<String>,
pub upstream_host: Option<String>,
}
impl IngressRule {
pub fn upstream(&self) -> Result<String> {
let host = self.upstream_host.as_deref().with_context(|| {
format!(
"slot [providers.{}] has no resolved upstream address for {}: the workload's \
mesh IP is allocated at deploy time (R599-F12), so it cannot come from the \
mirror. Either the fronting node's `GET /service-records?ready=true` must \
report a ready record exposing port {}, or the slot must pin \
`{UPSTREAM_HOST_FIELD} = \"<addr>\"` explicitly.",
self.slot, self.hostname, self.port
)
})?;
Ok(format!("{host}:{}", self.port))
}
pub fn service_url(&self) -> Result<String> {
Ok(format!("http://{}", self.upstream()?))
}
pub fn passway_upstream(&self) -> Result<String> {
Ok(format!("{}={}", self.hostname, self.upstream()?))
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IngressPlan {
pub provider: IngressProvider,
pub rules: Vec<IngressRule>,
}
impl IngressPlan {
pub fn passway_upstreams(&self) -> Result<Vec<String>> {
self.rules
.iter()
.map(IngressRule::passway_upstream)
.collect()
}
pub fn resolve_upstreams<F>(&mut self, mut discover: F) -> Result<()>
where
F: FnMut(&IngressRule) -> Result<Option<String>>,
{
for rule in &mut self.rules {
if rule.upstream_host.is_some() {
continue;
}
rule.upstream_host = discover(rule)?;
rule.upstream()?;
}
Ok(())
}
pub fn provider_id(&self) -> Option<&str> {
self.rules.iter().find_map(|r| r.provider_id.as_deref())
}
pub fn machine(&self) -> Option<&str> {
self.rules.iter().find_map(|r| r.machine.as_deref())
}
}
pub fn declared(mirror: &MirrorConfig) -> Option<IngressProvider> {
mirror.ingress.is_declared().then_some(mirror.ingress)
}
pub fn plan_ingress(mirror: &MirrorConfig) -> Result<Option<IngressPlan>> {
let Some(provider) = declared(mirror) else {
return Ok(None);
};
let mut rules = Vec::new();
for (role, slot) in &mirror.providers {
let fields = slot.fields();
let Some(port) = fields.get(PORT_FIELD).and_then(|v| v.as_integer()) else {
continue;
};
let Some(zone) = fields.get(ZONE_FIELD).and_then(|v| v.as_str()) else {
bail!(
"mirror declares ingress = {:?} and slot [providers.{role}] declares \
{PORT_FIELD} = {port}, but no `zone` — an ingress provider fans a public \
hostname in to a node-local port, so it has nothing to publish that port \
at. Add `zone = \"<hostname>\"` to the slot, or drop `{PORT_FIELD}` if \
this slot is not fronted.",
provider.as_str()
);
};
let port = u16::try_from(port).with_context(|| {
format!("slot [providers.{role}] {PORT_FIELD} = {port} is not a valid TCP port")
})?;
rules.push(IngressRule {
hostname: zone.to_string(),
port,
slot: role.clone(),
provider_id: slot.provider_id().map(str::to_string),
machine: fields
.get(MACHINE_FIELD)
.and_then(|v| v.as_str())
.map(str::to_string),
upstream_host: fields
.get(UPSTREAM_HOST_FIELD)
.and_then(|v| v.as_str())
.map(str::to_string),
});
}
rules.sort_by(|a, b| a.hostname.cmp(&b.hostname));
if rules.is_empty() {
bail!(
"mirror declares ingress = {:?} but no provider slot declares `{PORT_FIELD}` — \
there is nothing to publish. Either add `zone` + `{PORT_FIELD}` to the slot \
the front door fronts, or remove the `ingress` field.",
provider.as_str()
);
}
Ok(Some(IngressPlan { provider, rules }))
}
pub async fn ensure_tunnel_ingress(
cf: &CloudflareClient,
account_id: &str,
tunnel_id: &str,
plan: &IngressPlan,
) -> Result<TunnelIngressOutcome> {
let live = cf
.tunnel_configuration(account_id, tunnel_id)
.await
.with_context(|| format!("reading ingress configuration of tunnel {tunnel_id}"))?;
let live_ingress = live
.get("ingress")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default();
let merged = merge_tunnel_ingress(&live_ingress, &plan.rules)?;
if merged == live_ingress {
debug!(tunnel_id, "tunnel ingress already current — skipping PUT");
return Ok(TunnelIngressOutcome::AlreadyCurrent);
}
let mut config = live;
config
.as_object_mut()
.expect("tunnel_configuration always yields a JSON object")
.insert("ingress".into(), Value::Array(merged));
cf.put_tunnel_configuration(account_id, tunnel_id, &config)
.await
.with_context(|| format!("writing ingress configuration of tunnel {tunnel_id}"))?;
info!(
tunnel_id,
rules = plan.rules.len(),
"tunnel ingress configuration updated"
);
Ok(TunnelIngressOutcome::Updated)
}
pub async fn publish_tunnel_ingress(
workspace_root: &std::path::Path,
provider_id: &str,
tunnel_id: &str,
plan: &IngressPlan,
) -> Result<TunnelIngressOutcome> {
let cf_provider = super::cf_creds::CfProvider::resolve(workspace_root, provider_id)?;
let account_id = cf_provider.account_id.clone();
let cf = CloudflareClient::new(cf_provider.api_token()?);
ensure_tunnel_ingress(&cf, &account_id, tunnel_id, plan).await
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TunnelIngressOutcome {
AlreadyCurrent,
Updated,
}
fn merge_tunnel_ingress(live: &[Value], rules: &[IngressRule]) -> Result<Vec<Value>> {
let owned: std::collections::BTreeSet<&str> =
rules.iter().map(|r| r.hostname.as_str()).collect();
let mut out: Vec<Value> = Vec::with_capacity(live.len() + rules.len());
let mut catch_all: Option<Value> = None;
for rule in live {
match rule.get("hostname").and_then(Value::as_str) {
None | Some("") => catch_all = Some(rule.clone()),
Some(host) if owned.contains(host) => {} Some(_) => out.push(rule.clone()),
}
}
for rule in rules {
out.push(json!({
"hostname": rule.hostname,
"service": rule.service_url()?,
}));
}
out.push(catch_all.unwrap_or_else(|| json!({ "service": "http_status:404" })));
Ok(out)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::config::MirrorShape;
use std::collections::BTreeMap;
fn mirror(ingress: IngressProvider, slots: &str) -> MirrorConfig {
let providers: BTreeMap<String, crate::MirrorProviderSlot> =
toml::from_str(slots).expect("slot fixture parses");
MirrorConfig {
schema_version: 1,
shape: MirrorShape::SingleMachine,
ingress,
providers,
drivers: Default::default(),
asset_aliases: Default::default(),
}
}
#[test]
fn no_ingress_field_plans_nothing() {
let m = mirror(
IngressProvider::None,
"[compute]\nuse = \"hetzner\"\nzone = \"a.yah.dev\"\nport = 8080\n",
);
assert!(declared(&m).is_none());
assert_eq!(plan_ingress(&m).unwrap(), None);
}
#[test]
fn derives_a_rule_per_fronted_slot_sorted_by_hostname() {
let m = mirror(
IngressProvider::CloudflareTunnel,
"[compute]\nuse = \"hetzner\"\nmachine = \"us-east-001\"\nzone = \"z.yah.dev\"\nport = 8080\n\
[receiver]\nuse = \"cloudflare\"\nzone = \"a.yah.dev\"\nport = 9090\n",
);
let plan = plan_ingress(&m).unwrap().unwrap();
assert_eq!(plan.provider, IngressProvider::CloudflareTunnel);
assert_eq!(
plan.rules,
vec![
IngressRule {
hostname: "a.yah.dev".into(),
port: 9090,
slot: "receiver".into(),
provider_id: Some("cloudflare".into()),
machine: None,
upstream_host: None,
},
IngressRule {
hostname: "z.yah.dev".into(),
port: 8080,
slot: "compute".into(),
provider_id: Some("hetzner".into()),
machine: Some("us-east-001".into()),
upstream_host: None,
},
]
);
assert_eq!(plan.provider_id(), Some("cloudflare"));
assert_eq!(plan.machine(), Some("us-east-001"));
}
#[test]
fn slot_without_the_opt_in_marker_is_skipped() {
let m = mirror(
IngressProvider::Passway,
"[static]\nuse = \"cloudflare\"\nbucket = \"b\"\n\
[compute]\nuse = \"hetzner\"\nzone = \"a.yah.dev\"\nport = 8080\n",
);
let plan = plan_ingress(&m).unwrap().unwrap();
assert_eq!(plan.rules.len(), 1);
assert_eq!(plan.rules[0].slot, "compute");
}
#[test]
fn a_cdn_slots_zone_does_not_drag_it_into_the_plan() {
let m = mirror(
IngressProvider::CloudflareTunnel,
"[static]\nuse = \"cloudflare\"\nbucket = \"yah-app-dev\"\n\
zone = \"analytics.yah.dev\"\n\
[compute]\nuse = \"hetzner\"\nmachine = \"yah-cloud-1\"\n\
zone = \"analytics.yah.dev\"\nport = 8080\n",
);
let plan = plan_ingress(&m).unwrap().unwrap();
assert_eq!(plan.rules.len(), 1, "only the opted-in slot is fronted");
assert_eq!(plan.rules[0].slot, "compute");
}
#[test]
fn opt_in_without_zone_is_an_error_naming_the_slot() {
let m = mirror(
IngressProvider::CloudflareTunnel,
"[compute]\nuse = \"hetzner\"\nport = 8080\n",
);
let err = plan_ingress(&m).unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("providers.compute"), "got: {msg}");
assert!(msg.contains("zone"), "got: {msg}");
}
#[test]
fn ingress_with_no_fronted_slot_is_an_error() {
let m = mirror(
IngressProvider::Passway,
"[static]\nuse = \"cloudflare\"\nbucket = \"b\"\n",
);
let err = plan_ingress(&m).unwrap_err();
assert!(format!("{err:#}").contains("no provider slot declares `port`"));
}
#[test]
fn same_mirror_plans_identical_rules_under_either_provider() {
let slots = "[compute]\nuse = \"hetzner\"\nzone = \"a.yah.dev\"\nport = 8080\n";
let tunnel = plan_ingress(&mirror(IngressProvider::CloudflareTunnel, slots))
.unwrap()
.unwrap();
let passway = plan_ingress(&mirror(IngressProvider::Passway, slots))
.unwrap()
.unwrap();
assert_ne!(tunnel.provider, passway.provider);
assert_eq!(tunnel.rules, passway.rules);
}
#[test]
fn renders_both_provider_forms_from_one_rule() {
let r = rule("a.yah.dev", 8080);
assert_eq!(r.service_url().unwrap(), "http://100.64.0.5:8080");
assert_eq!(r.passway_upstream().unwrap(), "a.yah.dev=100.64.0.5:8080");
}
#[test]
fn an_unresolved_rule_renders_nothing_rather_than_dialing_loopback() {
let m = mirror(
IngressProvider::CloudflareTunnel,
"[compute]\nuse = \"hetzner\"\nzone = \"a.yah.dev\"\nport = 8080\n",
);
let plan = plan_ingress(&m).unwrap().unwrap();
let err = plan.rules[0].service_url().unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("service-records"), "got: {msg}");
assert!(msg.contains("upstream_host"), "got: {msg}");
}
#[test]
fn discovery_fills_the_upstream_from_placement() {
let m = mirror(
IngressProvider::Passway,
"[compute]\nuse = \"hetzner\"\nzone = \"a.yah.dev\"\nport = 8080\n",
);
let mut plan = plan_ingress(&m).unwrap().unwrap();
plan.resolve_upstreams(|r| {
assert_eq!(r.port, 8080);
Ok(Some("100.64.0.7".into()))
})
.unwrap();
assert_eq!(
plan.passway_upstreams().unwrap(),
vec!["a.yah.dev=100.64.0.7:8080"]
);
}
#[test]
fn an_explicit_upstream_host_wins_over_discovery() {
let m = mirror(
IngressProvider::Passway,
"[compute]\nuse = \"hetzner\"\nzone = \"a.yah.dev\"\nport = 8080\n\
upstream_host = \"127.0.0.1\"\n",
);
let mut plan = plan_ingress(&m).unwrap().unwrap();
plan.resolve_upstreams(|_| panic!("discovery must not run for a pinned slot"))
.unwrap();
assert_eq!(
plan.passway_upstreams().unwrap(),
vec!["a.yah.dev=127.0.0.1:8080"]
);
}
#[test]
fn no_ready_record_is_an_error_not_an_empty_upstream() {
let m = mirror(
IngressProvider::CloudflareTunnel,
"[compute]\nuse = \"hetzner\"\nzone = \"a.yah.dev\"\nport = 8080\n",
);
let mut plan = plan_ingress(&m).unwrap().unwrap();
let err = plan.resolve_upstreams(|_| Ok(None)).unwrap_err();
assert!(format!("{err:#}").contains("no resolved upstream address"));
}
fn rule(hostname: &str, port: u16) -> IngressRule {
IngressRule {
hostname: hostname.into(),
port,
slot: "compute".into(),
provider_id: None,
machine: None,
upstream_host: Some("100.64.0.5".into()),
}
}
#[test]
fn merge_appends_catch_all_when_tunnel_is_empty() {
let out = merge_tunnel_ingress(&[], &[rule("a.yah.dev", 8080)]).unwrap();
assert_eq!(
out,
vec![
json!({"hostname": "a.yah.dev", "service": "http://100.64.0.5:8080"}),
json!({"service": "http_status:404"}),
]
);
}
#[test]
fn merge_keeps_a_neighbours_rule_verbatim_including_unmodelled_fields() {
let live = vec![
json!({
"hostname": "other.yah.dev",
"service": "http://127.0.0.1:9999",
"path": "/api/*",
"originRequest": {"noTLSVerify": true}
}),
json!({"service": "http_status:404"}),
];
let out = merge_tunnel_ingress(&live, &[rule("a.yah.dev", 8080)]).unwrap();
assert_eq!(
out[0], live[0],
"neighbour rule must survive byte-identical"
);
assert_eq!(
out[1],
json!({"hostname": "a.yah.dev", "service": "http://100.64.0.5:8080"})
);
assert_eq!(out[2], json!({"service": "http_status:404"}));
}
#[test]
fn merge_replaces_an_owned_hostname_and_keeps_the_live_catch_all() {
let live = vec![
json!({"hostname": "a.yah.dev", "service": "http://127.0.0.1:1111"}),
json!({"service": "http_status:503"}),
];
let out = merge_tunnel_ingress(&live, &[rule("a.yah.dev", 8080)]).unwrap();
assert_eq!(
out,
vec![
json!({"hostname": "a.yah.dev", "service": "http://100.64.0.5:8080"}),
json!({"service": "http_status:503"}),
]
);
}
#[test]
fn merge_is_idempotent() {
let rules = vec![rule("a.yah.dev", 8080)];
let once = merge_tunnel_ingress(&[], &rules).unwrap();
let twice = merge_tunnel_ingress(&once, &rules).unwrap();
assert_eq!(once, twice);
}
#[test]
fn merge_treats_empty_hostname_as_the_catch_all() {
let live = vec![json!({"hostname": "", "service": "http_status:404"})];
let out = merge_tunnel_ingress(&live, &[rule("a.yah.dev", 8080)]).unwrap();
assert_eq!(out.len(), 2);
assert_eq!(
out[1],
json!({"hostname": "", "service": "http_status:404"})
);
}
}