1use 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
42pub struct CloudflareEnvoy {
51 client: Arc<CloudflareClient>,
52 account_id: String,
53}
54
55impl CloudflareEnvoy {
56 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 pub fn from_arc(client: Arc<CloudflareClient>, account_id: String) -> Self {
67 Self { client, account_id }
68 }
69
70 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 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 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 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 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 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 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 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}