1use anyhow::{bail, Context, Result};
22use async_trait::async_trait;
23use serde::Deserialize;
24use serde_json::Value;
25
26use super::floating_ip::{
27 reconcile_assignment, FloatingIpProvider, FloatingIpState, FloatingIpTarget,
28};
29use crate::config::MachineConfig;
30use crate::envoy::floating_ip::{
31 FloatingIpAssign, FloatingIpAssignInput, FloatingIpAssignOutput, FloatingIpStatus,
32 FloatingIpStatusInput, FloatingIpStatusOutput,
33};
34use crate::envoy::{AdapterFlavor, EnvoyAdapter, InternalVerb, Tier};
35
36const HETZNER_BASE: &str = "https://api.hetzner.cloud/v1";
37
38#[derive(Clone)]
41pub struct HetznerFloatingIp {
42 http: reqwest::Client,
43 token: String,
44 base_url: String,
45}
46
47impl HetznerFloatingIp {
48 pub fn new(token: impl Into<String>) -> Self {
50 Self {
51 http: reqwest::Client::new(),
52 token: token.into(),
53 base_url: HETZNER_BASE.to_string(),
54 }
55 }
56
57 pub fn with_base_url(mut self, url: impl Into<String>) -> Self {
60 self.base_url = url.into();
61 self
62 }
63
64 pub async fn floating_ip_assign(
67 &self,
68 input: FloatingIpAssignInput,
69 ) -> Result<FloatingIpAssignOutput> {
70 let target = FloatingIpTarget {
71 attach_id: input.attach_id,
72 zone: input.zone,
73 };
74 let outcome = reconcile_assignment(self, &input.ip_id, &target).await?;
75 Ok(FloatingIpAssignOutput {
76 reassigned: outcome.reassigned,
77 attached_to: outcome.attached_to,
78 })
79 }
80
81 pub async fn floating_ip_status(
83 &self,
84 input: FloatingIpStatusInput,
85 ) -> Result<FloatingIpStatusOutput> {
86 let state = self.current_assignment(&input.ip_id).await?;
87 Ok(FloatingIpStatusOutput {
88 zone: state.zone,
89 attached_to: state.attached_to,
90 })
91 }
92}
93
94#[async_trait]
95impl FloatingIpProvider for HetznerFloatingIp {
96 fn id(&self) -> &'static str {
97 "hetzner"
98 }
99
100 async fn resolve_target(&self, machine: &MachineConfig) -> Result<FloatingIpTarget> {
104 let zone = hetzner_network_zone(machine.location()).with_context(|| {
105 format!(
106 "floating_ip: resolving target for machine {:?}",
107 machine.name
108 )
109 })?;
110 let resp = self
111 .http
112 .get(format!("{}/servers", self.base_url))
113 .query(&[("name", machine.name.as_str())])
114 .bearer_auth(&self.token)
115 .send()
116 .await
117 .context("hetzner: GET /servers")?;
118 let status = resp.status();
119 if !status.is_success() {
120 let body = resp.text().await.unwrap_or_default();
121 bail!("hetzner GET /servers failed: {status} {body}");
122 }
123 let parsed: HetznerServersResponse = resp
124 .json()
125 .await
126 .context("hetzner: decode GET /servers response")?;
127 let server = parsed
128 .servers
129 .into_iter()
130 .find(|s| s.name.as_deref() == Some(machine.name.as_str()))
131 .with_context(|| format!("hetzner: no server named {:?}", machine.name))?;
132 Ok(FloatingIpTarget {
133 attach_id: server.id.to_string(),
134 zone: zone.to_string(),
135 })
136 }
137
138 async fn current_assignment(&self, ip_id: &str) -> Result<FloatingIpState> {
140 let resp = self
141 .http
142 .get(format!("{}/floating_ips/{}", self.base_url, ip_id))
143 .bearer_auth(&self.token)
144 .send()
145 .await
146 .context("hetzner: GET /floating_ips/{id}")?;
147 let status = resp.status();
148 if !status.is_success() {
149 let body = resp.text().await.unwrap_or_default();
150 bail!("hetzner GET /floating_ips/{ip_id} failed: {status} {body}");
151 }
152 let parsed: HetznerFloatingIpResponse = resp
153 .json()
154 .await
155 .context("hetzner: decode GET /floating_ips/{id} response")?;
156 Ok(FloatingIpState {
157 zone: parsed.floating_ip.home_location.network_zone,
158 attached_to: parsed.floating_ip.server.map(|id| id.to_string()),
159 })
160 }
161
162 async fn reassign(&self, ip_id: &str, target: &FloatingIpTarget) -> Result<()> {
164 let server_id: u64 = target.attach_id.parse().with_context(|| {
165 format!(
166 "hetzner: attach_id {:?} is not a numeric server id",
167 target.attach_id
168 )
169 })?;
170 let resp = self
171 .http
172 .post(format!(
173 "{}/floating_ips/{}/actions/assign",
174 self.base_url, ip_id
175 ))
176 .bearer_auth(&self.token)
177 .json(&serde_json::json!({ "server": server_id }))
178 .send()
179 .await
180 .context("hetzner: POST /floating_ips/{id}/actions/assign")?;
181 let status = resp.status();
182 if !status.is_success() {
183 let body = resp.text().await.unwrap_or_default();
184 bail!("hetzner POST /floating_ips/{ip_id}/actions/assign failed: {status} {body}");
185 }
186 Ok(())
187 }
188}
189
190#[async_trait]
191impl EnvoyAdapter for HetznerFloatingIp {
192 fn id(&self) -> &str {
193 "hetzner"
194 }
195 fn tier(&self) -> Tier {
196 Tier::S
197 }
198 fn flavor(&self) -> AdapterFlavor {
199 AdapterFlavor::Native
200 }
201 fn supported_verb_ids(&self) -> Vec<&'static str> {
202 vec![FloatingIpAssign::ID, FloatingIpStatus::ID]
203 }
204 async fn dispatch(&self, verb_id: &str, input: Value) -> Result<Value> {
205 match verb_id {
206 id if id == FloatingIpAssign::ID => {
207 let args: FloatingIpAssignInput =
208 serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
209 let out = self.floating_ip_assign(args).await?;
210 Ok(serde_json::to_value(out)?)
211 }
212 id if id == FloatingIpStatus::ID => {
213 let args: FloatingIpStatusInput =
214 serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
215 let out = self.floating_ip_status(args).await?;
216 Ok(serde_json::to_value(out)?)
217 }
218 other => bail!("hetzner floating-ip envoy does not support verb {other:?}"),
219 }
220 }
221}
222
223fn hetzner_network_zone(location: &str) -> Result<&'static str> {
228 match location {
229 "hil" => Ok("us-west"),
230 "ash" => Ok("us-east"),
231 "fsn1" | "nbg1" | "hel1" => Ok("eu-central"),
232 "sin" => Ok("ap-southeast"),
233 "" => bail!("hetzner: machine has no `location` set — required to derive its network zone"),
234 other => bail!("hetzner: unknown location {other:?}, cannot derive network zone"),
235 }
236}
237
238#[derive(Deserialize)]
239struct HetznerServersResponse {
240 servers: Vec<HetznerServerLite>,
241}
242
243#[derive(Deserialize)]
244struct HetznerServerLite {
245 id: u64,
246 #[serde(default)]
247 name: Option<String>,
248}
249
250#[derive(Deserialize)]
251struct HetznerFloatingIpResponse {
252 floating_ip: HetznerFloatingIpBody,
253}
254
255#[derive(Deserialize)]
256struct HetznerFloatingIpBody {
257 home_location: HetznerHomeLocation,
258 server: Option<u64>,
259}
260
261#[derive(Deserialize)]
262struct HetznerHomeLocation {
263 network_zone: String,
264}
265
266#[cfg(test)]
267mod tests {
268 use super::*;
269 use crate::provider::floating_ip::on_ingress_owner_changed;
270 use std::sync::atomic::{AtomicU32, Ordering};
271 use std::sync::{Arc, Mutex};
272
273 fn hil_machine(name: &str) -> MachineConfig {
274 MachineConfig {
275 name: name.into(),
276 provider: "hetzner".into(),
277 location: Some("hil".into()),
278 server_type: Some("cpx22".into()),
279 hosts_mirrors: vec![],
280 mesh_tags: vec![],
281 region: Some("us-west".into()),
282 zone: None,
283 arch: None,
284 bucket: None,
285 vendor: None,
286 nickname: None,
287 legacy_hostkey_fingerprint: None,
288 registration: Default::default(),
289 ssh_keys: vec![],
290 cloudflared: None,
291 hosts_operator_bridge: false,
292 connect: None,
293 allocatable: None,
294 taints: vec![],
295 sovereign_group: None,
296 sovereign_role: None,
297 }
298 }
299
300 async fn spawn_mock(
307 network_zone: &'static str,
308 initial_server: Option<u64>,
309 ) -> (String, Arc<AtomicU32>, tokio::task::JoinHandle<()>) {
310 let attached = Arc::new(Mutex::new(initial_server));
311 let assign_calls = Arc::new(AtomicU32::new(0));
312
313 let servers_route = {
314 axum::routing::get(move || async move {
315 axum::Json(serde_json::json!({ "servers": [ { "id": 555, "name": "edge-a" } ] }))
316 })
317 };
318
319 let floating_ip_get = {
320 let attached = attached.clone();
321 axum::routing::get(move || {
322 let attached = attached.clone();
323 async move {
324 let server = *attached.lock().unwrap();
325 axum::Json(serde_json::json!({
326 "floating_ip": {
327 "home_location": { "network_zone": network_zone },
328 "server": server,
329 }
330 }))
331 }
332 })
333 };
334
335 let assign_route = {
336 let attached = attached.clone();
337 let calls = assign_calls.clone();
338 axum::routing::post(move |axum::Json(body): axum::Json<serde_json::Value>| {
339 let attached = attached.clone();
340 let calls = calls.clone();
341 async move {
342 calls.fetch_add(1, Ordering::SeqCst);
343 let server = body.get("server").and_then(|v| v.as_u64());
344 *attached.lock().unwrap() = server;
345 axum::Json(serde_json::json!({ "action": { "status": "success" } }))
346 }
347 })
348 };
349
350 let app = axum::Router::new()
351 .route("/servers", servers_route)
352 .route("/floating_ips/{id}", floating_ip_get)
353 .route("/floating_ips/{id}/actions/assign", assign_route);
354
355 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
356 let addr = listener.local_addr().unwrap();
357 let handle = tokio::spawn(async move {
358 let _ = axum::serve(listener, app).await;
359 });
360
361 (format!("http://{addr}"), assign_calls, handle)
362 }
363
364 #[tokio::test]
365 async fn ingress_owner_flip_drives_exactly_one_reassign_call() {
366 let (base, calls, handle) = spawn_mock("us-west", Some(999)).await;
370 let client = HetznerFloatingIp::new("test-token").with_base_url(base);
371 let machine = hil_machine("edge-a");
372
373 let outcome = on_ingress_owner_changed(&client, &machine, "42")
374 .await
375 .unwrap();
376 assert!(outcome.reassigned, "owner flip must drive a reassign");
377 assert_eq!(outcome.attached_to, "555");
378 assert_eq!(calls.load(Ordering::SeqCst), 1);
379
380 handle.abort();
381 }
382
383 #[tokio::test]
384 async fn reapplying_the_same_owner_is_a_zero_call_noop() {
385 let (base, calls, handle) = spawn_mock("us-west", Some(555)).await;
387 let client = HetznerFloatingIp::new("test-token").with_base_url(base);
388 let machine = hil_machine("edge-a");
389
390 let outcome = on_ingress_owner_changed(&client, &machine, "42")
391 .await
392 .unwrap();
393 assert!(
394 !outcome.reassigned,
395 "re-applying the same owner must be a no-op"
396 );
397 assert_eq!(
398 calls.load(Ordering::SeqCst),
399 0,
400 "must not call the reassign endpoint"
401 );
402
403 handle.abort();
404 }
405
406 #[tokio::test]
407 async fn cross_zone_target_is_rejected_before_any_reassign_call() {
408 let (base, calls, handle) = spawn_mock("eu-central", None).await;
410 let client = HetznerFloatingIp::new("test-token").with_base_url(base);
411 let machine = hil_machine("edge-a"); let err = on_ingress_owner_changed(&client, &machine, "42")
414 .await
415 .unwrap_err();
416 let msg = format!("{err:#}");
417 assert!(
418 msg.contains("zone"),
419 "expected a zone-mismatch error, got: {msg}"
420 );
421 assert_eq!(
422 calls.load(Ordering::SeqCst),
423 0,
424 "zone mismatch must never call reassign"
425 );
426
427 handle.abort();
428 }
429
430 #[test]
431 fn hetzner_network_zone_maps_known_locations() {
432 assert_eq!(hetzner_network_zone("hil").unwrap(), "us-west");
433 assert_eq!(hetzner_network_zone("ash").unwrap(), "us-east");
434 assert_eq!(hetzner_network_zone("fsn1").unwrap(), "eu-central");
435 assert_eq!(hetzner_network_zone("nbg1").unwrap(), "eu-central");
436 assert_eq!(hetzner_network_zone("hel1").unwrap(), "eu-central");
437 }
438
439 #[test]
440 fn hetzner_network_zone_rejects_unknown_or_missing() {
441 assert!(hetzner_network_zone("mars1").is_err());
442 assert!(hetzner_network_zone("").is_err());
443 }
444}