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// ---------- public output types ----------
136
137/// A Cloudflare account the API token can access.
138#[derive(Debug, Clone, Serialize, Deserialize)]
139#[serde(rename_all = "camelCase")]
140pub struct CfAccountInfo {
141    pub id: String,
142    pub name: String,
143}
144
145/// One CNAME record a user needs to create in their external DNS registrar
146/// to route a hostname through a Cloudflare Tunnel.
147#[derive(Debug, Clone, Serialize, Deserialize)]
148#[serde(rename_all = "camelCase")]
149pub struct TunnelDnsRecord {
150    /// Human-readable tunnel name as entered in Cloudflare.
151    pub tunnel_name: String,
152    /// Hostname from the tunnel ingress rule (e.g. `yah.example.com`).
153    pub hostname: String,
154    /// CNAME target to enter in the registrar: `{tunnel_id}.cfargotunnel.com`.
155    pub cname_target: String,
156}
157
158/// Result of creating a Cloudflare Named Tunnel.
159#[derive(Debug, Clone, Serialize, Deserialize)]
160#[serde(rename_all = "camelCase")]
161pub struct CreateTunnelResult {
162    pub tunnel_id: String,
163    pub tunnel_name: String,
164    /// JWT token for `cloudflared tunnel run --token <TOKEN>`.
165    pub connector_token: String,
166    /// CNAME target: `{tunnel_id}.cfargotunnel.com`.
167    pub cname_target: String,
168}
169
170/// Result of creating a Cloudflare R2 bucket.
171#[derive(Debug, Clone, Serialize, Deserialize)]
172#[serde(rename_all = "camelCase")]
173pub struct CreateR2BucketResult {
174    pub name: String,
175    /// S3-compatible endpoint for object operations against this account.
176    pub endpoint: String,
177}
178
179/// Result of deploying a Cloudflare Worker script.
180#[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/// One binding to inject into a Worker's `env` at deploy time.
188///
189/// Each variant maps to a `{"type": …}` entry under `metadata.bindings` in the
190/// Workers upload multipart body. R2 buckets are bind-only — the Worker reads
191/// the bucket via `env.<name>` and no script-side change is needed beyond the
192/// metadata declaration.
193#[derive(Debug, Clone, Copy)]
194pub enum WorkerBinding<'a> {
195    /// `{"type":"plain_text","name":…,"text":…}` — runtime config string.
196    PlainText { name: &'a str, text: &'a str },
197    /// `{"type":"r2_bucket","name":…,"bucket_name":…}` — R2 bucket reference.
198    R2Bucket { name: &'a str, bucket_name: &'a str },
199}
200
201/// R2 bucket information from the list endpoint.
202#[derive(Debug, Clone, Serialize, Deserialize)]
203#[serde(rename_all = "camelCase")]
204pub struct R2BucketInfo {
205    pub name: String,
206    /// CF location hint, e.g. `"WEUR"`, `"ENAM"`, `"APAC"`. `None` when the
207    /// bucket was created without specifying a hint.
208    #[serde(default, skip_serializing_if = "Option::is_none")]
209    pub location: Option<String>,
210    /// ISO 8601 bucket creation timestamp.
211    #[serde(default, skip_serializing_if = "Option::is_none")]
212    pub creation_date: Option<String>,
213}
214
215/// One R2 custom-domain binding from
216/// `GET /accounts/{id}/r2/buckets/{bucket}/domains/custom`.
217///
218/// The CF response carries nested status + min_tls fields we don't act on
219/// today; the reconciler only needs the hostname + enabled flag to decide
220/// idempotency.
221#[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/// Live connection state of a Cloudflare Tunnel connector.
230///
231/// Derived from the `status` field returned by
232/// `GET /accounts/{id}/cfd_tunnel?is_deleted=false`. Falls back to
233/// [`TunnelConnState::Unknown`] when the field is absent or unrecognised.
234#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
235#[serde(rename_all = "kebab-case")]
236pub enum TunnelConnState {
237    /// At least one healthy `cloudflared` connector is active.
238    Active,
239    /// No connectors are running.
240    Inactive,
241    /// Connectors are running but unhealthy.
242    Degraded,
243    /// Status couldn't be determined (field absent or unrecognised value).
244    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/// Drift verdict for one tunnel ingress hostname: does live Cloudflare DNS
259/// route it to the tunnel's CNAME target?
260///
261/// Maps onto the designed Tunnels table (`infra-cloudflare.jsx`): `Synced`
262/// renders as a healthy pill, everything else as a `drift` pill (`Missing`
263/// is the "DNS record missing" case the design calls out).
264#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
265#[serde(rename_all = "kebab-case")]
266pub enum TunnelDriftState {
267    /// A live DNS record points the hostname at the expected tunnel target.
268    Synced,
269    /// No live DNS record exists for the hostname — the tunnel can't route to it.
270    Missing,
271    /// A live record exists but points somewhere other than the tunnel target.
272    Mismatch,
273    /// The hostname's zone couldn't be resolved or read (token lacks
274    /// `Zone: Read` / `DNS: Read`, or the apex isn't a zone in any accessible
275    /// account) — drift indeterminate, not a failure.
276    ZoneUnknown,
277}
278
279/// One row of the tunnel DNS-drift report: a tunnel ingress hostname paired
280/// with whether live Cloudflare DNS routes it to the tunnel, plus the live
281/// connector connection state.
282#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
283#[serde(rename_all = "camelCase")]
284pub struct TunnelDriftRow {
285    /// Human-readable tunnel name as entered in Cloudflare.
286    pub tunnel_name: String,
287    /// Hostname from the tunnel's ingress rule, e.g. `yubaba.yah.dev`.
288    pub hostname: String,
289    /// CNAME target the DNS record should point at: `{tunnel_id}.cfargotunnel.com`.
290    pub expected_target: String,
291    pub state: TunnelDriftState,
292    /// For [`TunnelDriftState::Mismatch`], the target the live record actually
293    /// points at. `None` for every other state.
294    #[serde(default, skip_serializing_if = "Option::is_none")]
295    pub live_target: Option<String>,
296    /// Live connector state from `GET /cfd_tunnel`. [`TunnelConnState::Unknown`]
297    /// when the field was absent or unrecognised.
298    pub conn_state: TunnelConnState,
299    /// RFC 3339 `conns_active_at` timestamp — when this tunnel last became
300    /// active. `None` when the tunnel has never connected or the field was
301    /// absent.
302    #[serde(default, skip_serializing_if = "Option::is_none")]
303    pub conn_since: Option<String>,
304}
305
306// ---------- API-token provisioning (R320-F12) ----------
307
308/// Resource scope a permission group applies at when building a token policy.
309#[derive(Debug, Clone, Copy, PartialEq, Eq)]
310pub enum GrantScope {
311    /// Granted on the whole account (`com.cloudflare.api.account.<id>`).
312    Account,
313    /// Granted on a single zone (`com.cloudflare.api.account.zone.<id>`).
314    Zone,
315}
316
317/// One permission to bake into a minted token: a Cloudflare permission-group
318/// display name, the scope it applies at, and a validated fallback ID.
319///
320/// Names are resolved to IDs against the live permission-groups catalog at
321/// create time ([`CloudflareClient::list_permission_group_ids`]); `fallback_id`
322/// is used only when the catalog lookup can't resolve the name. The IDs are
323/// global Cloudflare constants (validated 2026-05-26), not per-account.
324#[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
331/// Minimal permission set for a mesofact-static publish token: see the account,
332/// list/create R2 buckets, deploy Worker scripts, resolve the zone, manage
333/// Worker routes + the index-rewrite Transform Rule, and purge the CDN cache.
334///
335/// Account-scoped: `Account Settings Read`, `Workers R2 Storage Write`,
336/// `Workers Scripts Write`.
337/// Zone-scoped: `Zone Read`, `Zone Transform Rules Write`,
338/// `Workers Routes Write`, `Cache Purge`.
339///
340/// NB the Transform Rules group is the *zone-scoped* `Zone Transform Rules
341/// Write`, not the account-scoped `Transform Rules Write` —
342/// [`CloudflareClient::upsert_index_rewrite`] hits `/zones/{id}/rulesets`.
343/// Fallback IDs sourced from the global permission-groups catalog
344/// (validated 2026-05-26); catalog resolution at create-time takes precedence.
345pub 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/// Result of minting an account-owned API token. `value` is the secret and is
384/// returned by Cloudflare exactly once — store it immediately.
385#[derive(Debug, Clone, Serialize, Deserialize)]
386#[serde(rename_all = "camelCase")]
387pub struct CreateTokenResult {
388    pub id: String,
389    pub name: String,
390    /// The token secret. Shown once by Cloudflare; never retrievable again.
391    pub value: String,
392}
393
394// ---------- client ----------
395
396/// Cloudflare management API client.
397///
398/// Construct with [`CloudflareClient::new`] passing a pre-resolved API token.
399/// The token scope required per method is noted on each method.
400pub struct CloudflareClient {
401    token: String,
402    http: reqwest::Client,
403}
404
405impl CloudflareClient {
406    /// Create a client for the given API token.
407    pub fn new(token: String) -> Self {
408        Self {
409            token,
410            http: reqwest::Client::new(),
411        }
412    }
413
414    /// List accounts the token can access.
415    /// Requires: `Account: Read`.
416    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    /// List non-deleted Cloudflare Tunnels in `account_id` with connection state.
436    /// Requires: `Cloudflare Tunnel: Read`.
437    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    /// List non-deleted Cloudflare Tunnels in `account_id` as `(id, name)` pairs.
462    /// Requires: `Cloudflare Tunnel: Read`.
463    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    /// Collect CNAME records for all tunnels across all accessible accounts.
473    ///
474    /// Walks accounts → tunnels → ingress configurations. Returns an empty
475    /// vec when the token has no tunnels or no configured ingress hostnames.
476    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    /// Compute DNS drift for every tunnel ingress hostname, enriched with live
510    /// connector connection state.
511    ///
512    /// The declared side is the tunnel ingress config; the live side is the
513    /// zone's DNS records. Each ingress hostname is classified:
514    /// [`TunnelDriftState::Synced`] when a record points at the tunnel's CNAME
515    /// target, `Missing` when none exists, `Mismatch` when one points elsewhere.
516    ///
517    /// Connection state (`conn_state` / `conn_since`) comes from the `status`
518    /// and `conns_active_at` fields on the tunnel list response — fetched in the
519    /// same pass as the ingress configs to avoid an extra `list_accounts` round-trip.
520    ///
521    /// Degrades gracefully — a hostname whose zone can't be resolved or read is
522    /// reported `ZoneUnknown` rather than failing the whole report. Returns an
523    /// empty vec when the token has no tunnels or no ingress hostnames.
524    ///
525    /// Requires: `Cloudflare Tunnel: Read`, `Zone: Read`, `DNS: Read`.
526    pub async fn tunnel_dns_drift(&self) -> Result<Vec<TunnelDriftRow>> {
527        // Walk accounts → tunnels (with live connection state) → ingress configs
528        // in one pass, collecting both declared DNS records and conn-state in a
529        // single list_accounts round-trip.
530        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        // Resolve zones once. A permission failure leaves the set empty, so
573        // every hostname falls through to `ZoneUnknown` instead of erroring.
574        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    /// List zones the token can read, as `(zone_id, zone_name)` pairs.
602    /// Requires: `Zone: Read`.
603    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    /// Fetch DNS records in `zone_id` whose name exactly matches `name`.
620    /// Requires: `DNS: Read`.
621    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    /// Create a new Named Tunnel under `account_id` and return the connector token.
630    /// Requires: `Cloudflare Tunnel: Edit`.
631    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        // Fetch the connector JWT.
664        #[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    /// Read a tunnel's remotely-managed configuration body as raw JSON
690    /// (R594-F11).
691    ///
692    /// Returns the `result.config` object — the thing a PUT round-trips — or an
693    /// empty object when the tunnel has never been configured. Deliberately
694    /// untyped: the `ingress` list is the only key this crate owns, and every
695    /// sibling (`warp-routing`, `originRequest`, …) must survive a
696    /// read-modify-write untouched.
697    ///
698    /// Requires: `Cloudflare Tunnel: Read`.
699    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        // A tunnel configured with an explicit JSON `null` config reads back as
714        // Value::Null, which has no object to insert `ingress` into.
715        Ok(if config.is_object() {
716            config
717        } else {
718            serde_json::json!({})
719        })
720    }
721
722    /// Replace a tunnel's remotely-managed configuration (R594-F11).
723    ///
724    /// `config` is the whole config body, not a patch — Cloudflare replaces it
725    /// wholesale, which is why callers must GET-merge-PUT rather than PUT a
726    /// freshly-built list. See
727    /// [`reconciler::ingress::ensure_tunnel_ingress`](crate::reconciler::ingress::ensure_tunnel_ingress).
728    ///
729    /// Requires: `Cloudflare Tunnel: Edit`.
730    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    /// Create a new R2 bucket under `account_id`.
747    /// Requires: `Account: Cloudflare R2: Edit`.
748    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    /// Resolve a zone name (e.g. `"yah.dev"`) to its Cloudflare zone ID.
772    /// Requires: `Zone: Read`.
773    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    /// Purge content by cache tags from a zone.
790    ///
791    /// Cache tags must be applied to responses via the `Cache-Tag` header or
792    /// Cloudflare page rules. Returns `Ok(())` when all tags are queued for
793    /// purge. Requires: `Zone: Cache Purge`.
794    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    /// Upsert the Transform Rule that rewrites `GET /` → `/index.html` on the
812    /// zone, identified by the stable description tag `"yah:static-index"`.
813    ///
814    /// Idempotent: fetches the existing `http_request_transform` entrypoint,
815    /// drops any prior `"yah:static-index"` rule, appends the current one,
816    /// and PUTs the merged list back. Treats a missing entrypoint (no rules
817    /// yet) as an empty list.
818    ///
819    /// Requires: `Zone: Transform Rules: Edit`.
820    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        // Fetch existing rules; a missing entrypoint is not an error.
825        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        // Keep every rule except the one we manage, then append ours.
837        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    /// Deploy an ES-module Worker script with typed bindings for runtime config.
858    ///
859    /// Each entry in `bindings` becomes one `metadata.bindings[…]` declaration in
860    /// the upload payload — see [`WorkerBinding`] for the supported variants
861    /// (plain_text config, R2 bucket references).
862    ///
863    /// Uses a manual multipart/form-data upload (CF Workers API requires multipart
864    /// when metadata/bindings are attached). Idempotent: re-uploading the same script
865    /// is safe but costs one CF API round-trip — callers should hash-guard this.
866    ///
867    /// Requires: `Workers Scripts: Edit` (account-scoped).
868    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    /// Upsert a Worker route for `pattern` on `zone_id`, pointing at `script_name`.
897    ///
898    /// Idempotent: fetches existing routes, skips PUT/POST when the pattern already
899    /// points at the right script, updates an existing pattern pointing elsewhere,
900    /// or creates a new route entry.
901    ///
902    /// Requires: `Zone: Workers Routes: Edit` (zone-scoped).
903    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    /// Idempotently attach `hostname` (e.g. `cr.yah.dev`) as a Workers Custom
957    /// Domain on `script_name`. Custom Domains route every request for the
958    /// hostname into the Worker — distinct from a Worker Route, which only
959    /// matches a URL pattern within an already-proxied zone.
960    ///
961    /// Walks the existing Custom Domains list first; if `hostname` is already
962    /// bound to `script_name` on `zone_id`, returns Ok without an extra PUT.
963    /// Otherwise PUTs `/accounts/{account_id}/workers/domains`, which CF treats
964    /// as an upsert keyed on `(hostname, environment)`.
965    ///
966    /// Requires: `Workers Scripts: Edit` (account-scoped).
967    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    /// Delete an R2 bucket under `account_id`.
1018    ///
1019    /// Cloudflare's management API handles non-empty buckets — objects do not
1020    /// need to be drained first. Returns `Ok(())` on success, `Err` if the
1021    /// API returns a failure (including "bucket not found" — callers that
1022    /// need idempotency should probe [`Self::list_r2_buckets`] first).
1023    ///
1024    /// Requires: `Account: Cloudflare R2: Edit`.
1025    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    /// Idempotently upsert a DNS record in `zone_id`. Fetches existing records
1033    /// with the same name and type: updates the first match if found, creates
1034    /// a new record otherwise. Returns the provider-issued record ID.
1035    ///
1036    /// Requires: `DNS: Edit` (zone-scoped).
1037    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    /// Delete all DNS records in `zone_id` whose name matches `name` (and
1090    /// optionally `record_type`). Returns the count of records deleted.
1091    /// A count of 0 is not an error — the records may already have been absent.
1092    ///
1093    /// Requires: `DNS: Edit` (zone-scoped).
1094    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    /// Fetch the ID of the first DNS record matching `name` and `record_type`.
1130    /// Returns `None` when no matching record exists.
1131    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    /// List R2 buckets in `account_id`.
1156    /// Requires: `Account: Cloudflare R2: Read`.
1157    ///
1158    /// Unlike `/accounts` and `/zones`, the R2 list endpoint nests the array
1159    /// under `result.buckets` rather than returning `result` as a bare array,
1160    /// so it needs `CfSingle<R2ListResult>` and not `CfPage<BucketEntry>`.
1161    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    /// List R2 custom-domain bindings on `bucket_name`.
1193    ///
1194    /// Requires: `Workers R2 Storage: Read` (or Write, which implies Read).
1195    /// The response nests the array under `result.domains`, mirroring
1196    /// `list_r2_buckets`'s `result.buckets` shape.
1197    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    /// Bind a custom domain to an R2 bucket.
1232    ///
1233    /// `zone_id` names the zone that owns `domain` (resolve via
1234    /// [`Self::zone_id_for_name`]). CF requires it so the CNAME write into
1235    /// that zone is authorized — even though the caller is the bucket-side
1236    /// API. CF creates the CNAME automatically; no separate DNS-side call.
1237    /// Requires: `Workers R2 Storage: Edit` (account-scoped).
1238    ///
1239    /// `enabled: true` activates the binding immediately. CF still has to
1240    /// validate ownership + provision TLS in the background — the binding
1241    /// returns success the moment the record is queued, not when the
1242    /// hostname is fully resolvable. First-time DNS propagation is on the
1243    /// order of seconds to a minute.
1244    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    /// Fetch the account's permission-group catalog as a name → id map, used
1272    /// to resolve [`TokenGrant`] names before minting a token. Requires the
1273    /// calling token to carry `API Tokens: Read` (implied by `Write`).
1274    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    /// Mint an **account-owned** API token under `account_id` from `grants`
1298    /// scoped to `account_id` + `zone_id`.
1299    ///
1300    /// Resolves each grant's permission-group name against the live catalog
1301    /// (falling back to its baked-in ID), groups the IDs into account- and
1302    /// zone-scoped policy blocks, and POSTs to `/accounts/{id}/tokens`.
1303    /// Requires the calling token to carry `API Tokens: Write` — but the
1304    /// minted token is bounded by the *account's* access, not the calling
1305    /// token's, so the caller may hold only `API Tokens: Write`.
1306    ///
1307    /// The returned [`CreateTokenResult::value`] is the secret — Cloudflare
1308    /// reveals it only here.
1309    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        // A catalog lookup failure is non-fatal: build_token_body falls back
1317        // to the baked-in IDs for any name the (empty) catalog can't resolve.
1318        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    // ---------- HTTP helpers ----------
1333
1334    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
1418// ---------- worker-upload helpers (pure, network-free) ----------
1419
1420/// Build a `multipart/form-data` body for uploading a CF Worker script with
1421/// typed bindings for runtime config.
1422///
1423/// Each [`WorkerBinding`] becomes one entry in the Worker metadata's
1424/// `bindings` array — plain_text strings, R2 bucket references, etc.
1425///
1426/// Returns `(content_type_header_value, raw_body_bytes)`.
1427fn 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    // metadata part
1452    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    // script part
1461    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
1474// ---------- token-create helpers (pure, network-free) ----------
1475
1476/// Resolve a grant's permission-group ID: prefer the live catalog entry for its
1477/// name, fall back to the baked-in constant.
1478fn 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
1488/// Build the `POST /accounts/{id}/tokens` request body for `grants`, grouping
1489/// account- and zone-scoped permission groups into separate policy blocks
1490/// (Cloudflare rejects a single block mixing the two scopes). A scope with no
1491/// grants produces no block.
1492fn 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
1536// ---------- drift helpers (pure, network-free) ----------
1537
1538/// Normalise a DNS name/target for comparison: trim, drop a single trailing
1539/// dot, lowercase. So `yah.dev.` and `YAH.DEV` both compare equal to `yah.dev`.
1540fn norm_dns(s: &str) -> String {
1541    s.trim().trim_end_matches('.').to_ascii_lowercase()
1542}
1543
1544/// Pick the most specific accessible zone for `hostname` — the longest
1545/// zone-name suffix that is the apex of, or a parent of, the hostname.
1546/// Returns the matching zone id, or `None` when no accessible zone covers it.
1547fn 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
1559/// Classify drift for one ingress hostname against the live records fetched
1560/// for it. `live` is the zone's records (filtered by name upstream, but we
1561/// re-filter defensively so this stays a self-contained pure function).
1562fn 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        // Cloudflare returns FQDNs and CNAME content with/without trailing dots
1629        // and arbitrary case; normalisation must treat these as synced.
1630        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        // A label that merely ends in the zone string but isn't a subdomain
1650        // must NOT match: `notyah.dev` is not under `yah.dev`.
1651        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        // Catalog overrides Zone Read's id; everything else falls back to the
1664        // baked-in constant.
1665        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")); // Account Settings Read
1681        assert!(acct_ids.contains(&"bf7481a1826f439697cb59a20b22293e")); // R2 Storage Write
1682        assert!(acct_ids.contains(&"e086da7e2179491d91ee5f35b3ca210a")); // Workers Scripts Write
1683
1684        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")); // Zone Transform Rules Write
1700        assert!(zone_ids.contains(&"28f4b596e7d643029c524985477ae49a")); // Workers Routes Write
1701        assert!(zone_ids.contains(&"e17beae8b8cb423a99b1730f21238bed")); // Cache Purge
1702    }
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    /// Decode the multipart `metadata` part and return its parsed JSON.
1722    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}