Skip to main content

floating_ip_adapters/
vultr.rs

1//! [`VultrFloatingIp`] — Vultr reserved IPs (R594-F5).
2//!
3//! Vultr reserved IPs are **region-bound** (BGP-implemented inside AS20473) —
4//! verified 2026-07, W267 §Tier 1. Unlike Hetzner (network zone, a bucket of
5//! locations) or OVH (datacentre/country region, a bucket of datacenters),
6//! Vultr's own region code *is* the mobility unit — no further
7//! coarsening/mapping is needed, so this adapter uses
8//! [`floating_ip::FloatingIpMachine::location`] (Vultr's region slug, e.g.
9//! `"ewr"`) directly as the zone on both sides of the comparison in
10//! [`floating_ip::reconcile_assignment`].
11//!
12//! `"vultr"` already has an auto-provision driver bucket in `cloud`'s
13//! `provider_has_machine_driver`, so a Vultr-provider machine is expected to
14//! declare `location`.
15
16use anyhow::{bail, Context, Result};
17use async_trait::async_trait;
18use floating_ip::{FloatingIpMachine, FloatingIpProvider, FloatingIpState, FloatingIpTarget};
19use serde::Deserialize;
20
21const VULTR_BASE: &str = "https://api.vultr.com/v2";
22
23/// Vultr reserved-IP client. Bearer-token auth (Vultr's personal-access-token
24/// scheme).
25#[derive(Clone)]
26pub struct VultrFloatingIp {
27    http: reqwest::Client,
28    api_key: String,
29    base_url: String,
30}
31
32impl VultrFloatingIp {
33    /// Build against `api.vultr.com`.
34    pub fn new(api_key: impl Into<String>) -> Self {
35        Self {
36            http: reqwest::Client::new(),
37            api_key: api_key.into(),
38            base_url: VULTR_BASE.to_string(),
39        }
40    }
41
42    /// Override the base URL — used by the in-process axum mocks below.
43    pub fn with_base_url(mut self, url: impl Into<String>) -> Self {
44        self.base_url = url.into();
45        self
46    }
47}
48
49#[async_trait]
50impl FloatingIpProvider for VultrFloatingIp {
51    fn id(&self) -> &'static str {
52        "vultr"
53    }
54
55    /// `GET /v2/instances?label=<machine.name>` — Vultr instances are
56    /// looked up by their `label`, the same name-is-the-key convention
57    /// as Hetzner's `GET /servers?name=`.
58    async fn resolve_target(&self, machine: &FloatingIpMachine) -> Result<FloatingIpTarget> {
59        let zone = machine.location.clone().filter(|s| !s.is_empty()).with_context(|| {
60            format!(
61                "vultr: machine {:?} has no `location` set — required as its region (mobility zone)",
62                machine.name
63            )
64        })?;
65        let resp = self
66            .http
67            .get(format!("{}/instances", self.base_url))
68            .query(&[("label", machine.name.as_str())])
69            .bearer_auth(&self.api_key)
70            .send()
71            .await
72            .context("vultr: GET /instances")?;
73        let status = resp.status();
74        if !status.is_success() {
75            let body = resp.text().await.unwrap_or_default();
76            bail!("vultr GET /instances failed: {status} {body}");
77        }
78        let parsed: VultrInstancesResponse = resp
79            .json()
80            .await
81            .context("vultr: decode GET /instances response")?;
82        let instance = parsed
83            .instances
84            .into_iter()
85            .find(|i| i.label == machine.name)
86            .with_context(|| format!("vultr: no instance labeled {:?}", machine.name))?;
87        if instance.region != zone {
88            bail!(
89                "vultr: instance {:?} lives in region {:?}, but machine {:?} declares location {:?} — refusing to trust a mismatched declaration",
90                instance.id,
91                instance.region,
92                machine.name,
93                zone
94            );
95        }
96        Ok(FloatingIpTarget {
97            attach_id: instance.id,
98            zone,
99        })
100    }
101
102    /// `GET /v2/reserved-ips/{id}`.
103    async fn current_assignment(&self, ip_id: &str) -> Result<FloatingIpState> {
104        let resp = self
105            .http
106            .get(format!("{}/reserved-ips/{}", self.base_url, ip_id))
107            .bearer_auth(&self.api_key)
108            .send()
109            .await
110            .context("vultr: GET /reserved-ips/{id}")?;
111        let status = resp.status();
112        if !status.is_success() {
113            let body = resp.text().await.unwrap_or_default();
114            bail!("vultr GET /reserved-ips/{ip_id} failed: {status} {body}");
115        }
116        let parsed: VultrReservedIpResponse = resp
117            .json()
118            .await
119            .context("vultr: decode GET /reserved-ips/{id} response")?;
120        Ok(FloatingIpState {
121            zone: parsed.reserved_ip.region,
122            attached_to: parsed.reserved_ip.instance_id.filter(|s| !s.is_empty()),
123        })
124    }
125
126    /// `POST /v2/reserved-ips/{id}/attach`.
127    async fn reassign(&self, ip_id: &str, target: &FloatingIpTarget) -> Result<()> {
128        let resp = self
129            .http
130            .post(format!("{}/reserved-ips/{}/attach", self.base_url, ip_id))
131            .bearer_auth(&self.api_key)
132            .json(&serde_json::json!({ "instance_id": target.attach_id }))
133            .send()
134            .await
135            .context("vultr: POST /reserved-ips/{id}/attach")?;
136        let status = resp.status();
137        if !status.is_success() {
138            let body = resp.text().await.unwrap_or_default();
139            bail!("vultr POST /reserved-ips/{ip_id}/attach failed: {status} {body}");
140        }
141        Ok(())
142    }
143}
144
145#[derive(Deserialize)]
146struct VultrInstancesResponse {
147    instances: Vec<VultrInstanceLite>,
148}
149
150#[derive(Deserialize)]
151struct VultrInstanceLite {
152    id: String,
153    label: String,
154    region: String,
155}
156
157#[derive(Deserialize)]
158struct VultrReservedIpResponse {
159    reserved_ip: VultrReservedIpBody,
160}
161
162#[derive(Deserialize)]
163struct VultrReservedIpBody {
164    region: String,
165    #[serde(default)]
166    instance_id: Option<String>,
167}
168
169#[cfg(test)]
170mod tests {
171    use super::*;
172    use floating_ip::on_ingress_owner_changed;
173    use std::sync::atomic::{AtomicU32, Ordering};
174    use std::sync::{Arc, Mutex};
175
176    fn ewr_machine(name: &str) -> FloatingIpMachine {
177        FloatingIpMachine {
178            name: name.into(),
179            provider: "vultr".into(),
180            location: Some("ewr".into()),
181            region: None,
182            ingress_floating_ip: None,
183        }
184    }
185
186    /// Spin up an in-process axum server standing in for Vultr's API. See
187    /// `hetzner.rs`'s `spawn_mock` for the pattern this mirrors.
188    async fn spawn_mock(
189        instance_region: &'static str,
190        reserved_ip_region: &'static str,
191        initial_instance_id: Option<String>,
192    ) -> (String, Arc<AtomicU32>, tokio::task::JoinHandle<()>) {
193        let attached: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(initial_instance_id));
194        let attach_calls = Arc::new(AtomicU32::new(0));
195
196        let instances_route = axum::routing::get(move || async move {
197            axum::Json(serde_json::json!({
198                "instances": [ { "id": "inst-789", "label": "edge-c", "region": instance_region } ]
199            }))
200        });
201
202        let reserved_ip_get = {
203            let attached = attached.clone();
204            axum::routing::get(move || {
205                let attached = attached.clone();
206                async move {
207                    let instance_id = attached.lock().unwrap().clone();
208                    axum::Json(serde_json::json!({
209                        "reserved_ip": {
210                            "region": reserved_ip_region,
211                            "instance_id": instance_id,
212                        }
213                    }))
214                }
215            })
216        };
217
218        let attach_route = {
219            let attached = attached.clone();
220            let calls = attach_calls.clone();
221            axum::routing::post(move |axum::Json(body): axum::Json<serde_json::Value>| {
222                let attached = attached.clone();
223                let calls = calls.clone();
224                async move {
225                    calls.fetch_add(1, Ordering::SeqCst);
226                    let instance_id = body
227                        .get("instance_id")
228                        .and_then(|v| v.as_str())
229                        .map(String::from);
230                    *attached.lock().unwrap() = instance_id;
231                    axum::Json(serde_json::json!({}))
232                }
233            })
234        };
235
236        let app = axum::Router::new()
237            .route("/instances", instances_route)
238            .route("/reserved-ips/{id}", reserved_ip_get)
239            .route("/reserved-ips/{id}/attach", attach_route);
240
241        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
242        let addr = listener.local_addr().unwrap();
243        let handle = tokio::spawn(async move {
244            let _ = axum::serve(listener, app).await;
245        });
246
247        (format!("http://{addr}"), attach_calls, handle)
248    }
249
250    #[tokio::test]
251    async fn ingress_owner_flip_drives_exactly_one_reassign_call() {
252        let (base, calls, handle) = spawn_mock("ewr", "ewr", Some("old-instance".into())).await;
253        let client = VultrFloatingIp::new("test-api-key").with_base_url(base);
254        let machine = ewr_machine("edge-c");
255
256        let outcome = on_ingress_owner_changed(&client, &machine, "res-ip-1")
257            .await
258            .unwrap();
259        assert!(outcome.reassigned, "owner flip must drive a reassign");
260        assert_eq!(outcome.attached_to, "inst-789");
261        assert_eq!(calls.load(Ordering::SeqCst), 1);
262
263        handle.abort();
264    }
265
266    #[tokio::test]
267    async fn reapplying_the_same_owner_is_a_zero_call_noop() {
268        let (base, calls, handle) = spawn_mock("ewr", "ewr", Some("inst-789".into())).await;
269        let client = VultrFloatingIp::new("test-api-key").with_base_url(base);
270        let machine = ewr_machine("edge-c");
271
272        let outcome = on_ingress_owner_changed(&client, &machine, "res-ip-1")
273            .await
274            .unwrap();
275        assert!(
276            !outcome.reassigned,
277            "re-applying the same owner must be a no-op"
278        );
279        assert_eq!(calls.load(Ordering::SeqCst), 0, "must not call attach");
280
281        handle.abort();
282    }
283
284    #[tokio::test]
285    async fn cross_region_target_is_rejected_before_any_reassign_call() {
286        // Reserved IP is homed to lax; target instance lives in ewr.
287        let (base, calls, handle) = spawn_mock("ewr", "lax", None).await;
288        let client = VultrFloatingIp::new("test-api-key").with_base_url(base);
289        let machine = ewr_machine("edge-c");
290
291        let err = on_ingress_owner_changed(&client, &machine, "res-ip-1")
292            .await
293            .unwrap_err();
294        let msg = format!("{err:#}");
295        assert!(
296            msg.contains("zone"),
297            "expected a zone-mismatch error, got: {msg}"
298        );
299        assert_eq!(
300            calls.load(Ordering::SeqCst),
301            0,
302            "region mismatch must never call attach"
303        );
304
305        handle.abort();
306    }
307
308    #[tokio::test]
309    async fn instance_region_mismatching_declared_location_is_rejected() {
310        // The Vultr instance labeled "edge-c" actually lives in lax, but
311        // the machine TOML declares location=ewr — a data-consistency
312        // error the adapter must refuse to paper over.
313        let (base, calls, handle) = spawn_mock("lax", "lax", None).await;
314        let client = VultrFloatingIp::new("test-api-key").with_base_url(base);
315        let machine = ewr_machine("edge-c"); // declares ewr
316
317        let err = on_ingress_owner_changed(&client, &machine, "res-ip-1")
318            .await
319            .unwrap_err();
320        let msg = format!("{err:#}");
321        assert!(msg.contains("lax") && msg.contains("ewr"), "got: {msg}");
322        assert_eq!(calls.load(Ordering::SeqCst), 0);
323
324        handle.abort();
325    }
326}