Skip to main content

cloud/provider/
cloudflare_envoy.rs

1//! [`CloudflareEnvoy`] — tier-S adapter for Cloudflare R2 (`cloud.object.*`)
2//! and Cloudflare DNS (`dns.*`) verbs (R409-T6).
3//!
4//! Wraps [`CloudflareClient`] and exposes the envoy framework's verb dispatch
5//! interface. The client is held behind `Arc` so the adapter can be shared
6//! without copying the underlying token.
7//!
8//! ## Verb coverage
9//!
10//! | Verb | Cloudflare API |
11//! |---|---|
12//! | `cloud.object.bucket.create` | `POST /accounts/{id}/r2/buckets` |
13//! | `cloud.object.bucket.delete` | `DELETE /accounts/{id}/r2/buckets/{name}` |
14//! | `cloud.object.bucket.exists` | `GET /accounts/{id}/r2/buckets` + name scan |
15//! | `dns.record.upsert` | `GET` + `POST`/`PUT` `/zones/{id}/dns_records` |
16//! | `dns.record.list` | `GET /zones/{id}/dns_records` |
17//! | `dns.record.delete` | `GET` + `DELETE` `/zones/{id}/dns_records` |
18//! | `dns.zone.list` | `GET /zones` |
19//!
20//! `cloud.object.bucket.acl.set` is not wired — Cloudflare R2 does not have
21//! per-bucket ACLs in the S3 sense; public access is managed via custom
22//! domains and Workers.
23
24use std::sync::Arc;
25
26use anyhow::{bail, Context, Result};
27use async_trait::async_trait;
28use serde_json::Value;
29
30use super::cloudflare::CloudflareClient;
31use crate::envoy::cloud_object::{
32    CloudObjectBucketCreate, CloudObjectBucketCreateInput, CloudObjectBucketCreateOutput,
33    CloudObjectBucketDelete, CloudObjectBucketDeleteInput, CloudObjectBucketDeleteOutput,
34    CloudObjectBucketExists, CloudObjectBucketExistsInput, CloudObjectBucketExistsOutput,
35};
36use crate::envoy::dns_record::{
37    DnsRecordDelete, DnsRecordDeleteInput, DnsRecordDeleteOutput, DnsRecordEntry, DnsRecordList,
38    DnsRecordListInput, DnsRecordListOutput, DnsRecordUpsert, DnsRecordUpsertInput,
39    DnsRecordUpsertOutput, DnsZoneEntry, DnsZoneList, DnsZoneListInput, DnsZoneListOutput,
40};
41use crate::envoy::{AdapterFlavor, EnvoyAdapter, InternalVerb, Tier};
42
43/// Tier-S, native-flavored envoy adapter for Cloudflare.
44///
45/// Covers two verb categories via [`CloudflareClient`]:
46/// - `cloud.object.*` — R2 bucket lifecycle.
47/// - `dns.*` — DNS record management and zone listing.
48///
49/// `account_id` is required for R2 operations. DNS operations resolve zone
50/// names via `GET /zones` — no account_id needed.
51pub struct CloudflareEnvoy {
52    client: Arc<CloudflareClient>,
53    account_id: String,
54}
55
56impl CloudflareEnvoy {
57    /// Create an envoy adapter wrapping a fresh client for the given token
58    /// and Cloudflare account.
59    pub fn new(token: String, account_id: String) -> Self {
60        Self {
61            client: Arc::new(CloudflareClient::new(token)),
62            account_id,
63        }
64    }
65
66    /// Wrap an already-constructed (potentially shared) client.
67    pub fn from_arc(client: Arc<CloudflareClient>, account_id: String) -> Self {
68        Self { client, account_id }
69    }
70
71    // ── cloud.object.* ────────────────────────────────────────────────────
72
73    /// Typed handler for `cloud.object.bucket.create`.
74    ///
75    /// The `location_hint` input field is forwarded to Cloudflare as a
76    /// bucket-create body field when the CF API accepts it; today's
77    /// [`CloudflareClient::create_r2_bucket`] does not expose it, so the
78    /// hint is noted but not forwarded. Extend when bucket-placement control
79    /// becomes a requirement.
80    pub async fn cloud_object_bucket_create(
81        &self,
82        input: CloudObjectBucketCreateInput,
83    ) -> Result<CloudObjectBucketCreateOutput> {
84        let result = self
85            .client
86            .create_r2_bucket(&self.account_id, &input.name)
87            .await
88            .with_context(|| format!("cloud.object.bucket.create: create {:?}", input.name))?;
89        Ok(CloudObjectBucketCreateOutput {
90            endpoint: result.endpoint,
91        })
92    }
93
94    /// Typed handler for `cloud.object.bucket.delete`.
95    pub async fn cloud_object_bucket_delete(
96        &self,
97        input: CloudObjectBucketDeleteInput,
98    ) -> Result<CloudObjectBucketDeleteOutput> {
99        self.client
100            .delete_r2_bucket(&self.account_id, &input.name)
101            .await
102            .with_context(|| format!("cloud.object.bucket.delete: delete {:?}", input.name))?;
103        Ok(CloudObjectBucketDeleteOutput::default())
104    }
105
106    /// Typed handler for `cloud.object.bucket.exists`.
107    ///
108    /// Implemented by listing all R2 buckets and scanning for the name.
109    /// Cloudflare has no dedicated HEAD endpoint for R2 buckets.
110    pub async fn cloud_object_bucket_exists(
111        &self,
112        input: CloudObjectBucketExistsInput,
113    ) -> Result<CloudObjectBucketExistsOutput> {
114        let buckets = self
115            .client
116            .list_r2_buckets(&self.account_id)
117            .await
118            .context("cloud.object.bucket.exists: list buckets")?;
119        let found = buckets.iter().any(|b| b.name == input.name);
120        Ok(CloudObjectBucketExistsOutput {
121            exists: found,
122            endpoint: if found {
123                Some(format!(
124                    "https://{}.r2.cloudflarestorage.com",
125                    self.account_id
126                ))
127            } else {
128                None
129            },
130        })
131    }
132
133    // ── dns.* ─────────────────────────────────────────────────────────────
134
135    /// Typed handler for `dns.record.upsert`. Resolves the zone name to an
136    /// ID, then upserts the record via [`CloudflareClient::upsert_dns_record`].
137    pub async fn dns_record_upsert(
138        &self,
139        input: DnsRecordUpsertInput,
140    ) -> Result<DnsRecordUpsertOutput> {
141        let zone_id = self
142            .client
143            .zone_id_for_name(&input.zone)
144            .await
145            .with_context(|| format!("dns.record.upsert: resolve zone {:?}", input.zone))?;
146        let id = self
147            .client
148            .upsert_dns_record_matching(
149                &zone_id,
150                &input.name,
151                &input.record_type,
152                &input.content,
153                input.ttl,
154                input.proxied,
155                input.match_content,
156            )
157            .await
158            .with_context(|| {
159                format!(
160                    "dns.record.upsert: {}/{} {}",
161                    input.zone, input.name, input.record_type
162                )
163            })?;
164        Ok(DnsRecordUpsertOutput { id })
165    }
166
167    /// Typed handler for `dns.record.delete`.
168    pub async fn dns_record_delete(
169        &self,
170        input: DnsRecordDeleteInput,
171    ) -> Result<DnsRecordDeleteOutput> {
172        let zone_id = self
173            .client
174            .zone_id_for_name(&input.zone)
175            .await
176            .with_context(|| format!("dns.record.delete: resolve zone {:?}", input.zone))?;
177        let deleted = self
178            .client
179            .delete_dns_records_matching(
180                &zone_id,
181                &input.name,
182                input.record_type.as_deref(),
183                input.content.as_deref(),
184            )
185            .await
186            .with_context(|| format!("dns.record.delete: {}/{}", input.zone, input.name))?;
187        Ok(DnsRecordDeleteOutput { deleted })
188    }
189
190    /// Typed handler for `dns.record.list` (R859-F1).
191    pub async fn dns_record_list(&self, input: DnsRecordListInput) -> Result<DnsRecordListOutput> {
192        let zone_id = self
193            .client
194            .zone_id_for_name(&input.zone)
195            .await
196            .with_context(|| format!("dns.record.list: resolve zone {:?}", input.zone))?;
197        let records = self
198            .client
199            .list_dns_records(&zone_id, input.name.as_deref(), input.record_type.as_deref())
200            .await
201            .with_context(|| format!("dns.record.list: {}", input.zone))?;
202        Ok(DnsRecordListOutput {
203            records: records
204                .into_iter()
205                .map(|r| DnsRecordEntry {
206                    id: r.id,
207                    name: r.name,
208                    record_type: r.record_type,
209                    content: r.content,
210                    ttl: r.ttl,
211                    proxied: r.proxied,
212                })
213                .collect(),
214        })
215    }
216
217    /// Typed handler for `dns.zone.list`.
218    pub async fn dns_zone_list(&self, _input: DnsZoneListInput) -> Result<DnsZoneListOutput> {
219        let zones = self.client.list_zones().await.context("dns.zone.list")?;
220        Ok(DnsZoneListOutput {
221            zones: zones
222                .into_iter()
223                .map(|(id, name)| DnsZoneEntry { id, name })
224                .collect(),
225        })
226    }
227}
228
229#[async_trait]
230impl EnvoyAdapter for CloudflareEnvoy {
231    fn id(&self) -> &str {
232        "cloudflare"
233    }
234    fn tier(&self) -> Tier {
235        Tier::S
236    }
237    fn flavor(&self) -> AdapterFlavor {
238        AdapterFlavor::Native
239    }
240    fn supported_verb_ids(&self) -> Vec<&'static str> {
241        vec![
242            CloudObjectBucketCreate::ID,
243            CloudObjectBucketDelete::ID,
244            CloudObjectBucketExists::ID,
245            DnsRecordUpsert::ID,
246            DnsRecordList::ID,
247            DnsRecordDelete::ID,
248            DnsZoneList::ID,
249        ]
250    }
251    async fn dispatch(&self, verb_id: &str, input: Value) -> Result<Value> {
252        match verb_id {
253            id if id == CloudObjectBucketCreate::ID => {
254                let args: CloudObjectBucketCreateInput =
255                    serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
256                let out = self.cloud_object_bucket_create(args).await?;
257                Ok(serde_json::to_value(out)?)
258            }
259            id if id == CloudObjectBucketDelete::ID => {
260                let args: CloudObjectBucketDeleteInput =
261                    serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
262                let out = self.cloud_object_bucket_delete(args).await?;
263                Ok(serde_json::to_value(out)?)
264            }
265            id if id == CloudObjectBucketExists::ID => {
266                let args: CloudObjectBucketExistsInput =
267                    serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
268                let out = self.cloud_object_bucket_exists(args).await?;
269                Ok(serde_json::to_value(out)?)
270            }
271            id if id == DnsRecordUpsert::ID => {
272                let args: DnsRecordUpsertInput =
273                    serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
274                let out = self.dns_record_upsert(args).await?;
275                Ok(serde_json::to_value(out)?)
276            }
277            id if id == DnsRecordList::ID => {
278                let args: DnsRecordListInput =
279                    serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
280                let out = self.dns_record_list(args).await?;
281                Ok(serde_json::to_value(out)?)
282            }
283            id if id == DnsRecordDelete::ID => {
284                let args: DnsRecordDeleteInput =
285                    serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
286                let out = self.dns_record_delete(args).await?;
287                Ok(serde_json::to_value(out)?)
288            }
289            id if id == DnsZoneList::ID => {
290                let args: DnsZoneListInput =
291                    serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
292                let out = self.dns_zone_list(args).await?;
293                Ok(serde_json::to_value(out)?)
294            }
295            other => bail!("cloudflare envoy does not support verb {other:?}"),
296        }
297    }
298}
299
300#[cfg(test)]
301mod tests {
302    use super::*;
303
304    #[test]
305    fn cloudflare_envoy_metadata() {
306        let envoy = CloudflareEnvoy::new("tok".into(), "acct123".into());
307        assert_eq!(envoy.id(), "cloudflare");
308        assert_eq!(envoy.tier(), Tier::S);
309        assert_eq!(envoy.flavor(), AdapterFlavor::Native);
310    }
311
312    #[test]
313    fn supported_verb_ids_covers_both_categories() {
314        let envoy = CloudflareEnvoy::new("tok".into(), "acct123".into());
315        let ids = envoy.supported_verb_ids();
316        // cloud.object.*
317        assert!(ids.contains(&"cloud.object.bucket.create"), "{ids:?}");
318        assert!(ids.contains(&"cloud.object.bucket.delete"), "{ids:?}");
319        assert!(ids.contains(&"cloud.object.bucket.exists"), "{ids:?}");
320        // dns.*
321        assert!(ids.contains(&"dns.record.upsert"), "{ids:?}");
322        assert!(ids.contains(&"dns.record.list"), "{ids:?}");
323        assert!(ids.contains(&"dns.record.delete"), "{ids:?}");
324        assert!(ids.contains(&"dns.zone.list"), "{ids:?}");
325    }
326
327    #[tokio::test]
328    async fn dispatch_rejects_unsupported_verb() {
329        let envoy = CloudflareEnvoy::new("tok".into(), "acct123".into());
330        let err = envoy
331            .dispatch("cloud.vps.create", serde_json::json!({}))
332            .await
333            .unwrap_err();
334        assert!(
335            err.to_string().contains("cloud.vps.create"),
336            "error should name the rejected verb: {err}"
337        );
338    }
339}