Skip to main content

cloud/provider/
cloudflare.rs

1//! Cloudflare management API client — accounts, tunnels, R2 buckets, DNS.
2//!
3//! Shared by the desktop Tauri commands, the CLI, and the reconciler.
4//! Callers resolve the API token themselves (keychain, env, vault) and pass
5//! it to [`CloudflareClient::new`]; this module never reads credentials.
6//!
7//! Two distinct API surfaces:
8//! - **Management API** (this module) — accounts, R2-bucket CRUD, tunnels, DNS, cache purge.
9//! - **R2 object publish** — S3 SigV4, lives in `reconciler::r2_publish` and
10//!   reuses the existing `s3_sign` helper. Not part of this module.
11//!
12//! @yah:ticket(R320-F12, "yah cloud: mint scoped Cloudflare API token from a policy template (one-command onboarding for a new CF account)")
13//! @yah:assignee(agent:claude)
14//! @yah:at(2026-05-26T14:57:46Z)
15//! @yah:status(review)
16//! @yah:parent(R320)
17//! @arch:see(.yah/docs/working/W074-cloudflare-infra-provider.md)
18//! @yah:next("Resolve permission-group UUIDs at runtime from GET /accounts/{id}/tokens/permission_groups — CF references groups by ID not name; the policy template stores names and resolves to IDs at create time")
19//! @yah:next("Build the minimal mesofact-static policy: account-scoped block (Account Settings:Read + Workers R2 Storage:Edit) + zone-scoped block (Zone:Read + Transform Rules:Edit + Cache Purge). Keep scope blocks separate — mixing account- and zone-scoped groups in one block fails token-create validation")
20//! @yah:next("POST /accounts/{id}/tokens to create an ACCOUNT-OWNED token (DECIDED by user: survives the creating user, correct for a shared tool credential). Print the secret once and offer to store it in the cloudflare-api-token keystore slot")
21//! @yah:next("Expose as 'yah cloud cf token create --account <id> --zone <name>' so a fresh CF account is one command. The policy template is the checked-in artifact CF won't let you save dashboard-side")
22//! @yah:gotcha("Bootstrap credential needs ONLY API Tokens:Edit — VALIDATED live: the created token is bounded by the account's own access, not the bootstrap token's perms, so a minimal bootstrap suffices. Brand-new account still needs one such token minted manually (or Global Key) once.")
23//! @yah:gotcha("TRAP — two Transform Rules groups: 'Transform Rules Write' (ae16e88b…) is ACCOUNT-scoped = WRONG. upsert_index_rewrite hits /zones/{id}/rulesets so it needs ZONE-scoped 'Zone Transform Rules Write'. Zone Read + Cache Purge are also zone-scoped, not account-scoped.")
24//! @yah:gotcha("Permission-group IDs (global constants, validated 2026-05-26): acct-scoped Account Settings Read=c1fde68c7bcc44588cbb6ddbc16d6480, Workers R2 Storage Write=bf7481a1826f439697cb59a20b22293e; zone-scoped Zone Read=c8fed203ed3043cba015a93ad1616f1f, Zone Transform Rules Write=0ac90a90249747bca6b047d97f0803e9, Cache Purge=e17beae8b8cb423a99b1730f21238bed")
25//! @yah:gotcha("account-owned token verify = GET /accounts/{id}/tokens/verify; /user/tokens/verify returns success:false for account-owned tokens (NOT a failure — seen in live test)")
26//! @yah:assumes("VALIDATED live 2026-05-26: POST /accounts/{id}/tokens with the 2-block policy created 'yah-mesofact-static-yahdev' which lists R2 buckets successfully. R2 object upload still uses SEPARATE S3 keys (SigV4), not this management token.")
27//! @yah:handoff("Implemented + verified live end-to-end. CloudflareClient gained create_account_token (cloudflare.rs) + list_permission_group_ids; MESOFACT_STATIC_GRANTS const carries the 5 validated permission-group IDs (account: Account Settings Read, Workers R2 Storage Write; zone: Zone Read, Zone Transform Rules Write, Cache Purge) with group-name→id resolved against the live catalog at create time, baked-in IDs as fallback. Pure helpers build_token_body/resolve_grant_id split account- vs zone-scoped policy blocks (2 unit tests). CLI: 'yah cloud cf token create --zone <name> [--account <id>] [--store-slot <slot>] [--name <n>] [--bootstrap-slot <slot>]' in cloud.rs handle_cf_token_create — resolves account from --account or .yah/infra/providers/cloudflare.toml, resolves zone name→id, mints account-owned token via POST /accounts/{id}/tokens, stores to keystore (fail-fast on occupied slot, BEFORE minting) or prints once. Bootstrap defaults to cloudflare-api-token slot / $CLOUDFLARE_API_TOKEN.")
28//! @yah:next("Follow-ups (not in scope): 'yah cloud cf token revoke <id>' (DELETE endpoint already proven), and broader grant presets beyond MESOFACT_STATIC_GRANTS (e.g. + Tunnel:Read/DNS:Read for the Infra panel)")
29//! @yah:verify("cargo test -p cloud --lib — 207 passed (incl. token_body_splits_scopes_and_resolves_ids, token_body_omits_empty_scope_block)")
30//! @yah:verify("Live E2E 2026-05-26: 'yah cloud cf token create --zone yah.dev --store-slot cf-clitest' minted account-owned token, minted token listed R2 buckets successfully, then DELETE /accounts/{id}/tokens/{id} revoked it cleanly")
31//!
32//!
33//! @yah:ticket(R324-F5, "Tunnel connection status + uptime in the Tunnels table")
34//! @yah:assignee(agent:claude)
35//! @yah:at(2026-05-26T15:33:48Z)
36//! @yah:status(review)
37//! @yah:phase(P2)
38//! @yah:parent(R324)
39//! @yah:handoff("Tunnel connection state fully wired end-to-end. Rust: TunnelConnState enum (Active/Inactive/Degraded/Unknown) added to cloud crate; CfTunnel wire struct extended with status + conns_active_at; TunnelMeta private struct carries enriched data; list_tunnels_meta() fetches status in one pass; TunnelDriftRow gained conn_state + conn_since (RFC3339 optional); tunnel_dns_drift() refactored to walk accounts→tunnels_meta→configs in a single list_accounts pass instead of calling tunnel_dns_records() separately (avoids duplicate API round-trip). Re-exported TunnelConnState through provider/mod.rs + cloud/src/lib.rs + desktop cloudflare.rs. TS: TunnelConnState type + connState/connSince on TunnelDriftRow in types.ts. UI: DriftPill now shows 'connected · <uptime>' (pulsing forest dot) for synced+active, 'inactive' neutral pill for synced+inactive, 'degraded' warn for synced+degraded. fmtUptime() formats ISO → '<1m'/'45m'/'12h'/'3d'. TunnelsSection right label changed to 'N active' (or 'N active · M drift'). collectCfSlots test drift() fixture updated with connState: 'unknown'.")
40//! @yah:verify("cargo test -p cloud --lib cloudflare — 15 passed (incl. 8 drift unit tests)")
41//! @yah:verify("cargo check -p desktop — clean (TunnelConnState re-exported + TunnelDriftRow extended)")
42//! @yah:verify("cd packages/yah/ui && bun test src/components/infra/CloudflarePanel.collectCfSlots.test.ts — 8 pass")
43//! @yah:verify("cd packages/yah/ui && bun run typecheck — no errors in infra/ or env/ (pre-existing failures elsewhere unchanged)")
44//!
45//! @yah:ticket(R324-F6, "R2 bucket size + object count + region in Accounts section")
46//! @yah:assignee(agent:claude)
47//! @yah:at(2026-05-26T15:33:49Z)
48//! @yah:status(review)
49//! @yah:phase(P2)
50//! @yah:parent(R324)
51//! @yah:next("Extend R2BucketInfo + list_r2_buckets with size/object-count/region; design AccountBlock shows '412 MB · 142 obj'. Needs extra R2 (or S3 list) calls per bucket.")
52//! @yah:next("Render the new fields in CloudflarePanel AccountBlock bucket rows.")
53//! @yah:handoff("Added location + creation_date to R2BucketInfo (Rust struct + TS interface). Both fields come from the existing GET /accounts/{id}/r2/buckets list call — no extra round-trips. fmtR2Location() maps CF location codes (WEUR/EEUR/WNAM/ENAM/APAC) to short labels. AccountBlock bucket cards now show the region badge on the right when present. Note: bucket size + object count are not available from the CF management API without per-bucket S3 calls (paginated ListObjectsV2 + S3 credentials); deferred to a future ticket.")
54//! @yah:verify("cargo check -p cloud -p desktop — clean (R2BucketInfo extended, deserialized from BucketEntry, re-exported unchanged)")
55//! @yah:verify("cd packages/yah/ui && bun run typecheck — no new errors in infra/ or env/")
56//!
57//!
58//! @yah:ticket(R419-F1, "Extend deploy_worker_script for r2_bucket bindings")
59//! @yah:assignee(agent:claude)
60//! @yah:at(2026-06-03T08:02:38Z)
61//! @yah:status(review)
62//! @yah:parent(R419)
63//! @yah:handoff("Widened deploy_worker_script + build_worker_multipart to typed bindings. New pub enum WorkerBinding<'a> { PlainText { name, text }, R2Bucket { name, bucket_name } } encodes both shapes; multipart metadata.bindings now emits the matching CF wire JSON. Single existing caller (mesofact_static.rs:505) maps its (String,String) plain_text vec into WorkerBinding::PlainText refs — runtime behavior unchanged. F2's CloudflareWorkerReconciler now has the surface it needs: pass [WorkerBinding::R2Bucket { name: binding_name_from_workload_toml, bucket_name: from_mirror_providers_cache }, ...].")
64//! @yah:verify("cargo check -p cloud --lib — clean")
65//! @yah:verify("cargo test -p cloud --lib provider::cloudflare — 12 passed, incl. multipart_includes_r2_bucket_binding_metadata + multipart_mixes_plain_text_and_r2_bindings")
66//! @yah:next("F2 pickup: bind via WorkerBinding::R2Bucket { name: <workload.toml [[bindings]].name>, bucket_name: <mirror providers.cache.bucket> } after fail-fast on workload<->mirror binding-name drift.")
67
68use anyhow::{anyhow, Result};
69use serde::{Deserialize, Serialize};
70
71const CF_API: &str = "https://api.cloudflare.com/client/v4";
72
73// ---------- internal wire types ----------
74
75#[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    /// Cloudflare's live status: `"active"` | `"inactive"` | `"degraded"` | `"unknown"`.
94    #[serde(default)]
95    status: Option<String>,
96    /// RFC 3339 timestamp — when the tunnel last became active. `null` when
97    /// never active or the field is absent.
98    #[serde(default)]
99    conns_active_at: Option<String>,
100}
101
102/// Enriched tunnel record used internally when walking accounts for drift.
103struct 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/// One live DNS record from `GET /zones/{id}/dns_records`. We only read the
126/// name + content (the target); record type and proxy flags are ignored —
127/// drift is decided by whether *some* record routes the hostname to the
128/// expected tunnel target.
129#[derive(Debug, Clone, Deserialize)]
130struct CfDnsRecord {
131    name: String,
132    content: String,
133}
134
135/// One live DNS record with everything a reconciler needs to decide what to
136/// change — R859-F1, the read side of [`CloudflareClient::list_dns_records`].
137///
138/// Distinct from the private [`CfDnsRecord`] above, which deliberately reads
139/// only name + content because tunnel-drift detection asks a narrower
140/// question. Reconciling a record set needs the `id` (to delete one member of
141/// a multi-valued RRset) and `proxied` (an orange-clouded record at an apex
142/// that should be grey is drift, even when the content is right).
143#[derive(Debug, Clone, PartialEq, Eq, Deserialize)]
144pub struct DnsRecordDetail {
145    pub id: String,
146    pub name: String,
147    #[serde(rename = "type")]
148    pub record_type: String,
149    pub content: String,
150    #[serde(default)]
151    pub ttl: u32,
152    #[serde(default)]
153    pub proxied: bool,
154}
155
156// ---------- public output types ----------
157
158/// A Cloudflare account the API token can access.
159#[derive(Debug, Clone, Serialize, Deserialize)]
160#[serde(rename_all = "camelCase")]
161pub struct CfAccountInfo {
162    pub id: String,
163    pub name: String,
164}
165
166/// One CNAME record a user needs to create in their external DNS registrar
167/// to route a hostname through a Cloudflare Tunnel.
168#[derive(Debug, Clone, Serialize, Deserialize)]
169#[serde(rename_all = "camelCase")]
170pub struct TunnelDnsRecord {
171    /// Human-readable tunnel name as entered in Cloudflare.
172    pub tunnel_name: String,
173    /// Hostname from the tunnel ingress rule (e.g. `yah.example.com`).
174    pub hostname: String,
175    /// CNAME target to enter in the registrar: `{tunnel_id}.cfargotunnel.com`.
176    pub cname_target: String,
177}
178
179/// Result of creating a Cloudflare Named Tunnel.
180#[derive(Debug, Clone, Serialize, Deserialize)]
181#[serde(rename_all = "camelCase")]
182pub struct CreateTunnelResult {
183    pub tunnel_id: String,
184    pub tunnel_name: String,
185    /// JWT token for `cloudflared tunnel run --token <TOKEN>`.
186    pub connector_token: String,
187    /// CNAME target: `{tunnel_id}.cfargotunnel.com`.
188    pub cname_target: String,
189}
190
191/// Result of creating a Cloudflare R2 bucket.
192#[derive(Debug, Clone, Serialize, Deserialize)]
193#[serde(rename_all = "camelCase")]
194pub struct CreateR2BucketResult {
195    pub name: String,
196    /// S3-compatible endpoint for object operations against this account.
197    pub endpoint: String,
198}
199
200/// Result of deploying a Cloudflare Worker script.
201#[derive(Debug, Clone, Serialize, Deserialize)]
202#[serde(rename_all = "camelCase")]
203pub struct WorkerDeployResult {
204    pub id: String,
205    pub etag: Option<String>,
206}
207
208/// One binding to inject into a Worker's `env` at deploy time.
209///
210/// Each variant maps to a `{"type": …}` entry under `metadata.bindings` in the
211/// Workers upload multipart body. R2 buckets are bind-only — the Worker reads
212/// the bucket via `env.<name>` and no script-side change is needed beyond the
213/// metadata declaration.
214#[derive(Debug, Clone, Copy)]
215pub enum WorkerBinding<'a> {
216    /// `{"type":"plain_text","name":…,"text":…}` — runtime config string.
217    PlainText { name: &'a str, text: &'a str },
218    /// `{"type":"r2_bucket","name":…,"bucket_name":…}` — R2 bucket reference.
219    R2Bucket { name: &'a str, bucket_name: &'a str },
220}
221
222/// R2 bucket information from the list endpoint.
223#[derive(Debug, Clone, Serialize, Deserialize)]
224#[serde(rename_all = "camelCase")]
225pub struct R2BucketInfo {
226    pub name: String,
227    /// CF location hint, e.g. `"WEUR"`, `"ENAM"`, `"APAC"`. `None` when the
228    /// bucket was created without specifying a hint.
229    #[serde(default, skip_serializing_if = "Option::is_none")]
230    pub location: Option<String>,
231    /// ISO 8601 bucket creation timestamp.
232    #[serde(default, skip_serializing_if = "Option::is_none")]
233    pub creation_date: Option<String>,
234}
235
236/// One R2 custom-domain binding from
237/// `GET /accounts/{id}/r2/buckets/{bucket}/domains/custom`.
238///
239/// The CF response carries nested status + min_tls fields we don't act on
240/// today; the reconciler only needs the hostname + enabled flag to decide
241/// idempotency.
242#[derive(Debug, Clone, Serialize, Deserialize)]
243#[serde(rename_all = "camelCase")]
244pub struct R2CustomDomain {
245    pub domain: String,
246    #[serde(default)]
247    pub enabled: bool,
248}
249
250/// Live connection state of a Cloudflare Tunnel connector.
251///
252/// Derived from the `status` field returned by
253/// `GET /accounts/{id}/cfd_tunnel?is_deleted=false`. Falls back to
254/// [`TunnelConnState::Unknown`] when the field is absent or unrecognised.
255#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
256#[serde(rename_all = "kebab-case")]
257pub enum TunnelConnState {
258    /// At least one healthy `cloudflared` connector is active.
259    Active,
260    /// No connectors are running.
261    Inactive,
262    /// Connectors are running but unhealthy.
263    Degraded,
264    /// Status couldn't be determined (field absent or unrecognised value).
265    Unknown,
266}
267
268impl TunnelConnState {
269    fn from_cf_status(s: &str) -> Self {
270        match s {
271            "active" | "healthy" => Self::Active,
272            "inactive" | "down" => Self::Inactive,
273            "degraded" | "unhealthy" => Self::Degraded,
274            _ => Self::Unknown,
275        }
276    }
277}
278
279/// Drift verdict for one tunnel ingress hostname: does live Cloudflare DNS
280/// route it to the tunnel's CNAME target?
281///
282/// Maps onto the designed Tunnels table (`infra-cloudflare.jsx`): `Synced`
283/// renders as a healthy pill, everything else as a `drift` pill (`Missing`
284/// is the "DNS record missing" case the design calls out).
285#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
286#[serde(rename_all = "kebab-case")]
287pub enum TunnelDriftState {
288    /// A live DNS record points the hostname at the expected tunnel target.
289    Synced,
290    /// No live DNS record exists for the hostname — the tunnel can't route to it.
291    Missing,
292    /// A live record exists but points somewhere other than the tunnel target.
293    Mismatch,
294    /// The hostname's zone couldn't be resolved or read (token lacks
295    /// `Zone: Read` / `DNS: Read`, or the apex isn't a zone in any accessible
296    /// account) — drift indeterminate, not a failure.
297    ZoneUnknown,
298}
299
300/// One row of the tunnel DNS-drift report: a tunnel ingress hostname paired
301/// with whether live Cloudflare DNS routes it to the tunnel, plus the live
302/// connector connection state.
303#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
304#[serde(rename_all = "camelCase")]
305pub struct TunnelDriftRow {
306    /// Human-readable tunnel name as entered in Cloudflare.
307    pub tunnel_name: String,
308    /// Hostname from the tunnel's ingress rule, e.g. `yubaba.yah.dev`.
309    pub hostname: String,
310    /// CNAME target the DNS record should point at: `{tunnel_id}.cfargotunnel.com`.
311    pub expected_target: String,
312    pub state: TunnelDriftState,
313    /// For [`TunnelDriftState::Mismatch`], the target the live record actually
314    /// points at. `None` for every other state.
315    #[serde(default, skip_serializing_if = "Option::is_none")]
316    pub live_target: Option<String>,
317    /// Live connector state from `GET /cfd_tunnel`. [`TunnelConnState::Unknown`]
318    /// when the field was absent or unrecognised.
319    pub conn_state: TunnelConnState,
320    /// RFC 3339 `conns_active_at` timestamp — when this tunnel last became
321    /// active. `None` when the tunnel has never connected or the field was
322    /// absent.
323    #[serde(default, skip_serializing_if = "Option::is_none")]
324    pub conn_since: Option<String>,
325}
326
327// ---------- API-token provisioning (R320-F12) ----------
328
329/// Resource scope a permission group applies at when building a token policy.
330#[derive(Debug, Clone, Copy, PartialEq, Eq)]
331pub enum GrantScope {
332    /// Granted on the whole account (`com.cloudflare.api.account.<id>`).
333    Account,
334    /// Granted on a single zone (`com.cloudflare.api.account.zone.<id>`).
335    Zone,
336}
337
338/// One permission to bake into a minted token: a Cloudflare permission-group
339/// display name, the scope it applies at, and a validated fallback ID.
340///
341/// Names are resolved to IDs against the live permission-groups catalog at
342/// create time ([`CloudflareClient::list_permission_group_ids`]); `fallback_id`
343/// is used only when the catalog lookup can't resolve the name. The IDs are
344/// global Cloudflare constants (validated 2026-05-26), not per-account.
345#[derive(Debug, Clone, Copy)]
346pub struct TokenGrant {
347    pub group_name: &'static str,
348    pub scope: GrantScope,
349    pub fallback_id: &'static str,
350}
351
352/// Minimal permission set for a mesofact-static publish token: see the account,
353/// list/create R2 buckets, deploy Worker scripts, resolve the zone, manage
354/// Worker routes + the index-rewrite Transform Rule, and purge the CDN cache.
355///
356/// Account-scoped: `Account Settings Read`, `Workers R2 Storage Write`,
357/// `Workers Scripts Write`.
358/// Zone-scoped: `Zone Read`, `Zone Transform Rules Write`,
359/// `Workers Routes Write`, `Cache Purge`.
360///
361/// NB the Transform Rules group is the *zone-scoped* `Zone Transform Rules
362/// Write`, not the account-scoped `Transform Rules Write` —
363/// [`CloudflareClient::upsert_index_rewrite`] hits `/zones/{id}/rulesets`.
364/// Fallback IDs sourced from the global permission-groups catalog
365/// (validated 2026-05-26); catalog resolution at create-time takes precedence.
366pub const MESOFACT_STATIC_GRANTS: &[TokenGrant] = &[
367    TokenGrant {
368        group_name: "Account Settings Read",
369        scope: GrantScope::Account,
370        fallback_id: "c1fde68c7bcc44588cbb6ddbc16d6480",
371    },
372    TokenGrant {
373        group_name: "Workers R2 Storage Write",
374        scope: GrantScope::Account,
375        fallback_id: "bf7481a1826f439697cb59a20b22293e",
376    },
377    TokenGrant {
378        group_name: "Workers Scripts Write",
379        scope: GrantScope::Account,
380        fallback_id: "e086da7e2179491d91ee5f35b3ca210a",
381    },
382    TokenGrant {
383        group_name: "Zone Read",
384        scope: GrantScope::Zone,
385        fallback_id: "c8fed203ed3043cba015a93ad1616f1f",
386    },
387    TokenGrant {
388        group_name: "Zone Transform Rules Write",
389        scope: GrantScope::Zone,
390        fallback_id: "0ac90a90249747bca6b047d97f0803e9",
391    },
392    TokenGrant {
393        group_name: "Workers Routes Write",
394        scope: GrantScope::Zone,
395        fallback_id: "28f4b596e7d643029c524985477ae49a",
396    },
397    TokenGrant {
398        group_name: "Cache Purge",
399        scope: GrantScope::Zone,
400        fallback_id: "e17beae8b8cb423a99b1730f21238bed",
401    },
402    // R859-F1: deploy_domain_passway (the sovereign-apex A-record reconciler)
403    // is the first production consumer of the dns.* envoy verbs against this
404    // token, and it both lists and upserts/prunes A records — hence both
405    // Read and Write, not just Read as the R324-F5 @yah:next follow-up noted.
406    TokenGrant {
407        group_name: "DNS Read",
408        scope: GrantScope::Zone,
409        fallback_id: "82e64a83756745bbbb1c9c2701bf816b",
410    },
411    TokenGrant {
412        group_name: "DNS Write",
413        scope: GrantScope::Zone,
414        fallback_id: "4755a26eedb94da69e1066d98aa820be",
415    },
416];
417
418/// Result of minting an account-owned API token. `value` is the secret and is
419/// returned by Cloudflare exactly once — store it immediately.
420#[derive(Debug, Clone, Serialize, Deserialize)]
421#[serde(rename_all = "camelCase")]
422pub struct CreateTokenResult {
423    pub id: String,
424    pub name: String,
425    /// The token secret. Shown once by Cloudflare; never retrievable again.
426    pub value: String,
427}
428
429// ---------- client ----------
430
431/// Cloudflare management API client.
432///
433/// Construct with [`CloudflareClient::new`] passing a pre-resolved API token.
434/// The token scope required per method is noted on each method.
435pub struct CloudflareClient {
436    token: String,
437    http: reqwest::Client,
438}
439
440impl CloudflareClient {
441    /// Create a client for the given API token.
442    pub fn new(token: String) -> Self {
443        Self {
444            token,
445            http: reqwest::Client::new(),
446        }
447    }
448
449    /// List accounts the token can access.
450    /// Requires: `Account: Read`.
451    pub async fn list_accounts(&self) -> Result<Vec<CfAccountInfo>> {
452        #[derive(Deserialize)]
453        struct Entry {
454            id: String,
455            name: String,
456        }
457        let resp: CfPage<Entry> = self.cf_get("/accounts").await?;
458        self.ok(&resp.success, &resp.errors)?;
459        Ok(resp
460            .result
461            .unwrap_or_default()
462            .into_iter()
463            .map(|a| CfAccountInfo {
464                id: a.id,
465                name: a.name,
466            })
467            .collect())
468    }
469
470    /// List non-deleted Cloudflare Tunnels in `account_id` with connection state.
471    /// Requires: `Cloudflare Tunnel: Read`.
472    async fn list_tunnels_meta(&self, account_id: &str) -> Result<Vec<TunnelMeta>> {
473        let resp: CfPage<CfTunnel> = self
474            .cf_get(&format!(
475                "/accounts/{account_id}/cfd_tunnel?is_deleted=false"
476            ))
477            .await?;
478        self.ok(&resp.success, &resp.errors)?;
479        Ok(resp
480            .result
481            .unwrap_or_default()
482            .into_iter()
483            .map(|t| TunnelMeta {
484                conn_state: t
485                    .status
486                    .as_deref()
487                    .map(TunnelConnState::from_cf_status)
488                    .unwrap_or(TunnelConnState::Unknown),
489                conn_since: t.conns_active_at,
490                id: t.id,
491                name: t.name,
492            })
493            .collect())
494    }
495
496    /// List non-deleted Cloudflare Tunnels in `account_id` as `(id, name)` pairs.
497    /// Requires: `Cloudflare Tunnel: Read`.
498    pub async fn list_tunnels(&self, account_id: &str) -> Result<Vec<(String, String)>> {
499        Ok(self
500            .list_tunnels_meta(account_id)
501            .await?
502            .into_iter()
503            .map(|m| (m.id, m.name))
504            .collect())
505    }
506
507    /// Collect CNAME records for all tunnels across all accessible accounts.
508    ///
509    /// Walks accounts → tunnels → ingress configurations. Returns an empty
510    /// vec when the token has no tunnels or no configured ingress hostnames.
511    pub async fn tunnel_dns_records(&self) -> Result<Vec<TunnelDnsRecord>> {
512        let mut records = Vec::new();
513        let accounts = self.list_accounts().await?;
514
515        for account in &accounts {
516            let tunnels = self.list_tunnels(&account.id).await?;
517            for (tunnel_id, tunnel_name) in &tunnels {
518                let cname_target = format!("{tunnel_id}.cfargotunnel.com");
519                let path = format!(
520                    "/accounts/{}/cfd_tunnel/{tunnel_id}/configurations",
521                    account.id
522                );
523                let config_resp: CfSingle<CfTunnelConfig> = self.cf_get(&path).await?;
524                self.ok(&config_resp.success, &config_resp.errors)?;
525                let ingress = config_resp
526                    .result
527                    .and_then(|c| c.config)
528                    .and_then(|c| c.ingress)
529                    .unwrap_or_default();
530                for rule in ingress {
531                    if let Some(hostname) = rule.hostname.filter(|h| !h.is_empty()) {
532                        records.push(TunnelDnsRecord {
533                            tunnel_name: tunnel_name.clone(),
534                            hostname,
535                            cname_target: cname_target.clone(),
536                        });
537                    }
538                }
539            }
540        }
541        Ok(records)
542    }
543
544    /// Compute DNS drift for every tunnel ingress hostname, enriched with live
545    /// connector connection state.
546    ///
547    /// The declared side is the tunnel ingress config; the live side is the
548    /// zone's DNS records. Each ingress hostname is classified:
549    /// [`TunnelDriftState::Synced`] when a record points at the tunnel's CNAME
550    /// target, `Missing` when none exists, `Mismatch` when one points elsewhere.
551    ///
552    /// Connection state (`conn_state` / `conn_since`) comes from the `status`
553    /// and `conns_active_at` fields on the tunnel list response — fetched in the
554    /// same pass as the ingress configs to avoid an extra `list_accounts` round-trip.
555    ///
556    /// Degrades gracefully — a hostname whose zone can't be resolved or read is
557    /// reported `ZoneUnknown` rather than failing the whole report. Returns an
558    /// empty vec when the token has no tunnels or no ingress hostnames.
559    ///
560    /// Requires: `Cloudflare Tunnel: Read`, `Zone: Read`, `DNS: Read`.
561    pub async fn tunnel_dns_drift(&self) -> Result<Vec<TunnelDriftRow>> {
562        // Walk accounts → tunnels (with live connection state) → ingress configs
563        // in one pass, collecting both declared DNS records and conn-state in a
564        // single list_accounts round-trip.
565        let accounts = self.list_accounts().await?;
566        let mut declared: Vec<TunnelDnsRecord> = Vec::new();
567        let mut conn_by_tunnel: std::collections::HashMap<
568            String,
569            (TunnelConnState, Option<String>),
570        > = Default::default();
571
572        for account in &accounts {
573            let metas = self.list_tunnels_meta(&account.id).await?;
574            for meta in &metas {
575                conn_by_tunnel.insert(
576                    meta.name.clone(),
577                    (meta.conn_state, meta.conn_since.clone()),
578                );
579                let cname_target = format!("{}.cfargotunnel.com", meta.id);
580                let path = format!(
581                    "/accounts/{}/cfd_tunnel/{}/configurations",
582                    account.id, meta.id
583                );
584                let config_resp: CfSingle<CfTunnelConfig> = self.cf_get(&path).await?;
585                self.ok(&config_resp.success, &config_resp.errors)?;
586                let ingress = config_resp
587                    .result
588                    .and_then(|c| c.config)
589                    .and_then(|c| c.ingress)
590                    .unwrap_or_default();
591                for rule in ingress {
592                    if let Some(hostname) = rule.hostname.filter(|h| !h.is_empty()) {
593                        declared.push(TunnelDnsRecord {
594                            tunnel_name: meta.name.clone(),
595                            hostname,
596                            cname_target: cname_target.clone(),
597                        });
598                    }
599                }
600            }
601        }
602
603        if declared.is_empty() {
604            return Ok(Vec::new());
605        }
606
607        // Resolve zones once. A permission failure leaves the set empty, so
608        // every hostname falls through to `ZoneUnknown` instead of erroring.
609        let zones = self.list_zones().await.unwrap_or_default();
610
611        let mut rows = Vec::with_capacity(declared.len());
612        for rec in declared {
613            let (conn_state, conn_since) = conn_by_tunnel
614                .remove(&rec.tunnel_name)
615                .unwrap_or((TunnelConnState::Unknown, None));
616            let (state, live_target) = match best_zone_for(&rec.hostname, &zones) {
617                None => (TunnelDriftState::ZoneUnknown, None),
618                Some(zone_id) => match self.dns_records_named(zone_id, &rec.hostname).await {
619                    Ok(live) => classify_tunnel_drift(&rec.hostname, &rec.cname_target, &live),
620                    Err(_) => (TunnelDriftState::ZoneUnknown, None),
621                },
622            };
623            rows.push(TunnelDriftRow {
624                tunnel_name: rec.tunnel_name,
625                hostname: rec.hostname,
626                expected_target: rec.cname_target,
627                state,
628                live_target,
629                conn_state,
630                conn_since,
631            });
632        }
633        Ok(rows)
634    }
635
636    /// List zones the token can read, as `(zone_id, zone_name)` pairs.
637    /// Requires: `Zone: Read`.
638    pub async fn list_zones(&self) -> Result<Vec<(String, String)>> {
639        #[derive(Deserialize)]
640        struct ZoneEntry {
641            id: String,
642            name: String,
643        }
644        let resp: CfPage<ZoneEntry> = self.cf_get("/zones?per_page=50").await?;
645        self.ok(&resp.success, &resp.errors)?;
646        Ok(resp
647            .result
648            .unwrap_or_default()
649            .into_iter()
650            .map(|z| (z.id, z.name))
651            .collect())
652    }
653
654    /// Fetch DNS records in `zone_id` whose name exactly matches `name`.
655    /// Requires: `DNS: Read`.
656    async fn dns_records_named(&self, zone_id: &str, name: &str) -> Result<Vec<CfDnsRecord>> {
657        let resp: CfPage<CfDnsRecord> = self
658            .cf_get(&format!("/zones/{zone_id}/dns_records?name={name}"))
659            .await?;
660        self.ok(&resp.success, &resp.errors)?;
661        Ok(resp.result.unwrap_or_default())
662    }
663
664    /// Create a new Named Tunnel under `account_id` and return the connector token.
665    /// Requires: `Cloudflare Tunnel: Edit`.
666    pub async fn create_tunnel(&self, account_id: &str, name: &str) -> Result<CreateTunnelResult> {
667        use base64::Engine as _;
668
669        let mut secret_bytes = [0u8; 32];
670        getrandom::getrandom(&mut secret_bytes)
671            .map_err(|e| anyhow!("generate tunnel secret: {e}"))?;
672        let tunnel_secret = base64::engine::general_purpose::STANDARD.encode(secret_bytes);
673
674        #[derive(Serialize)]
675        struct CreateBody<'a> {
676            name: &'a str,
677            tunnel_secret: String,
678        }
679        #[derive(Deserialize)]
680        struct CreatedTunnel {
681            id: String,
682            name: String,
683        }
684        let create_resp: CfSingle<CreatedTunnel> = self
685            .cf_post(
686                &format!("/accounts/{account_id}/cfd_tunnel"),
687                &CreateBody {
688                    name,
689                    tunnel_secret,
690                },
691            )
692            .await?;
693        self.ok(&create_resp.success, &create_resp.errors)?;
694        let created = create_resp
695            .result
696            .ok_or_else(|| anyhow!("tunnel create: no result in response"))?;
697
698        // Fetch the connector JWT.
699        #[derive(Deserialize)]
700        struct TokenResp {
701            success: bool,
702            result: Option<String>,
703            errors: Option<Vec<serde_json::Value>>,
704        }
705        let token_resp: TokenResp = self
706            .cf_get(&format!(
707                "/accounts/{account_id}/cfd_tunnel/{}/token",
708                created.id
709            ))
710            .await?;
711        self.ok(&token_resp.success, &token_resp.errors)?;
712        let connector_token = token_resp
713            .result
714            .ok_or_else(|| anyhow!("no connector token in response"))?;
715
716        Ok(CreateTunnelResult {
717            cname_target: format!("{}.cfargotunnel.com", created.id),
718            tunnel_id: created.id,
719            tunnel_name: created.name,
720            connector_token,
721        })
722    }
723
724    /// Read a tunnel's remotely-managed configuration body as raw JSON
725    /// (R594-F11).
726    ///
727    /// Returns the `result.config` object — the thing a PUT round-trips — or an
728    /// empty object when the tunnel has never been configured. Deliberately
729    /// untyped: the `ingress` list is the only key this crate owns, and every
730    /// sibling (`warp-routing`, `originRequest`, …) must survive a
731    /// read-modify-write untouched.
732    ///
733    /// Requires: `Cloudflare Tunnel: Read`.
734    pub async fn tunnel_configuration(
735        &self,
736        account_id: &str,
737        tunnel_id: &str,
738    ) -> Result<serde_json::Value> {
739        let path = format!("/accounts/{account_id}/cfd_tunnel/{tunnel_id}/configurations");
740        let resp: CfSingle<serde_json::Value> = self.cf_get(&path).await?;
741        self.ok(&resp.success, &resp.errors)?;
742        let config = resp
743            .result
744            .as_ref()
745            .and_then(|r| r.get("config"))
746            .cloned()
747            .unwrap_or_else(|| serde_json::json!({}));
748        // A tunnel configured with an explicit JSON `null` config reads back as
749        // Value::Null, which has no object to insert `ingress` into.
750        Ok(if config.is_object() {
751            config
752        } else {
753            serde_json::json!({})
754        })
755    }
756
757    /// Replace a tunnel's remotely-managed configuration (R594-F11).
758    ///
759    /// `config` is the whole config body, not a patch — Cloudflare replaces it
760    /// wholesale, which is why callers must GET-merge-PUT rather than PUT a
761    /// freshly-built list. See
762    /// [`reconciler::ingress::ensure_tunnel_ingress`](crate::reconciler::ingress::ensure_tunnel_ingress).
763    ///
764    /// Requires: `Cloudflare Tunnel: Edit`.
765    pub async fn put_tunnel_configuration(
766        &self,
767        account_id: &str,
768        tunnel_id: &str,
769        config: &serde_json::Value,
770    ) -> Result<()> {
771        #[derive(Serialize)]
772        struct ConfigBody<'a> {
773            config: &'a serde_json::Value,
774        }
775        let path = format!("/accounts/{account_id}/cfd_tunnel/{tunnel_id}/configurations");
776        let resp: CfSingle<serde_json::Value> = self.cf_put(&path, &ConfigBody { config }).await?;
777        self.ok(&resp.success, &resp.errors)?;
778        Ok(())
779    }
780
781    /// Create a new R2 bucket under `account_id`.
782    /// Requires: `Account: Cloudflare R2: Edit`.
783    pub async fn create_r2_bucket(
784        &self,
785        account_id: &str,
786        bucket_name: &str,
787    ) -> Result<CreateR2BucketResult> {
788        #[derive(Serialize)]
789        struct CreateBody<'a> {
790            name: &'a str,
791        }
792        let resp: CfSingle<serde_json::Value> = self
793            .cf_post(
794                &format!("/accounts/{account_id}/r2/buckets"),
795                &CreateBody { name: bucket_name },
796            )
797            .await?;
798        self.ok(&resp.success, &resp.errors)?;
799
800        Ok(CreateR2BucketResult {
801            endpoint: format!("https://{account_id}.r2.cloudflarestorage.com"),
802            name: bucket_name.to_string(),
803        })
804    }
805
806    /// Resolve a zone name (e.g. `"yah.dev"`) to its Cloudflare zone ID.
807    /// Requires: `Zone: Read`.
808    pub async fn zone_id_for_name(&self, zone_name: &str) -> Result<String> {
809        #[derive(Deserialize)]
810        struct ZoneEntry {
811            id: String,
812            name: String,
813        }
814        let resp: CfPage<ZoneEntry> = self.cf_get(&format!("/zones?name={zone_name}")).await?;
815        self.ok(&resp.success, &resp.errors)?;
816        resp.result
817            .unwrap_or_default()
818            .into_iter()
819            .find(|z| z.name == zone_name)
820            .map(|z| z.id)
821            .ok_or_else(|| anyhow!("no Cloudflare zone found for name {zone_name:?}"))
822    }
823
824    /// Purge content by cache tags from a zone.
825    ///
826    /// Cache tags must be applied to responses via the `Cache-Tag` header or
827    /// Cloudflare page rules. Returns `Ok(())` when all tags are queued for
828    /// purge. Requires: `Zone: Cache Purge`.
829    pub async fn purge_cache_tags(&self, zone_id: &str, tags: &[String]) -> Result<()> {
830        if tags.is_empty() {
831            return Ok(());
832        }
833        #[derive(Serialize)]
834        struct PurgeBody<'a> {
835            tags: &'a [String],
836        }
837        let resp: CfSingle<serde_json::Value> = self
838            .cf_post(
839                &format!("/zones/{zone_id}/purge_cache"),
840                &PurgeBody { tags },
841            )
842            .await?;
843        self.ok(&resp.success, &resp.errors)
844    }
845
846    /// Upsert the Transform Rule that rewrites `GET /` → `/index.html` on the
847    /// zone, identified by the stable description tag `"yah:static-index"`.
848    ///
849    /// Idempotent: fetches the existing `http_request_transform` entrypoint,
850    /// drops any prior `"yah:static-index"` rule, appends the current one,
851    /// and PUTs the merged list back. Treats a missing entrypoint (no rules
852    /// yet) as an empty list.
853    ///
854    /// Requires: `Zone: Transform Rules: Edit`.
855    pub async fn upsert_index_rewrite(&self, zone_id: &str) -> Result<()> {
856        const RULE_DESC: &str = "yah:static-index";
857        let path = format!("/zones/{zone_id}/rulesets/phases/http_request_transform/entrypoint");
858
859        // Fetch existing rules; a missing entrypoint is not an error.
860        let existing: Vec<serde_json::Value> = {
861            #[derive(Deserialize)]
862            struct Rs {
863                rules: Option<Vec<serde_json::Value>>,
864            }
865            match self.cf_get::<CfSingle<Rs>>(&path).await {
866                Ok(resp) if resp.success => resp.result.and_then(|r| r.rules).unwrap_or_default(),
867                _ => Vec::new(),
868            }
869        };
870
871        // Keep every rule except the one we manage, then append ours.
872        let mut rules: Vec<serde_json::Value> = existing
873            .into_iter()
874            .filter(|r| r.get("description").and_then(|v| v.as_str()) != Some(RULE_DESC))
875            .collect();
876        rules.push(serde_json::json!({
877            "action": "rewrite",
878            "description": RULE_DESC,
879            "expression": "(http.request.uri.path eq \"/\")",
880            "action_parameters": {
881                "uri": { "path": { "value": "/index.html" } }
882            },
883            "enabled": true
884        }));
885
886        let resp: CfSingle<serde_json::Value> = self
887            .cf_put(&path, &serde_json::json!({ "rules": rules }))
888            .await?;
889        self.ok(&resp.success, &resp.errors)
890    }
891
892    /// Deploy an ES-module Worker script with typed bindings for runtime config.
893    ///
894    /// Each entry in `bindings` becomes one `metadata.bindings[…]` declaration in
895    /// the upload payload — see [`WorkerBinding`] for the supported variants
896    /// (plain_text config, R2 bucket references).
897    ///
898    /// Uses a manual multipart/form-data upload (CF Workers API requires multipart
899    /// when metadata/bindings are attached). Idempotent: re-uploading the same script
900    /// is safe but costs one CF API round-trip — callers should hash-guard this.
901    ///
902    /// Requires: `Workers Scripts: Edit` (account-scoped).
903    pub async fn deploy_worker_script(
904        &self,
905        account_id: &str,
906        script_name: &str,
907        script_js: &str,
908        bindings: &[WorkerBinding<'_>],
909    ) -> Result<WorkerDeployResult> {
910        let url = format!("{CF_API}/accounts/{account_id}/workers/scripts/{script_name}");
911        let (content_type, body) = build_worker_multipart(script_js, bindings);
912        let resp = self
913            .http
914            .put(&url)
915            .header("Authorization", format!("Bearer {}", self.token))
916            .header("Content-Type", content_type)
917            .body(body)
918            .send()
919            .await
920            .map_err(|e| anyhow!("PUT {url}: {e}"))?;
921        let result: CfSingle<WorkerDeployResult> = resp
922            .json()
923            .await
924            .map_err(|e| anyhow!("PUT {url} parse: {e}"))?;
925        self.ok(&result.success, &result.errors)?;
926        result
927            .result
928            .ok_or_else(|| anyhow!("deploy worker: no result in response"))
929    }
930
931    /// Upsert a Worker route for `pattern` on `zone_id`, pointing at `script_name`.
932    ///
933    /// Idempotent: fetches existing routes, skips PUT/POST when the pattern already
934    /// points at the right script, updates an existing pattern pointing elsewhere,
935    /// or creates a new route entry.
936    ///
937    /// Requires: `Zone: Workers Routes: Edit` (zone-scoped).
938    pub async fn upsert_worker_route(
939        &self,
940        zone_id: &str,
941        pattern: &str,
942        script_name: &str,
943    ) -> Result<()> {
944        let list_path = format!("/zones/{zone_id}/workers/routes");
945
946        #[derive(Deserialize)]
947        struct RouteEntry {
948            id: String,
949            pattern: String,
950            #[serde(default)]
951            script: Option<String>,
952        }
953        let list: CfPage<RouteEntry> = self.cf_get(&list_path).await?;
954        self.ok(&list.success, &list.errors)?;
955        let routes = list.result.unwrap_or_default();
956
957        #[derive(Serialize)]
958        struct RouteBody<'a> {
959            pattern: &'a str,
960            script: &'a str,
961        }
962
963        if let Some(existing) = routes.iter().find(|r| r.pattern == pattern) {
964            if existing.script.as_deref() == Some(script_name) {
965                return Ok(());
966            }
967            let resp: CfSingle<serde_json::Value> = self
968                .cf_put(
969                    &format!("/zones/{zone_id}/workers/routes/{}", existing.id),
970                    &RouteBody {
971                        pattern,
972                        script: script_name,
973                    },
974                )
975                .await?;
976            self.ok(&resp.success, &resp.errors)
977        } else {
978            let resp: CfSingle<serde_json::Value> = self
979                .cf_post(
980                    &list_path,
981                    &RouteBody {
982                        pattern,
983                        script: script_name,
984                    },
985                )
986                .await?;
987            self.ok(&resp.success, &resp.errors)
988        }
989    }
990
991    /// Idempotently attach `hostname` (e.g. `cr.yah.dev`) as a Workers Custom
992    /// Domain on `script_name`. Custom Domains route every request for the
993    /// hostname into the Worker — distinct from a Worker Route, which only
994    /// matches a URL pattern within an already-proxied zone.
995    ///
996    /// Walks the existing Custom Domains list first; if `hostname` is already
997    /// bound to `script_name` on `zone_id`, returns Ok without an extra PUT.
998    /// Otherwise PUTs `/accounts/{account_id}/workers/domains`, which CF treats
999    /// as an upsert keyed on `(hostname, environment)`.
1000    ///
1001    /// Requires: `Workers Scripts: Edit` (account-scoped).
1002    pub async fn upsert_worker_custom_domain(
1003        &self,
1004        account_id: &str,
1005        zone_id: &str,
1006        hostname: &str,
1007        script_name: &str,
1008    ) -> Result<()> {
1009        let list_path = format!("/accounts/{account_id}/workers/domains");
1010
1011        #[derive(Deserialize)]
1012        struct DomainEntry {
1013            #[serde(default)]
1014            hostname: Option<String>,
1015            #[serde(default)]
1016            service: Option<String>,
1017            #[serde(default, rename = "zone_id")]
1018            zone_id: Option<String>,
1019        }
1020        let list: CfPage<DomainEntry> = self.cf_get(&list_path).await?;
1021        self.ok(&list.success, &list.errors)?;
1022        let domains = list.result.unwrap_or_default();
1023        if domains.iter().any(|d| {
1024            d.hostname.as_deref() == Some(hostname)
1025                && d.service.as_deref() == Some(script_name)
1026                && d.zone_id.as_deref() == Some(zone_id)
1027        }) {
1028            return Ok(());
1029        }
1030
1031        #[derive(Serialize)]
1032        struct DomainBody<'a> {
1033            environment: &'a str,
1034            hostname: &'a str,
1035            service: &'a str,
1036            zone_id: &'a str,
1037        }
1038        let resp: CfSingle<serde_json::Value> = self
1039            .cf_put(
1040                &list_path,
1041                &DomainBody {
1042                    environment: "production",
1043                    hostname,
1044                    service: script_name,
1045                    zone_id,
1046                },
1047            )
1048            .await?;
1049        self.ok(&resp.success, &resp.errors)
1050    }
1051
1052    /// Delete an R2 bucket under `account_id`.
1053    ///
1054    /// Cloudflare's management API handles non-empty buckets — objects do not
1055    /// need to be drained first. Returns `Ok(())` on success, `Err` if the
1056    /// API returns a failure (including "bucket not found" — callers that
1057    /// need idempotency should probe [`Self::list_r2_buckets`] first).
1058    ///
1059    /// Requires: `Account: Cloudflare R2: Edit`.
1060    pub async fn delete_r2_bucket(&self, account_id: &str, bucket_name: &str) -> Result<()> {
1061        let resp: CfSingle<serde_json::Value> = self
1062            .cf_delete(&format!("/accounts/{account_id}/r2/buckets/{bucket_name}"))
1063            .await?;
1064        self.ok(&resp.success, &resp.errors)
1065    }
1066
1067    /// Idempotently upsert a DNS record in `zone_id`. Fetches existing records
1068    /// with the same name and type: updates the first match if found, creates
1069    /// a new record otherwise. Returns the provider-issued record ID.
1070    ///
1071    /// Requires: `DNS: Edit` (zone-scoped).
1072    pub async fn upsert_dns_record(
1073        &self,
1074        zone_id: &str,
1075        name: &str,
1076        record_type: &str,
1077        content: &str,
1078        ttl: u32,
1079        proxied: bool,
1080    ) -> Result<String> {
1081        self.upsert_dns_record_matching(zone_id, name, record_type, content, ttl, proxied, false)
1082            .await
1083    }
1084
1085    /// [`upsert_dns_record`](Self::upsert_dns_record) with control over what
1086    /// counts as "the existing record" — R859-F1.
1087    ///
1088    /// `match_content = false` reproduces the original behaviour: the first
1089    /// record sharing `name` + `record_type` is updated in place. That is right
1090    /// for a single-valued name (one CNAME at `cdn.yah.dev`) and **wrong for a
1091    /// multi-valued RRset**: adding the second A record of a round-robin apex
1092    /// would rewrite the first one's content, silently halving the origin set
1093    /// to one box.
1094    ///
1095    /// `match_content = true` keys the lookup on `(name, type, content)`, so
1096    /// the call means "ensure exactly this record exists" — a no-op update when
1097    /// it already does, a create when it does not, and never a mutation of a
1098    /// sibling record at the same name.
1099    #[allow(clippy::too_many_arguments)]
1100    pub async fn upsert_dns_record_matching(
1101        &self,
1102        zone_id: &str,
1103        name: &str,
1104        record_type: &str,
1105        content: &str,
1106        ttl: u32,
1107        proxied: bool,
1108        match_content: bool,
1109    ) -> Result<String> {
1110        let existing_id = if match_content {
1111            self.list_dns_records(zone_id, Some(name), Some(record_type))
1112                .await?
1113                .into_iter()
1114                .find(|r| r.content == content)
1115                .map(|r| r.id)
1116        } else {
1117            self.find_dns_record(zone_id, name, record_type).await?
1118        };
1119
1120        #[derive(Serialize)]
1121        struct RecordBody<'a> {
1122            name: &'a str,
1123            #[serde(rename = "type")]
1124            record_type: &'a str,
1125            content: &'a str,
1126            ttl: u32,
1127            proxied: bool,
1128        }
1129        let body = RecordBody {
1130            name,
1131            record_type,
1132            content,
1133            ttl,
1134            proxied,
1135        };
1136
1137        #[derive(Deserialize)]
1138        struct RecordResult {
1139            id: String,
1140        }
1141
1142        if let Some(id) = existing_id {
1143            let resp: CfSingle<RecordResult> = self
1144                .cf_put(&format!("/zones/{zone_id}/dns_records/{id}"), &body)
1145                .await?;
1146            self.ok(&resp.success, &resp.errors)?;
1147            resp.result
1148                .map(|r| r.id)
1149                .ok_or_else(|| anyhow!("dns record update: no id in response"))
1150        } else {
1151            let resp: CfSingle<RecordResult> = self
1152                .cf_post(&format!("/zones/{zone_id}/dns_records"), &body)
1153                .await?;
1154            self.ok(&resp.success, &resp.errors)?;
1155            resp.result
1156                .map(|r| r.id)
1157                .ok_or_else(|| anyhow!("dns record create: no id in response"))
1158        }
1159    }
1160
1161    /// Read the DNS records in `zone_id`, optionally narrowed to one `name`
1162    /// and/or one `record_type` — R859-F1, the read half the `dns.*` catalog
1163    /// was missing.
1164    ///
1165    /// Unlike [`dns_records_named`](Self::dns_records_named) (drift-detection
1166    /// only, name + content) this returns the record **id**, type, ttl and
1167    /// proxy flag, which is what a reconciler needs to decide what to change.
1168    ///
1169    /// Requires: `DNS: Read` (zone-scoped).
1170    pub async fn list_dns_records(
1171        &self,
1172        zone_id: &str,
1173        name: Option<&str>,
1174        record_type: Option<&str>,
1175    ) -> Result<Vec<DnsRecordDetail>> {
1176        let mut query = format!("/zones/{zone_id}/dns_records?per_page=100");
1177        if let Some(n) = name {
1178            query.push_str(&format!("&name={n}"));
1179        }
1180        if let Some(t) = record_type {
1181            query.push_str(&format!("&type={t}"));
1182        }
1183        let resp: CfPage<DnsRecordDetail> = self.cf_get(&query).await?;
1184        self.ok(&resp.success, &resp.errors)?;
1185        Ok(resp.result.unwrap_or_default())
1186    }
1187
1188    /// Delete all DNS records in `zone_id` whose name matches `name` (and
1189    /// optionally `record_type`). Returns the count of records deleted.
1190    /// A count of 0 is not an error — the records may already have been absent.
1191    ///
1192    /// Requires: `DNS: Edit` (zone-scoped).
1193    pub async fn delete_dns_records(
1194        &self,
1195        zone_id: &str,
1196        name: &str,
1197        record_type: Option<&str>,
1198    ) -> Result<u32> {
1199        self.delete_dns_records_matching(zone_id, name, record_type, None)
1200            .await
1201    }
1202
1203    /// [`delete_dns_records`](Self::delete_dns_records) narrowed to records
1204    /// carrying one exact value — R859-F1.
1205    ///
1206    /// A round-robin apex holds several A records under one name, so
1207    /// "delete the A records at `yah.dev`" is not a way to withdraw *one*
1208    /// origin: it takes the live ones with it. `content = Some(ip)` deletes
1209    /// only the withdrawn member.
1210    ///
1211    /// Requires: `DNS: Edit` (zone-scoped).
1212    pub async fn delete_dns_records_matching(
1213        &self,
1214        zone_id: &str,
1215        name: &str,
1216        record_type: Option<&str>,
1217        content: Option<&str>,
1218    ) -> Result<u32> {
1219        let ids: Vec<String> = self
1220            .list_dns_records(zone_id, Some(name), record_type)
1221            .await?
1222            .into_iter()
1223            .filter(|r| content.is_none_or(|c| r.content == c))
1224            .map(|r| r.id)
1225            .collect();
1226        let mut deleted = 0u32;
1227        for id in &ids {
1228            let del: CfSingle<serde_json::Value> = self
1229                .cf_delete(&format!("/zones/{zone_id}/dns_records/{id}"))
1230                .await?;
1231            self.ok(&del.success, &del.errors)?;
1232            deleted += 1;
1233        }
1234        Ok(deleted)
1235    }
1236
1237    /// Fetch the ID of the first DNS record matching `name` and `record_type`.
1238    /// Returns `None` when no matching record exists.
1239    async fn find_dns_record(
1240        &self,
1241        zone_id: &str,
1242        name: &str,
1243        record_type: &str,
1244    ) -> Result<Option<String>> {
1245        #[derive(Deserialize)]
1246        struct RecordEntry {
1247            id: String,
1248        }
1249        let resp: CfPage<RecordEntry> = self
1250            .cf_get(&format!(
1251                "/zones/{zone_id}/dns_records?name={name}&type={record_type}"
1252            ))
1253            .await?;
1254        self.ok(&resp.success, &resp.errors)?;
1255        Ok(resp
1256            .result
1257            .unwrap_or_default()
1258            .into_iter()
1259            .next()
1260            .map(|r| r.id))
1261    }
1262
1263    /// List R2 buckets in `account_id`.
1264    /// Requires: `Account: Cloudflare R2: Read`.
1265    ///
1266    /// Unlike `/accounts` and `/zones`, the R2 list endpoint nests the array
1267    /// under `result.buckets` rather than returning `result` as a bare array,
1268    /// so it needs `CfSingle<R2ListResult>` and not `CfPage<BucketEntry>`.
1269    pub async fn list_r2_buckets(&self, account_id: &str) -> Result<Vec<R2BucketInfo>> {
1270        #[derive(Deserialize)]
1271        struct BucketEntry {
1272            name: String,
1273            #[serde(default)]
1274            location: Option<String>,
1275            #[serde(default)]
1276            creation_date: Option<String>,
1277        }
1278        #[derive(Deserialize)]
1279        struct R2ListResult {
1280            #[serde(default)]
1281            buckets: Vec<BucketEntry>,
1282        }
1283        let resp: CfSingle<R2ListResult> = self
1284            .cf_get(&format!("/accounts/{account_id}/r2/buckets"))
1285            .await?;
1286        self.ok(&resp.success, &resp.errors)?;
1287        Ok(resp
1288            .result
1289            .map(|r| r.buckets)
1290            .unwrap_or_default()
1291            .into_iter()
1292            .map(|b| R2BucketInfo {
1293                name: b.name,
1294                location: b.location,
1295                creation_date: b.creation_date,
1296            })
1297            .collect())
1298    }
1299
1300    /// List R2 custom-domain bindings on `bucket_name`.
1301    ///
1302    /// Requires: `Workers R2 Storage: Read` (or Write, which implies Read).
1303    /// The response nests the array under `result.domains`, mirroring
1304    /// `list_r2_buckets`'s `result.buckets` shape.
1305    pub async fn list_r2_custom_domains(
1306        &self,
1307        account_id: &str,
1308        bucket_name: &str,
1309    ) -> Result<Vec<R2CustomDomain>> {
1310        #[derive(Deserialize)]
1311        struct DomainEntry {
1312            domain: String,
1313            #[serde(default)]
1314            enabled: bool,
1315        }
1316        #[derive(Deserialize)]
1317        struct R2DomainListResult {
1318            #[serde(default)]
1319            domains: Vec<DomainEntry>,
1320        }
1321        let resp: CfSingle<R2DomainListResult> = self
1322            .cf_get(&format!(
1323                "/accounts/{account_id}/r2/buckets/{bucket_name}/domains/custom"
1324            ))
1325            .await?;
1326        self.ok(&resp.success, &resp.errors)?;
1327        Ok(resp
1328            .result
1329            .map(|r| r.domains)
1330            .unwrap_or_default()
1331            .into_iter()
1332            .map(|d| R2CustomDomain {
1333                domain: d.domain,
1334                enabled: d.enabled,
1335            })
1336            .collect())
1337    }
1338
1339    /// Bind a custom domain to an R2 bucket.
1340    ///
1341    /// `zone_id` names the zone that owns `domain` (resolve via
1342    /// [`Self::zone_id_for_name`]). CF requires it so the CNAME write into
1343    /// that zone is authorized — even though the caller is the bucket-side
1344    /// API. CF creates the CNAME automatically; no separate DNS-side call.
1345    /// Requires: `Workers R2 Storage: Edit` (account-scoped).
1346    ///
1347    /// `enabled: true` activates the binding immediately. CF still has to
1348    /// validate ownership + provision TLS in the background — the binding
1349    /// returns success the moment the record is queued, not when the
1350    /// hostname is fully resolvable. First-time DNS propagation is on the
1351    /// order of seconds to a minute.
1352    pub async fn add_r2_custom_domain(
1353        &self,
1354        account_id: &str,
1355        bucket_name: &str,
1356        domain: &str,
1357        zone_id: &str,
1358    ) -> Result<()> {
1359        #[derive(Serialize)]
1360        #[serde(rename_all = "camelCase")]
1361        struct AddBody<'a> {
1362            domain: &'a str,
1363            enabled: bool,
1364            zone_id: &'a str,
1365        }
1366        let resp: CfSingle<serde_json::Value> = self
1367            .cf_post(
1368                &format!("/accounts/{account_id}/r2/buckets/{bucket_name}/domains/custom"),
1369                &AddBody {
1370                    domain,
1371                    enabled: true,
1372                    zone_id,
1373                },
1374            )
1375            .await?;
1376        self.ok(&resp.success, &resp.errors)
1377    }
1378
1379    /// Fetch the account's permission-group catalog as a name → id map, used
1380    /// to resolve [`TokenGrant`] names before minting a token. Requires the
1381    /// calling token to carry `API Tokens: Read` (implied by `Write`).
1382    pub async fn list_permission_group_ids(
1383        &self,
1384        account_id: &str,
1385    ) -> Result<std::collections::BTreeMap<String, String>> {
1386        #[derive(Deserialize)]
1387        struct PgEntry {
1388            id: String,
1389            name: String,
1390        }
1391        let resp: CfPage<PgEntry> = self
1392            .cf_get(&format!(
1393                "/accounts/{account_id}/tokens/permission_groups?per_page=500"
1394            ))
1395            .await?;
1396        self.ok(&resp.success, &resp.errors)?;
1397        Ok(resp
1398            .result
1399            .unwrap_or_default()
1400            .into_iter()
1401            .map(|p| (p.name, p.id))
1402            .collect())
1403    }
1404
1405    /// Mint an **account-owned** API token under `account_id` from `grants`
1406    /// scoped to `account_id` + `zone_id`.
1407    ///
1408    /// Resolves each grant's permission-group name against the live catalog
1409    /// (falling back to its baked-in ID), groups the IDs into account- and
1410    /// zone-scoped policy blocks, and POSTs to `/accounts/{id}/tokens`.
1411    /// Requires the calling token to carry `API Tokens: Write` — but the
1412    /// minted token is bounded by the *account's* access, not the calling
1413    /// token's, so the caller may hold only `API Tokens: Write`.
1414    ///
1415    /// The returned [`CreateTokenResult::value`] is the secret — Cloudflare
1416    /// reveals it only here.
1417    pub async fn create_account_token(
1418        &self,
1419        account_id: &str,
1420        zone_id: &str,
1421        token_name: &str,
1422        grants: &[TokenGrant],
1423    ) -> Result<CreateTokenResult> {
1424        // A catalog lookup failure is non-fatal: build_token_body falls back
1425        // to the baked-in IDs for any name the (empty) catalog can't resolve.
1426        let catalog = self
1427            .list_permission_group_ids(account_id)
1428            .await
1429            .unwrap_or_default();
1430        let body = build_token_body(token_name, account_id, zone_id, grants, &catalog);
1431
1432        let resp: CfSingle<CreateTokenResult> = self
1433            .cf_post(&format!("/accounts/{account_id}/tokens"), &body)
1434            .await?;
1435        self.ok(&resp.success, &resp.errors)?;
1436        resp.result
1437            .ok_or_else(|| anyhow!("token create: no result in response"))
1438    }
1439
1440    // ---------- HTTP helpers ----------
1441
1442    async fn cf_get<T: serde::de::DeserializeOwned>(&self, path: &str) -> Result<T> {
1443        let url = format!("{CF_API}{path}");
1444        let resp = self
1445            .http
1446            .get(&url)
1447            .header("Authorization", format!("Bearer {}", self.token))
1448            .header("Content-Type", "application/json")
1449            .send()
1450            .await
1451            .map_err(|e| anyhow!("GET {url}: {e}"))?;
1452        resp.json::<T>()
1453            .await
1454            .map_err(|e| anyhow!("GET {url} parse: {e}"))
1455    }
1456
1457    async fn cf_post<B: Serialize, T: serde::de::DeserializeOwned>(
1458        &self,
1459        path: &str,
1460        body: &B,
1461    ) -> Result<T> {
1462        let url = format!("{CF_API}{path}");
1463        let resp = self
1464            .http
1465            .post(&url)
1466            .header("Authorization", format!("Bearer {}", self.token))
1467            .header("Content-Type", "application/json")
1468            .json(body)
1469            .send()
1470            .await
1471            .map_err(|e| anyhow!("POST {url}: {e}"))?;
1472        resp.json::<T>()
1473            .await
1474            .map_err(|e| anyhow!("POST {url} parse: {e}"))
1475    }
1476
1477    async fn cf_put<B: Serialize, T: serde::de::DeserializeOwned>(
1478        &self,
1479        path: &str,
1480        body: &B,
1481    ) -> Result<T> {
1482        let url = format!("{CF_API}{path}");
1483        let resp = self
1484            .http
1485            .put(&url)
1486            .header("Authorization", format!("Bearer {}", self.token))
1487            .header("Content-Type", "application/json")
1488            .json(body)
1489            .send()
1490            .await
1491            .map_err(|e| anyhow!("PUT {url}: {e}"))?;
1492        resp.json::<T>()
1493            .await
1494            .map_err(|e| anyhow!("PUT {url} parse: {e}"))
1495    }
1496
1497    async fn cf_delete<T: serde::de::DeserializeOwned>(&self, path: &str) -> Result<T> {
1498        let url = format!("{CF_API}{path}");
1499        let resp = self
1500            .http
1501            .delete(&url)
1502            .header("Authorization", format!("Bearer {}", self.token))
1503            .header("Content-Type", "application/json")
1504            .send()
1505            .await
1506            .map_err(|e| anyhow!("DELETE {url}: {e}"))?;
1507        resp.json::<T>()
1508            .await
1509            .map_err(|e| anyhow!("DELETE {url} parse: {e}"))
1510    }
1511
1512    fn ok(&self, success: &bool, errors: &Option<Vec<serde_json::Value>>) -> Result<()> {
1513        if *success {
1514            return Ok(());
1515        }
1516        let msg = errors
1517            .as_ref()
1518            .and_then(|e| e.first())
1519            .and_then(|e| e.get("message"))
1520            .and_then(|m| m.as_str())
1521            .unwrap_or("Cloudflare returned an error");
1522        Err(anyhow!("{msg}"))
1523    }
1524}
1525
1526// ---------- worker-upload helpers (pure, network-free) ----------
1527
1528/// Build a `multipart/form-data` body for uploading a CF Worker script with
1529/// typed bindings for runtime config.
1530///
1531/// Each [`WorkerBinding`] becomes one entry in the Worker metadata's
1532/// `bindings` array — plain_text strings, R2 bucket references, etc.
1533///
1534/// Returns `(content_type_header_value, raw_body_bytes)`.
1535fn build_worker_multipart(script_js: &str, bindings: &[WorkerBinding<'_>]) -> (String, Vec<u8>) {
1536    const BOUNDARY: &str = "yahWorkerUpload0";
1537    const SCRIPT_FILENAME: &str = "worker.js";
1538    let binding_json: Vec<serde_json::Value> = bindings
1539        .iter()
1540        .map(|b| match *b {
1541            WorkerBinding::PlainText { name, text } => serde_json::json!({
1542                "type": "plain_text",
1543                "name": name,
1544                "text": text,
1545            }),
1546            WorkerBinding::R2Bucket { name, bucket_name } => serde_json::json!({
1547                "type": "r2_bucket",
1548                "name": name,
1549                "bucket_name": bucket_name,
1550            }),
1551        })
1552        .collect();
1553    let metadata = serde_json::json!({
1554        "main_module": SCRIPT_FILENAME,
1555        "bindings": binding_json,
1556    });
1557    let mut body: Vec<u8> = Vec::new();
1558    let push = |v: &mut Vec<u8>, s: &str| v.extend_from_slice(s.as_bytes());
1559    // metadata part
1560    push(&mut body, &format!("--{BOUNDARY}\r\n"));
1561    push(
1562        &mut body,
1563        "Content-Disposition: form-data; name=\"metadata\"\r\n",
1564    );
1565    push(&mut body, "Content-Type: application/json\r\n\r\n");
1566    push(&mut body, &metadata.to_string());
1567    push(&mut body, "\r\n");
1568    // script part
1569    push(&mut body, &format!("--{BOUNDARY}\r\n"));
1570    push(&mut body, &format!(
1571        "Content-Disposition: form-data; name=\"{SCRIPT_FILENAME}\"; filename=\"{SCRIPT_FILENAME}\"\r\n"
1572    ));
1573    push(
1574        &mut body,
1575        "Content-Type: application/javascript+module\r\n\r\n",
1576    );
1577    push(&mut body, script_js);
1578    push(&mut body, &format!("\r\n--{BOUNDARY}--\r\n"));
1579    (format!("multipart/form-data; boundary={BOUNDARY}"), body)
1580}
1581
1582// ---------- token-create helpers (pure, network-free) ----------
1583
1584/// Resolve a grant's permission-group ID: prefer the live catalog entry for its
1585/// name, fall back to the baked-in constant.
1586fn resolve_grant_id(
1587    grant: &TokenGrant,
1588    catalog: &std::collections::BTreeMap<String, String>,
1589) -> String {
1590    catalog
1591        .get(grant.group_name)
1592        .cloned()
1593        .unwrap_or_else(|| grant.fallback_id.to_string())
1594}
1595
1596/// Build the `POST /accounts/{id}/tokens` request body for `grants`, grouping
1597/// account- and zone-scoped permission groups into separate policy blocks
1598/// (Cloudflare rejects a single block mixing the two scopes). A scope with no
1599/// grants produces no block.
1600fn build_token_body(
1601    token_name: &str,
1602    account_id: &str,
1603    zone_id: &str,
1604    grants: &[TokenGrant],
1605    catalog: &std::collections::BTreeMap<String, String>,
1606) -> serde_json::Value {
1607    let ids_for = |scope: GrantScope| -> Vec<serde_json::Value> {
1608        grants
1609            .iter()
1610            .filter(|g| g.scope == scope)
1611            .map(|g| serde_json::json!({ "id": resolve_grant_id(g, catalog) }))
1612            .collect()
1613    };
1614    let block = |resource: String, groups: Vec<serde_json::Value>| -> Option<serde_json::Value> {
1615        if groups.is_empty() {
1616            return None;
1617        }
1618        let mut resources = serde_json::Map::new();
1619        resources.insert(resource, serde_json::Value::String("*".into()));
1620        Some(serde_json::json!({
1621            "effect": "allow",
1622            "resources": serde_json::Value::Object(resources),
1623            "permission_groups": groups,
1624        }))
1625    };
1626
1627    let policies: Vec<serde_json::Value> = [
1628        block(
1629            format!("com.cloudflare.api.account.{account_id}"),
1630            ids_for(GrantScope::Account),
1631        ),
1632        block(
1633            format!("com.cloudflare.api.account.zone.{zone_id}"),
1634            ids_for(GrantScope::Zone),
1635        ),
1636    ]
1637    .into_iter()
1638    .flatten()
1639    .collect();
1640
1641    serde_json::json!({ "name": token_name, "policies": policies })
1642}
1643
1644// ---------- drift helpers (pure, network-free) ----------
1645
1646/// Normalise a DNS name/target for comparison: trim, drop a single trailing
1647/// dot, lowercase. So `yah.dev.` and `YAH.DEV` both compare equal to `yah.dev`.
1648fn norm_dns(s: &str) -> String {
1649    s.trim().trim_end_matches('.').to_ascii_lowercase()
1650}
1651
1652/// Pick the most specific accessible zone for `hostname` — the longest
1653/// zone-name suffix that is the apex of, or a parent of, the hostname.
1654/// Returns the matching zone id, or `None` when no accessible zone covers it.
1655fn best_zone_for<'a>(hostname: &str, zones: &'a [(String, String)]) -> Option<&'a str> {
1656    let h = norm_dns(hostname);
1657    zones
1658        .iter()
1659        .filter(|(_, name)| {
1660            let z = norm_dns(name);
1661            h == z || h.ends_with(&format!(".{z}"))
1662        })
1663        .max_by_key(|(_, name)| name.len())
1664        .map(|(id, _)| id.as_str())
1665}
1666
1667/// Classify drift for one ingress hostname against the live records fetched
1668/// for it. `live` is the zone's records (filtered by name upstream, but we
1669/// re-filter defensively so this stays a self-contained pure function).
1670fn classify_tunnel_drift(
1671    hostname: &str,
1672    expected_target: &str,
1673    live: &[CfDnsRecord],
1674) -> (TunnelDriftState, Option<String>) {
1675    let h = norm_dns(hostname);
1676    let matching: Vec<&CfDnsRecord> = live.iter().filter(|r| norm_dns(&r.name) == h).collect();
1677    if matching.is_empty() {
1678        return (TunnelDriftState::Missing, None);
1679    }
1680    let want = norm_dns(expected_target);
1681    if matching.iter().any(|r| norm_dns(&r.content) == want) {
1682        return (TunnelDriftState::Synced, None);
1683    }
1684    (
1685        TunnelDriftState::Mismatch,
1686        Some(matching[0].content.clone()),
1687    )
1688}
1689
1690#[cfg(test)]
1691mod tests {
1692    use super::*;
1693
1694    fn rec(name: &str, content: &str) -> CfDnsRecord {
1695        CfDnsRecord {
1696            name: name.into(),
1697            content: content.into(),
1698        }
1699    }
1700
1701    #[test]
1702    fn drift_synced_when_live_matches() {
1703        let live = vec![rec("yubaba.yah.dev", "9e4d.cfargotunnel.com")];
1704        let (state, target) =
1705            classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
1706        assert_eq!(state, TunnelDriftState::Synced);
1707        assert!(target.is_none());
1708    }
1709
1710    #[test]
1711    fn drift_missing_when_no_matching_record() {
1712        let live = vec![rec("other.yah.dev", "x.cfargotunnel.com")];
1713        let (state, target) =
1714            classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
1715        assert_eq!(state, TunnelDriftState::Missing);
1716        assert!(target.is_none());
1717    }
1718
1719    #[test]
1720    fn drift_missing_when_zone_empty() {
1721        let (state, _) = classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &[]);
1722        assert_eq!(state, TunnelDriftState::Missing);
1723    }
1724
1725    #[test]
1726    fn drift_mismatch_surfaces_live_target() {
1727        let live = vec![rec("yubaba.yah.dev", "stale.cfargotunnel.com")];
1728        let (state, target) =
1729            classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
1730        assert_eq!(state, TunnelDriftState::Mismatch);
1731        assert_eq!(target.as_deref(), Some("stale.cfargotunnel.com"));
1732    }
1733
1734    #[test]
1735    fn drift_normalises_trailing_dot_and_case() {
1736        // Cloudflare returns FQDNs and CNAME content with/without trailing dots
1737        // and arbitrary case; normalisation must treat these as synced.
1738        let live = vec![rec("Yubaba.YAH.dev.", "9E4D.cfargotunnel.com.")];
1739        let (state, _) = classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
1740        assert_eq!(state, TunnelDriftState::Synced);
1741    }
1742
1743    #[test]
1744    fn best_zone_picks_longest_suffix() {
1745        let zones = vec![
1746            ("z_apex".to_string(), "dev".to_string()),
1747            ("z_zone".to_string(), "yah.dev".to_string()),
1748        ];
1749        assert_eq!(best_zone_for("yubaba.yah.dev", &zones), Some("z_zone"));
1750        assert_eq!(best_zone_for("yah.dev", &zones), Some("z_zone"));
1751    }
1752
1753    #[test]
1754    fn best_zone_none_when_no_suffix_covers() {
1755        let zones = vec![("z1".to_string(), "yah.dev".to_string())];
1756        assert_eq!(best_zone_for("example.com", &zones), None);
1757        // A label that merely ends in the zone string but isn't a subdomain
1758        // must NOT match: `notyah.dev` is not under `yah.dev`.
1759        assert_eq!(best_zone_for("notyah.dev", &zones), None);
1760    }
1761
1762    #[test]
1763    fn best_zone_apex_self_match() {
1764        let zones = vec![("z1".to_string(), "yah.dev".to_string())];
1765        assert_eq!(best_zone_for("yah.dev", &zones), Some("z1"));
1766    }
1767
1768    #[test]
1769    fn token_body_splits_scopes_and_resolves_ids() {
1770        use std::collections::BTreeMap;
1771        // Catalog overrides Zone Read's id; everything else falls back to the
1772        // baked-in constant.
1773        let mut catalog = BTreeMap::new();
1774        catalog.insert("Zone Read".to_string(), "CATALOG_ZONE_READ".to_string());
1775
1776        let body = build_token_body("t", "ACCT", "ZONE", MESOFACT_STATIC_GRANTS, &catalog);
1777        let policies = body["policies"].as_array().unwrap();
1778        assert_eq!(policies.len(), 2, "one account block + one zone block");
1779
1780        let acct = &policies[0];
1781        assert_eq!(acct["resources"]["com.cloudflare.api.account.ACCT"], "*");
1782        let acct_ids: Vec<&str> = acct["permission_groups"]
1783            .as_array()
1784            .unwrap()
1785            .iter()
1786            .map(|g| g["id"].as_str().unwrap())
1787            .collect();
1788        assert!(acct_ids.contains(&"c1fde68c7bcc44588cbb6ddbc16d6480")); // Account Settings Read
1789        assert!(acct_ids.contains(&"bf7481a1826f439697cb59a20b22293e")); // R2 Storage Write
1790        assert!(acct_ids.contains(&"e086da7e2179491d91ee5f35b3ca210a")); // Workers Scripts Write
1791
1792        let zone = &policies[1];
1793        assert_eq!(
1794            zone["resources"]["com.cloudflare.api.account.zone.ZONE"],
1795            "*"
1796        );
1797        let zone_ids: Vec<&str> = zone["permission_groups"]
1798            .as_array()
1799            .unwrap()
1800            .iter()
1801            .map(|g| g["id"].as_str().unwrap())
1802            .collect();
1803        assert!(
1804            zone_ids.contains(&"CATALOG_ZONE_READ"),
1805            "catalog id wins over fallback"
1806        );
1807        assert!(zone_ids.contains(&"0ac90a90249747bca6b047d97f0803e9")); // Zone Transform Rules Write
1808        assert!(zone_ids.contains(&"28f4b596e7d643029c524985477ae49a")); // Workers Routes Write
1809        assert!(zone_ids.contains(&"e17beae8b8cb423a99b1730f21238bed")); // Cache Purge
1810    }
1811
1812    #[test]
1813    fn token_body_omits_empty_scope_block() {
1814        use std::collections::BTreeMap;
1815        let only_zone = &[TokenGrant {
1816            group_name: "Zone Read",
1817            scope: GrantScope::Zone,
1818            fallback_id: "ZR",
1819        }];
1820        let body = build_token_body("t", "A", "Z", only_zone, &BTreeMap::new());
1821        let policies = body["policies"].as_array().unwrap();
1822        assert_eq!(policies.len(), 1, "no account block when no account grants");
1823        assert_eq!(
1824            policies[0]["resources"]["com.cloudflare.api.account.zone.Z"],
1825            "*"
1826        );
1827    }
1828
1829    /// Decode the multipart `metadata` part and return its parsed JSON.
1830    fn extract_metadata_json(body: &[u8]) -> serde_json::Value {
1831        let s = std::str::from_utf8(body).expect("multipart body is utf-8 for these tests");
1832        let (_, after) = s
1833            .split_once("name=\"metadata\"")
1834            .expect("metadata part present");
1835        let (_, after) = after
1836            .split_once("\r\n\r\n")
1837            .expect("metadata body delimited");
1838        let (json, _) = after
1839            .split_once("\r\n--")
1840            .expect("metadata terminated by boundary");
1841        serde_json::from_str(json).expect("metadata JSON parses")
1842    }
1843
1844    #[test]
1845    fn multipart_includes_r2_bucket_binding_metadata() {
1846        let bindings = [WorkerBinding::R2Bucket {
1847            name: "CACHE",
1848            bucket_name: "yah-cr-cache",
1849        }];
1850        let (content_type, body) = build_worker_multipart("export default {}", &bindings);
1851
1852        assert!(
1853            content_type.starts_with("multipart/form-data; boundary="),
1854            "content-type advertises multipart with boundary: got {content_type}",
1855        );
1856
1857        let metadata = extract_metadata_json(&body);
1858        assert_eq!(metadata["main_module"], "worker.js");
1859        let bindings = metadata["bindings"].as_array().expect("bindings array");
1860        assert_eq!(bindings.len(), 1);
1861        assert_eq!(bindings[0]["type"], "r2_bucket");
1862        assert_eq!(bindings[0]["name"], "CACHE");
1863        assert_eq!(bindings[0]["bucket_name"], "yah-cr-cache");
1864    }
1865
1866    #[test]
1867    fn multipart_mixes_plain_text_and_r2_bindings() {
1868        let bindings = [
1869            WorkerBinding::PlainText {
1870                name: "MODE",
1871                text: "cache",
1872            },
1873            WorkerBinding::R2Bucket {
1874                name: "CACHE",
1875                bucket_name: "yah-cr-cache",
1876            },
1877        ];
1878        let (_, body) = build_worker_multipart("export default {}", &bindings);
1879
1880        let metadata = extract_metadata_json(&body);
1881        let entries = metadata["bindings"].as_array().unwrap();
1882        assert_eq!(entries.len(), 2);
1883        assert_eq!(entries[0]["type"], "plain_text");
1884        assert_eq!(entries[0]["text"], "cache");
1885        assert_eq!(entries[1]["type"], "r2_bucket");
1886        assert_eq!(entries[1]["bucket_name"], "yah-cr-cache");
1887    }
1888}