1use anyhow::{bail, Context, Result};
26use async_trait::async_trait;
27use serde::Deserialize;
28use serde_json::Value;
29
30use super::floating_ip::{
31 reconcile_assignment, FloatingIpProvider, FloatingIpState, FloatingIpTarget,
32};
33use crate::config::MachineConfig;
34use crate::envoy::floating_ip::{
35 FloatingIpAssign, FloatingIpAssignInput, FloatingIpAssignOutput, FloatingIpStatus,
36 FloatingIpStatusInput, FloatingIpStatusOutput,
37};
38use crate::envoy::{AdapterFlavor, EnvoyAdapter, InternalVerb, Tier};
39
40const OVH_BASE: &str = "https://api.ovh.com/1.0";
41
42#[derive(Clone)]
44pub struct OvhFloatingIp {
45 http: reqwest::Client,
46 consumer_key: String,
47 base_url: String,
48}
49
50impl OvhFloatingIp {
51 pub fn new(consumer_key: impl Into<String>) -> Self {
53 Self {
54 http: reqwest::Client::new(),
55 consumer_key: consumer_key.into(),
56 base_url: OVH_BASE.to_string(),
57 }
58 }
59
60 pub fn with_base_url(mut self, url: impl Into<String>) -> Self {
62 self.base_url = url.into();
63 self
64 }
65
66 pub async fn floating_ip_assign(
68 &self,
69 input: FloatingIpAssignInput,
70 ) -> Result<FloatingIpAssignOutput> {
71 let target = FloatingIpTarget {
72 attach_id: input.attach_id,
73 zone: input.zone,
74 };
75 let outcome = reconcile_assignment(self, &input.ip_id, &target).await?;
76 Ok(FloatingIpAssignOutput {
77 reassigned: outcome.reassigned,
78 attached_to: outcome.attached_to,
79 })
80 }
81
82 pub async fn floating_ip_status(
84 &self,
85 input: FloatingIpStatusInput,
86 ) -> Result<FloatingIpStatusOutput> {
87 let state = self.current_assignment(&input.ip_id).await?;
88 Ok(FloatingIpStatusOutput {
89 zone: state.zone,
90 attached_to: state.attached_to,
91 })
92 }
93}
94
95#[async_trait]
96impl FloatingIpProvider for OvhFloatingIp {
97 fn id(&self) -> &'static str {
98 "ovh"
99 }
100
101 async fn resolve_target(&self, machine: &MachineConfig) -> Result<FloatingIpTarget> {
106 let zone = ovh_zone_for(machine)?;
107 let resp = self
108 .http
109 .get(format!(
110 "{}/dedicated/server/{}",
111 self.base_url, machine.name
112 ))
113 .header("X-Ovh-Consumer", &self.consumer_key)
114 .send()
115 .await
116 .context("ovh: GET /dedicated/server/{serviceName}")?;
117 let status = resp.status();
118 if status == reqwest::StatusCode::NOT_FOUND {
119 bail!("ovh: no dedicated server named {:?}", machine.name);
120 }
121 if !status.is_success() {
122 let body = resp.text().await.unwrap_or_default();
123 bail!(
124 "ovh GET /dedicated/server/{} failed: {status} {body}",
125 machine.name
126 );
127 }
128 Ok(FloatingIpTarget {
129 attach_id: machine.name.clone(),
130 zone,
131 })
132 }
133
134 async fn current_assignment(&self, ip_id: &str) -> Result<FloatingIpState> {
137 let resp = self
138 .http
139 .get(format!("{}/ip/{}", self.base_url, ip_id))
140 .header("X-Ovh-Consumer", &self.consumer_key)
141 .send()
142 .await
143 .context("ovh: GET /ip/{ip}")?;
144 let status = resp.status();
145 if !status.is_success() {
146 let body = resp.text().await.unwrap_or_default();
147 bail!("ovh GET /ip/{ip_id} failed: {status} {body}");
148 }
149 let parsed: OvhIpInfo = resp
150 .json()
151 .await
152 .context("ovh: decode GET /ip/{ip} response")?;
153 Ok(FloatingIpState {
154 zone: ovh_datacenter_region(&parsed.datacenter)?.to_string(),
155 attached_to: parsed.routed_to.filter(|s| !s.is_empty()),
156 })
157 }
158
159 async fn reassign(&self, ip_id: &str, target: &FloatingIpTarget) -> Result<()> {
162 let resp = self
163 .http
164 .post(format!(
165 "{}/dedicated/server/{}/ipMove",
166 self.base_url, target.attach_id
167 ))
168 .header("X-Ovh-Consumer", &self.consumer_key)
169 .json(&serde_json::json!({ "ip": ip_id }))
170 .send()
171 .await
172 .context("ovh: POST /dedicated/server/{serviceName}/ipMove")?;
173 let status = resp.status();
174 if !status.is_success() {
175 let body = resp.text().await.unwrap_or_default();
176 bail!(
177 "ovh POST /dedicated/server/{}/ipMove failed: {status} {body}",
178 target.attach_id
179 );
180 }
181 Ok(())
182 }
183}
184
185#[async_trait]
186impl EnvoyAdapter for OvhFloatingIp {
187 fn id(&self) -> &str {
188 "ovh"
189 }
190 fn tier(&self) -> Tier {
191 Tier::A
192 }
193 fn flavor(&self) -> AdapterFlavor {
194 AdapterFlavor::Native
195 }
196 fn supported_verb_ids(&self) -> Vec<&'static str> {
197 vec![FloatingIpAssign::ID, FloatingIpStatus::ID]
198 }
199 async fn dispatch(&self, verb_id: &str, input: Value) -> Result<Value> {
200 match verb_id {
201 id if id == FloatingIpAssign::ID => {
202 let args: FloatingIpAssignInput =
203 serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
204 let out = self.floating_ip_assign(args).await?;
205 Ok(serde_json::to_value(out)?)
206 }
207 id if id == FloatingIpStatus::ID => {
208 let args: FloatingIpStatusInput =
209 serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
210 let out = self.floating_ip_status(args).await?;
211 Ok(serde_json::to_value(out)?)
212 }
213 other => bail!("ovh floating-ip envoy does not support verb {other:?}"),
214 }
215 }
216}
217
218fn ovh_datacenter_region(dc: &str) -> Result<&'static str> {
223 match dc.to_ascii_lowercase().as_str() {
224 "gra" | "rbx" | "sbg" => Ok("eu-west"),
225 "bhs" => Ok("ca-east"),
226 "waw" => Ok("eu-central-pl"),
227 "syd" => Ok("au-east"),
228 "sgp" => Ok("ap-southeast"),
229 other => bail!("ovh: unknown datacenter {other:?}, cannot derive mobility region"),
230 }
231}
232
233fn ovh_zone_for(machine: &MachineConfig) -> Result<String> {
243 if let Some(dc) = machine.location.as_deref().filter(|s| !s.is_empty()) {
244 return Ok(ovh_datacenter_region(dc)?.to_string());
245 }
246 machine.region.clone().with_context(|| {
247 format!(
248 "ovh: machine {:?} has neither `location` nor `region` set — cannot derive its mobility zone",
249 machine.name
250 )
251 })
252}
253
254#[derive(Deserialize)]
255struct OvhIpInfo {
256 datacenter: String,
257 #[serde(rename = "routedTo", default)]
258 routed_to: Option<String>,
259}
260
261#[cfg(test)]
262mod tests {
263 use super::*;
264 use crate::provider::floating_ip::on_ingress_owner_changed;
265 use std::sync::atomic::{AtomicU32, Ordering};
266 use std::sync::{Arc, Mutex};
267
268 fn gra_machine(name: &str) -> MachineConfig {
269 MachineConfig {
270 name: name.into(),
271 provider: "ovh".into(),
272 location: None,
273 server_type: None,
274 hosts_mirrors: vec![],
275 mesh_tags: vec![],
276 region: Some("eu-west".into()),
277 zone: None,
278 arch: None,
279 bucket: None,
280 vendor: None,
281 nickname: None,
282 legacy_hostkey_fingerprint: None,
283 registration: Default::default(),
284 ssh_keys: vec![],
285 cloudflared: None,
286 hosts_operator_bridge: false,
287 connect: None,
288 allocatable: None,
289 taints: vec![],
290 }
291 }
292
293 async fn spawn_mock(
297 datacenter: &'static str,
298 initial_routed_to: Option<String>,
299 ) -> (String, Arc<AtomicU32>, tokio::task::JoinHandle<()>) {
300 let routed_to: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(initial_routed_to));
301 let move_calls = Arc::new(AtomicU32::new(0));
302
303 let server_exists_route =
304 axum::routing::get(|| async { axum::Json(serde_json::json!({ "datacenter": "gra" })) });
305
306 let ip_get_route = {
307 let routed_to = routed_to.clone();
308 axum::routing::get(move || {
309 let routed_to = routed_to.clone();
310 async move {
311 let current = routed_to.lock().unwrap().clone();
312 axum::Json(serde_json::json!({
313 "datacenter": datacenter,
314 "routedTo": current,
315 }))
316 }
317 })
318 };
319
320 let ip_move_route = {
321 let routed_to = routed_to.clone();
322 let calls = move_calls.clone();
323 axum::routing::post(
324 move |axum::extract::Path(service_name): axum::extract::Path<String>,
325 axum::Json(_body): axum::Json<serde_json::Value>| {
326 let routed_to = routed_to.clone();
327 let calls = calls.clone();
328 async move {
329 calls.fetch_add(1, Ordering::SeqCst);
330 *routed_to.lock().unwrap() = Some(service_name);
331 axum::Json(serde_json::json!({}))
332 }
333 },
334 )
335 };
336
337 let app = axum::Router::new()
338 .route("/dedicated/server/{service_name}", server_exists_route)
339 .route("/ip/{ip}", ip_get_route)
340 .route("/dedicated/server/{service_name}/ipMove", ip_move_route);
341
342 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
343 let addr = listener.local_addr().unwrap();
344 let handle = tokio::spawn(async move {
345 let _ = axum::serve(listener, app).await;
346 });
347
348 (format!("http://{addr}"), move_calls, handle)
349 }
350
351 #[tokio::test]
352 async fn ingress_owner_flip_drives_exactly_one_reassign_call() {
353 let (base, calls, handle) = spawn_mock("gra", Some("old-server".into())).await;
354 let client = OvhFloatingIp::new("test-consumer-key").with_base_url(base);
355 let machine = gra_machine("edge-b");
356
357 let outcome = on_ingress_owner_changed(&client, &machine, "51.81.85.200")
358 .await
359 .unwrap();
360 assert!(outcome.reassigned, "owner flip must drive a reassign");
361 assert_eq!(outcome.attached_to, "edge-b");
362 assert_eq!(calls.load(Ordering::SeqCst), 1);
363
364 handle.abort();
365 }
366
367 #[tokio::test]
368 async fn reapplying_the_same_owner_is_a_zero_call_noop() {
369 let (base, calls, handle) = spawn_mock("gra", Some("edge-b".into())).await;
370 let client = OvhFloatingIp::new("test-consumer-key").with_base_url(base);
371 let machine = gra_machine("edge-b");
372
373 let outcome = on_ingress_owner_changed(&client, &machine, "51.81.85.200")
374 .await
375 .unwrap();
376 assert!(
377 !outcome.reassigned,
378 "re-applying the same owner must be a no-op"
379 );
380 assert_eq!(calls.load(Ordering::SeqCst), 0, "must not call ipMove");
381
382 handle.abort();
383 }
384
385 #[tokio::test]
386 async fn cross_region_target_is_rejected_before_any_reassign_call() {
387 let (base, calls, handle) = spawn_mock("bhs", None).await;
389 let client = OvhFloatingIp::new("test-consumer-key").with_base_url(base);
390 let machine = gra_machine("edge-b"); let err = on_ingress_owner_changed(&client, &machine, "51.81.85.200")
393 .await
394 .unwrap_err();
395 let msg = format!("{err:#}");
396 assert!(
397 msg.contains("zone"),
398 "expected a zone-mismatch error, got: {msg}"
399 );
400 assert_eq!(
401 calls.load(Ordering::SeqCst),
402 0,
403 "region mismatch must never call ipMove"
404 );
405
406 handle.abort();
407 }
408
409 #[test]
410 fn ovh_datacenter_region_maps_eu_west_trio() {
411 assert_eq!(ovh_datacenter_region("gra").unwrap(), "eu-west");
412 assert_eq!(ovh_datacenter_region("rbx").unwrap(), "eu-west");
413 assert_eq!(ovh_datacenter_region("sbg").unwrap(), "eu-west");
414 assert_eq!(ovh_datacenter_region("bhs").unwrap(), "ca-east");
415 }
416
417 #[test]
418 fn ovh_datacenter_region_rejects_unknown() {
419 assert!(ovh_datacenter_region("xyz").is_err());
420 }
421
422 #[test]
423 fn ovh_zone_for_prefers_location_over_region() {
424 let mut m = gra_machine("edge-b");
425 m.location = Some("bhs".into());
426 assert_eq!(ovh_zone_for(&m).unwrap(), "ca-east");
427 }
428
429 #[test]
430 fn ovh_zone_for_falls_back_to_region() {
431 let m = gra_machine("edge-b"); assert_eq!(ovh_zone_for(&m).unwrap(), "eu-west");
433 }
434
435 #[test]
436 fn ovh_zone_for_errors_with_neither_field() {
437 let mut m = gra_machine("edge-b");
438 m.region = None;
439 assert!(ovh_zone_for(&m).is_err());
440 }
441}