use std::collections::HashMap;
use anyhow::{bail, Context, Result};
use serde_json::{json, Value};
use tracing::{debug, info};
use crate::config::{resolve_machine_among, IngressEdge, IngressProvider};
use crate::{CloudflareClient, MachineConfig, MirrorConfig};
const ZONE_FIELD: &str = "zone";
const PORT_FIELD: &str = "port";
const MACHINE_FIELD: &str = "machine";
const MACHINES_FIELD: &str = "machines";
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>,
pub front_doors: Vec<String>,
pub tunnel_id: Option<String>,
}
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 workload_machine(&self) -> Option<&str> {
self.rules.iter().find_map(|r| r.machine.as_deref())
}
}
pub fn declared(mirror: &MirrorConfig) -> Result<Vec<IngressEdge>> {
mirror.ingress_edges()
}
pub fn resolve_ingress_placements(
machines: &[MachineConfig],
mirror: &MirrorConfig,
) -> Result<HashMap<String, String>> {
let mut placements = HashMap::new();
for (role, slot) in &mirror.providers {
let fields = slot.fields();
if fields.contains_key(MACHINE_FIELD) || fields.contains_key(MACHINES_FIELD) {
continue;
}
let Some(required) = slot.required() else {
continue;
};
let machine = resolve_machine_among(machines, &required).with_context(|| {
format!("resolving placement for [providers.{role}] required = {{ … }}")
})?;
placements.insert(role.clone(), machine.name.clone());
}
Ok(placements)
}
pub fn plan_ingress(
mirror: &MirrorConfig,
placements: &HashMap<String, String>,
) -> Result<Vec<IngressPlan>> {
let edges = declared(mirror)?;
if edges.is_empty() {
return Ok(Vec::new());
}
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 edge(s) 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.",
edges.len()
);
};
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)
.or_else(|| {
fields
.get(MACHINES_FIELD)
.and_then(|v| v.as_array())
.and_then(|a| a.first())
.and_then(|v| v.as_str())
.map(str::to_string)
})
.or_else(|| placements.get(role).cloned()),
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 edge(s) 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` declaration.",
edges.len()
);
}
if edges.len() > 1 {
if let Some(edge) = edges.iter().find(|e| !e.has_selector()) {
bail!(
"{}: a mirror with {} edges needs every edge to name what it fronts. Add \
`slots = [...]` or `hostnames = [...]`. Mixing front doors is the whole point \
of declaring several, and an implicit catch-all would publish a service \
through whichever edge happened to be written first.",
edge.label(),
edges.len()
);
}
}
partition(&edges, rules)
}
fn partition(edges: &[IngressEdge], rules: Vec<IngressRule>) -> Result<Vec<IngressPlan>> {
let mut buckets: Vec<Vec<IngressRule>> = vec![Vec::new(); edges.len()];
for rule in rules {
let claimants: Vec<usize> = edges
.iter()
.enumerate()
.filter(|(_, e)| e.claims(&rule.slot, &rule.hostname))
.map(|(i, _)| i)
.collect();
match claimants.as_slice() {
[i] => buckets[*i].push(rule),
[] => bail!(
"slot [providers.{}] is fronted at {:?} but no `[[ingress]]` edge claims it — \
the partition has a hole, so that hostname would be published by nothing while \
the declaration says otherwise. Add {:?} to an edge's `slots`, or {:?} to its \
`hostnames`.",
rule.slot,
rule.hostname,
rule.slot,
rule.hostname
),
many => bail!(
"slot [providers.{}] ({:?}) is claimed by {} edges — {}. One hostname cannot be \
published through two front doors: DNS points one way, so the second is dead \
config that looks live. Narrow the selectors so exactly one claims it.",
rule.slot,
rule.hostname,
many.len(),
many.iter()
.map(|i| edges[*i].label())
.collect::<Vec<_>>()
.join(" and ")
),
}
}
let mut plans = Vec::with_capacity(edges.len());
for (edge, rules) in edges.iter().zip(buckets) {
if rules.is_empty() {
bail!(
"{}: fronts nothing. Its selector matches no slot that declares `zone` + \
`{PORT_FIELD}` — a typo'd slot name is otherwise invisible, because the front \
door still deploys and simply publishes an empty rule set.",
edge.label()
);
}
let front_doors = if edge.machines.is_empty() {
rules
.iter()
.find_map(|r| r.machine.clone())
.into_iter()
.collect()
} else {
edge.machines.clone()
};
plans.push(IngressPlan {
provider: edge.provider,
rules,
front_doors,
tunnel_id: edge.tunnel_id.clone(),
});
}
Ok(plans)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PlannedEdge {
pub service: String,
pub env: String,
pub plan: IngressPlan,
}
impl PlannedEdge {
pub fn label(&self) -> String {
format!("{}/{}", self.service, self.env)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NodeFrontDoor {
pub machine: String,
pub provider: IngressProvider,
pub tunnel_id: Option<String>,
pub rules: Vec<IngressRule>,
pub sources: Vec<String>,
}
impl NodeFrontDoor {
pub fn passway_upstreams(&self) -> Result<Vec<String>> {
self.rules.iter().map(IngressRule::passway_upstream).collect()
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct Collation {
pub front_doors: Vec<NodeFrontDoor>,
pub unplaced: Vec<String>,
}
pub fn collate_front_doors(planned: &[PlannedEdge]) -> Result<Collation> {
let mut owner: std::collections::BTreeMap<&str, (&PlannedEdge, &IngressRule)> =
std::collections::BTreeMap::new();
for edge in planned {
for rule in &edge.plan.rules {
match owner.get(rule.hostname.as_str()) {
None => {
owner.insert(rule.hostname.as_str(), (edge, rule));
}
Some((first_edge, first_rule)) => {
if first_edge.plan.provider != edge.plan.provider {
bail!(
"hostname {:?} is fronted by two different providers — {} declares \
{:?} and {} declares {:?}. DNS points one way, so one of them is \
dead config that still reads as live. Pick one front door for that \
hostname.",
rule.hostname,
first_edge.label(),
first_edge.plan.provider.as_str(),
edge.label(),
edge.plan.provider.as_str()
);
}
if (first_rule.port, &first_rule.upstream_host)
!= (rule.port, &rule.upstream_host)
{
bail!(
"hostname {:?} is fronted at two different upstreams — {} \
(providers.{}, port {}) and {} (providers.{}, port {}). One front \
door publishes one rule per hostname, so whichever service applies \
last wins on the box.",
rule.hostname,
first_edge.label(),
first_rule.slot,
first_rule.port,
edge.label(),
rule.slot,
rule.port
);
}
}
}
}
}
type Key = (String, &'static str, Option<String>);
let mut grouped: std::collections::BTreeMap<Key, NodeFrontDoor> =
std::collections::BTreeMap::new();
let mut unplaced = Vec::new();
for edge in planned {
if edge.plan.front_doors.is_empty() {
unplaced.push(edge.label());
continue;
}
for machine in &edge.plan.front_doors {
let key: Key = (
machine.clone(),
edge.plan.provider.as_str(),
edge.plan.tunnel_id.clone(),
);
let door = grouped.entry(key).or_insert_with(|| NodeFrontDoor {
machine: machine.clone(),
provider: edge.plan.provider,
tunnel_id: edge.plan.tunnel_id.clone(),
rules: Vec::new(),
sources: Vec::new(),
});
for rule in &edge.plan.rules {
if !door.rules.iter().any(|r| r.hostname == rule.hostname) {
door.rules.push(rule.clone());
}
}
let label = edge.label();
if !door.sources.contains(&label) {
door.sources.push(label);
}
}
}
let mut front_doors: Vec<NodeFrontDoor> = grouped.into_values().collect();
for door in &mut front_doors {
door.rules.sort_by(|a, b| a.hostname.cmp(&b.hostname));
door.sources.sort();
}
unplaced.sort();
unplaced.dedup();
Ok(Collation {
front_doors,
unplaced,
})
}
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 {
mirror_placed(ingress, &[], slots)
}
fn mirror_placed(
ingress: IngressProvider,
ingress_machines: &[&str],
slots: &str,
) -> MirrorConfig {
let mut m = mirror_edges(vec![], slots);
m.ingress = ingress.into();
m.ingress_machines = ingress_machines.iter().map(|s| s.to_string()).collect();
m
}
fn mirror_edges(edges: Vec<IngressEdge>, 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: crate::config::IngressDecl::Edges(edges),
ingress_machines: Vec::new(),
providers,
drivers: Default::default(),
asset_aliases: Default::default(),
}
}
fn edge(provider: IngressProvider, machines: &[&str], slots: &[&str]) -> IngressEdge {
IngressEdge {
provider,
machines: machines.iter().map(|s| s.to_string()).collect(),
slots: slots.iter().map(|s| s.to_string()).collect(),
hostnames: Vec::new(),
tunnel_id: None,
}
}
fn only_plan(m: &MirrorConfig) -> IngressPlan {
let mut plans = plan_ingress(m, &HashMap::new()).expect("mirror plans");
assert_eq!(plans.len(), 1, "fixture declares exactly one edge");
plans.remove(0)
}
#[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).unwrap().is_empty());
assert_eq!(plan_ingress(&m, &HashMap::new()).unwrap(), vec![]);
}
#[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 = only_plan(&m);
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.workload_machine(), Some("us-east-001"));
assert_eq!(plan.front_doors, vec!["us-east-001".to_string()]);
}
#[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 = only_plan(&m);
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 = only_plan(&m);
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, &HashMap::new()).unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("providers.compute"), "got: {msg}");
assert!(msg.contains("zone"), "got: {msg}");
}
#[test]
fn front_doors_are_independent_of_where_the_fronted_workload_runs() {
let m = mirror_placed(
IngressProvider::Passway,
&["us-east-001", "us-west-001"],
"[bundle]\nuse = \"cloudflare\"\nmachines = [\"us-east-001\"]\n\
zone = \"yah.dev\"\nport = 8080\nupstream_host = \"100.64.0.3\"\n",
);
let plan = only_plan(&m);
assert_eq!(
plan.front_doors,
vec!["us-east-001".to_string(), "us-west-001".to_string()]
);
assert_eq!(plan.workload_machine(), Some("us-east-001"));
}
#[test]
fn a_bundle_slots_machines_list_is_read_as_placement() {
let m = mirror(
IngressProvider::Passway,
"[bundle]\nuse = \"cloudflare\"\nmachines = [\"us-east-001\", \"us-west-001\"]\n\
zone = \"yah.dev\"\nport = 8080\n",
);
let plan = only_plan(&m);
assert_eq!(plan.rules[0].machine.as_deref(), Some("us-east-001"));
assert_eq!(plan.front_doors, vec!["us-east-001".to_string()]);
}
#[test]
fn a_singular_machine_field_still_wins_over_the_plural_one() {
let m = mirror(
IngressProvider::Passway,
"[compute]\nuse = \"hetzner\"\nmachine = \"pinned\"\n\
machines = [\"ignored\"]\nzone = \"a.yah.dev\"\nport = 8080\n",
);
let plan = only_plan(&m);
assert_eq!(plan.rules[0].machine.as_deref(), Some("pinned"));
}
#[test]
fn front_door_placement_without_a_front_door_is_an_error() {
let m = mirror_placed(
IngressProvider::None,
&["us-west-001"],
"[compute]\nuse = \"hetzner\"\nzone = \"a.yah.dev\"\nport = 8080\n",
);
let err = plan_ingress(&m, &HashMap::new()).unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("ingress_machines"), "got: {msg}");
assert!(msg.contains("us-west-001"), "got: {msg}");
}
#[test]
fn every_front_door_gets_the_same_upstream_set() {
let m = mirror_placed(
IngressProvider::Passway,
&["us-east-001", "us-west-001", "us-south-001"],
"[bundle]\nuse = \"cloudflare\"\nmachines = [\"us-east-001\"]\n\
zone = \"yah.dev\"\nport = 8080\nupstream_host = \"100.64.0.3\"\n",
);
let plan = only_plan(&m);
assert_eq!(plan.front_doors.len(), 3);
assert_eq!(
plan.passway_upstreams().unwrap(),
vec!["yah.dev=100.64.0.3:8080".to_string()]
);
}
#[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, &HashMap::new()).unwrap_err();
assert!(format!("{err:#}").contains("no provider slot declares `port`"));
}
const MIXED_SLOTS: &str = "[bundle]\nuse = \"cloudflare\"\nmachines = [\"us-east-001\"]\n\
zone = \"yah.dev\"\nport = 8080\nupstream_host = \"100.64.0.3\"\n\
[internal]\nuse = \"hetzner\"\nmachine = \"us-west-001\"\n\
zone = \"admin.yah.dev\"\nport = 9443\nupstream_host = \"100.64.0.4\"\n";
#[test]
fn two_edges_split_the_rules_and_keep_their_own_placement() {
let m = mirror_edges(
vec![
edge(IngressProvider::Passway, &["us-east-001"], &["bundle"]),
edge(
IngressProvider::CloudflareTunnel,
&["us-west-001"],
&["internal"],
),
],
MIXED_SLOTS,
);
let plans = plan_ingress(&m, &HashMap::new()).unwrap();
assert_eq!(plans.len(), 2, "one plan per declared edge");
assert_eq!(plans[0].provider, IngressProvider::Passway);
assert_eq!(
plans[0].passway_upstreams().unwrap(),
vec!["yah.dev=100.64.0.3:8080"]
);
assert_eq!(plans[0].front_doors, vec!["us-east-001".to_string()]);
assert_eq!(plans[1].provider, IngressProvider::CloudflareTunnel);
assert_eq!(plans[1].rules.len(), 1);
assert_eq!(plans[1].rules[0].hostname, "admin.yah.dev");
assert_eq!(plans[1].front_doors, vec!["us-west-001".to_string()]);
}
#[test]
fn an_edge_may_select_by_hostname_instead_of_by_slot() {
let m = mirror_edges(
vec![
IngressEdge {
hostnames: vec!["yah.dev".into()],
..edge(IngressProvider::Passway, &["us-east-001"], &[])
},
IngressEdge {
hostnames: vec!["admin.yah.dev".into()],
..edge(IngressProvider::CloudflareTunnel, &["us-west-001"], &[])
},
],
MIXED_SLOTS,
);
let plans = plan_ingress(&m, &HashMap::new()).unwrap();
assert_eq!(plans[0].rules[0].hostname, "yah.dev");
assert_eq!(plans[1].rules[0].hostname, "admin.yah.dev");
}
#[test]
fn a_slot_no_edge_claims_is_an_error_naming_it() {
let m = mirror_edges(
vec![
edge(IngressProvider::Passway, &["us-east-001"], &["bundle"]),
edge(IngressProvider::Passway, &["us-west-001"], &["nonexistent"]),
],
MIXED_SLOTS,
);
let msg = format!("{:#}", plan_ingress(&m, &HashMap::new()).unwrap_err());
assert!(msg.contains("internal"), "names the unclaimed slot: {msg}");
assert!(msg.contains("admin.yah.dev"), "got: {msg}");
}
#[test]
fn a_slot_two_edges_claim_is_an_error_naming_both() {
let m = mirror_edges(
vec![
edge(IngressProvider::Passway, &["us-east-001"], &["bundle"]),
IngressEdge {
hostnames: vec!["yah.dev".into()],
..edge(IngressProvider::Passway, &["us-west-001"], &["internal"])
},
],
MIXED_SLOTS,
);
let msg = format!("{:#}", plan_ingress(&m, &HashMap::new()).unwrap_err());
assert!(msg.contains("claimed by 2 edges"), "got: {msg}");
assert!(msg.contains("slots = [\"bundle\"]"), "got: {msg}");
}
#[test]
fn an_edge_whose_selector_matches_nothing_is_an_error() {
let m = mirror_edges(
vec![
edge(IngressProvider::Passway, &["us-east-001"], &["bundle"]),
edge(IngressProvider::Passway, &["us-west-001"], &["internl"]),
edge(IngressProvider::Passway, &["us-south-001"], &["internal"]),
],
MIXED_SLOTS,
);
let msg = format!("{:#}", plan_ingress(&m, &HashMap::new()).unwrap_err());
assert!(msg.contains("fronts nothing"), "got: {msg}");
assert!(msg.contains("internl"), "names the typo: {msg}");
}
#[test]
fn several_edges_require_every_one_to_say_what_it_fronts() {
let m = mirror_edges(
vec![
edge(IngressProvider::Passway, &["us-east-001"], &[]),
edge(
IngressProvider::CloudflareTunnel,
&["us-west-001"],
&["internal"],
),
],
MIXED_SLOTS,
);
let msg = format!("{:#}", plan_ingress(&m, &HashMap::new()).unwrap_err());
assert!(msg.contains("needs every edge to name what it fronts"), "got: {msg}");
}
#[test]
fn one_selectorless_edge_still_fronts_everything() {
let m = mirror_edges(
vec![edge(IngressProvider::Passway, &["us-east-001"], &[])],
MIXED_SLOTS,
);
let plan = only_plan(&m);
assert_eq!(plan.rules.len(), 2);
}
#[test]
fn the_legacy_scalar_spelling_plans_the_same_edge_as_the_list_form() {
let slots = "[compute]\nuse = \"hetzner\"\nmachine = \"us-east-001\"\n\
zone = \"a.yah.dev\"\nport = 8080\n";
let scalar = only_plan(&mirror_placed(
IngressProvider::Passway,
&["us-east-001", "us-south-001"],
slots,
));
let listed = only_plan(&mirror_edges(
vec![edge(
IngressProvider::Passway,
&["us-east-001", "us-south-001"],
&[],
)],
slots,
));
assert_eq!(scalar, listed);
}
#[test]
fn placement_declared_twice_is_an_error_not_a_precedence_rule() {
let mut m = mirror_edges(
vec![edge(IngressProvider::Passway, &["us-east-001"], &[])],
"[compute]\nuse = \"hetzner\"\nzone = \"a.yah.dev\"\nport = 8080\n",
);
m.ingress_machines = vec!["us-west-001".into()];
let msg = format!("{:#}", plan_ingress(&m, &HashMap::new()).unwrap_err());
assert!(msg.contains("stated twice"), "got: {msg}");
assert!(msg.contains("us-west-001"), "got: {msg}");
}
#[test]
fn an_edge_may_name_its_own_tunnel_overriding_the_nodes() {
let m = mirror_edges(
vec![IngressEdge {
tunnel_id: Some("cohort-b-tunnel".into()),
..edge(IngressProvider::CloudflareTunnel, &["us-east-001"], &[])
}],
"[compute]\nuse = \"cloudflare\"\nzone = \"a.yah.dev\"\nport = 8080\n",
);
assert_eq!(only_plan(&m).tunnel_id.as_deref(), Some("cohort-b-tunnel"));
}
#[test]
fn an_edge_that_fronts_with_nothing_is_rejected() {
let m = mirror_edges(
vec![edge(IngressProvider::None, &["us-east-001"], &[])],
"[compute]\nuse = \"hetzner\"\nzone = \"a.yah.dev\"\nport = 8080\n",
);
let msg = format!("{:#}", plan_ingress(&m, &HashMap::new()).unwrap_err());
assert!(msg.contains("fronts nothing"), "got: {msg}");
}
#[test]
fn both_spellings_round_trip_through_toml() {
let scalar: MirrorConfig = toml::from_str(
"schema_version = 1\nshape = \"single-machine\"\ningress = \"passway\"\n\
ingress_machines = [\"us-east-001\"]\n",
)
.expect("scalar spelling parses");
assert_eq!(scalar.ingress_edges().unwrap().len(), 1);
let listed: MirrorConfig = toml::from_str(
"schema_version = 1\nshape = \"single-machine\"\n\
[[ingress]]\nprovider = \"passway\"\nslots = [\"bundle\"]\n\
[[ingress]]\nprovider = \"cloudflare-tunnel\"\nhostnames = [\"x.yah.dev\"]\n",
)
.expect("list spelling parses");
let edges = listed.ingress_edges().unwrap();
assert_eq!(edges.len(), 2);
assert_eq!(edges[0].slots, vec!["bundle".to_string()]);
assert_eq!(edges[1].provider, IngressProvider::CloudflareTunnel);
let back: MirrorConfig = toml::from_str(&toml::to_string_pretty(&listed).unwrap())
.expect("list spelling round-trips");
assert_eq!(back.ingress_edges().unwrap(), edges);
}
#[test]
fn a_misspelled_provider_says_what_the_legal_values_are() {
let err = toml::from_str::<MirrorConfig>(
"schema_version = 1\nshape = \"single-machine\"\n\
[[ingress]]\nprovider = \"passwya\"\n",
)
.unwrap_err();
let msg = err.to_string();
assert!(msg.contains("passwya"), "names the bad value: {msg}");
assert!(msg.contains("passway"), "names the legal ones: {msg}");
let err = toml::from_str::<MirrorConfig>(
"schema_version = 1\nshape = \"single-machine\"\ningress = \"pasway\"\n",
)
.unwrap_err();
assert!(err.to_string().contains("pasway"), "got: {err}");
}
fn planned(service: &str, plan: IngressPlan) -> PlannedEdge {
PlannedEdge {
service: service.into(),
env: "cloud".into(),
plan,
}
}
fn plan(provider: IngressProvider, machines: &[&str], hostname: &str, port: u16) -> IngressPlan {
IngressPlan {
provider,
rules: vec![IngressRule {
hostname: hostname.into(),
port,
slot: "compute".into(),
provider_id: None,
machine: None,
upstream_host: Some("100.64.0.5".into()),
}],
front_doors: machines.iter().map(|s| s.to_string()).collect(),
tunnel_id: None,
}
}
#[test]
fn two_services_fronting_one_node_collate_into_one_appliance() {
let c = collate_front_doors(&[
planned(
"yah-marketing",
plan(IngressProvider::Passway, &["us-east-001"], "yah.dev", 8080),
),
planned(
"yah-issues",
plan(
IngressProvider::Passway,
&["us-east-001"],
"issues.yah.dev",
8731,
),
),
])
.unwrap();
assert_eq!(c.front_doors.len(), 1, "one appliance, not one per service");
let door = &c.front_doors[0];
assert_eq!(door.machine, "us-east-001");
assert_eq!(
door.passway_upstreams().unwrap(),
vec!["issues.yah.dev=100.64.0.5:8731", "yah.dev=100.64.0.5:8080"]
);
assert_eq!(
door.sources,
vec!["yah-issues/cloud".to_string(), "yah-marketing/cloud".to_string()],
"provenance answers `why is this hostname on this box`"
);
assert!(c.unplaced.is_empty());
}
#[test]
fn one_service_across_three_origins_collates_to_three_appliances() {
let c = collate_front_doors(&[planned(
"yah-marketing",
plan(
IngressProvider::Passway,
&["us-east-001", "us-south-001", "us-west-001"],
"yah.dev",
8080,
),
)])
.unwrap();
assert_eq!(c.front_doors.len(), 3);
for door in &c.front_doors {
assert_eq!(
door.passway_upstreams().unwrap(),
vec!["yah.dev=100.64.0.5:8080"]
);
}
}
#[test]
fn two_cohorts_on_one_node_are_two_connectors_not_a_conflict() {
let mut a = plan(
IngressProvider::CloudflareTunnel,
&["us-east-001"],
"a.yah.dev",
8080,
);
a.tunnel_id = Some("cohort-a".into());
let mut b = plan(
IngressProvider::CloudflareTunnel,
&["us-east-001"],
"b.yah.dev",
8081,
);
b.tunnel_id = Some("cohort-b".into());
let c = collate_front_doors(&[planned("svc-a", a), planned("svc-b", b)]).unwrap();
assert_eq!(c.front_doors.len(), 2);
assert_eq!(
c.front_doors
.iter()
.map(|d| d.tunnel_id.clone())
.collect::<Vec<_>>(),
vec![Some("cohort-a".into()), Some("cohort-b".into())]
);
}
#[test]
fn one_hostname_through_two_providers_is_a_conflict_naming_both_services() {
let err = collate_front_doors(&[
planned(
"yah-marketing",
plan(IngressProvider::Passway, &["us-east-001"], "yah.dev", 8080),
),
planned(
"yah-legacy",
plan(
IngressProvider::CloudflareTunnel,
&["us-west-001"],
"yah.dev",
8080,
),
),
])
.unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("two different providers"), "got: {msg}");
assert!(msg.contains("yah-marketing/cloud"), "got: {msg}");
assert!(msg.contains("yah-legacy/cloud"), "got: {msg}");
}
#[test]
fn one_hostname_at_two_upstreams_is_a_conflict() {
let err = collate_front_doors(&[
planned(
"svc-a",
plan(IngressProvider::Passway, &["us-east-001"], "yah.dev", 8080),
),
planned(
"svc-b",
plan(IngressProvider::Passway, &["us-east-001"], "yah.dev", 9090),
),
])
.unwrap_err();
assert!(
format!("{err:#}").contains("two different upstreams"),
"got: {err:#}"
);
}
#[test]
fn an_edge_with_nowhere_to_run_is_reported_not_dropped() {
let c = collate_front_doors(&[planned(
"yah-marketing",
plan(IngressProvider::Passway, &[], "yah.dev", 8080),
)])
.unwrap();
assert!(c.front_doors.is_empty());
assert_eq!(c.unplaced, vec!["yah-marketing/cloud".to_string()]);
}
#[test]
fn same_mirror_plans_identical_rules_under_either_provider() {
let slots = "[compute]\nuse = \"hetzner\"\nzone = \"a.yah.dev\"\nport = 8080\n";
let tunnel = only_plan(&mirror(IngressProvider::CloudflareTunnel, slots));
let passway = only_plan(&mirror(IngressProvider::Passway, slots));
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 = only_plan(&m);
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 = only_plan(&m);
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 = only_plan(&m);
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 = only_plan(&m);
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"})
);
}
}