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