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