use std::sync::Arc;
use anyhow::{bail, Context, Result};
use async_trait::async_trait;
use serde_json::Value;
use super::cloudflare::CloudflareClient;
use crate::envoy::cloud_object::{
CloudObjectBucketCreate, CloudObjectBucketCreateInput, CloudObjectBucketCreateOutput,
CloudObjectBucketDelete, CloudObjectBucketDeleteInput, CloudObjectBucketDeleteOutput,
CloudObjectBucketExists, CloudObjectBucketExistsInput, CloudObjectBucketExistsOutput,
};
use crate::envoy::dns_record::{
DnsRecordDelete, DnsRecordDeleteInput, DnsRecordDeleteOutput, DnsRecordEntry, DnsRecordList,
DnsRecordListInput, DnsRecordListOutput, DnsRecordUpsert, DnsRecordUpsertInput,
DnsRecordUpsertOutput, DnsZoneEntry, DnsZoneList, DnsZoneListInput, DnsZoneListOutput,
};
use crate::envoy::{AdapterFlavor, EnvoyAdapter, InternalVerb, Tier};
pub struct CloudflareEnvoy {
client: Arc<CloudflareClient>,
account_id: String,
}
impl CloudflareEnvoy {
pub fn new(token: String, account_id: String) -> Self {
Self {
client: Arc::new(CloudflareClient::new(token)),
account_id,
}
}
pub fn from_arc(client: Arc<CloudflareClient>, account_id: String) -> Self {
Self { client, account_id }
}
pub async fn cloud_object_bucket_create(
&self,
input: CloudObjectBucketCreateInput,
) -> Result<CloudObjectBucketCreateOutput> {
let result = self
.client
.create_r2_bucket(&self.account_id, &input.name)
.await
.with_context(|| format!("cloud.object.bucket.create: create {:?}", input.name))?;
Ok(CloudObjectBucketCreateOutput {
endpoint: result.endpoint,
})
}
pub async fn cloud_object_bucket_delete(
&self,
input: CloudObjectBucketDeleteInput,
) -> Result<CloudObjectBucketDeleteOutput> {
self.client
.delete_r2_bucket(&self.account_id, &input.name)
.await
.with_context(|| format!("cloud.object.bucket.delete: delete {:?}", input.name))?;
Ok(CloudObjectBucketDeleteOutput::default())
}
pub async fn cloud_object_bucket_exists(
&self,
input: CloudObjectBucketExistsInput,
) -> Result<CloudObjectBucketExistsOutput> {
let buckets = self
.client
.list_r2_buckets(&self.account_id)
.await
.context("cloud.object.bucket.exists: list buckets")?;
let found = buckets.iter().any(|b| b.name == input.name);
Ok(CloudObjectBucketExistsOutput {
exists: found,
endpoint: if found {
Some(format!(
"https://{}.r2.cloudflarestorage.com",
self.account_id
))
} else {
None
},
})
}
pub async fn dns_record_upsert(
&self,
input: DnsRecordUpsertInput,
) -> Result<DnsRecordUpsertOutput> {
let zone_id = self
.client
.zone_id_for_name(&input.zone)
.await
.with_context(|| format!("dns.record.upsert: resolve zone {:?}", input.zone))?;
let id = self
.client
.upsert_dns_record_matching(
&zone_id,
&input.name,
&input.record_type,
&input.content,
input.ttl,
input.proxied,
input.match_content,
)
.await
.with_context(|| {
format!(
"dns.record.upsert: {}/{} {}",
input.zone, input.name, input.record_type
)
})?;
Ok(DnsRecordUpsertOutput { id })
}
pub async fn dns_record_delete(
&self,
input: DnsRecordDeleteInput,
) -> Result<DnsRecordDeleteOutput> {
let zone_id = self
.client
.zone_id_for_name(&input.zone)
.await
.with_context(|| format!("dns.record.delete: resolve zone {:?}", input.zone))?;
let deleted = self
.client
.delete_dns_records_matching(
&zone_id,
&input.name,
input.record_type.as_deref(),
input.content.as_deref(),
)
.await
.with_context(|| format!("dns.record.delete: {}/{}", input.zone, input.name))?;
Ok(DnsRecordDeleteOutput { deleted })
}
pub async fn dns_record_list(&self, input: DnsRecordListInput) -> Result<DnsRecordListOutput> {
let zone_id = self
.client
.zone_id_for_name(&input.zone)
.await
.with_context(|| format!("dns.record.list: resolve zone {:?}", input.zone))?;
let records = self
.client
.list_dns_records(&zone_id, input.name.as_deref(), input.record_type.as_deref())
.await
.with_context(|| format!("dns.record.list: {}", input.zone))?;
Ok(DnsRecordListOutput {
records: records
.into_iter()
.map(|r| DnsRecordEntry {
id: r.id,
name: r.name,
record_type: r.record_type,
content: r.content,
ttl: r.ttl,
proxied: r.proxied,
})
.collect(),
})
}
pub async fn dns_zone_list(&self, _input: DnsZoneListInput) -> Result<DnsZoneListOutput> {
let zones = self.client.list_zones().await.context("dns.zone.list")?;
Ok(DnsZoneListOutput {
zones: zones
.into_iter()
.map(|(id, name)| DnsZoneEntry { id, name })
.collect(),
})
}
}
#[async_trait]
impl EnvoyAdapter for CloudflareEnvoy {
fn id(&self) -> &str {
"cloudflare"
}
fn tier(&self) -> Tier {
Tier::S
}
fn flavor(&self) -> AdapterFlavor {
AdapterFlavor::Native
}
fn supported_verb_ids(&self) -> Vec<&'static str> {
vec![
CloudObjectBucketCreate::ID,
CloudObjectBucketDelete::ID,
CloudObjectBucketExists::ID,
DnsRecordUpsert::ID,
DnsRecordList::ID,
DnsRecordDelete::ID,
DnsZoneList::ID,
]
}
async fn dispatch(&self, verb_id: &str, input: Value) -> Result<Value> {
match verb_id {
id if id == CloudObjectBucketCreate::ID => {
let args: CloudObjectBucketCreateInput =
serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
let out = self.cloud_object_bucket_create(args).await?;
Ok(serde_json::to_value(out)?)
}
id if id == CloudObjectBucketDelete::ID => {
let args: CloudObjectBucketDeleteInput =
serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
let out = self.cloud_object_bucket_delete(args).await?;
Ok(serde_json::to_value(out)?)
}
id if id == CloudObjectBucketExists::ID => {
let args: CloudObjectBucketExistsInput =
serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
let out = self.cloud_object_bucket_exists(args).await?;
Ok(serde_json::to_value(out)?)
}
id if id == DnsRecordUpsert::ID => {
let args: DnsRecordUpsertInput =
serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
let out = self.dns_record_upsert(args).await?;
Ok(serde_json::to_value(out)?)
}
id if id == DnsRecordList::ID => {
let args: DnsRecordListInput =
serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
let out = self.dns_record_list(args).await?;
Ok(serde_json::to_value(out)?)
}
id if id == DnsRecordDelete::ID => {
let args: DnsRecordDeleteInput =
serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
let out = self.dns_record_delete(args).await?;
Ok(serde_json::to_value(out)?)
}
id if id == DnsZoneList::ID => {
let args: DnsZoneListInput =
serde_json::from_value(input).with_context(|| format!("{id}: decode input"))?;
let out = self.dns_zone_list(args).await?;
Ok(serde_json::to_value(out)?)
}
other => bail!("cloudflare envoy does not support verb {other:?}"),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn cloudflare_envoy_metadata() {
let envoy = CloudflareEnvoy::new("tok".into(), "acct123".into());
assert_eq!(envoy.id(), "cloudflare");
assert_eq!(envoy.tier(), Tier::S);
assert_eq!(envoy.flavor(), AdapterFlavor::Native);
}
#[test]
fn supported_verb_ids_covers_both_categories() {
let envoy = CloudflareEnvoy::new("tok".into(), "acct123".into());
let ids = envoy.supported_verb_ids();
assert!(ids.contains(&"cloud.object.bucket.create"), "{ids:?}");
assert!(ids.contains(&"cloud.object.bucket.delete"), "{ids:?}");
assert!(ids.contains(&"cloud.object.bucket.exists"), "{ids:?}");
assert!(ids.contains(&"dns.record.upsert"), "{ids:?}");
assert!(ids.contains(&"dns.record.list"), "{ids:?}");
assert!(ids.contains(&"dns.record.delete"), "{ids:?}");
assert!(ids.contains(&"dns.zone.list"), "{ids:?}");
}
#[tokio::test]
async fn dispatch_rejects_unsupported_verb() {
let envoy = CloudflareEnvoy::new("tok".into(), "acct123".into());
let err = envoy
.dispatch("cloud.vps.create", serde_json::json!({}))
.await
.unwrap_err();
assert!(
err.to_string().contains("cloud.vps.create"),
"error should name the rejected verb: {err}"
);
}
}