Skip to main content

cloud/provider/
digitalocean.rs

1//! [`DigitalOceanEnvoy`] — second native `cloud.vps.*` adapter, spike scope
2//! (R409-T10). Together with [`HetznerEnvoy`] this is the catalog-shape
3//! validator R409-T11's postmortem decides on.
4//!
5//! Layering differs from Hetzner: there is no separate `MachineProvider` impl
6//! in the way. Hetzner has a long-lived [`HetznerDriver`] that the
7//! orchestration code (`provision.rs`, `cloud_init.rs`) calls into, and the
8//! envoy was retrofitted on top. DigitalOcean has never had such a driver, so
9//! [`DigitalOceanClient`] is built thin and lives only behind
10//! [`DigitalOceanEnvoy`] — no parallel `MachineProvider` impl. If T11 says
11//! "expand", the client lifts into its own `crates/yah/digitalocean` crate
12//! the way Hetzner did in R040-F14.
13//!
14//! The spike is **pre-envoy-host**: nothing wires `DigitalOceanEnvoy` through
15//! `KgToolRegistry` (that's R409-T9). Verification is unit tests of the
16//! conversion layer — the network methods are not exercised in tests.
17//!
18//! ## Post-T11 shape
19//!
20//! Three findings drove revisions in the R409-T11 postmortem (see W144
21//! §"Catalog-shape postmortem"):
22//!
23//! - **D5**: wire `location` is now a coarse region tag (`na-west`,
24//!   `na-east`, `eu-central`); [`do_region`] maps each tag to the nearest
25//!   DO slug (`sfo3`, `nyc3`, `fra1`).
26//! - **D6**: `CloudVpsCreateInput.project` is gone — DO Projects remain
27//!   available as a future `cloud.project.*` verb tree but are no longer
28//!   a free parameter on `cloud.vps.create`.
29//! - **D7**: `ssh_keys` is now `Vec<String>`; DO accepts both numeric IDs
30//!   and SHA-256 fingerprints as JSON strings, so the spec passes the
31//!   values through verbatim.
32
33use std::sync::Arc;
34
35use anyhow::{bail, Context, Result};
36use async_trait::async_trait;
37use reqwest::StatusCode;
38use serde::{Deserialize, Serialize};
39use serde_json::Value;
40
41use crate::envoy::cloud_vps::{
42    CloudVpsCreate, CloudVpsCreateInput, CloudVpsCreateOutput, CloudVpsDestroy,
43    CloudVpsDestroyInput, CloudVpsDestroyOutput, CloudVpsStatus, CloudVpsStatusInput,
44    CloudVpsStatusOutput, VpsPhase,
45};
46use crate::envoy::{AdapterFlavor, EnvoyAdapter, InternalVerb, Tier};
47
48const DO_BASE: &str = "https://api.digitalocean.com/v2";
49
50/// Minimal DigitalOcean Cloud API client — droplet lifecycle only.
51///
52/// Token-source policy stays out of the client: callers obtain the bearer
53/// token (env, keystore, vault) and pass it in. Same convention as
54/// [`hetzner::HetznerClient`].
55#[derive(Clone)]
56pub struct DigitalOceanClient {
57    http: reqwest::Client,
58    token: String,
59    base_url: String,
60}
61
62impl DigitalOceanClient {
63    /// Build a client against `api.digitalocean.com`.
64    pub fn new(token: impl Into<String>) -> Self {
65        Self {
66            http: reqwest::Client::new(),
67            token: token.into(),
68            base_url: DO_BASE.to_string(),
69        }
70    }
71
72    /// Override the base URL — useful for integration tests against a
73    /// mocked endpoint. Production callers shouldn't need this.
74    pub fn with_base_url(mut self, url: impl Into<String>) -> Self {
75        self.base_url = url.into();
76        self
77    }
78
79    /// `POST /v2/droplets`. Returns the new droplet's numeric id.
80    pub async fn create_droplet(&self, spec: &DoCreateDropletSpec) -> Result<u64> {
81        let resp = self
82            .http
83            .post(format!("{}/droplets", self.base_url))
84            .bearer_auth(&self.token)
85            .json(spec)
86            .send()
87            .await
88            .context("digitalocean: POST /droplets")?;
89        let status = resp.status();
90        if !status.is_success() {
91            let body = resp.text().await.unwrap_or_default();
92            bail!("digitalocean POST /droplets failed: {status} {body}");
93        }
94        let parsed: CreateDropletResponse = resp
95            .json()
96            .await
97            .context("digitalocean: decode POST /droplets response")?;
98        Ok(parsed.droplet.id)
99    }
100
101    /// `GET /v2/droplets/{id}`. Returns the raw droplet `status` string
102    /// (one of `new`/`active`/`off`/`archive`, or something newer DO has
103    /// added that we have not yet mapped).
104    pub async fn droplet_status(&self, id: u64) -> Result<String> {
105        let resp = self
106            .http
107            .get(format!("{}/droplets/{}", self.base_url, id))
108            .bearer_auth(&self.token)
109            .send()
110            .await
111            .context("digitalocean: GET /droplets/{id}")?;
112        let status = resp.status();
113        if !status.is_success() {
114            let body = resp.text().await.unwrap_or_default();
115            bail!("digitalocean GET /droplets/{id} failed: {status} {body}");
116        }
117        let parsed: DropletResponse = resp
118            .json()
119            .await
120            .context("digitalocean: decode GET /droplets/{id} response")?;
121        Ok(parsed.droplet.status)
122    }
123
124    /// `DELETE /v2/droplets/{id}`. Idempotent — a 404 returns `Ok(())`.
125    pub async fn destroy_droplet(&self, id: u64) -> Result<()> {
126        let resp = self
127            .http
128            .delete(format!("{}/droplets/{}", self.base_url, id))
129            .bearer_auth(&self.token)
130            .send()
131            .await
132            .context("digitalocean: DELETE /droplets/{id}")?;
133        let status = resp.status();
134        if status == StatusCode::NOT_FOUND {
135            return Ok(());
136        }
137        if !status.is_success() {
138            let body = resp.text().await.unwrap_or_default();
139            bail!("digitalocean DELETE /droplets/{id} failed: {status} {body}");
140        }
141        Ok(())
142    }
143}
144
145/// Request body for `POST /v2/droplets`. Field names match DO's API
146/// exactly so this serializes straight to the wire.
147#[derive(Debug, Clone, Serialize)]
148pub struct DoCreateDropletSpec {
149    pub name: String,
150    /// DO region slug (e.g. `"sfo3"`, `"fra1"`).
151    pub region: String,
152    /// DO size slug (e.g. `"s-1vcpu-1gb"`, `"s-2vcpu-4gb"`).
153    pub size: String,
154    /// DO image slug (e.g. `"debian-12-x64"`).
155    pub image: String,
156    /// SSH keys to authorize for `root`. DO accepts numeric ids or SHA-256
157    /// fingerprints — W144 D7 lets both flow through as strings.
158    #[serde(skip_serializing_if = "Vec::is_empty")]
159    pub ssh_keys: Vec<String>,
160    /// cloud-init `user_data`. Omitted from the body when `None`.
161    #[serde(skip_serializing_if = "Option::is_none")]
162    pub user_data: Option<String>,
163}
164
165#[derive(Deserialize)]
166struct CreateDropletResponse {
167    droplet: RawDroplet,
168}
169
170#[derive(Deserialize)]
171struct DropletResponse {
172    droplet: RawDroplet,
173}
174
175#[derive(Deserialize)]
176struct RawDroplet {
177    id: u64,
178    status: String,
179}
180
181/// Tier-S, native-flavored envoy adapter that bridges the DigitalOcean
182/// droplet API to the `cloud.vps.*` verb framework. Spike scope (R409-T10):
183/// only the three verbs from [`crate::envoy::cloud_vps`] are wired.
184pub struct DigitalOceanEnvoy {
185    client: Arc<DigitalOceanClient>,
186}
187
188impl DigitalOceanEnvoy {
189    /// Wrap an owned client.
190    pub fn new(client: DigitalOceanClient) -> Self {
191        Self {
192            client: Arc::new(client),
193        }
194    }
195
196    /// Wrap a shared client.
197    pub fn from_arc(client: Arc<DigitalOceanClient>) -> Self {
198        Self { client }
199    }
200
201    /// Typed handler for `cloud.vps.create`. Public so integration tests
202    /// can bypass [`EnvoyAdapter::dispatch`].
203    pub async fn cloud_vps_create(
204        &self,
205        input: CloudVpsCreateInput,
206    ) -> Result<CloudVpsCreateOutput> {
207        let region = do_region(&input.location)
208            .with_context(|| format!("cloud.vps.create: unknown location {:?}", input.location))?;
209        let spec = DoCreateDropletSpec {
210            name: input.name,
211            region: region.to_string(),
212            size: input.server_type,
213            image: input.image,
214            ssh_keys: input.ssh_keys,
215            user_data: if input.user_data.is_empty() {
216                None
217            } else {
218                Some(input.user_data)
219            },
220        };
221        let id = self
222            .client
223            .create_droplet(&spec)
224            .await
225            .context("cloud.vps.create: create_droplet")?;
226        Ok(CloudVpsCreateOutput { id: id.to_string() })
227    }
228
229    /// Typed handler for `cloud.vps.destroy`. Idempotent.
230    pub async fn cloud_vps_destroy(
231        &self,
232        input: CloudVpsDestroyInput,
233    ) -> Result<CloudVpsDestroyOutput> {
234        let id = parse_droplet_id(&input.id, "cloud.vps.destroy")?;
235        self.client
236            .destroy_droplet(id)
237            .await
238            .context("cloud.vps.destroy: destroy_droplet")?;
239        Ok(CloudVpsDestroyOutput::default())
240    }
241
242    /// Typed handler for `cloud.vps.status`. Maps DO's small status
243    /// taxonomy onto the wire [`VpsPhase`] via [`do_status_to_output`].
244    pub async fn cloud_vps_status(
245        &self,
246        input: CloudVpsStatusInput,
247    ) -> Result<CloudVpsStatusOutput> {
248        let id = parse_droplet_id(&input.id, "cloud.vps.status")?;
249        let raw = self
250            .client
251            .droplet_status(id)
252            .await
253            .context("cloud.vps.status: droplet_status")?;
254        Ok(do_status_to_output(&raw))
255    }
256}
257
258#[async_trait]
259impl EnvoyAdapter for DigitalOceanEnvoy {
260    fn id(&self) -> &str {
261        "digitalocean"
262    }
263    fn tier(&self) -> Tier {
264        Tier::S
265    }
266    fn flavor(&self) -> AdapterFlavor {
267        AdapterFlavor::Native
268    }
269    fn supported_verb_ids(&self) -> Vec<&'static str> {
270        vec![CloudVpsCreate::ID, CloudVpsDestroy::ID, CloudVpsStatus::ID]
271    }
272    async fn dispatch(&self, verb_id: &str, input: Value) -> Result<Value> {
273        match verb_id {
274            id if id == CloudVpsCreate::ID => {
275                let args: CloudVpsCreateInput =
276                    serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
277                let out = self.cloud_vps_create(args).await?;
278                Ok(serde_json::to_value(out)?)
279            }
280            id if id == CloudVpsDestroy::ID => {
281                let args: CloudVpsDestroyInput =
282                    serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
283                let out = self.cloud_vps_destroy(args).await?;
284                Ok(serde_json::to_value(out)?)
285            }
286            id if id == CloudVpsStatus::ID => {
287                let args: CloudVpsStatusInput =
288                    serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
289                let out = self.cloud_vps_status(args).await?;
290                Ok(serde_json::to_value(out)?)
291            }
292            other => bail!("digitalocean envoy does not support verb {other:?}"),
293        }
294    }
295}
296
297/// Pure conversion: W144-D5 coarse region tag → DO region slug.
298///
299/// The tags promise *broad region*, not specific city: `na-west` lands in
300/// the DO San Francisco region; `na-east` in New York; `eu-central` in
301/// Frankfurt. New tags follow `<continent>-<direction>` and are added here
302/// monotonically as DO opens more regions.
303fn do_region(loc: &str) -> Result<&'static str> {
304    match loc {
305        "na-west" => Ok("sfo3"),
306        "na-east" => Ok("nyc3"),
307        "eu-central" => Ok("fra1"),
308        other => Err(anyhow::anyhow!("unknown location: {other}")),
309    }
310}
311
312/// Pure conversion: DO droplet status string → wire [`CloudVpsStatusOutput`].
313///
314/// DO has four documented statuses (`new`, `active`, `off`, `archive`). Newer
315/// or vendor-internal values fall through to [`VpsPhase::Unknown`] with the
316/// raw string preserved in `detail` so the caller can see what happened.
317fn do_status_to_output(raw: &str) -> CloudVpsStatusOutput {
318    let (phase, detail) = match raw {
319        "new" => (VpsPhase::Initializing, None),
320        "active" => (VpsPhase::Running, None),
321        "off" => (VpsPhase::Off, None),
322        "archive" => (VpsPhase::Deleting, None),
323        other => (VpsPhase::Unknown, Some(other.to_string())),
324    };
325    CloudVpsStatusOutput { phase, detail }
326}
327
328fn parse_droplet_id(raw: &str, verb: &str) -> Result<u64> {
329    raw.parse()
330        .with_context(|| format!("{verb}: id {raw:?} is not a u64 droplet id"))
331}
332
333#[cfg(test)]
334mod tests {
335    use super::*;
336
337    #[test]
338    fn do_region_maps_known_locations() {
339        assert_eq!(do_region("na-west").unwrap(), "sfo3");
340        assert_eq!(do_region("na-east").unwrap(), "nyc3");
341        assert_eq!(do_region("eu-central").unwrap(), "fra1");
342    }
343
344    #[test]
345    fn do_region_rejects_unknown() {
346        let err = do_region("hil").unwrap_err();
347        assert!(err.to_string().contains("hil"));
348    }
349
350    #[test]
351    fn do_region_rejects_legacy_airport_codes() {
352        // W144 D5: the old Hetzner-centric airport codes (pdx/iad/fsn) are
353        // no longer wire-valid. Adapters must surface a clear error so the
354        // caller updates their input to the new coarse region tag.
355        for legacy in ["pdx", "iad", "fsn"] {
356            let err = do_region(legacy).unwrap_err();
357            assert!(
358                err.to_string().contains(legacy),
359                "expected error to name {legacy}"
360            );
361        }
362    }
363
364    #[test]
365    fn do_status_new_maps_to_initializing() {
366        let out = do_status_to_output("new");
367        assert_eq!(out.phase, VpsPhase::Initializing);
368        assert!(out.detail.is_none());
369    }
370
371    #[test]
372    fn do_status_active_maps_to_running() {
373        let out = do_status_to_output("active");
374        assert_eq!(out.phase, VpsPhase::Running);
375        assert!(out.detail.is_none());
376    }
377
378    #[test]
379    fn do_status_off_maps_to_off() {
380        let out = do_status_to_output("off");
381        assert_eq!(out.phase, VpsPhase::Off);
382        assert!(out.detail.is_none());
383    }
384
385    #[test]
386    fn do_status_archive_maps_to_deleting() {
387        let out = do_status_to_output("archive");
388        assert_eq!(out.phase, VpsPhase::Deleting);
389        assert!(out.detail.is_none());
390    }
391
392    #[test]
393    fn do_status_unknown_carries_raw_string_as_detail() {
394        // T11 input: DO may add states (e.g. `provisioning`) without notice.
395        // The wire shape lets unknown values pass through with the raw string
396        // so the caller can see what happened.
397        let out = do_status_to_output("provisioning");
398        assert_eq!(out.phase, VpsPhase::Unknown);
399        assert_eq!(out.detail.as_deref(), Some("provisioning"));
400    }
401
402    #[test]
403    fn create_spec_omits_optional_fields_when_empty() {
404        let spec = DoCreateDropletSpec {
405            name: "x".into(),
406            region: "fra1".into(),
407            size: "s-1vcpu-1gb".into(),
408            image: "debian-12-x64".into(),
409            ssh_keys: vec![],
410            user_data: None,
411        };
412        let wire = serde_json::to_value(&spec).unwrap();
413        assert!(wire.get("user_data").is_none());
414        assert!(wire.get("ssh_keys").is_none());
415        assert_eq!(wire["name"], "x");
416        assert_eq!(wire["region"], "fra1");
417        assert_eq!(wire["size"], "s-1vcpu-1gb");
418        assert_eq!(wire["image"], "debian-12-x64");
419    }
420
421    #[test]
422    fn create_spec_serializes_user_data_and_ssh_keys() {
423        let spec = DoCreateDropletSpec {
424            name: "x".into(),
425            region: "fra1".into(),
426            size: "s-1vcpu-1gb".into(),
427            image: "debian-12-x64".into(),
428            ssh_keys: vec!["123".into(), "e0:7a:1b".into()],
429            user_data: Some("#cloud-config\n".into()),
430        };
431        let wire = serde_json::to_value(&spec).unwrap();
432        assert_eq!(wire["user_data"], "#cloud-config\n");
433        // W144 D7: DO accepts both numeric IDs and fingerprints — the wire
434        // payload preserves whatever string form the caller supplied.
435        assert_eq!(wire["ssh_keys"], serde_json::json!(["123", "e0:7a:1b"]));
436    }
437
438    #[test]
439    fn envoy_advertises_three_cloud_vps_verbs() {
440        let envoy = DigitalOceanEnvoy::new(DigitalOceanClient::new("stub-token"));
441        let ids = envoy.supported_verb_ids();
442        assert_eq!(ids.len(), 3);
443        assert!(ids.contains(&"cloud.vps.create"));
444        assert!(ids.contains(&"cloud.vps.destroy"));
445        assert!(ids.contains(&"cloud.vps.status"));
446        assert_eq!(envoy.id(), "digitalocean");
447        assert_eq!(envoy.tier(), Tier::S);
448        assert_eq!(envoy.flavor(), AdapterFlavor::Native);
449    }
450
451    #[tokio::test]
452    async fn dispatch_unknown_verb_errors() {
453        let envoy = DigitalOceanEnvoy::new(DigitalOceanClient::new("stub-token"));
454        let err = envoy
455            .dispatch("cloud.dns.upsert", serde_json::json!({}))
456            .await
457            .unwrap_err();
458        assert!(err.to_string().contains("does not support"));
459    }
460
461    #[tokio::test]
462    async fn destroy_with_non_numeric_id_errors_before_http() {
463        // The id parse fails before any HTTP call — so this exercises the
464        // input-validation path against a stub-token client without needing
465        // a mocked DO endpoint.
466        let envoy = DigitalOceanEnvoy::new(DigitalOceanClient::new("stub-token"));
467        let err = envoy
468            .cloud_vps_destroy(CloudVpsDestroyInput {
469                id: "not-a-number".into(),
470            })
471            .await
472            .unwrap_err();
473        let msg = format!("{err:#}");
474        assert!(msg.contains("not-a-number"), "{msg}");
475    }
476
477    #[tokio::test]
478    async fn status_with_non_numeric_id_errors_before_http() {
479        let envoy = DigitalOceanEnvoy::new(DigitalOceanClient::new("stub-token"));
480        let err = envoy
481            .cloud_vps_status(CloudVpsStatusInput {
482                id: "droplet-abc".into(),
483            })
484            .await
485            .unwrap_err();
486        let msg = format!("{err:#}");
487        assert!(msg.contains("droplet-abc"), "{msg}");
488    }
489}