Skip to main content

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}