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 sovereign_group: None,
291 sovereign_role: None,
292 }
293 }
294
295 async fn spawn_mock(
299 datacenter: &'static str,
300 initial_routed_to: Option<String>,
301 ) -> (String, Arc<AtomicU32>, tokio::task::JoinHandle<()>) {
302 let routed_to: Arc<Mutex<Option<String>>> = Arc::new(Mutex::new(initial_routed_to));
303 let move_calls = Arc::new(AtomicU32::new(0));
304
305 let server_exists_route =
306 axum::routing::get(|| async { axum::Json(serde_json::json!({ "datacenter": "gra" })) });
307
308 let ip_get_route = {
309 let routed_to = routed_to.clone();
310 axum::routing::get(move || {
311 let routed_to = routed_to.clone();
312 async move {
313 let current = routed_to.lock().unwrap().clone();
314 axum::Json(serde_json::json!({
315 "datacenter": datacenter,
316 "routedTo": current,
317 }))
318 }
319 })
320 };
321
322 let ip_move_route = {
323 let routed_to = routed_to.clone();
324 let calls = move_calls.clone();
325 axum::routing::post(
326 move |axum::extract::Path(service_name): axum::extract::Path<String>,
327 axum::Json(_body): axum::Json<serde_json::Value>| {
328 let routed_to = routed_to.clone();
329 let calls = calls.clone();
330 async move {
331 calls.fetch_add(1, Ordering::SeqCst);
332 *routed_to.lock().unwrap() = Some(service_name);
333 axum::Json(serde_json::json!({}))
334 }
335 },
336 )
337 };
338
339 let app = axum::Router::new()
340 .route("/dedicated/server/{service_name}", server_exists_route)
341 .route("/ip/{ip}", ip_get_route)
342 .route("/dedicated/server/{service_name}/ipMove", ip_move_route);
343
344 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
345 let addr = listener.local_addr().unwrap();
346 let handle = tokio::spawn(async move {
347 let _ = axum::serve(listener, app).await;
348 });
349
350 (format!("http://{addr}"), move_calls, handle)
351 }
352
353 #[tokio::test]
354 async fn ingress_owner_flip_drives_exactly_one_reassign_call() {
355 let (base, calls, handle) = spawn_mock("gra", Some("old-server".into())).await;
356 let client = OvhFloatingIp::new("test-consumer-key").with_base_url(base);
357 let machine = gra_machine("edge-b");
358
359 let outcome = on_ingress_owner_changed(&client, &machine, "51.81.85.200")
360 .await
361 .unwrap();
362 assert!(outcome.reassigned, "owner flip must drive a reassign");
363 assert_eq!(outcome.attached_to, "edge-b");
364 assert_eq!(calls.load(Ordering::SeqCst), 1);
365
366 handle.abort();
367 }
368
369 #[tokio::test]
370 async fn reapplying_the_same_owner_is_a_zero_call_noop() {
371 let (base, calls, handle) = spawn_mock("gra", Some("edge-b".into())).await;
372 let client = OvhFloatingIp::new("test-consumer-key").with_base_url(base);
373 let machine = gra_machine("edge-b");
374
375 let outcome = on_ingress_owner_changed(&client, &machine, "51.81.85.200")
376 .await
377 .unwrap();
378 assert!(
379 !outcome.reassigned,
380 "re-applying the same owner must be a no-op"
381 );
382 assert_eq!(calls.load(Ordering::SeqCst), 0, "must not call ipMove");
383
384 handle.abort();
385 }
386
387 #[tokio::test]
388 async fn cross_region_target_is_rejected_before_any_reassign_call() {
389 let (base, calls, handle) = spawn_mock("bhs", None).await;
391 let client = OvhFloatingIp::new("test-consumer-key").with_base_url(base);
392 let machine = gra_machine("edge-b"); let err = on_ingress_owner_changed(&client, &machine, "51.81.85.200")
395 .await
396 .unwrap_err();
397 let msg = format!("{err:#}");
398 assert!(
399 msg.contains("zone"),
400 "expected a zone-mismatch error, got: {msg}"
401 );
402 assert_eq!(
403 calls.load(Ordering::SeqCst),
404 0,
405 "region mismatch must never call ipMove"
406 );
407
408 handle.abort();
409 }
410
411 #[test]
412 fn ovh_datacenter_region_maps_eu_west_trio() {
413 assert_eq!(ovh_datacenter_region("gra").unwrap(), "eu-west");
414 assert_eq!(ovh_datacenter_region("rbx").unwrap(), "eu-west");
415 assert_eq!(ovh_datacenter_region("sbg").unwrap(), "eu-west");
416 assert_eq!(ovh_datacenter_region("bhs").unwrap(), "ca-east");
417 }
418
419 #[test]
420 fn ovh_datacenter_region_rejects_unknown() {
421 assert!(ovh_datacenter_region("xyz").is_err());
422 }
423
424 #[test]
425 fn ovh_zone_for_prefers_location_over_region() {
426 let mut m = gra_machine("edge-b");
427 m.location = Some("bhs".into());
428 assert_eq!(ovh_zone_for(&m).unwrap(), "ca-east");
429 }
430
431 #[test]
432 fn ovh_zone_for_falls_back_to_region() {
433 let m = gra_machine("edge-b"); assert_eq!(ovh_zone_for(&m).unwrap(), "eu-west");
435 }
436
437 #[test]
438 fn ovh_zone_for_errors_with_neither_field() {
439 let mut m = gra_machine("edge-b");
440 m.region = None;
441 assert!(ovh_zone_for(&m).is_err());
442 }
443}