1use 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#[derive(Clone)]
37pub struct VultrFloatingIp {
38 http: reqwest::Client,
39 api_key: String,
40 base_url: String,
41}
42
43impl VultrFloatingIp {
44 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 pub fn with_base_url(mut self, url: impl Into<String>) -> Self {
55 self.base_url = url.into();
56 self
57 }
58
59 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 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 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 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 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 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 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 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"); 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}