1use 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#[derive(Clone)]
56pub struct DigitalOceanClient {
57 http: reqwest::Client,
58 token: String,
59 base_url: String,
60}
61
62impl DigitalOceanClient {
63 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 pub fn with_base_url(mut self, url: impl Into<String>) -> Self {
75 self.base_url = url.into();
76 self
77 }
78
79 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 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 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#[derive(Debug, Clone, Serialize)]
148pub struct DoCreateDropletSpec {
149 pub name: String,
150 pub region: String,
152 pub size: String,
154 pub image: String,
156 #[serde(skip_serializing_if = "Vec::is_empty")]
159 pub ssh_keys: Vec<String>,
160 #[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
181pub struct DigitalOceanEnvoy {
185 client: Arc<DigitalOceanClient>,
186}
187
188impl DigitalOceanEnvoy {
189 pub fn new(client: DigitalOceanClient) -> Self {
191 Self {
192 client: Arc::new(client),
193 }
194 }
195
196 pub fn from_arc(client: Arc<DigitalOceanClient>) -> Self {
198 Self { client }
199 }
200
201 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 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 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
297fn 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
312fn 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 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 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 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 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}