Skip to main content

cloud/provider/
vultr_floating_ip.rs

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