1use anyhow::{anyhow, Result};
69use serde::{Deserialize, Serialize};
70
71const CF_API: &str = "https://api.cloudflare.com/client/v4";
72
73#[derive(Deserialize)]
76struct CfPage<T> {
77 success: bool,
78 result: Option<Vec<T>>,
79 errors: Option<Vec<serde_json::Value>>,
80}
81
82#[derive(Deserialize)]
83struct CfSingle<T> {
84 success: bool,
85 result: Option<T>,
86 errors: Option<Vec<serde_json::Value>>,
87}
88
89#[derive(Deserialize)]
90struct CfTunnel {
91 id: String,
92 name: String,
93 #[serde(default)]
95 status: Option<String>,
96 #[serde(default)]
99 conns_active_at: Option<String>,
100}
101
102struct TunnelMeta {
104 id: String,
105 name: String,
106 conn_state: TunnelConnState,
107 conn_since: Option<String>,
108}
109
110#[derive(Deserialize)]
111struct CfTunnelConfig {
112 config: Option<CfIngressConfig>,
113}
114
115#[derive(Deserialize)]
116struct CfIngressConfig {
117 ingress: Option<Vec<CfIngressRule>>,
118}
119
120#[derive(Deserialize)]
121struct CfIngressRule {
122 hostname: Option<String>,
123}
124
125#[derive(Debug, Clone, Deserialize)]
130struct CfDnsRecord {
131 name: String,
132 content: String,
133}
134
135#[derive(Debug, Clone, Serialize, Deserialize)]
139#[serde(rename_all = "camelCase")]
140pub struct CfAccountInfo {
141 pub id: String,
142 pub name: String,
143}
144
145#[derive(Debug, Clone, Serialize, Deserialize)]
148#[serde(rename_all = "camelCase")]
149pub struct TunnelDnsRecord {
150 pub tunnel_name: String,
152 pub hostname: String,
154 pub cname_target: String,
156}
157
158#[derive(Debug, Clone, Serialize, Deserialize)]
160#[serde(rename_all = "camelCase")]
161pub struct CreateTunnelResult {
162 pub tunnel_id: String,
163 pub tunnel_name: String,
164 pub connector_token: String,
166 pub cname_target: String,
168}
169
170#[derive(Debug, Clone, Serialize, Deserialize)]
172#[serde(rename_all = "camelCase")]
173pub struct CreateR2BucketResult {
174 pub name: String,
175 pub endpoint: String,
177}
178
179#[derive(Debug, Clone, Serialize, Deserialize)]
181#[serde(rename_all = "camelCase")]
182pub struct WorkerDeployResult {
183 pub id: String,
184 pub etag: Option<String>,
185}
186
187#[derive(Debug, Clone, Copy)]
194pub enum WorkerBinding<'a> {
195 PlainText { name: &'a str, text: &'a str },
197 R2Bucket { name: &'a str, bucket_name: &'a str },
199}
200
201#[derive(Debug, Clone, Serialize, Deserialize)]
203#[serde(rename_all = "camelCase")]
204pub struct R2BucketInfo {
205 pub name: String,
206 #[serde(default, skip_serializing_if = "Option::is_none")]
209 pub location: Option<String>,
210 #[serde(default, skip_serializing_if = "Option::is_none")]
212 pub creation_date: Option<String>,
213}
214
215#[derive(Debug, Clone, Serialize, Deserialize)]
222#[serde(rename_all = "camelCase")]
223pub struct R2CustomDomain {
224 pub domain: String,
225 #[serde(default)]
226 pub enabled: bool,
227}
228
229#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
235#[serde(rename_all = "kebab-case")]
236pub enum TunnelConnState {
237 Active,
239 Inactive,
241 Degraded,
243 Unknown,
245}
246
247impl TunnelConnState {
248 fn from_cf_status(s: &str) -> Self {
249 match s {
250 "active" | "healthy" => Self::Active,
251 "inactive" | "down" => Self::Inactive,
252 "degraded" | "unhealthy" => Self::Degraded,
253 _ => Self::Unknown,
254 }
255 }
256}
257
258#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
265#[serde(rename_all = "kebab-case")]
266pub enum TunnelDriftState {
267 Synced,
269 Missing,
271 Mismatch,
273 ZoneUnknown,
277}
278
279#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
283#[serde(rename_all = "camelCase")]
284pub struct TunnelDriftRow {
285 pub tunnel_name: String,
287 pub hostname: String,
289 pub expected_target: String,
291 pub state: TunnelDriftState,
292 #[serde(default, skip_serializing_if = "Option::is_none")]
295 pub live_target: Option<String>,
296 pub conn_state: TunnelConnState,
299 #[serde(default, skip_serializing_if = "Option::is_none")]
303 pub conn_since: Option<String>,
304}
305
306#[derive(Debug, Clone, Copy, PartialEq, Eq)]
310pub enum GrantScope {
311 Account,
313 Zone,
315}
316
317#[derive(Debug, Clone, Copy)]
325pub struct TokenGrant {
326 pub group_name: &'static str,
327 pub scope: GrantScope,
328 pub fallback_id: &'static str,
329}
330
331pub const MESOFACT_STATIC_GRANTS: &[TokenGrant] = &[
346 TokenGrant {
347 group_name: "Account Settings Read",
348 scope: GrantScope::Account,
349 fallback_id: "c1fde68c7bcc44588cbb6ddbc16d6480",
350 },
351 TokenGrant {
352 group_name: "Workers R2 Storage Write",
353 scope: GrantScope::Account,
354 fallback_id: "bf7481a1826f439697cb59a20b22293e",
355 },
356 TokenGrant {
357 group_name: "Workers Scripts Write",
358 scope: GrantScope::Account,
359 fallback_id: "e086da7e2179491d91ee5f35b3ca210a",
360 },
361 TokenGrant {
362 group_name: "Zone Read",
363 scope: GrantScope::Zone,
364 fallback_id: "c8fed203ed3043cba015a93ad1616f1f",
365 },
366 TokenGrant {
367 group_name: "Zone Transform Rules Write",
368 scope: GrantScope::Zone,
369 fallback_id: "0ac90a90249747bca6b047d97f0803e9",
370 },
371 TokenGrant {
372 group_name: "Workers Routes Write",
373 scope: GrantScope::Zone,
374 fallback_id: "28f4b596e7d643029c524985477ae49a",
375 },
376 TokenGrant {
377 group_name: "Cache Purge",
378 scope: GrantScope::Zone,
379 fallback_id: "e17beae8b8cb423a99b1730f21238bed",
380 },
381];
382
383#[derive(Debug, Clone, Serialize, Deserialize)]
386#[serde(rename_all = "camelCase")]
387pub struct CreateTokenResult {
388 pub id: String,
389 pub name: String,
390 pub value: String,
392}
393
394pub struct CloudflareClient {
401 token: String,
402 http: reqwest::Client,
403}
404
405impl CloudflareClient {
406 pub fn new(token: String) -> Self {
408 Self {
409 token,
410 http: reqwest::Client::new(),
411 }
412 }
413
414 pub async fn list_accounts(&self) -> Result<Vec<CfAccountInfo>> {
417 #[derive(Deserialize)]
418 struct Entry {
419 id: String,
420 name: String,
421 }
422 let resp: CfPage<Entry> = self.cf_get("/accounts").await?;
423 self.ok(&resp.success, &resp.errors)?;
424 Ok(resp
425 .result
426 .unwrap_or_default()
427 .into_iter()
428 .map(|a| CfAccountInfo {
429 id: a.id,
430 name: a.name,
431 })
432 .collect())
433 }
434
435 async fn list_tunnels_meta(&self, account_id: &str) -> Result<Vec<TunnelMeta>> {
438 let resp: CfPage<CfTunnel> = self
439 .cf_get(&format!(
440 "/accounts/{account_id}/cfd_tunnel?is_deleted=false"
441 ))
442 .await?;
443 self.ok(&resp.success, &resp.errors)?;
444 Ok(resp
445 .result
446 .unwrap_or_default()
447 .into_iter()
448 .map(|t| TunnelMeta {
449 conn_state: t
450 .status
451 .as_deref()
452 .map(TunnelConnState::from_cf_status)
453 .unwrap_or(TunnelConnState::Unknown),
454 conn_since: t.conns_active_at,
455 id: t.id,
456 name: t.name,
457 })
458 .collect())
459 }
460
461 pub async fn list_tunnels(&self, account_id: &str) -> Result<Vec<(String, String)>> {
464 Ok(self
465 .list_tunnels_meta(account_id)
466 .await?
467 .into_iter()
468 .map(|m| (m.id, m.name))
469 .collect())
470 }
471
472 pub async fn tunnel_dns_records(&self) -> Result<Vec<TunnelDnsRecord>> {
477 let mut records = Vec::new();
478 let accounts = self.list_accounts().await?;
479
480 for account in &accounts {
481 let tunnels = self.list_tunnels(&account.id).await?;
482 for (tunnel_id, tunnel_name) in &tunnels {
483 let cname_target = format!("{tunnel_id}.cfargotunnel.com");
484 let path = format!(
485 "/accounts/{}/cfd_tunnel/{tunnel_id}/configurations",
486 account.id
487 );
488 let config_resp: CfSingle<CfTunnelConfig> = self.cf_get(&path).await?;
489 self.ok(&config_resp.success, &config_resp.errors)?;
490 let ingress = config_resp
491 .result
492 .and_then(|c| c.config)
493 .and_then(|c| c.ingress)
494 .unwrap_or_default();
495 for rule in ingress {
496 if let Some(hostname) = rule.hostname.filter(|h| !h.is_empty()) {
497 records.push(TunnelDnsRecord {
498 tunnel_name: tunnel_name.clone(),
499 hostname,
500 cname_target: cname_target.clone(),
501 });
502 }
503 }
504 }
505 }
506 Ok(records)
507 }
508
509 pub async fn tunnel_dns_drift(&self) -> Result<Vec<TunnelDriftRow>> {
527 let accounts = self.list_accounts().await?;
531 let mut declared: Vec<TunnelDnsRecord> = Vec::new();
532 let mut conn_by_tunnel: std::collections::HashMap<
533 String,
534 (TunnelConnState, Option<String>),
535 > = Default::default();
536
537 for account in &accounts {
538 let metas = self.list_tunnels_meta(&account.id).await?;
539 for meta in &metas {
540 conn_by_tunnel.insert(
541 meta.name.clone(),
542 (meta.conn_state, meta.conn_since.clone()),
543 );
544 let cname_target = format!("{}.cfargotunnel.com", meta.id);
545 let path = format!(
546 "/accounts/{}/cfd_tunnel/{}/configurations",
547 account.id, meta.id
548 );
549 let config_resp: CfSingle<CfTunnelConfig> = self.cf_get(&path).await?;
550 self.ok(&config_resp.success, &config_resp.errors)?;
551 let ingress = config_resp
552 .result
553 .and_then(|c| c.config)
554 .and_then(|c| c.ingress)
555 .unwrap_or_default();
556 for rule in ingress {
557 if let Some(hostname) = rule.hostname.filter(|h| !h.is_empty()) {
558 declared.push(TunnelDnsRecord {
559 tunnel_name: meta.name.clone(),
560 hostname,
561 cname_target: cname_target.clone(),
562 });
563 }
564 }
565 }
566 }
567
568 if declared.is_empty() {
569 return Ok(Vec::new());
570 }
571
572 let zones = self.list_zones().await.unwrap_or_default();
575
576 let mut rows = Vec::with_capacity(declared.len());
577 for rec in declared {
578 let (conn_state, conn_since) = conn_by_tunnel
579 .remove(&rec.tunnel_name)
580 .unwrap_or((TunnelConnState::Unknown, None));
581 let (state, live_target) = match best_zone_for(&rec.hostname, &zones) {
582 None => (TunnelDriftState::ZoneUnknown, None),
583 Some(zone_id) => match self.dns_records_named(zone_id, &rec.hostname).await {
584 Ok(live) => classify_tunnel_drift(&rec.hostname, &rec.cname_target, &live),
585 Err(_) => (TunnelDriftState::ZoneUnknown, None),
586 },
587 };
588 rows.push(TunnelDriftRow {
589 tunnel_name: rec.tunnel_name,
590 hostname: rec.hostname,
591 expected_target: rec.cname_target,
592 state,
593 live_target,
594 conn_state,
595 conn_since,
596 });
597 }
598 Ok(rows)
599 }
600
601 pub async fn list_zones(&self) -> Result<Vec<(String, String)>> {
604 #[derive(Deserialize)]
605 struct ZoneEntry {
606 id: String,
607 name: String,
608 }
609 let resp: CfPage<ZoneEntry> = self.cf_get("/zones?per_page=50").await?;
610 self.ok(&resp.success, &resp.errors)?;
611 Ok(resp
612 .result
613 .unwrap_or_default()
614 .into_iter()
615 .map(|z| (z.id, z.name))
616 .collect())
617 }
618
619 async fn dns_records_named(&self, zone_id: &str, name: &str) -> Result<Vec<CfDnsRecord>> {
622 let resp: CfPage<CfDnsRecord> = self
623 .cf_get(&format!("/zones/{zone_id}/dns_records?name={name}"))
624 .await?;
625 self.ok(&resp.success, &resp.errors)?;
626 Ok(resp.result.unwrap_or_default())
627 }
628
629 pub async fn create_tunnel(&self, account_id: &str, name: &str) -> Result<CreateTunnelResult> {
632 use base64::Engine as _;
633
634 let mut secret_bytes = [0u8; 32];
635 getrandom::getrandom(&mut secret_bytes)
636 .map_err(|e| anyhow!("generate tunnel secret: {e}"))?;
637 let tunnel_secret = base64::engine::general_purpose::STANDARD.encode(secret_bytes);
638
639 #[derive(Serialize)]
640 struct CreateBody<'a> {
641 name: &'a str,
642 tunnel_secret: String,
643 }
644 #[derive(Deserialize)]
645 struct CreatedTunnel {
646 id: String,
647 name: String,
648 }
649 let create_resp: CfSingle<CreatedTunnel> = self
650 .cf_post(
651 &format!("/accounts/{account_id}/cfd_tunnel"),
652 &CreateBody {
653 name,
654 tunnel_secret,
655 },
656 )
657 .await?;
658 self.ok(&create_resp.success, &create_resp.errors)?;
659 let created = create_resp
660 .result
661 .ok_or_else(|| anyhow!("tunnel create: no result in response"))?;
662
663 #[derive(Deserialize)]
665 struct TokenResp {
666 success: bool,
667 result: Option<String>,
668 errors: Option<Vec<serde_json::Value>>,
669 }
670 let token_resp: TokenResp = self
671 .cf_get(&format!(
672 "/accounts/{account_id}/cfd_tunnel/{}/token",
673 created.id
674 ))
675 .await?;
676 self.ok(&token_resp.success, &token_resp.errors)?;
677 let connector_token = token_resp
678 .result
679 .ok_or_else(|| anyhow!("no connector token in response"))?;
680
681 Ok(CreateTunnelResult {
682 cname_target: format!("{}.cfargotunnel.com", created.id),
683 tunnel_id: created.id,
684 tunnel_name: created.name,
685 connector_token,
686 })
687 }
688
689 pub async fn tunnel_configuration(
700 &self,
701 account_id: &str,
702 tunnel_id: &str,
703 ) -> Result<serde_json::Value> {
704 let path = format!("/accounts/{account_id}/cfd_tunnel/{tunnel_id}/configurations");
705 let resp: CfSingle<serde_json::Value> = self.cf_get(&path).await?;
706 self.ok(&resp.success, &resp.errors)?;
707 let config = resp
708 .result
709 .as_ref()
710 .and_then(|r| r.get("config"))
711 .cloned()
712 .unwrap_or_else(|| serde_json::json!({}));
713 Ok(if config.is_object() {
716 config
717 } else {
718 serde_json::json!({})
719 })
720 }
721
722 pub async fn put_tunnel_configuration(
731 &self,
732 account_id: &str,
733 tunnel_id: &str,
734 config: &serde_json::Value,
735 ) -> Result<()> {
736 #[derive(Serialize)]
737 struct ConfigBody<'a> {
738 config: &'a serde_json::Value,
739 }
740 let path = format!("/accounts/{account_id}/cfd_tunnel/{tunnel_id}/configurations");
741 let resp: CfSingle<serde_json::Value> = self.cf_put(&path, &ConfigBody { config }).await?;
742 self.ok(&resp.success, &resp.errors)?;
743 Ok(())
744 }
745
746 pub async fn create_r2_bucket(
749 &self,
750 account_id: &str,
751 bucket_name: &str,
752 ) -> Result<CreateR2BucketResult> {
753 #[derive(Serialize)]
754 struct CreateBody<'a> {
755 name: &'a str,
756 }
757 let resp: CfSingle<serde_json::Value> = self
758 .cf_post(
759 &format!("/accounts/{account_id}/r2/buckets"),
760 &CreateBody { name: bucket_name },
761 )
762 .await?;
763 self.ok(&resp.success, &resp.errors)?;
764
765 Ok(CreateR2BucketResult {
766 endpoint: format!("https://{account_id}.r2.cloudflarestorage.com"),
767 name: bucket_name.to_string(),
768 })
769 }
770
771 pub async fn zone_id_for_name(&self, zone_name: &str) -> Result<String> {
774 #[derive(Deserialize)]
775 struct ZoneEntry {
776 id: String,
777 name: String,
778 }
779 let resp: CfPage<ZoneEntry> = self.cf_get(&format!("/zones?name={zone_name}")).await?;
780 self.ok(&resp.success, &resp.errors)?;
781 resp.result
782 .unwrap_or_default()
783 .into_iter()
784 .find(|z| z.name == zone_name)
785 .map(|z| z.id)
786 .ok_or_else(|| anyhow!("no Cloudflare zone found for name {zone_name:?}"))
787 }
788
789 pub async fn purge_cache_tags(&self, zone_id: &str, tags: &[String]) -> Result<()> {
795 if tags.is_empty() {
796 return Ok(());
797 }
798 #[derive(Serialize)]
799 struct PurgeBody<'a> {
800 tags: &'a [String],
801 }
802 let resp: CfSingle<serde_json::Value> = self
803 .cf_post(
804 &format!("/zones/{zone_id}/purge_cache"),
805 &PurgeBody { tags },
806 )
807 .await?;
808 self.ok(&resp.success, &resp.errors)
809 }
810
811 pub async fn upsert_index_rewrite(&self, zone_id: &str) -> Result<()> {
821 const RULE_DESC: &str = "yah:static-index";
822 let path = format!("/zones/{zone_id}/rulesets/phases/http_request_transform/entrypoint");
823
824 let existing: Vec<serde_json::Value> = {
826 #[derive(Deserialize)]
827 struct Rs {
828 rules: Option<Vec<serde_json::Value>>,
829 }
830 match self.cf_get::<CfSingle<Rs>>(&path).await {
831 Ok(resp) if resp.success => resp.result.and_then(|r| r.rules).unwrap_or_default(),
832 _ => Vec::new(),
833 }
834 };
835
836 let mut rules: Vec<serde_json::Value> = existing
838 .into_iter()
839 .filter(|r| r.get("description").and_then(|v| v.as_str()) != Some(RULE_DESC))
840 .collect();
841 rules.push(serde_json::json!({
842 "action": "rewrite",
843 "description": RULE_DESC,
844 "expression": "(http.request.uri.path eq \"/\")",
845 "action_parameters": {
846 "uri": { "path": { "value": "/index.html" } }
847 },
848 "enabled": true
849 }));
850
851 let resp: CfSingle<serde_json::Value> = self
852 .cf_put(&path, &serde_json::json!({ "rules": rules }))
853 .await?;
854 self.ok(&resp.success, &resp.errors)
855 }
856
857 pub async fn deploy_worker_script(
869 &self,
870 account_id: &str,
871 script_name: &str,
872 script_js: &str,
873 bindings: &[WorkerBinding<'_>],
874 ) -> Result<WorkerDeployResult> {
875 let url = format!("{CF_API}/accounts/{account_id}/workers/scripts/{script_name}");
876 let (content_type, body) = build_worker_multipart(script_js, bindings);
877 let resp = self
878 .http
879 .put(&url)
880 .header("Authorization", format!("Bearer {}", self.token))
881 .header("Content-Type", content_type)
882 .body(body)
883 .send()
884 .await
885 .map_err(|e| anyhow!("PUT {url}: {e}"))?;
886 let result: CfSingle<WorkerDeployResult> = resp
887 .json()
888 .await
889 .map_err(|e| anyhow!("PUT {url} parse: {e}"))?;
890 self.ok(&result.success, &result.errors)?;
891 result
892 .result
893 .ok_or_else(|| anyhow!("deploy worker: no result in response"))
894 }
895
896 pub async fn upsert_worker_route(
904 &self,
905 zone_id: &str,
906 pattern: &str,
907 script_name: &str,
908 ) -> Result<()> {
909 let list_path = format!("/zones/{zone_id}/workers/routes");
910
911 #[derive(Deserialize)]
912 struct RouteEntry {
913 id: String,
914 pattern: String,
915 #[serde(default)]
916 script: Option<String>,
917 }
918 let list: CfPage<RouteEntry> = self.cf_get(&list_path).await?;
919 self.ok(&list.success, &list.errors)?;
920 let routes = list.result.unwrap_or_default();
921
922 #[derive(Serialize)]
923 struct RouteBody<'a> {
924 pattern: &'a str,
925 script: &'a str,
926 }
927
928 if let Some(existing) = routes.iter().find(|r| r.pattern == pattern) {
929 if existing.script.as_deref() == Some(script_name) {
930 return Ok(());
931 }
932 let resp: CfSingle<serde_json::Value> = self
933 .cf_put(
934 &format!("/zones/{zone_id}/workers/routes/{}", existing.id),
935 &RouteBody {
936 pattern,
937 script: script_name,
938 },
939 )
940 .await?;
941 self.ok(&resp.success, &resp.errors)
942 } else {
943 let resp: CfSingle<serde_json::Value> = self
944 .cf_post(
945 &list_path,
946 &RouteBody {
947 pattern,
948 script: script_name,
949 },
950 )
951 .await?;
952 self.ok(&resp.success, &resp.errors)
953 }
954 }
955
956 pub async fn upsert_worker_custom_domain(
968 &self,
969 account_id: &str,
970 zone_id: &str,
971 hostname: &str,
972 script_name: &str,
973 ) -> Result<()> {
974 let list_path = format!("/accounts/{account_id}/workers/domains");
975
976 #[derive(Deserialize)]
977 struct DomainEntry {
978 #[serde(default)]
979 hostname: Option<String>,
980 #[serde(default)]
981 service: Option<String>,
982 #[serde(default, rename = "zone_id")]
983 zone_id: Option<String>,
984 }
985 let list: CfPage<DomainEntry> = self.cf_get(&list_path).await?;
986 self.ok(&list.success, &list.errors)?;
987 let domains = list.result.unwrap_or_default();
988 if domains.iter().any(|d| {
989 d.hostname.as_deref() == Some(hostname)
990 && d.service.as_deref() == Some(script_name)
991 && d.zone_id.as_deref() == Some(zone_id)
992 }) {
993 return Ok(());
994 }
995
996 #[derive(Serialize)]
997 struct DomainBody<'a> {
998 environment: &'a str,
999 hostname: &'a str,
1000 service: &'a str,
1001 zone_id: &'a str,
1002 }
1003 let resp: CfSingle<serde_json::Value> = self
1004 .cf_put(
1005 &list_path,
1006 &DomainBody {
1007 environment: "production",
1008 hostname,
1009 service: script_name,
1010 zone_id,
1011 },
1012 )
1013 .await?;
1014 self.ok(&resp.success, &resp.errors)
1015 }
1016
1017 pub async fn delete_r2_bucket(&self, account_id: &str, bucket_name: &str) -> Result<()> {
1026 let resp: CfSingle<serde_json::Value> = self
1027 .cf_delete(&format!("/accounts/{account_id}/r2/buckets/{bucket_name}"))
1028 .await?;
1029 self.ok(&resp.success, &resp.errors)
1030 }
1031
1032 pub async fn upsert_dns_record(
1038 &self,
1039 zone_id: &str,
1040 name: &str,
1041 record_type: &str,
1042 content: &str,
1043 ttl: u32,
1044 proxied: bool,
1045 ) -> Result<String> {
1046 let existing_id = self.find_dns_record(zone_id, name, record_type).await?;
1047
1048 #[derive(Serialize)]
1049 struct RecordBody<'a> {
1050 name: &'a str,
1051 #[serde(rename = "type")]
1052 record_type: &'a str,
1053 content: &'a str,
1054 ttl: u32,
1055 proxied: bool,
1056 }
1057 let body = RecordBody {
1058 name,
1059 record_type,
1060 content,
1061 ttl,
1062 proxied,
1063 };
1064
1065 #[derive(Deserialize)]
1066 struct RecordResult {
1067 id: String,
1068 }
1069
1070 if let Some(id) = existing_id {
1071 let resp: CfSingle<RecordResult> = self
1072 .cf_put(&format!("/zones/{zone_id}/dns_records/{id}"), &body)
1073 .await?;
1074 self.ok(&resp.success, &resp.errors)?;
1075 resp.result
1076 .map(|r| r.id)
1077 .ok_or_else(|| anyhow!("dns record update: no id in response"))
1078 } else {
1079 let resp: CfSingle<RecordResult> = self
1080 .cf_post(&format!("/zones/{zone_id}/dns_records"), &body)
1081 .await?;
1082 self.ok(&resp.success, &resp.errors)?;
1083 resp.result
1084 .map(|r| r.id)
1085 .ok_or_else(|| anyhow!("dns record create: no id in response"))
1086 }
1087 }
1088
1089 pub async fn delete_dns_records(
1095 &self,
1096 zone_id: &str,
1097 name: &str,
1098 record_type: Option<&str>,
1099 ) -> Result<u32> {
1100 #[derive(Deserialize)]
1101 struct RecordEntry {
1102 id: String,
1103 }
1104 let query = match record_type {
1105 Some(t) => {
1106 format!("/zones/{zone_id}/dns_records?name={name}&type={t}")
1107 }
1108 None => format!("/zones/{zone_id}/dns_records?name={name}"),
1109 };
1110 let resp: CfPage<RecordEntry> = self.cf_get(&query).await?;
1111 self.ok(&resp.success, &resp.errors)?;
1112 let ids: Vec<String> = resp
1113 .result
1114 .unwrap_or_default()
1115 .into_iter()
1116 .map(|r| r.id)
1117 .collect();
1118 let mut deleted = 0u32;
1119 for id in &ids {
1120 let del: CfSingle<serde_json::Value> = self
1121 .cf_delete(&format!("/zones/{zone_id}/dns_records/{id}"))
1122 .await?;
1123 self.ok(&del.success, &del.errors)?;
1124 deleted += 1;
1125 }
1126 Ok(deleted)
1127 }
1128
1129 async fn find_dns_record(
1132 &self,
1133 zone_id: &str,
1134 name: &str,
1135 record_type: &str,
1136 ) -> Result<Option<String>> {
1137 #[derive(Deserialize)]
1138 struct RecordEntry {
1139 id: String,
1140 }
1141 let resp: CfPage<RecordEntry> = self
1142 .cf_get(&format!(
1143 "/zones/{zone_id}/dns_records?name={name}&type={record_type}"
1144 ))
1145 .await?;
1146 self.ok(&resp.success, &resp.errors)?;
1147 Ok(resp
1148 .result
1149 .unwrap_or_default()
1150 .into_iter()
1151 .next()
1152 .map(|r| r.id))
1153 }
1154
1155 pub async fn list_r2_buckets(&self, account_id: &str) -> Result<Vec<R2BucketInfo>> {
1162 #[derive(Deserialize)]
1163 struct BucketEntry {
1164 name: String,
1165 #[serde(default)]
1166 location: Option<String>,
1167 #[serde(default)]
1168 creation_date: Option<String>,
1169 }
1170 #[derive(Deserialize)]
1171 struct R2ListResult {
1172 #[serde(default)]
1173 buckets: Vec<BucketEntry>,
1174 }
1175 let resp: CfSingle<R2ListResult> = self
1176 .cf_get(&format!("/accounts/{account_id}/r2/buckets"))
1177 .await?;
1178 self.ok(&resp.success, &resp.errors)?;
1179 Ok(resp
1180 .result
1181 .map(|r| r.buckets)
1182 .unwrap_or_default()
1183 .into_iter()
1184 .map(|b| R2BucketInfo {
1185 name: b.name,
1186 location: b.location,
1187 creation_date: b.creation_date,
1188 })
1189 .collect())
1190 }
1191
1192 pub async fn list_r2_custom_domains(
1198 &self,
1199 account_id: &str,
1200 bucket_name: &str,
1201 ) -> Result<Vec<R2CustomDomain>> {
1202 #[derive(Deserialize)]
1203 struct DomainEntry {
1204 domain: String,
1205 #[serde(default)]
1206 enabled: bool,
1207 }
1208 #[derive(Deserialize)]
1209 struct R2DomainListResult {
1210 #[serde(default)]
1211 domains: Vec<DomainEntry>,
1212 }
1213 let resp: CfSingle<R2DomainListResult> = self
1214 .cf_get(&format!(
1215 "/accounts/{account_id}/r2/buckets/{bucket_name}/domains/custom"
1216 ))
1217 .await?;
1218 self.ok(&resp.success, &resp.errors)?;
1219 Ok(resp
1220 .result
1221 .map(|r| r.domains)
1222 .unwrap_or_default()
1223 .into_iter()
1224 .map(|d| R2CustomDomain {
1225 domain: d.domain,
1226 enabled: d.enabled,
1227 })
1228 .collect())
1229 }
1230
1231 pub async fn add_r2_custom_domain(
1245 &self,
1246 account_id: &str,
1247 bucket_name: &str,
1248 domain: &str,
1249 zone_id: &str,
1250 ) -> Result<()> {
1251 #[derive(Serialize)]
1252 #[serde(rename_all = "camelCase")]
1253 struct AddBody<'a> {
1254 domain: &'a str,
1255 enabled: bool,
1256 zone_id: &'a str,
1257 }
1258 let resp: CfSingle<serde_json::Value> = self
1259 .cf_post(
1260 &format!("/accounts/{account_id}/r2/buckets/{bucket_name}/domains/custom"),
1261 &AddBody {
1262 domain,
1263 enabled: true,
1264 zone_id,
1265 },
1266 )
1267 .await?;
1268 self.ok(&resp.success, &resp.errors)
1269 }
1270
1271 pub async fn list_permission_group_ids(
1275 &self,
1276 account_id: &str,
1277 ) -> Result<std::collections::BTreeMap<String, String>> {
1278 #[derive(Deserialize)]
1279 struct PgEntry {
1280 id: String,
1281 name: String,
1282 }
1283 let resp: CfPage<PgEntry> = self
1284 .cf_get(&format!(
1285 "/accounts/{account_id}/tokens/permission_groups?per_page=500"
1286 ))
1287 .await?;
1288 self.ok(&resp.success, &resp.errors)?;
1289 Ok(resp
1290 .result
1291 .unwrap_or_default()
1292 .into_iter()
1293 .map(|p| (p.name, p.id))
1294 .collect())
1295 }
1296
1297 pub async fn create_account_token(
1310 &self,
1311 account_id: &str,
1312 zone_id: &str,
1313 token_name: &str,
1314 grants: &[TokenGrant],
1315 ) -> Result<CreateTokenResult> {
1316 let catalog = self
1319 .list_permission_group_ids(account_id)
1320 .await
1321 .unwrap_or_default();
1322 let body = build_token_body(token_name, account_id, zone_id, grants, &catalog);
1323
1324 let resp: CfSingle<CreateTokenResult> = self
1325 .cf_post(&format!("/accounts/{account_id}/tokens"), &body)
1326 .await?;
1327 self.ok(&resp.success, &resp.errors)?;
1328 resp.result
1329 .ok_or_else(|| anyhow!("token create: no result in response"))
1330 }
1331
1332 async fn cf_get<T: serde::de::DeserializeOwned>(&self, path: &str) -> Result<T> {
1335 let url = format!("{CF_API}{path}");
1336 let resp = self
1337 .http
1338 .get(&url)
1339 .header("Authorization", format!("Bearer {}", self.token))
1340 .header("Content-Type", "application/json")
1341 .send()
1342 .await
1343 .map_err(|e| anyhow!("GET {url}: {e}"))?;
1344 resp.json::<T>()
1345 .await
1346 .map_err(|e| anyhow!("GET {url} parse: {e}"))
1347 }
1348
1349 async fn cf_post<B: Serialize, T: serde::de::DeserializeOwned>(
1350 &self,
1351 path: &str,
1352 body: &B,
1353 ) -> Result<T> {
1354 let url = format!("{CF_API}{path}");
1355 let resp = self
1356 .http
1357 .post(&url)
1358 .header("Authorization", format!("Bearer {}", self.token))
1359 .header("Content-Type", "application/json")
1360 .json(body)
1361 .send()
1362 .await
1363 .map_err(|e| anyhow!("POST {url}: {e}"))?;
1364 resp.json::<T>()
1365 .await
1366 .map_err(|e| anyhow!("POST {url} parse: {e}"))
1367 }
1368
1369 async fn cf_put<B: Serialize, T: serde::de::DeserializeOwned>(
1370 &self,
1371 path: &str,
1372 body: &B,
1373 ) -> Result<T> {
1374 let url = format!("{CF_API}{path}");
1375 let resp = self
1376 .http
1377 .put(&url)
1378 .header("Authorization", format!("Bearer {}", self.token))
1379 .header("Content-Type", "application/json")
1380 .json(body)
1381 .send()
1382 .await
1383 .map_err(|e| anyhow!("PUT {url}: {e}"))?;
1384 resp.json::<T>()
1385 .await
1386 .map_err(|e| anyhow!("PUT {url} parse: {e}"))
1387 }
1388
1389 async fn cf_delete<T: serde::de::DeserializeOwned>(&self, path: &str) -> Result<T> {
1390 let url = format!("{CF_API}{path}");
1391 let resp = self
1392 .http
1393 .delete(&url)
1394 .header("Authorization", format!("Bearer {}", self.token))
1395 .header("Content-Type", "application/json")
1396 .send()
1397 .await
1398 .map_err(|e| anyhow!("DELETE {url}: {e}"))?;
1399 resp.json::<T>()
1400 .await
1401 .map_err(|e| anyhow!("DELETE {url} parse: {e}"))
1402 }
1403
1404 fn ok(&self, success: &bool, errors: &Option<Vec<serde_json::Value>>) -> Result<()> {
1405 if *success {
1406 return Ok(());
1407 }
1408 let msg = errors
1409 .as_ref()
1410 .and_then(|e| e.first())
1411 .and_then(|e| e.get("message"))
1412 .and_then(|m| m.as_str())
1413 .unwrap_or("Cloudflare returned an error");
1414 Err(anyhow!("{msg}"))
1415 }
1416}
1417
1418fn build_worker_multipart(script_js: &str, bindings: &[WorkerBinding<'_>]) -> (String, Vec<u8>) {
1428 const BOUNDARY: &str = "yahWorkerUpload0";
1429 const SCRIPT_FILENAME: &str = "worker.js";
1430 let binding_json: Vec<serde_json::Value> = bindings
1431 .iter()
1432 .map(|b| match *b {
1433 WorkerBinding::PlainText { name, text } => serde_json::json!({
1434 "type": "plain_text",
1435 "name": name,
1436 "text": text,
1437 }),
1438 WorkerBinding::R2Bucket { name, bucket_name } => serde_json::json!({
1439 "type": "r2_bucket",
1440 "name": name,
1441 "bucket_name": bucket_name,
1442 }),
1443 })
1444 .collect();
1445 let metadata = serde_json::json!({
1446 "main_module": SCRIPT_FILENAME,
1447 "bindings": binding_json,
1448 });
1449 let mut body: Vec<u8> = Vec::new();
1450 let push = |v: &mut Vec<u8>, s: &str| v.extend_from_slice(s.as_bytes());
1451 push(&mut body, &format!("--{BOUNDARY}\r\n"));
1453 push(
1454 &mut body,
1455 "Content-Disposition: form-data; name=\"metadata\"\r\n",
1456 );
1457 push(&mut body, "Content-Type: application/json\r\n\r\n");
1458 push(&mut body, &metadata.to_string());
1459 push(&mut body, "\r\n");
1460 push(&mut body, &format!("--{BOUNDARY}\r\n"));
1462 push(&mut body, &format!(
1463 "Content-Disposition: form-data; name=\"{SCRIPT_FILENAME}\"; filename=\"{SCRIPT_FILENAME}\"\r\n"
1464 ));
1465 push(
1466 &mut body,
1467 "Content-Type: application/javascript+module\r\n\r\n",
1468 );
1469 push(&mut body, script_js);
1470 push(&mut body, &format!("\r\n--{BOUNDARY}--\r\n"));
1471 (format!("multipart/form-data; boundary={BOUNDARY}"), body)
1472}
1473
1474fn resolve_grant_id(
1479 grant: &TokenGrant,
1480 catalog: &std::collections::BTreeMap<String, String>,
1481) -> String {
1482 catalog
1483 .get(grant.group_name)
1484 .cloned()
1485 .unwrap_or_else(|| grant.fallback_id.to_string())
1486}
1487
1488fn build_token_body(
1493 token_name: &str,
1494 account_id: &str,
1495 zone_id: &str,
1496 grants: &[TokenGrant],
1497 catalog: &std::collections::BTreeMap<String, String>,
1498) -> serde_json::Value {
1499 let ids_for = |scope: GrantScope| -> Vec<serde_json::Value> {
1500 grants
1501 .iter()
1502 .filter(|g| g.scope == scope)
1503 .map(|g| serde_json::json!({ "id": resolve_grant_id(g, catalog) }))
1504 .collect()
1505 };
1506 let block = |resource: String, groups: Vec<serde_json::Value>| -> Option<serde_json::Value> {
1507 if groups.is_empty() {
1508 return None;
1509 }
1510 let mut resources = serde_json::Map::new();
1511 resources.insert(resource, serde_json::Value::String("*".into()));
1512 Some(serde_json::json!({
1513 "effect": "allow",
1514 "resources": serde_json::Value::Object(resources),
1515 "permission_groups": groups,
1516 }))
1517 };
1518
1519 let policies: Vec<serde_json::Value> = [
1520 block(
1521 format!("com.cloudflare.api.account.{account_id}"),
1522 ids_for(GrantScope::Account),
1523 ),
1524 block(
1525 format!("com.cloudflare.api.account.zone.{zone_id}"),
1526 ids_for(GrantScope::Zone),
1527 ),
1528 ]
1529 .into_iter()
1530 .flatten()
1531 .collect();
1532
1533 serde_json::json!({ "name": token_name, "policies": policies })
1534}
1535
1536fn norm_dns(s: &str) -> String {
1541 s.trim().trim_end_matches('.').to_ascii_lowercase()
1542}
1543
1544fn best_zone_for<'a>(hostname: &str, zones: &'a [(String, String)]) -> Option<&'a str> {
1548 let h = norm_dns(hostname);
1549 zones
1550 .iter()
1551 .filter(|(_, name)| {
1552 let z = norm_dns(name);
1553 h == z || h.ends_with(&format!(".{z}"))
1554 })
1555 .max_by_key(|(_, name)| name.len())
1556 .map(|(id, _)| id.as_str())
1557}
1558
1559fn classify_tunnel_drift(
1563 hostname: &str,
1564 expected_target: &str,
1565 live: &[CfDnsRecord],
1566) -> (TunnelDriftState, Option<String>) {
1567 let h = norm_dns(hostname);
1568 let matching: Vec<&CfDnsRecord> = live.iter().filter(|r| norm_dns(&r.name) == h).collect();
1569 if matching.is_empty() {
1570 return (TunnelDriftState::Missing, None);
1571 }
1572 let want = norm_dns(expected_target);
1573 if matching.iter().any(|r| norm_dns(&r.content) == want) {
1574 return (TunnelDriftState::Synced, None);
1575 }
1576 (
1577 TunnelDriftState::Mismatch,
1578 Some(matching[0].content.clone()),
1579 )
1580}
1581
1582#[cfg(test)]
1583mod tests {
1584 use super::*;
1585
1586 fn rec(name: &str, content: &str) -> CfDnsRecord {
1587 CfDnsRecord {
1588 name: name.into(),
1589 content: content.into(),
1590 }
1591 }
1592
1593 #[test]
1594 fn drift_synced_when_live_matches() {
1595 let live = vec![rec("yubaba.yah.dev", "9e4d.cfargotunnel.com")];
1596 let (state, target) =
1597 classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
1598 assert_eq!(state, TunnelDriftState::Synced);
1599 assert!(target.is_none());
1600 }
1601
1602 #[test]
1603 fn drift_missing_when_no_matching_record() {
1604 let live = vec![rec("other.yah.dev", "x.cfargotunnel.com")];
1605 let (state, target) =
1606 classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
1607 assert_eq!(state, TunnelDriftState::Missing);
1608 assert!(target.is_none());
1609 }
1610
1611 #[test]
1612 fn drift_missing_when_zone_empty() {
1613 let (state, _) = classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &[]);
1614 assert_eq!(state, TunnelDriftState::Missing);
1615 }
1616
1617 #[test]
1618 fn drift_mismatch_surfaces_live_target() {
1619 let live = vec![rec("yubaba.yah.dev", "stale.cfargotunnel.com")];
1620 let (state, target) =
1621 classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
1622 assert_eq!(state, TunnelDriftState::Mismatch);
1623 assert_eq!(target.as_deref(), Some("stale.cfargotunnel.com"));
1624 }
1625
1626 #[test]
1627 fn drift_normalises_trailing_dot_and_case() {
1628 let live = vec![rec("Yubaba.YAH.dev.", "9E4D.cfargotunnel.com.")];
1631 let (state, _) = classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
1632 assert_eq!(state, TunnelDriftState::Synced);
1633 }
1634
1635 #[test]
1636 fn best_zone_picks_longest_suffix() {
1637 let zones = vec![
1638 ("z_apex".to_string(), "dev".to_string()),
1639 ("z_zone".to_string(), "yah.dev".to_string()),
1640 ];
1641 assert_eq!(best_zone_for("yubaba.yah.dev", &zones), Some("z_zone"));
1642 assert_eq!(best_zone_for("yah.dev", &zones), Some("z_zone"));
1643 }
1644
1645 #[test]
1646 fn best_zone_none_when_no_suffix_covers() {
1647 let zones = vec![("z1".to_string(), "yah.dev".to_string())];
1648 assert_eq!(best_zone_for("example.com", &zones), None);
1649 assert_eq!(best_zone_for("notyah.dev", &zones), None);
1652 }
1653
1654 #[test]
1655 fn best_zone_apex_self_match() {
1656 let zones = vec![("z1".to_string(), "yah.dev".to_string())];
1657 assert_eq!(best_zone_for("yah.dev", &zones), Some("z1"));
1658 }
1659
1660 #[test]
1661 fn token_body_splits_scopes_and_resolves_ids() {
1662 use std::collections::BTreeMap;
1663 let mut catalog = BTreeMap::new();
1666 catalog.insert("Zone Read".to_string(), "CATALOG_ZONE_READ".to_string());
1667
1668 let body = build_token_body("t", "ACCT", "ZONE", MESOFACT_STATIC_GRANTS, &catalog);
1669 let policies = body["policies"].as_array().unwrap();
1670 assert_eq!(policies.len(), 2, "one account block + one zone block");
1671
1672 let acct = &policies[0];
1673 assert_eq!(acct["resources"]["com.cloudflare.api.account.ACCT"], "*");
1674 let acct_ids: Vec<&str> = acct["permission_groups"]
1675 .as_array()
1676 .unwrap()
1677 .iter()
1678 .map(|g| g["id"].as_str().unwrap())
1679 .collect();
1680 assert!(acct_ids.contains(&"c1fde68c7bcc44588cbb6ddbc16d6480")); assert!(acct_ids.contains(&"bf7481a1826f439697cb59a20b22293e")); assert!(acct_ids.contains(&"e086da7e2179491d91ee5f35b3ca210a")); let zone = &policies[1];
1685 assert_eq!(
1686 zone["resources"]["com.cloudflare.api.account.zone.ZONE"],
1687 "*"
1688 );
1689 let zone_ids: Vec<&str> = zone["permission_groups"]
1690 .as_array()
1691 .unwrap()
1692 .iter()
1693 .map(|g| g["id"].as_str().unwrap())
1694 .collect();
1695 assert!(
1696 zone_ids.contains(&"CATALOG_ZONE_READ"),
1697 "catalog id wins over fallback"
1698 );
1699 assert!(zone_ids.contains(&"0ac90a90249747bca6b047d97f0803e9")); assert!(zone_ids.contains(&"28f4b596e7d643029c524985477ae49a")); assert!(zone_ids.contains(&"e17beae8b8cb423a99b1730f21238bed")); }
1703
1704 #[test]
1705 fn token_body_omits_empty_scope_block() {
1706 use std::collections::BTreeMap;
1707 let only_zone = &[TokenGrant {
1708 group_name: "Zone Read",
1709 scope: GrantScope::Zone,
1710 fallback_id: "ZR",
1711 }];
1712 let body = build_token_body("t", "A", "Z", only_zone, &BTreeMap::new());
1713 let policies = body["policies"].as_array().unwrap();
1714 assert_eq!(policies.len(), 1, "no account block when no account grants");
1715 assert_eq!(
1716 policies[0]["resources"]["com.cloudflare.api.account.zone.Z"],
1717 "*"
1718 );
1719 }
1720
1721 fn extract_metadata_json(body: &[u8]) -> serde_json::Value {
1723 let s = std::str::from_utf8(body).expect("multipart body is utf-8 for these tests");
1724 let (_, after) = s
1725 .split_once("name=\"metadata\"")
1726 .expect("metadata part present");
1727 let (_, after) = after
1728 .split_once("\r\n\r\n")
1729 .expect("metadata body delimited");
1730 let (json, _) = after
1731 .split_once("\r\n--")
1732 .expect("metadata terminated by boundary");
1733 serde_json::from_str(json).expect("metadata JSON parses")
1734 }
1735
1736 #[test]
1737 fn multipart_includes_r2_bucket_binding_metadata() {
1738 let bindings = [WorkerBinding::R2Bucket {
1739 name: "CACHE",
1740 bucket_name: "yah-cr-cache",
1741 }];
1742 let (content_type, body) = build_worker_multipart("export default {}", &bindings);
1743
1744 assert!(
1745 content_type.starts_with("multipart/form-data; boundary="),
1746 "content-type advertises multipart with boundary: got {content_type}",
1747 );
1748
1749 let metadata = extract_metadata_json(&body);
1750 assert_eq!(metadata["main_module"], "worker.js");
1751 let bindings = metadata["bindings"].as_array().expect("bindings array");
1752 assert_eq!(bindings.len(), 1);
1753 assert_eq!(bindings[0]["type"], "r2_bucket");
1754 assert_eq!(bindings[0]["name"], "CACHE");
1755 assert_eq!(bindings[0]["bucket_name"], "yah-cr-cache");
1756 }
1757
1758 #[test]
1759 fn multipart_mixes_plain_text_and_r2_bindings() {
1760 let bindings = [
1761 WorkerBinding::PlainText {
1762 name: "MODE",
1763 text: "cache",
1764 },
1765 WorkerBinding::R2Bucket {
1766 name: "CACHE",
1767 bucket_name: "yah-cr-cache",
1768 },
1769 ];
1770 let (_, body) = build_worker_multipart("export default {}", &bindings);
1771
1772 let metadata = extract_metadata_json(&body);
1773 let entries = metadata["bindings"].as_array().unwrap();
1774 assert_eq!(entries.len(), 2);
1775 assert_eq!(entries[0]["type"], "plain_text");
1776 assert_eq!(entries[0]["text"], "cache");
1777 assert_eq!(entries[1]["type"], "r2_bucket");
1778 assert_eq!(entries[1]["bucket_name"], "yah-cr-cache");
1779 }
1780}