1use 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#[derive(Clone)]
26pub struct VultrFloatingIp {
27 http: reqwest::Client,
28 api_key: String,
29 base_url: String,
30}
31
32impl VultrFloatingIp {
33 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 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 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 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 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 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 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 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"); 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}