cloud/provider/hetzner.rs
1//! @yah:ticket(R040-F9, "Lift KeysStore into shared crate; cloud reads vault then env")
2//! @yah:at(2026-05-05T00:33:17Z)
3//! @yah:status(review)
4//! @yah:assignee(agent:claude)
5//! @yah:parent(R040)
6//! @yah:handoff("DRY-up landed: app/yah/cli/src/keys.rs lifted to crates/yah/keys (own Cargo.toml, ProjectDirs::data_dir() unchanged so the existing on-disk vault keeps working). app/yah/cli now `keys = { path = ... }` deps the new crate; aes-gcm and rand drops out of the CLI's direct deps (transitive now). main.rs/agent.rs/agentd.rs swapped `mod keys;`/`crate::keys::` for the external crate path. cloud crate gained `keys` dep + `HetznerDriver::from_default_sources()` that tries KeysStore::open().get(slot) per-key then falls back to env: hetzner-api-token↔HETZNER_API_TOKEN, hetzner-s3-access-key↔HETZNER_S3_ACCESS_KEY, hetzner-s3-secret-key↔HETZNER_S3_SECRET_KEY. Vault open errors are swallowed (no vault → fall back to env). yah cloud callsites (app/yah/cli/src/cloud.rs:257 + :423) flipped to from_default_sources(). Tests: 6/6 green for keys (moved from CLI), 26/26 green for cloud (no test changes — all existing). cargo check -p keys + -p cloud + -p yah --bin yah-agentd clean.")
7//! @yah:next("Follow-up scoped as R043 (relay) with phases F1 bridge / F2 naming / F3 cleanup — unify desktop api_keys with this vault, KeysStore as canonical, drop keyring dep once soaked.")
8//! @yah:next("Optional: yah keys CLI could grow `--from-keychain <provider>` flag to one-shot import a desktop-vault token without typing it. Not urgent; the user can already pipe via `security find-generic-password ... | yah keys set --from-stdin <slot>`.")
9//! @yah:verify("cargo test -p keys")
10//! @yah:verify("cargo test -p cloud")
11//! @yah:verify("cargo check -p yah --bin yah-agentd")
12//! @yah:gotcha("cargo check --workspace currently fails on app/yah/cli/src/cloud.rs:210 because handle_agent is referenced but not yet defined — that's parallel R040-F7 WIP (yah cloud agent ping/services/logs against yah-yubaba), not this refactor. The `keys = ...` and HetznerDriver::from_default_sources additions compile clean on their own.")
13//!
14//! @yah:ticket(R040-F10, "Pass ssh_keys through to Hetzner create_server (pre-mesh SSH access)")
15//! @yah:at(2026-05-05T00:33:17Z)
16//! @yah:status(review)
17//! @yah:assignee(agent:claude)
18//! @yah:parent(R040)
19//! @yah:handoff("Threaded ssh_keys end-to-end. MachineConfig.ssh_keys: Vec<u64> with #[serde(default)] (existing toml files unaffected). ServerSpec.ssh_keys mirrors. hetzner.rs::create_server attaches `\"ssh_keys\": [...]` to the POST body when non-empty (Hetzner ignores empty arrays anyway). provision.rs threads MachineConfig.ssh_keys → ProvisionRequest → ServerSpec. 6 test sites updated to add the field; 29/29 cloud tests green.")
20//! @yah:handoff("Live-verified: provisioned yah-cloud-1 (cpx11, hil) via direct API (cloud-crate provision flow blew the 32KiB user_data cap embedding yubaba — separate ticket). With ssh_keys=[111513970, 111525493] the box accepted SSH on first boot via ~/.ssh/yah, no out-of-band root password recovery needed. The agentd round-trip succeeded over the socket: agent.list_sessions returned {sessions:[]}, bogus method returned -32601, stop-unknown returned {stopped:false}.")
21//! @yah:next("Follow-up R040-Tx: 32KiB user_data cap blocks the canonical `yah cloud machine provision` path because the yubaba binary embeds at ~3MB after base64. Refactor cloud-init to fetch yubaba from a URL (GitHub release artifact during cloud-init runcmd) instead of inlining the bytes. Until then, provision flows that need yubaba need a different transport (post-boot scp + register-hostkey).")
22//! @yah:verify("cargo test -p cloud")
23//! @yah:verify("cargo run -p yah --bin yah -- cloud machine status (sees ssh_keys field via the new yah-cloud-1.toml in .yah/cloud/machines/)")
24//!
25//! @yah:ticket(R040-F14, "Extract shared crates/yah/hetzner: lift transport+DTOs out of desktop and cloud parallel impls")
26//! @yah:at(2026-05-05T00:33:17Z)
27//! @yah:assignee(agent:claude)
28//! @yah:status(review)
29//! @yah:parent(R040)
30//! @yah:handoff("app/yah/desktop/src/hetzner.rs (~641 LOC) and crates/yah/cloud/src/provider/hetzner.rs (~560 LOC) are parallel HTTP clients. Both ship their own auth_client/check_status, RawServer/HetznerServer DTOs, error types. Lift the shared core into crates/yah/hetzner; both call sites become thin layers over it. Surfaced in R040-F12 + F13 (destroy-button work) where the obvious DRY move would have ballooned the change set.")
31//! @yah:next("New crate crates/yah/hetzner: HetznerClient (reqwest+bearer + check_status), wire DTOs (RawServer/RawSshKey/RawLocation/RawImage), HetznerError")
32//! @yah:next("Desktop hetzner.rs becomes UI-flow operations (list_server_types/_locations/_images stay desktop-only — catalog browse for the form) over the shared client")
33//! @yah:next("Cloud hetzner.rs becomes reconcile-flow operations (find_server_by_name, server_status, destroy_server, bucket APIs) over the shared client")
34//! @yah:next("Token-source convergence is a separate ticket: shared client takes &str token, desktop reads keychain blob, cloud reads KeysStore vault — unifying those vaults is its own refactor")
35//! @yah:next("Estimate: 2-4 hours; tests on both sides pass without functional change")
36//!
37//! @yah:ticket(R040-T17, "Decision-recorded: do NOT build floating-IP plumbing (superseded by Cloudflare Tunnels + Headscale mesh)")
38//! @yah:at(2026-05-05T00:33:17Z)
39//! @yah:assignee(agent:claude)
40//! @yah:status(review)
41//! @yah:parent(R040)
42//! @yah:kind(task)
43//! @yah:handoff("Question raised this session: 'should yah cloud have one floating IP per region as an ingress target?' Answer recorded so future-self doesn't redo the analysis. NO — combination of R040-F15 (Cloudflare Tunnel for public ingress) + R040-F16 (Headscale mesh for inter-node TCP) means the entire stable-IP question evaporates. Public DNS points at <tunnel-id>.cfargotunnel.com (CF-managed). Inter-node uses 100.64.x.x mesh IPs that survive box replacement. Hetzner inline IPs can churn weekly without consequence. Floating IPs in this design would just be ~€0.50/mo per IP burned to solve a problem we don't have. Future-self: do NOT re-litigate without a concrete TCP/UDP public-ingress requirement that CF Spectrum (paid) can't cover.")
44//! @yah:next("Archive condition: this task is purely a design-decision marker. Archive it once R040-F15 lands (the architectural commitment is real). Until then it stays here as a tripwire so a future claim of 'we should add primary-IP support' surfaces the prior reasoning before the work starts.")
45//! @yah:next("Reopen condition: a concrete service requirement that needs raw TCP/UDP from the public internet AND can't justify CF Spectrum's pricing. If reopened, the natural shape is MachineConfig.primary_ip: Option<u64> + ServerSpec.primary_ip + a `yah cloud ip {create,list,destroy}` subcommand — but DON'T pre-build any of that until reopened.")
46//!
47//!
48
49use super::{
50 BucketAcl, BucketRef, Location, MachineProvider, ProjectId, ServerId, ServerSpec, ServerStatus,
51 ServerSummary,
52};
53use anyhow::{bail, Context, Result};
54use async_trait::async_trait;
55use local_driver::s3_sign::{
56 sign_s3_delete_bucket, sign_s3_head_bucket, sign_s3_put_bucket, sign_s3_put_bucket_acl,
57};
58use reqwest::StatusCode;
59use yah_hetzner::{HetznerClient, HetznerCreateServerSpec};
60
61/// Hetzner Cloud + Object Storage driver for Phase-1 mirror bootstrap.
62///
63/// Cloud operations (server lifecycle) use `HETZNER_API_TOKEN`.
64/// Bucket operations use the Hetzner Object Storage S3-compat API
65/// (`HETZNER_S3_ACCESS_KEY` + `HETZNER_S3_SECRET_KEY`).
66/// Run `yah cloud secrets` for the canonical contract (vault slots + env).
67#[derive(Clone)]
68pub struct HetznerDriver {
69 hclient: HetznerClient,
70 s3_access_key: Option<String>,
71 s3_secret_key: Option<String>,
72 /// Overrides the computed S3 base endpoint (useful for integration tests).
73 s3_endpoint_override: Option<String>,
74}
75
76impl HetznerDriver {
77 /// Build with a Cloud API token only. Bucket operations will fail until
78 /// S3 credentials are added via [`with_storage`].
79 pub fn new(token: impl Into<String>) -> Self {
80 Self {
81 hclient: HetznerClient::new(token),
82 s3_access_key: None,
83 s3_secret_key: None,
84 s3_endpoint_override: None,
85 }
86 }
87
88 /// Attach Hetzner Object Storage S3 credentials.
89 pub fn with_storage(
90 mut self,
91 access_key: impl Into<String>,
92 secret_key: impl Into<String>,
93 ) -> Self {
94 self.s3_access_key = Some(access_key.into());
95 self.s3_secret_key = Some(secret_key.into());
96 self
97 }
98
99 /// Override the S3 base endpoint (e.g. `"http://localhost:9000"` for MinIO in tests).
100 pub fn with_s3_endpoint(mut self, endpoint: impl Into<String>) -> Self {
101 self.s3_endpoint_override = Some(endpoint.into());
102 self
103 }
104
105 /// Build from environment variables.
106 ///
107 /// Required: `HETZNER_API_TOKEN`.
108 /// Optional: `HETZNER_S3_ACCESS_KEY` + `HETZNER_S3_SECRET_KEY` (needed for `create_bucket`).
109 pub fn from_env() -> Result<Self> {
110 let token = std::env::var("HETZNER_API_TOKEN")
111 .context("HETZNER_API_TOKEN not set — run `yah cloud secrets` for the contract")?;
112 let mut driver = Self::new(token);
113 if let (Ok(ak), Ok(sk)) = (
114 std::env::var("HETZNER_S3_ACCESS_KEY"),
115 std::env::var("HETZNER_S3_SECRET_KEY"),
116 ) {
117 driver = driver.with_storage(ak, sk);
118 }
119 Ok(driver)
120 }
121
122 /// Build from the shared `keys` vault first, then fall back to env vars.
123 ///
124 /// Slot ↔ env mapping:
125 /// - `hetzner-api-token` ↔ `HETZNER_API_TOKEN` (required)
126 /// - `hetzner-s3-access-key` ↔ `HETZNER_S3_ACCESS_KEY` (optional pair)
127 /// - `hetzner-s3-secret-key` ↔ `HETZNER_S3_SECRET_KEY` (optional pair)
128 ///
129 /// The vault path is the production-default for `yah` callers; env is
130 /// kept as a fallback so CI / one-shot scripts (and hosts that don't
131 /// have a vault yet) keep working. The two are checked per-key, so a
132 /// vault-stored API token combined with env-supplied S3 creds is fine.
133 /// Vault open errors (missing machine.key, etc.) are swallowed — they
134 /// just mean "no vault here, try env."
135 pub fn from_default_sources() -> Result<Self> {
136 let token = fob::get_or_env("hetzner-api-token", "HETZNER_API_TOKEN")?.context(
137 "no Hetzner API token — set one with `yah keys set hetzner-api-token` \
138 or export HETZNER_API_TOKEN; run `yah cloud secrets` for the full contract",
139 )?;
140 let mut driver = Self::new(token);
141 if let (Some(ak), Some(sk)) = (
142 fob::get_or_env("hetzner-s3-access-key", "HETZNER_S3_ACCESS_KEY")?,
143 fob::get_or_env("hetzner-s3-secret-key", "HETZNER_S3_SECRET_KEY")?,
144 ) {
145 driver = driver.with_storage(ak, sk);
146 }
147 Ok(driver)
148 }
149
150 fn s3_endpoint(&self, location: &Location) -> String {
151 self.s3_endpoint_override
152 .clone()
153 .unwrap_or_else(|| location.hetzner_storage_endpoint().to_string())
154 }
155}
156
157// ─── MachineProvider ─────────────────────────────────────────────────────────
158
159#[async_trait]
160impl MachineProvider for HetznerDriver {
161 async fn ensure_project(&self, name: &str) -> Result<ProjectId> {
162 // Hetzner Cloud API tokens are already project-scoped; no API call needed.
163 Ok(ProjectId(name.to_string()))
164 }
165
166 async fn create_server(
167 &self,
168 _project: &ProjectId,
169 spec: &ServerSpec,
170 user_data: &str,
171 ) -> Result<ServerId> {
172 let hspec = HetznerCreateServerSpec {
173 name: spec.name.clone(),
174 server_type: spec.server_type.clone(),
175 location: spec.location.hetzner_cloud_id().to_string(),
176 image: spec.image.clone(),
177 ssh_keys: spec.ssh_keys.clone(),
178 user_data: if user_data.is_empty() {
179 None
180 } else {
181 Some(user_data.to_string())
182 },
183 };
184 let server = self
185 .hclient
186 .create_server(&hspec)
187 .await
188 .context("create_server")?;
189 Ok(ServerId(server.id.to_string()))
190 }
191
192 async fn server_status(&self, id: &ServerId) -> Result<ServerStatus> {
193 let server_id: u64 = id.0.parse().context("invalid server id")?;
194 match self
195 .hclient
196 .get_server(server_id)
197 .await
198 .context("server_status")?
199 {
200 None => Ok(ServerStatus::Unknown("not-found".into())),
201 Some(s) => Ok(parse_server_status(&s.status)),
202 }
203 }
204
205 async fn find_server_by_name(&self, name: &str) -> Result<Option<ServerSummary>> {
206 let servers = self
207 .hclient
208 .find_servers_by_name(name)
209 .await
210 .context("find_server_by_name")?;
211 Ok(servers.into_iter().next().map(|s| ServerSummary {
212 id: ServerId(s.id.to_string()),
213 server_type: s.server_type,
214 status: parse_server_status(&s.status),
215 public_ipv4: s.ipv4,
216 location: s.location,
217 }))
218 }
219
220 async fn bucket_exists(&self, name: &str, location: Location) -> Result<bool> {
221 let (ak, sk) = match (&self.s3_access_key, &self.s3_secret_key) {
222 (Some(a), Some(s)) => (a.as_str(), s.as_str()),
223 _ => bail!(
224 "S3 credentials not configured — set HETZNER_S3_ACCESS_KEY and \
225 HETZNER_S3_SECRET_KEY (run `yah cloud secrets` for the contract)"
226 ),
227 };
228
229 let endpoint = self.s3_endpoint(&location);
230 let region = location.hetzner_storage_region();
231 let url = format!("{endpoint}/{name}");
232
233 let headers = sign_s3_head_bucket(&url, region, ak, sk)?;
234
235 let resp = self
236 .hclient
237 .raw_client()
238 .head(&url)
239 .headers(headers)
240 .send()
241 .await
242 .context("HEAD bucket")?;
243
244 let http_status = resp.status();
245 match http_status {
246 StatusCode::OK | StatusCode::NO_CONTENT => Ok(true),
247 StatusCode::NOT_FOUND => Ok(false),
248 _ => {
249 let text = resp.text().await.unwrap_or_default();
250 bail!("bucket_exists failed ({http_status}): {text}")
251 }
252 }
253 }
254
255 async fn destroy_server(&self, id: &ServerId) -> Result<()> {
256 let server_id: u64 = id.0.parse().context("invalid server id")?;
257 self.hclient
258 .destroy_server(server_id)
259 .await
260 .context("destroy_server")?;
261 Ok(())
262 }
263
264 async fn delete_bucket(&self, name: &str, location: Location) -> Result<()> {
265 let (ak, sk) = match (&self.s3_access_key, &self.s3_secret_key) {
266 (Some(a), Some(s)) => (a.as_str(), s.as_str()),
267 _ => bail!(
268 "S3 credentials not configured — set HETZNER_S3_ACCESS_KEY and \
269 HETZNER_S3_SECRET_KEY (run `yah cloud secrets` for the contract)"
270 ),
271 };
272
273 let endpoint = self.s3_endpoint(&location);
274 let region = location.hetzner_storage_region();
275 let url = format!("{endpoint}/{name}");
276
277 let headers = sign_s3_delete_bucket(&url, region, ak, sk)?;
278
279 let resp = self
280 .hclient
281 .raw_client()
282 .delete(&url)
283 .headers(headers)
284 .send()
285 .await
286 .context("DELETE bucket")?;
287
288 let http_status = resp.status();
289 match http_status {
290 StatusCode::NO_CONTENT | StatusCode::OK | StatusCode::NOT_FOUND => Ok(()),
291 StatusCode::CONFLICT => {
292 let text = resp.text().await.unwrap_or_default();
293 bail!(
294 "delete_bucket failed ({http_status} likely BucketNotEmpty): {text}\n\
295 yah doesn't yet list-and-delete objects — empty the bucket first \
296 via aws-cli (`aws s3 rm --recursive --endpoint-url <region-endpoint> \
297 s3://{name}`) and retry."
298 )
299 }
300 _ => {
301 let text = resp.text().await.unwrap_or_default();
302 bail!("delete_bucket failed ({http_status}): {text}")
303 }
304 }
305 }
306
307 async fn create_bucket(&self, name: &str, location: Location) -> Result<BucketRef> {
308 let (ak, sk) = match (&self.s3_access_key, &self.s3_secret_key) {
309 (Some(a), Some(s)) => (a.as_str(), s.as_str()),
310 _ => bail!(
311 "S3 credentials not configured — set HETZNER_S3_ACCESS_KEY and \
312 HETZNER_S3_SECRET_KEY (run `yah cloud secrets` for the contract)"
313 ),
314 };
315
316 let endpoint = self.s3_endpoint(&location);
317 let region = location.hetzner_storage_region();
318 // Path-style: https://<endpoint>/<bucket>
319 let url = format!("{endpoint}/{name}");
320
321 let headers = sign_s3_put_bucket(&url, region, ak, sk)?;
322
323 let resp = self
324 .hclient
325 .raw_client()
326 .put(&url)
327 .headers(headers)
328 .body("")
329 .send()
330 .await
331 .context("PUT bucket")?;
332
333 let http_status = resp.status();
334 // 200 OK or 409 Conflict (bucket already exists and belongs to caller) are both fine.
335 if http_status.is_success() || http_status == StatusCode::CONFLICT {
336 return Ok(BucketRef {
337 name: name.to_string(),
338 endpoint,
339 });
340 }
341 let text = resp.text().await.unwrap_or_default();
342 bail!("create_bucket failed ({http_status}): {text}");
343 }
344
345 async fn set_bucket_acl(&self, name: &str, location: Location, acl: BucketAcl) -> Result<()> {
346 let (ak, sk) = match (&self.s3_access_key, &self.s3_secret_key) {
347 (Some(a), Some(s)) => (a.as_str(), s.as_str()),
348 _ => bail!(
349 "S3 credentials not configured — set HETZNER_S3_ACCESS_KEY and \
350 HETZNER_S3_SECRET_KEY (run `yah cloud secrets` for the contract)"
351 ),
352 };
353
354 let endpoint = self.s3_endpoint(&location);
355 let region = location.hetzner_storage_region();
356 let url = format!("{endpoint}/{name}?acl");
357
358 let headers = sign_s3_put_bucket_acl(&url, region, ak, sk, acl.as_canned())?;
359
360 let resp = self
361 .hclient
362 .raw_client()
363 .put(&url)
364 .headers(headers)
365 .body("")
366 .send()
367 .await
368 .context("PUT bucket?acl")?;
369
370 let http_status = resp.status();
371 if http_status.is_success() {
372 return Ok(());
373 }
374 let text = resp.text().await.unwrap_or_default();
375 bail!("set_bucket_acl failed ({http_status}): {text}");
376 }
377}
378
379fn parse_server_status(s: &str) -> ServerStatus {
380 match s {
381 "initializing" => ServerStatus::Initializing,
382 "starting" => ServerStatus::Starting,
383 "running" => ServerStatus::Running,
384 "stopping" => ServerStatus::Stopping,
385 "off" => ServerStatus::Off,
386 "deleting" => ServerStatus::Deleting,
387 other => ServerStatus::Unknown(other.to_string()),
388 }
389}
390
391// ─── Tests ────────────────────────────────────────────────────────────────────
392
393#[cfg(test)]
394mod tests {
395 use super::*;
396
397 #[test]
398 fn parse_status_known_values() {
399 assert_eq!(parse_server_status("running"), ServerStatus::Running);
400 assert_eq!(parse_server_status("off"), ServerStatus::Off);
401 assert_eq!(
402 parse_server_status("initializing"),
403 ServerStatus::Initializing
404 );
405 assert_eq!(
406 parse_server_status("banana"),
407 ServerStatus::Unknown("banana".into())
408 );
409 }
410
411 #[test]
412 fn location_ids_correct() {
413 assert_eq!(Location::Pdx.hetzner_cloud_id(), "hil");
414 assert_eq!(Location::Iad.hetzner_cloud_id(), "ash");
415 assert_eq!(Location::Fsn.hetzner_cloud_id(), "fsn1");
416 }
417
418 #[tokio::test]
419 #[ignore = "requires HETZNER_API_TOKEN"]
420 async fn integration_create_and_destroy_server() {
421 let driver = HetznerDriver::from_env().unwrap();
422 let project = driver.ensure_project("yah-test").await.unwrap();
423 let spec = ServerSpec {
424 name: "yah-test-ephemeral".into(),
425 server_type: "cpx11".into(),
426 image: "debian-12".into(),
427 location: Location::Fsn,
428 ssh_keys: vec![],
429 };
430 let id = driver
431 .create_server(&project, &spec, "#cloud-config\n")
432 .await
433 .unwrap();
434 println!("created server {}", id.0);
435
436 let status = driver.server_status(&id).await.unwrap();
437 println!("initial status: {status:?}");
438 assert!(matches!(
439 status,
440 ServerStatus::Running | ServerStatus::Initializing | ServerStatus::Starting
441 ));
442
443 driver.destroy_server(&id).await.unwrap();
444 println!("destroyed");
445 }
446
447 #[tokio::test]
448 #[ignore = "requires HETZNER_API_TOKEN + HETZNER_S3_ACCESS_KEY + HETZNER_S3_SECRET_KEY"]
449 async fn integration_create_bucket() {
450 let driver = HetznerDriver::from_env().unwrap();
451 let bucket = driver
452 .create_bucket("yah-ci-test-bucket-fsn1", Location::Fsn)
453 .await
454 .unwrap();
455 println!("bucket: {} @ {}", bucket.name, bucket.endpoint);
456 }
457}