use std::collections::{BTreeMap, HashMap};
use std::net::Ipv4Addr;
use serde::{Deserialize, Serialize};
use workload_spec::{
EnvValue, EnvVar, ExposeSpec, HealthProbe, Healthcheck, ImageRef, LifecycleArchetype,
MeshExpose, MeshIdent, Millis, NamespaceId, ResourceLimits, RestartPolicy, SchemaVersion,
SecretMount, SecretRef, SecretTarget, StopPolicy, TenantId, TierTag, VolumeMount, Workload,
WorkloadSpec, HOST_NETWORK_ANNOTATION, HOST_NETWORK_VALUE, PUBLIC_IP_TAINT,
REQUIRES_TAINT_ANNOTATION,
};
pub const INGRESS_WORKLOAD_NAME: &str = "passway-ingress";
pub const DEFAULT_LISTEN: &str = "0.0.0.0:443";
const TLS_PORT: u16 = 443;
const CERT_MOUNT_PATH: &str = "/run/secrets/tls.crt";
const KEY_MOUNT_PATH: &str = "/run/secrets/tls.key";
const PID_FILE: &str = "/run/passway/pingora.pid";
const UPGRADE_SOCK: &str = "/run/passway/upgrade.sock";
const DEFAULT_COMMAND: &str = "/usr/local/bin/passway";
pub const CATCH_ALL_LABEL: &str = "*";
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(untagged)]
pub enum HostScoped<T> {
Global(T),
PerHost(BTreeMap<String, T>),
}
impl<T: Default> Default for HostScoped<T> {
fn default() -> Self {
Self::Global(T::default())
}
}
impl<T> HostScoped<T> {
pub fn render(&self, value_str: impl Fn(&T) -> String) -> String {
match self {
Self::Global(v) => value_str(v),
Self::PerHost(map) => map
.iter()
.map(|(host, v)| format!("{host}={}", value_str(v)))
.collect::<Vec<_>>()
.join(","),
}
}
pub fn validate_keys(&self, var: &str) -> Result<(), String> {
let Self::PerHost(map) = self else {
return Ok(());
};
for host in map.keys() {
if host.trim().is_empty() {
return Err(format!(
"{var} has an entry with an empty hostname; write \
{CATCH_ALL_LABEL:?} if you meant the catch-all set"
));
}
if host.contains(',') || host.contains('=') {
return Err(format!(
"{var} hostname {host:?} contains ',' or '=', which the wire grammar uses \
as separators — it would be re-split into entries that parse as something \
else"
));
}
}
Ok(())
}
pub fn with_host(self, hostname: impl Into<String>, value: T) -> Self
where
T: Clone,
{
let mut map = match self {
Self::PerHost(map) => map,
Self::Global(global) => {
BTreeMap::from([(CATCH_ALL_LABEL.to_string(), global)])
}
};
map.insert(hostname.into(), value);
Self::PerHost(map)
}
}
pub fn front_door_upstream_rule(
hostname: &str,
placed_at: Ipv4Addr,
this_door: Ipv4Addr,
port: u16,
) -> String {
let upstream = if placed_at == this_door {
Ipv4Addr::LOCALHOST
} else {
placed_at
};
format!("{hostname}={upstream}:{port}")
}
fn cert_secret_name(domain: &str) -> String {
format!("tls/{domain}/cert")
}
fn key_secret_name(domain: &str) -> String {
format!("tls/{domain}/key")
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PasswayIngressSpec {
pub domain: String,
#[serde(default = "default_listen")]
pub listen: String,
#[serde(default)]
pub upstreams: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub discover_from: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub discover_ident: Option<String>,
#[serde(default)]
pub upstream_tls: HostScoped<bool>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub upstream_sni: Option<HostScoped<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub command: Option<Vec<String>>,
}
fn default_listen() -> String {
DEFAULT_LISTEN.to_string()
}
impl PasswayIngressSpec {
pub fn validate(&self) -> Result<(), String> {
self.upstream_tls.validate_keys("PASSWAY_UPSTREAM_TLS")?;
if let Some(sni) = &self.upstream_sni {
sni.validate_keys("PASSWAY_UPSTREAM_SNI")?;
if let HostScoped::PerHost(map) = sni {
for (host, value) in map {
if value.trim().is_empty() {
return Err(format!(
"PASSWAY_UPSTREAM_SNI entry for {host:?} names no SNI; omit the \
entry to inherit the process-wide default"
));
}
}
}
}
Ok(())
}
pub fn into_container_workload(&self, image: ImageRef) -> Workload {
let mut env = vec![
literal_env("PASSWAY_LISTEN", self.listen.clone()),
literal_env("PASSWAY_TLS_MODE", "manual".into()),
literal_env("PASSWAY_TLS_CERT", CERT_MOUNT_PATH.into()),
literal_env("PASSWAY_TLS_KEY", KEY_MOUNT_PATH.into()),
literal_env(
"PASSWAY_UPSTREAM_TLS",
self.upstream_tls.render(bool::to_string),
),
literal_env("PASSWAY_PID_FILE", PID_FILE.into()),
literal_env("PASSWAY_UPGRADE_SOCK", UPGRADE_SOCK.into()),
];
if let Some(sni) = &self.upstream_sni {
env.push(literal_env("PASSWAY_UPSTREAM_SNI", sni.render(String::clone)));
}
match &self.discover_from {
Some(base_url) => {
env.push(literal_env("PASSWAY_UPSTREAM_SOURCE", "yubaba".into()));
env.push(literal_env("PASSWAY_YUBABA_URL", base_url.clone()));
if let Some(ident) = &self.discover_ident {
env.push(literal_env("PASSWAY_YUBABA_IDENT", ident.clone()));
}
if !self.upstreams.is_empty() {
env.push(literal_env("PASSWAY_UPSTREAMS", self.upstreams.join(",")));
}
}
None => {
env.push(literal_env("PASSWAY_UPSTREAM_SOURCE", "static".into()));
env.push(literal_env("PASSWAY_UPSTREAMS", self.upstreams.join(",")));
}
}
let secrets = vec![
cluster_file_secret(cert_secret_name(&self.domain), CERT_MOUNT_PATH),
cluster_file_secret(key_secret_name(&self.domain), KEY_MOUNT_PATH),
];
let mut annotations = HashMap::new();
annotations.insert(
HOST_NETWORK_ANNOTATION.to_string(),
HOST_NETWORK_VALUE.to_string(),
);
annotations.insert(
REQUIRES_TAINT_ANNOTATION.to_string(),
PUBLIC_IP_TAINT.to_string(),
);
let grace = Millis::from_secs(5);
let spec = WorkloadSpec {
schema_version: SchemaVersion::V1,
name: INGRESS_WORKLOAD_NAME.into(),
image,
tier: TierTag("infra".into()),
tenant: TenantId::singleton(),
namespace: NamespaceId::singleton(),
replicas: 1,
command: Some(
self.command
.clone()
.unwrap_or_else(|| vec![DEFAULT_COMMAND.into()]),
),
entrypoint: None,
workdir: None,
user: None,
env,
secrets,
volumes: Vec::<VolumeMount>::new(),
resources: ResourceLimits {
memory_mb: 256,
cpu_millis: 512,
ephemeral_storage_mb: 256,
},
depends_on: vec![],
requires: vec![],
healthcheck: Some(Healthcheck {
probe: HealthProbe::TcpConnect { port: TLS_PORT },
interval: Millis::from_secs(10),
timeout: Millis::from_secs(2),
initial_delay: Millis::from_secs(10),
failure_threshold: 3,
}),
restart_policy: RestartPolicy::Always,
archetype: Some(LifecycleArchetype::Appliance),
stop_policy: StopPolicy {
signal: 15,
grace_period: grace,
},
expose: ExposeSpec {
mesh: MeshExpose {
identity: MeshIdent(INGRESS_WORKLOAD_NAME.into()),
ports: MeshExpose::anonymous_ports([TLS_PORT]),
allow_from: vec![],
},
public: None,
operator: None,
},
labels: HashMap::new(),
annotations,
};
Workload::container(spec)
}
}
fn literal_env(name: &str, value: String) -> EnvVar {
EnvVar {
name: name.into(),
value: EnvValue::Literal { value },
}
}
fn cluster_file_secret(name: String, mount_path: &str) -> SecretMount {
SecretMount {
source: SecretRef::Cluster { name },
target: SecretTarget::File {
path: mount_path.into(),
mode: 0o400,
},
}
}
#[cfg(test)]
mod tests {
use super::*;
fn sample_image() -> ImageRef {
ImageRef {
registry: "localhost".into(),
repository: "passway".into(),
tag: "r600f5".into(),
digest: workload_spec::testing::test_digest(),
}
}
fn sample_spec() -> PasswayIngressSpec {
PasswayIngressSpec {
domain: "yah.dev".into(),
listen: DEFAULT_LISTEN.into(),
upstreams: vec!["yah-marketing.pdx:8080".into(), "yah-dashboard.pdx:8080".into()],
upstream_tls: HostScoped::Global(false),
upstream_sni: None,
discover_from: None,
discover_ident: None,
command: None,
}
}
fn discovery_spec_with_pin() -> PasswayIngressSpec {
let mut s = sample_spec();
s.discover_from = Some("http://100.64.0.2:7443".into());
s.discover_ident = Some("yah.dev=yah-marketing".into());
s.upstreams = vec![front_door_upstream_rule(
"cloud.mesh.yah.dev",
Ipv4Addr::new(100, 64, 0, 1),
Ipv4Addr::new(100, 64, 0, 2),
443,
)];
s.upstream_tls = HostScoped::Global(false).with_host("cloud.mesh.yah.dev", true);
s.upstream_sni = Some(HostScoped::PerHost(BTreeMap::from([(
"cloud.mesh.yah.dev".to_string(),
"cloud.mesh.yah.dev".to_string(),
)])));
s
}
fn lower(spec: &PasswayIngressSpec) -> WorkloadSpec {
spec.into_container_workload(sample_image())
.container_spec()
.expect("expected a container reference workload")
.clone()
}
fn env_val<'a>(spec: &'a WorkloadSpec, name: &str) -> Option<&'a str> {
spec.env.iter().find(|e| e.name == name).and_then(|e| match &e.value {
EnvValue::Literal { value } => Some(value.as_str()),
_ => None,
})
}
#[test]
fn consumes_shared_cert_via_two_cluster_file_mounts() {
let spec = lower(&sample_spec());
assert_eq!(spec.secrets.len(), 2, "cert + key");
let cert = &spec.secrets[0];
assert_eq!(
cert.source,
SecretRef::Cluster {
name: "tls/yah.dev/cert".into()
}
);
assert_eq!(
cert.target,
SecretTarget::File {
path: CERT_MOUNT_PATH.into(),
mode: 0o400
}
);
let key = &spec.secrets[1];
assert_eq!(
key.source,
SecretRef::Cluster {
name: "tls/yah.dev/key".into()
}
);
assert_eq!(
key.target,
SecretTarget::File {
path: KEY_MOUNT_PATH.into(),
mode: 0o400
}
);
}
#[test]
fn runs_manual_tls_pointing_at_the_mounts() {
let spec = lower(&sample_spec());
assert_eq!(env_val(&spec, "PASSWAY_TLS_MODE"), Some("manual"));
assert_eq!(env_val(&spec, "PASSWAY_TLS_CERT"), Some(CERT_MOUNT_PATH));
assert_eq!(env_val(&spec, "PASSWAY_TLS_KEY"), Some(KEY_MOUNT_PATH));
}
#[test]
fn drops_self_issuance_no_acme_env() {
let spec = lower(&sample_spec());
assert!(
!spec.env.iter().any(|e| e.name.starts_with("PASSWAY_ACME")),
"ingress must carry no PASSWAY_ACME_* env"
);
assert_ne!(env_val(&spec, "PASSWAY_TLS_MODE"), Some("acme"));
}
#[test]
fn upstreams_join_into_one_env() {
let spec = lower(&sample_spec());
assert_eq!(env_val(&spec, "PASSWAY_UPSTREAM_SOURCE"), Some("static"));
assert_eq!(
env_val(&spec, "PASSWAY_UPSTREAMS"),
Some("yah-marketing.pdx:8080,yah-dashboard.pdx:8080")
);
assert_eq!(env_val(&spec, "PASSWAY_UPSTREAM_TLS"), Some("false"));
}
#[test]
fn discover_from_selects_the_yubaba_upstream_source() {
let mut s = sample_spec();
s.discover_from = Some("http://100.64.0.2:7443".into());
s.discover_ident = Some("yah-marketing".into());
let spec = lower(&s);
assert_eq!(env_val(&spec, "PASSWAY_UPSTREAM_SOURCE"), Some("yubaba"));
assert_eq!(
env_val(&spec, "PASSWAY_YUBABA_URL"),
Some("http://100.64.0.2:7443")
);
}
#[test]
fn discovery_mode_renders_the_workload_ident() {
let mut s = sample_spec();
s.discover_from = Some("http://100.64.0.2:7443".into());
s.discover_ident = Some("yah-marketing".into());
let spec = lower(&s);
assert_eq!(
env_val(&spec, "PASSWAY_YUBABA_IDENT"),
Some("yah-marketing")
);
}
#[test]
fn static_mode_renders_no_workload_ident() {
let spec = lower(&sample_spec());
assert_eq!(env_val(&spec, "PASSWAY_YUBABA_IDENT"), None);
}
#[test]
fn discovery_mode_renders_static_pins_only_when_they_are_given() {
let mut s = sample_spec();
s.discover_from = Some("http://100.64.0.2:7443".into());
s.upstreams = vec![];
let spec = lower(&s);
assert_eq!(
env_val(&spec, "PASSWAY_UPSTREAMS"),
None,
"a discovery-mode ingress with nothing pinned must not carry an empty list"
);
let spec = lower(&discovery_spec_with_pin());
assert_eq!(env_val(&spec, "PASSWAY_UPSTREAM_SOURCE"), Some("yubaba"));
assert_eq!(
env_val(&spec, "PASSWAY_UPSTREAMS"),
Some("cloud.mesh.yah.dev=100.64.0.1:443"),
"the pin rides alongside discovery — it names a node this door's \
local yubaba cannot see (R844-B11)"
);
assert_eq!(
env_val(&spec, "PASSWAY_YUBABA_URL"),
Some("http://100.64.0.2:7443"),
"and discovery stays live for every other hostname"
);
assert_eq!(
env_val(&spec, "PASSWAY_YUBABA_IDENT"),
Some("yah.dev=yah-marketing")
);
}
#[test]
fn per_host_upstream_scheme_renders_passways_fan_in_grammar() {
let spec = lower(&discovery_spec_with_pin());
assert_eq!(
env_val(&spec, "PASSWAY_UPSTREAM_TLS"),
Some("*=false,cloud.mesh.yah.dev=true"),
"the prior global is carried over as the explicit catch-all, so no \
other hostname's scheme changes"
);
assert_eq!(
env_val(&spec, "PASSWAY_UPSTREAM_SNI"),
Some("cloud.mesh.yah.dev=cloud.mesh.yah.dev")
);
}
#[test]
fn bare_upstream_scheme_still_renders_the_pre_r858_shape() {
let mut s = sample_spec();
s.upstream_tls = HostScoped::Global(true);
s.upstream_sni = Some(HostScoped::Global("edge.example.com".into()));
let spec = lower(&s);
assert_eq!(env_val(&spec, "PASSWAY_UPSTREAM_TLS"), Some("true"));
assert_eq!(
env_val(&spec, "PASSWAY_UPSTREAM_SNI"),
Some("edge.example.com")
);
}
#[test]
fn omitted_sni_renders_no_variable_at_all() {
let spec = lower(&sample_spec());
assert_eq!(env_val(&spec, "PASSWAY_UPSTREAM_SNI"), None);
}
#[test]
fn front_door_rule_sends_the_co_located_door_to_loopback() {
const WEST: Ipv4Addr = Ipv4Addr::new(100, 64, 0, 1);
const EAST: Ipv4Addr = Ipv4Addr::new(100, 64, 0, 2);
assert_eq!(
front_door_upstream_rule("cloud.mesh.yah.dev", WEST, EAST, 443),
"cloud.mesh.yah.dev=100.64.0.1:443"
);
assert_eq!(
front_door_upstream_rule("cloud.mesh.yah.dev", WEST, WEST, 443),
"cloud.mesh.yah.dev=127.0.0.1:443"
);
}
#[test]
fn validate_rejects_the_keys_passway_would_panic_on() {
let mut s = sample_spec();
s.upstream_tls = HostScoped::PerHost(BTreeMap::from([("".to_string(), true)]));
let err = s.validate().expect_err("empty hostname");
assert!(err.contains("PASSWAY_UPSTREAM_TLS"), "names the var: {err}");
let mut s = sample_spec();
s.upstream_sni = Some(HostScoped::PerHost(BTreeMap::from([(
"cloud.mesh.yah.dev".to_string(),
String::new(),
)])));
let err = s.validate().expect_err("empty SNI value");
assert!(err.contains("names no SNI"), "explains why: {err}");
assert!(discovery_spec_with_pin().validate().is_ok());
assert!(sample_spec().validate().is_ok());
}
#[test]
fn static_mode_renders_no_discovery_url() {
let spec = lower(&sample_spec());
assert_eq!(env_val(&spec, "PASSWAY_YUBABA_URL"), None);
}
#[test]
fn upgrade_wiring_present_for_graceful_reload() {
let spec = lower(&sample_spec());
assert_eq!(env_val(&spec, "PASSWAY_PID_FILE"), Some(PID_FILE));
assert_eq!(env_val(&spec, "PASSWAY_UPGRADE_SOCK"), Some(UPGRADE_SOCK));
assert!(!spec.env.iter().any(|e| e.name == "PASSWAY_UPGRADE"));
}
#[test]
fn is_public_ip_appliance_on_host_network() {
let spec = lower(&sample_spec());
assert_eq!(spec.archetype, Some(LifecycleArchetype::Appliance));
assert_eq!(spec.requires_taint(), Some(PUBLIC_IP_TAINT));
assert!(spec.wants_host_network());
assert_eq!(spec.tier.0, "infra");
}
#[test]
fn ha_diagnose_can_find_it_by_ident() {
let spec = lower(&sample_spec());
assert!(spec.expose.mesh.identity.0.to_lowercase().contains("passway"));
assert_eq!(spec.expose.mesh.numbers(), vec![443]);
}
#[test]
fn command_defaults_and_overrides() {
let spec = lower(&sample_spec());
assert_eq!(spec.command.as_deref().unwrap(), [DEFAULT_COMMAND]);
let mut custom = sample_spec();
custom.command = Some(vec!["/usr/bin/tini".into(), "--".into(), "/bin/passway".into()]);
let spec = lower(&custom);
assert_eq!(
spec.command.as_deref().unwrap(),
["/usr/bin/tini", "--", "/bin/passway"]
);
}
#[test]
fn lowered_spec_passes_shape_validation() {
let spec = lower(&sample_spec());
workload_spec::validate::shape(&spec).expect("ingress spec passes shape validation");
}
#[test]
fn spec_round_trips_through_serde_with_defaults() {
let json = r#"{ "domain": "yah.dev", "upstreams": ["a.pdx:8080"] }"#;
let spec: PasswayIngressSpec = serde_json::from_str(json).unwrap();
assert_eq!(spec.listen, DEFAULT_LISTEN);
assert_eq!(spec.upstream_tls, HostScoped::Global(false));
assert!(spec.upstream_sni.is_none(), "no SNI is passway's default");
assert!(spec.command.is_none());
assert!(spec.discover_from.is_none(), "static is the default source");
assert_eq!(spec.upstreams, vec!["a.pdx:8080".to_string()]);
}
#[test]
fn host_scoped_settings_round_trip_in_both_forms() {
let bare: PasswayIngressSpec =
serde_json::from_str(r#"{ "domain": "yah.dev", "upstream_tls": true }"#).unwrap();
assert_eq!(bare.upstream_tls, HostScoped::Global(true));
assert_eq!(serde_json::to_value(&bare.upstream_tls).unwrap(), true);
let fanned: PasswayIngressSpec = serde_json::from_str(
r#"{ "domain": "yah.dev",
"upstream_tls": { "cloud.mesh.yah.dev": true, "*": false },
"upstream_sni": { "cloud.mesh.yah.dev": "cloud.mesh.yah.dev" } }"#,
)
.unwrap();
assert_eq!(
fanned.upstream_tls.render(bool::to_string),
"*=false,cloud.mesh.yah.dev=true"
);
assert_eq!(
fanned.upstream_sni.as_ref().unwrap().render(String::clone),
"cloud.mesh.yah.dev=cloud.mesh.yah.dev"
);
}
}