use anyhow::{bail, Context, Result};
use async_trait::async_trait;
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct FloatingIpMachine {
pub name: String,
pub provider: String,
pub location: Option<String>,
pub region: Option<String>,
pub ingress_floating_ip: Option<String>,
}
impl FloatingIpMachine {
pub fn location(&self) -> &str {
self.location.as_deref().unwrap_or("")
}
}
#[async_trait]
pub trait FloatingIpProvider: Send + Sync {
fn id(&self) -> &'static str;
async fn resolve_target(&self, machine: &FloatingIpMachine) -> Result<FloatingIpTarget>;
async fn current_assignment(&self, ip_id: &str) -> Result<FloatingIpState>;
async fn reassign(&self, ip_id: &str, target: &FloatingIpTarget) -> Result<()>;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FloatingIpTarget {
pub attach_id: String,
pub zone: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FloatingIpState {
pub zone: String,
pub attached_to: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FloatingIpAssignOutcome {
pub reassigned: bool,
pub attached_to: String,
}
pub async fn reconcile_assignment<P: FloatingIpProvider + ?Sized>(
provider: &P,
ip_id: &str,
target: &FloatingIpTarget,
) -> Result<FloatingIpAssignOutcome> {
let current = provider.current_assignment(ip_id).await?;
if current.zone != target.zone {
bail!(
"floating_ip.assign: {} ip {ip_id:?} is homed to zone {:?}, cannot move it into zone {:?} (target attach id {:?}) — {} floating/reserved IPs are not mobile across zones (W267 §Tier 1)",
provider.id(),
current.zone,
target.zone,
target.attach_id,
provider.id(),
);
}
if current.attached_to.as_deref() == Some(target.attach_id.as_str()) {
return Ok(FloatingIpAssignOutcome {
reassigned: false,
attached_to: target.attach_id.clone(),
});
}
provider.reassign(ip_id, target).await?;
Ok(FloatingIpAssignOutcome {
reassigned: true,
attached_to: target.attach_id.clone(),
})
}
pub async fn on_ingress_owner_changed<P: FloatingIpProvider + ?Sized>(
provider: &P,
machine: &FloatingIpMachine,
ip_id: &str,
) -> Result<FloatingIpAssignOutcome> {
let target = provider.resolve_target(machine).await?;
reconcile_assignment(provider, ip_id, &target).await
}
pub const FLOATING_IP_PROVIDERS: &[(&str, &str, &str)] = &[
("hetzner", "hetzner-api-token", "HETZNER_API_TOKEN"),
("ovh", "ovh-consumer-key", "OVH_CONSUMER_KEY"),
("vultr", "vultr-api-key", "VULTR_API_KEY"),
];
pub fn provider_has_floating_ip_adapter(provider: &str) -> bool {
FLOATING_IP_PROVIDERS.iter().any(|(id, _, _)| *id == provider)
}
pub fn supported_floating_ip_providers() -> String {
FLOATING_IP_PROVIDERS
.iter()
.map(|(id, _, _)| *id)
.collect::<Vec<_>>()
.join(", ")
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OwnerLiveness {
ConfirmedUp,
ConfirmedDown,
Unconfirmed,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum QuorumHealth {
Healthy,
Degraded {
reason: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum IngressOwnerEffect {
Reassign {
machine: String,
ip_id: String,
},
Withdraw {
machine: String,
reason: String,
},
Refuse { reason: String },
NoOp { reason: String },
}
impl IngressOwnerEffect {
pub fn is_action(&self) -> bool {
matches!(self, Self::Reassign { .. } | Self::Withdraw { .. })
}
pub fn reason(&self) -> String {
match self {
Self::Reassign { machine, ip_id } => {
format!("reassign floating IP {ip_id} to {machine}")
}
Self::Withdraw { machine, reason } => {
format!("withdraw {machine} from the apex: {reason}")
}
Self::Refuse { reason } | Self::NoOp { reason } => reason.clone(),
}
}
}
pub fn resolve_ingress_owner<'a>(
owner: &str,
machines: &'a [FloatingIpMachine],
) -> Result<&'a FloatingIpMachine> {
machines
.iter()
.find(|m| m.name == owner)
.with_context(|| {
format!(
"raft names {owner:?} as the ingress owner, but no .yah/infra/machines/*.toml \
declares a machine with that name (declared: {}). Note `ingress_owner` carries \
the node's /etc/hostname, which is not always its machine name — R841 saw \
`vps-4c1efa56` recorded for the box declared as `us-west-001`. Rename the box's \
hostname to match its machine name, or this mapping cannot be made safely.",
machines
.iter()
.map(|m| m.name.as_str())
.collect::<Vec<_>>()
.join(", "),
)
})
}
pub fn plan_ingress_owner_effect(
previous_owner: Option<&str>,
current_owner: Option<&str>,
current_owner_liveness: OwnerLiveness,
quorum: &QuorumHealth,
machines: &[FloatingIpMachine],
) -> IngressOwnerEffect {
let Some(owner) = current_owner else {
return IngressOwnerEffect::NoOp {
reason: match previous_owner {
Some(prev) => format!(
"ingress owner cleared (was {prev}) — leaving the floating IP on the \
last-known-good node; there is no detach verb and no specified \
safe-unassigned state at Tier 1"
),
None => "no ingress owner recorded".to_string(),
},
};
};
let machine = match resolve_ingress_owner(owner, machines) {
Ok(m) => m,
Err(e) => return IngressOwnerEffect::Refuse { reason: format!("{e:#}") },
};
let owner_changed = previous_owner != Some(owner);
if owner_changed {
let Some(ip_id) = machine.ingress_floating_ip.as_deref() else {
return IngressOwnerEffect::NoOp {
reason: format!(
"ingress owner moved to {owner}, which declares no `ingress_floating_ip` — \
this machine has no floating-IP path"
),
};
};
if current_owner_liveness == OwnerLiveness::ConfirmedDown {
return IngressOwnerEffect::Refuse {
reason: format!(
"ingress owner moved to {owner}, but liveness has confirmed it DOWN — \
refusing to point the public IP at a box we have positive evidence is dead"
),
};
}
if let QuorumHealth::Degraded { reason } = quorum {
return IngressOwnerEffect::Refuse {
reason: format!(
"ingress owner moved to {owner} but the reassign is refused: {reason} \
(yubaba-failover.md pre-check 1). A reassign takes the IP off the old \
owner, so it is a withdrawal and fails closed."
),
};
}
return IngressOwnerEffect::Reassign {
machine: owner.to_string(),
ip_id: ip_id.to_string(),
};
}
if current_owner_liveness == OwnerLiveness::ConfirmedDown {
if let QuorumHealth::Degraded { reason } = quorum {
return IngressOwnerEffect::Refuse {
reason: format!(
"ingress owner {owner} is confirmed down, but the withdrawal is refused: \
{reason} (yubaba-failover.md pre-check 1)"
),
};
}
return IngressOwnerEffect::Withdraw {
machine: owner.to_string(),
reason: format!("ingress owner {owner} is confirmed down by the lease channel"),
};
}
IngressOwnerEffect::NoOp {
reason: format!("ingress owner unchanged ({owner}) and not confirmed down"),
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::Mutex;
struct FakeProvider {
zone: &'static str,
attached_to: Mutex<Option<String>>,
reassign_calls: AtomicU32,
}
#[async_trait]
impl FloatingIpProvider for FakeProvider {
fn id(&self) -> &'static str {
"fake"
}
async fn resolve_target(&self, machine: &FloatingIpMachine) -> Result<FloatingIpTarget> {
Ok(FloatingIpTarget {
attach_id: machine.name.clone(),
zone: self.zone.to_string(),
})
}
async fn current_assignment(&self, _ip_id: &str) -> Result<FloatingIpState> {
Ok(FloatingIpState {
zone: self.zone.to_string(),
attached_to: self.attached_to.lock().unwrap().clone(),
})
}
async fn reassign(&self, _ip_id: &str, target: &FloatingIpTarget) -> Result<()> {
self.reassign_calls.fetch_add(1, Ordering::SeqCst);
*self.attached_to.lock().unwrap() = Some(target.attach_id.clone());
Ok(())
}
}
fn machine(name: &str) -> FloatingIpMachine {
FloatingIpMachine {
name: name.into(),
provider: "fake".into(),
..Default::default()
}
}
#[tokio::test]
async fn ownership_flip_drives_exactly_one_reassign_call() {
let provider = FakeProvider {
zone: "us-west",
attached_to: Mutex::new(Some("old-node".into())),
reassign_calls: AtomicU32::new(0),
};
let outcome = on_ingress_owner_changed(&provider, &machine("new-node"), "ip-1")
.await
.unwrap();
assert!(outcome.reassigned);
assert_eq!(outcome.attached_to, "new-node");
assert_eq!(provider.reassign_calls.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn reapplying_the_same_owner_is_a_zero_call_noop() {
let provider = FakeProvider {
zone: "us-west",
attached_to: Mutex::new(Some("new-node".into())),
reassign_calls: AtomicU32::new(0),
};
let outcome = on_ingress_owner_changed(&provider, &machine("new-node"), "ip-1")
.await
.unwrap();
assert!(!outcome.reassigned);
assert_eq!(outcome.attached_to, "new-node");
assert_eq!(
provider.reassign_calls.load(Ordering::SeqCst),
0,
"idempotent re-apply must not call reassign"
);
}
#[tokio::test]
async fn never_assigned_ip_gets_a_first_assign_call() {
let provider = FakeProvider {
zone: "us-west",
attached_to: Mutex::new(None),
reassign_calls: AtomicU32::new(0),
};
let outcome = on_ingress_owner_changed(&provider, &machine("new-node"), "ip-1")
.await
.unwrap();
assert!(outcome.reassigned);
assert_eq!(provider.reassign_calls.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn cross_zone_target_is_rejected_before_any_reassign_call() {
let provider = FakeProvider {
zone: "eu-central",
attached_to: Mutex::new(None),
reassign_calls: AtomicU32::new(0),
};
let target = FloatingIpTarget {
attach_id: "new-node".into(),
zone: "us-west".into(),
};
let err = reconcile_assignment(&provider, "ip-1", &target)
.await
.unwrap_err();
let msg = format!("{err:#}");
assert!(
msg.contains("zone"),
"expected a zone-mismatch message, got: {msg}"
);
assert_eq!(
provider.reassign_calls.load(Ordering::SeqCst),
0,
"zone mismatch must never call reassign"
);
}
#[test]
fn the_three_shipped_adapters_are_all_reachable_by_provider_id() {
for id in ["hetzner", "ovh", "vultr"] {
assert!(
provider_has_floating_ip_adapter(id),
"{id} ships a FloatingIpProvider impl but the registry cannot reach it"
);
}
for id in ["digitalocean", "static", "local-docker", ""] {
assert!(!provider_has_floating_ip_adapter(id), "{id}");
}
}
#[test]
fn every_registered_provider_is_named_in_the_supported_list() {
let supported = supported_floating_ip_providers();
for (id, _, _) in FLOATING_IP_PROVIDERS {
assert!(supported.contains(id), "{id} missing from {supported:?}");
}
}
fn fleet() -> Vec<FloatingIpMachine> {
let mut west = machine("us-west-001");
west.provider = "hetzner".into();
west.ingress_floating_ip = Some("fip-42".into());
let mut east = machine("us-east-001");
east.provider = "hetzner".into();
east.ingress_floating_ip = Some("fip-42".into());
let mesh_only = machine("us-west-002");
vec![west, east, mesh_only]
}
fn degraded() -> QuorumHealth {
QuorumHealth::Degraded {
reason: "quorum AT RISK: 2/3 voters available".into(),
}
}
#[test]
fn an_ownership_flip_onto_a_machine_with_a_floating_ip_reassigns_it() {
let effect = plan_ingress_owner_effect(
Some("us-west-001"),
Some("us-east-001"),
OwnerLiveness::ConfirmedUp,
&QuorumHealth::Healthy,
&fleet(),
);
assert_eq!(
effect,
IngressOwnerEffect::Reassign {
machine: "us-east-001".into(),
ip_id: "fip-42".into(),
}
);
assert!(effect.is_action());
}
#[test]
fn an_unconfirmed_new_owner_still_reassigns_because_liveness_may_only_veto() {
assert!(matches!(
plan_ingress_owner_effect(
Some("us-west-001"),
Some("us-east-001"),
OwnerLiveness::Unconfirmed,
&QuorumHealth::Healthy,
&fleet(),
),
IngressOwnerEffect::Reassign { .. }
));
}
#[test]
fn a_new_owner_confirmed_down_is_refused_rather_than_pointed_at() {
let effect = plan_ingress_owner_effect(
Some("us-west-001"),
Some("us-east-001"),
OwnerLiveness::ConfirmedDown,
&QuorumHealth::Healthy,
&fleet(),
);
assert!(matches!(effect, IngressOwnerEffect::Refuse { .. }), "{effect:?}");
assert!(effect.reason().contains("DOWN"), "{}", effect.reason());
}
#[test]
fn a_degraded_quorum_refuses_the_reassign_and_carries_the_verdicts_reason() {
let effect = plan_ingress_owner_effect(
Some("us-west-001"),
Some("us-east-001"),
OwnerLiveness::ConfirmedUp,
°raded(),
&fleet(),
);
assert!(matches!(effect, IngressOwnerEffect::Refuse { .. }), "{effect:?}");
assert!(
effect.reason().contains("2/3 voters available"),
"the refusal must carry the quorum verdict's own reason, got: {}",
effect.reason()
);
}
#[test]
fn a_machine_with_no_floating_ip_is_a_clean_skip_not_an_error() {
let effect = plan_ingress_owner_effect(
Some("us-west-001"),
Some("us-west-002"),
OwnerLiveness::ConfirmedUp,
&QuorumHealth::Healthy,
&fleet(),
);
assert!(matches!(effect, IngressOwnerEffect::NoOp { .. }), "{effect:?}");
assert!(!effect.is_action());
assert!(
effect.reason().contains("no floating-IP path"),
"{}",
effect.reason()
);
}
#[test]
fn a_steady_healthy_owner_does_nothing() {
let effect = plan_ingress_owner_effect(
Some("us-east-001"),
Some("us-east-001"),
OwnerLiveness::ConfirmedUp,
&QuorumHealth::Healthy,
&fleet(),
);
assert!(matches!(effect, IngressOwnerEffect::NoOp { .. }), "{effect:?}");
}
#[test]
fn a_steady_owner_confirmed_down_is_withdrawn_from_the_apex() {
let effect = plan_ingress_owner_effect(
Some("us-east-001"),
Some("us-east-001"),
OwnerLiveness::ConfirmedDown,
&QuorumHealth::Healthy,
&fleet(),
);
assert_eq!(
effect,
IngressOwnerEffect::Withdraw {
machine: "us-east-001".into(),
reason: "ingress owner us-east-001 is confirmed down by the lease channel".into(),
}
);
}
#[test]
fn a_degraded_quorum_refuses_the_withdrawal_too() {
let effect = plan_ingress_owner_effect(
Some("us-east-001"),
Some("us-east-001"),
OwnerLiveness::ConfirmedDown,
°raded(),
&fleet(),
);
assert!(matches!(effect, IngressOwnerEffect::Refuse { .. }), "{effect:?}");
assert!(effect.reason().contains("2/3 voters available"), "{}", effect.reason());
}
#[test]
fn an_ingress_owner_that_names_no_declared_machine_refuses_loudly() {
let effect = plan_ingress_owner_effect(
Some("us-west-001"),
Some("vps-4c1efa56"),
OwnerLiveness::ConfirmedUp,
&QuorumHealth::Healthy,
&fleet(),
);
assert!(matches!(effect, IngressOwnerEffect::Refuse { .. }), "{effect:?}");
let reason = effect.reason();
assert!(reason.contains("vps-4c1efa56"), "{reason}");
assert!(
reason.contains("us-west-001") && reason.contains("us-east-001"),
"the refusal must name the declared machines it compared against: {reason}"
);
assert!(
reason.contains("hostname"),
"and must explain WHY the two spaces differ: {reason}"
);
}
#[test]
fn clearing_the_ingress_owner_leaves_the_ip_where_it_is() {
let effect = plan_ingress_owner_effect(
Some("us-east-001"),
None,
OwnerLiveness::ConfirmedDown,
&QuorumHealth::Healthy,
&fleet(),
);
assert!(matches!(effect, IngressOwnerEffect::NoOp { .. }), "{effect:?}");
assert!(
effect.reason().contains("last-known-good"),
"{}",
effect.reason()
);
}
#[test]
fn no_ingress_owner_at_all_is_a_no_op() {
assert!(matches!(
plan_ingress_owner_effect(
None,
None,
OwnerLiveness::Unconfirmed,
&QuorumHealth::Healthy,
&fleet(),
),
IngressOwnerEffect::NoOp { .. }
));
}
#[test]
fn a_first_observation_of_an_existing_owner_converges_the_ip() {
assert_eq!(
plan_ingress_owner_effect(
None,
Some("us-east-001"),
OwnerLiveness::ConfirmedUp,
&QuorumHealth::Healthy,
&fleet(),
),
IngressOwnerEffect::Reassign {
machine: "us-east-001".into(),
ip_id: "fip-42".into(),
},
"reconcile_assignment is idempotent, so a redundant converge costs zero \
provider calls — but skipping it would leave a stale IP unfixed forever"
);
}
}