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//!
68//! @yah:ticket(R893-T22, "MESOFACT_STATIC_GRANTS gains account + zone Analytics Read, so the re-minted token can serve R893-F8")
69//! @yah:at(2026-09-12T22:27:26Z)
70//! @yah:status(review)
71//! @yah:assignee(agent:bundle-anthropic-ashguard)
72//! @yah:phase(P4)
73//! @yah:parent(R893)
74//! @yah:next("Tier: Cleric — a small, well-specified addition to one const array, but the permission-group ids must be resolved against the live Cloudflare catalog rather than guessed, and three stale prose enumerations need correcting in the same pass.")
75//! @yah:handoff("Added Account Analytics Read (account, id b89a480218d04ceb98b4fe57ca29dc1f) and Analytics Read (zone, id 9c88f9c5bce24ce7af9a958ba9c504db) to MESOFACT_STATIC_GRANTS in oss/yubaba/crates/cloud/src/provider/cloudflare.rs (now 11 grants: 4 account + 7 zone). Both group names + ids resolved live 2026-09-12 via GET /accounts/3948dc292e724e71b0deefde0ea95999/tokens/permission_groups using the cloudflare-legacy-yah keystore slot (the token the operator granted Analytics on) — never printed, only piped through curl+python filtering.")
76//! @yah:handoff("Updated the three stale enumerations named in the ticket: app/yah/cli/src/cloud.rs:1212-1218 doc comment now lists Account Analytics:Read + Analytics:Read; .yah/infra/providers/cloudflare.toml:37-52 gained a third dated paragraph ('eleven total') rather than editing the historical seven/nine paragraphs in place; oss/yah-base/crates/keys/src/spec.rs:695 now says 'eleven' and its cloudflare.toml consumer line reference corrected from :47 to :53.")
77//! @yah:handoff("No mint performed — scope fence respected. cargo test -p yah-cloud --lib: 1205 passed both before and after (baseline established fresh, same run). cargo check -p desktop --lib, -p yah --lib, -p fob --lib all clean (pre-existing warnings only, no new errors).")
78//! @yah:verify("cd oss/yubaba && cargo test -p yah-cloud --lib — 1205 passed, 0 failed (matches pre-change baseline)")
79//! @yah:verify("cargo check -p desktop --lib — clean")
80//! @yah:verify("cargo check -p yah --lib -p fob --lib — clean")
81//! @yah:assumes("cloudflare-legacy-yah keystore slot is the same bootstrap token the operator granted Analytics on 2026-09-12 per the ticket description; its successful GET on /tokens/permission_groups (API Tokens:Read/Edit scope) is what proved that grant reached the account, not an independent verification of the Analytics grant itself.")
82//! @yah:verify("Leader re-ran the gate independently: `cd oss/yubaba && cargo test -p yah-cloud --lib` = 1205 passed / 0 failed / 4 ignored, identical to the leader's own pre-change baseline measured before dispatch. The camp build rail reported \"input closure unchanged across the whole run: no skew\" on that baseline, so both numbers describe the same tree. Leader also read the grants array directly: `Account Analytics Read` (account scope) and `Analytics Read` (zone scope) are present with the live-resolved ids, and the doc comment at cloudflare.rs:382-388 records that the two groups are genuinely distinct and were validated against the catalog rather than inferred from their names.")
83
84use anyhow::{anyhow, Result};
85use serde::{Deserialize, Serialize};
86
87const CF_API: &str = "https://api.cloudflare.com/client/v4";
88
89// ---------- internal wire types ----------
90
91#[derive(Deserialize)]
92struct CfPage<T> {
93    success: bool,
94    result: Option<Vec<T>>,
95    errors: Option<Vec<serde_json::Value>>,
96}
97
98#[derive(Deserialize)]
99struct CfSingle<T> {
100    success: bool,
101    result: Option<T>,
102    errors: Option<Vec<serde_json::Value>>,
103}
104
105#[derive(Deserialize)]
106struct CfTunnel {
107    id: String,
108    name: String,
109    /// Cloudflare's live status: `"active"` | `"inactive"` | `"degraded"` | `"unknown"`.
110    #[serde(default)]
111    status: Option<String>,
112    /// RFC 3339 timestamp — when the tunnel last became active. `null` when
113    /// never active or the field is absent.
114    #[serde(default)]
115    conns_active_at: Option<String>,
116}
117
118/// Enriched tunnel record used internally when walking accounts for drift.
119struct TunnelMeta {
120    id: String,
121    name: String,
122    conn_state: TunnelConnState,
123    conn_since: Option<String>,
124}
125
126#[derive(Deserialize)]
127struct CfTunnelConfig {
128    config: Option<CfIngressConfig>,
129}
130
131#[derive(Deserialize)]
132struct CfIngressConfig {
133    ingress: Option<Vec<CfIngressRule>>,
134}
135
136#[derive(Deserialize)]
137struct CfIngressRule {
138    hostname: Option<String>,
139}
140
141/// One live DNS record from `GET /zones/{id}/dns_records`. We only read the
142/// name + content (the target); record type and proxy flags are ignored —
143/// drift is decided by whether *some* record routes the hostname to the
144/// expected tunnel target.
145#[derive(Debug, Clone, Deserialize)]
146struct CfDnsRecord {
147    name: String,
148    content: String,
149}
150
151/// One live DNS record with everything a reconciler needs to decide what to
152/// change — R859-F1, the read side of [`CloudflareClient::list_dns_records`].
153///
154/// Distinct from the private [`CfDnsRecord`] above, which deliberately reads
155/// only name + content because tunnel-drift detection asks a narrower
156/// question. Reconciling a record set needs the `id` (to delete one member of
157/// a multi-valued RRset) and `proxied` (an orange-clouded record at an apex
158/// that should be grey is drift, even when the content is right).
159#[derive(Debug, Clone, PartialEq, Eq, Deserialize)]
160pub struct DnsRecordDetail {
161    pub id: String,
162    pub name: String,
163    #[serde(rename = "type")]
164    pub record_type: String,
165    pub content: String,
166    #[serde(default)]
167    pub ttl: u32,
168    #[serde(default)]
169    pub proxied: bool,
170}
171
172// ---------- public output types ----------
173
174/// A Cloudflare account the API token can access.
175#[derive(Debug, Clone, Serialize, Deserialize)]
176#[serde(rename_all = "camelCase")]
177pub struct CfAccountInfo {
178    pub id: String,
179    pub name: String,
180}
181
182/// One CNAME record a user needs to create in their external DNS registrar
183/// to route a hostname through a Cloudflare Tunnel.
184#[derive(Debug, Clone, Serialize, Deserialize)]
185#[serde(rename_all = "camelCase")]
186pub struct TunnelDnsRecord {
187    /// Human-readable tunnel name as entered in Cloudflare.
188    pub tunnel_name: String,
189    /// Hostname from the tunnel ingress rule (e.g. `yah.example.com`).
190    pub hostname: String,
191    /// CNAME target to enter in the registrar: `{tunnel_id}.cfargotunnel.com`.
192    pub cname_target: String,
193}
194
195/// Result of creating a Cloudflare Named Tunnel.
196#[derive(Debug, Clone, Serialize, Deserialize)]
197#[serde(rename_all = "camelCase")]
198pub struct CreateTunnelResult {
199    pub tunnel_id: String,
200    pub tunnel_name: String,
201    /// JWT token for `cloudflared tunnel run --token <TOKEN>`.
202    pub connector_token: String,
203    /// CNAME target: `{tunnel_id}.cfargotunnel.com`.
204    pub cname_target: String,
205}
206
207/// Result of creating a Cloudflare R2 bucket.
208#[derive(Debug, Clone, Serialize, Deserialize)]
209#[serde(rename_all = "camelCase")]
210pub struct CreateR2BucketResult {
211    pub name: String,
212    /// S3-compatible endpoint for object operations against this account.
213    pub endpoint: String,
214}
215
216/// Result of deploying a Cloudflare Worker script.
217#[derive(Debug, Clone, Serialize, Deserialize)]
218#[serde(rename_all = "camelCase")]
219pub struct WorkerDeployResult {
220    pub id: String,
221    pub etag: Option<String>,
222}
223
224/// One binding to inject into a Worker's `env` at deploy time.
225///
226/// Each variant maps to a `{"type": …}` entry under `metadata.bindings` in the
227/// Workers upload multipart body. R2 buckets are bind-only — the Worker reads
228/// the bucket via `env.<name>` and no script-side change is needed beyond the
229/// metadata declaration.
230#[derive(Debug, Clone, Copy)]
231pub enum WorkerBinding<'a> {
232    /// `{"type":"plain_text","name":…,"text":…}` — runtime config string.
233    PlainText { name: &'a str, text: &'a str },
234    /// `{"type":"r2_bucket","name":…,"bucket_name":…}` — R2 bucket reference.
235    R2Bucket { name: &'a str, bucket_name: &'a str },
236}
237
238/// R2 bucket information from the list endpoint.
239#[derive(Debug, Clone, Serialize, Deserialize)]
240#[serde(rename_all = "camelCase")]
241pub struct R2BucketInfo {
242    pub name: String,
243    /// CF location hint, e.g. `"WEUR"`, `"ENAM"`, `"APAC"`. `None` when the
244    /// bucket was created without specifying a hint.
245    #[serde(default, skip_serializing_if = "Option::is_none")]
246    pub location: Option<String>,
247    /// ISO 8601 bucket creation timestamp.
248    #[serde(default, skip_serializing_if = "Option::is_none")]
249    pub creation_date: Option<String>,
250}
251
252/// One R2 custom-domain binding from
253/// `GET /accounts/{id}/r2/buckets/{bucket}/domains/custom`.
254///
255/// The CF response carries nested status + min_tls fields we don't act on
256/// today; the reconciler only needs the hostname + enabled flag to decide
257/// idempotency.
258#[derive(Debug, Clone, Serialize, Deserialize)]
259#[serde(rename_all = "camelCase")]
260pub struct R2CustomDomain {
261    pub domain: String,
262    #[serde(default)]
263    pub enabled: bool,
264}
265
266/// Live connection state of a Cloudflare Tunnel connector.
267///
268/// Derived from the `status` field returned by
269/// `GET /accounts/{id}/cfd_tunnel?is_deleted=false`. Falls back to
270/// [`TunnelConnState::Unknown`] when the field is absent or unrecognised.
271#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
272#[serde(rename_all = "kebab-case")]
273pub enum TunnelConnState {
274    /// At least one healthy `cloudflared` connector is active.
275    Active,
276    /// No connectors are running.
277    Inactive,
278    /// Connectors are running but unhealthy.
279    Degraded,
280    /// Status couldn't be determined (field absent or unrecognised value).
281    Unknown,
282}
283
284impl TunnelConnState {
285    fn from_cf_status(s: &str) -> Self {
286        match s {
287            "active" | "healthy" => Self::Active,
288            "inactive" | "down" => Self::Inactive,
289            "degraded" | "unhealthy" => Self::Degraded,
290            _ => Self::Unknown,
291        }
292    }
293}
294
295/// Drift verdict for one tunnel ingress hostname: does live Cloudflare DNS
296/// route it to the tunnel's CNAME target?
297///
298/// Maps onto the designed Tunnels table (`infra-cloudflare.jsx`): `Synced`
299/// renders as a healthy pill, everything else as a `drift` pill (`Missing`
300/// is the "DNS record missing" case the design calls out).
301#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
302#[serde(rename_all = "kebab-case")]
303pub enum TunnelDriftState {
304    /// A live DNS record points the hostname at the expected tunnel target.
305    Synced,
306    /// No live DNS record exists for the hostname — the tunnel can't route to it.
307    Missing,
308    /// A live record exists but points somewhere other than the tunnel target.
309    Mismatch,
310    /// The hostname's zone couldn't be resolved or read (token lacks
311    /// `Zone: Read` / `DNS: Read`, or the apex isn't a zone in any accessible
312    /// account) — drift indeterminate, not a failure.
313    ZoneUnknown,
314}
315
316/// One row of the tunnel DNS-drift report: a tunnel ingress hostname paired
317/// with whether live Cloudflare DNS routes it to the tunnel, plus the live
318/// connector connection state.
319#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
320#[serde(rename_all = "camelCase")]
321pub struct TunnelDriftRow {
322    /// Human-readable tunnel name as entered in Cloudflare.
323    pub tunnel_name: String,
324    /// Hostname from the tunnel's ingress rule, e.g. `yubaba.yah.dev`.
325    pub hostname: String,
326    /// CNAME target the DNS record should point at: `{tunnel_id}.cfargotunnel.com`.
327    pub expected_target: String,
328    pub state: TunnelDriftState,
329    /// For [`TunnelDriftState::Mismatch`], the target the live record actually
330    /// points at. `None` for every other state.
331    #[serde(default, skip_serializing_if = "Option::is_none")]
332    pub live_target: Option<String>,
333    /// Live connector state from `GET /cfd_tunnel`. [`TunnelConnState::Unknown`]
334    /// when the field was absent or unrecognised.
335    pub conn_state: TunnelConnState,
336    /// RFC 3339 `conns_active_at` timestamp — when this tunnel last became
337    /// active. `None` when the tunnel has never connected or the field was
338    /// absent.
339    #[serde(default, skip_serializing_if = "Option::is_none")]
340    pub conn_since: Option<String>,
341}
342
343// ---------- API-token provisioning (R320-F12) ----------
344
345/// Resource scope a permission group applies at when building a token policy.
346#[derive(Debug, Clone, Copy, PartialEq, Eq)]
347pub enum GrantScope {
348    /// Granted on the whole account (`com.cloudflare.api.account.<id>`).
349    Account,
350    /// Granted on a single zone (`com.cloudflare.api.account.zone.<id>`).
351    Zone,
352}
353
354/// One permission to bake into a minted token: a Cloudflare permission-group
355/// display name, the scope it applies at, and a validated fallback ID.
356///
357/// Names are resolved to IDs against the live permission-groups catalog at
358/// create time ([`CloudflareClient::list_permission_group_ids`]); `fallback_id`
359/// is used only when the catalog lookup can't resolve the name. The IDs are
360/// global Cloudflare constants (validated 2026-05-26), not per-account.
361#[derive(Debug, Clone, Copy)]
362pub struct TokenGrant {
363    pub group_name: &'static str,
364    pub scope: GrantScope,
365    pub fallback_id: &'static str,
366}
367
368/// Minimal permission set for a mesofact-static publish token: see the account,
369/// list/create R2 buckets, deploy Worker scripts, resolve the zone, manage
370/// Worker routes + the index-rewrite Transform Rule, purge the CDN cache, and
371/// read GraphQL Analytics (account + zone, R893-T22 — backs the front-door
372/// latency panel in `app/yah/desktop/src/front_door.rs`).
373///
374/// Account-scoped: `Account Settings Read`, `Workers R2 Storage Write`,
375/// `Workers Scripts Write`, `Account Analytics Read`.
376/// Zone-scoped: `Zone Read`, `Zone Transform Rules Write`,
377/// `Workers Routes Write`, `Cache Purge`, `DNS Read`, `DNS Write`,
378/// `Analytics Read`.
379///
380/// NB the Transform Rules group is the *zone-scoped* `Zone Transform Rules
381/// Write`, not the account-scoped `Transform Rules Write` —
382/// [`CloudflareClient::upsert_index_rewrite`] hits `/zones/{id}/rulesets`.
383/// NB likewise there are two distinct Analytics groups, resolved live
384/// 2026-09-12 against `GET /accounts/{id}/tokens/permission_groups`:
385/// `Account Analytics Read` (`com.cloudflare.api.account`) and `Analytics
386/// Read` (`com.cloudflare.api.account.zone`) — same trap shape as Transform
387/// Rules, a scoped and an unscoped group sharing a name fragment.
388/// Fallback IDs sourced from the global permission-groups catalog
389/// (validated 2026-05-26, Analytics pair added 2026-09-12); catalog
390/// resolution at create-time takes precedence.
391pub const MESOFACT_STATIC_GRANTS: &[TokenGrant] = &[
392    TokenGrant {
393        group_name: "Account Settings Read",
394        scope: GrantScope::Account,
395        fallback_id: "c1fde68c7bcc44588cbb6ddbc16d6480",
396    },
397    TokenGrant {
398        group_name: "Workers R2 Storage Write",
399        scope: GrantScope::Account,
400        fallback_id: "bf7481a1826f439697cb59a20b22293e",
401    },
402    TokenGrant {
403        group_name: "Workers Scripts Write",
404        scope: GrantScope::Account,
405        fallback_id: "e086da7e2179491d91ee5f35b3ca210a",
406    },
407    // R893-T22: backs the front-door latency panel's GraphQL Analytics query
408    // (app/yah/desktop/src/front_door.rs). Resolved live 2026-09-12 against
409    // GET /accounts/3948dc292e724e71b0deefde0ea95999/tokens/permission_groups.
410    TokenGrant {
411        group_name: "Account Analytics Read",
412        scope: GrantScope::Account,
413        fallback_id: "b89a480218d04ceb98b4fe57ca29dc1f",
414    },
415    TokenGrant {
416        group_name: "Zone Read",
417        scope: GrantScope::Zone,
418        fallback_id: "c8fed203ed3043cba015a93ad1616f1f",
419    },
420    TokenGrant {
421        group_name: "Zone Transform Rules Write",
422        scope: GrantScope::Zone,
423        fallback_id: "0ac90a90249747bca6b047d97f0803e9",
424    },
425    TokenGrant {
426        group_name: "Workers Routes Write",
427        scope: GrantScope::Zone,
428        fallback_id: "28f4b596e7d643029c524985477ae49a",
429    },
430    TokenGrant {
431        group_name: "Cache Purge",
432        scope: GrantScope::Zone,
433        fallback_id: "e17beae8b8cb423a99b1730f21238bed",
434    },
435    // R859-F1: deploy_domain_passway (the sovereign-apex A-record reconciler)
436    // is the first production consumer of the dns.* envoy verbs against this
437    // token, and it both lists and upserts/prunes A records — hence both
438    // Read and Write, not just Read as the R324-F5 @yah:next follow-up noted.
439    TokenGrant {
440        group_name: "DNS Read",
441        scope: GrantScope::Zone,
442        fallback_id: "82e64a83756745bbbb1c9c2701bf816b",
443    },
444    TokenGrant {
445        group_name: "DNS Write",
446        scope: GrantScope::Zone,
447        fallback_id: "4755a26eedb94da69e1066d98aa820be",
448    },
449    // R893-T22: the zone-scoped counterpart to Account Analytics Read above —
450    // NOT the same permission group despite the near-identical name.
451    TokenGrant {
452        group_name: "Analytics Read",
453        scope: GrantScope::Zone,
454        fallback_id: "9c88f9c5bce24ce7af9a958ba9c504db",
455    },
456];
457
458/// Account-scoped grant for a token that only needs to publish and read back
459/// a Cloudflare Tunnel's ingress configuration (`ensure_tunnel_ingress` in
460/// `reconciler/ingress.rs` — a PUT, so it needs the Write/"Edit" group, not
461/// the Read one). Group name + fallback ID resolved live 2026-09-15 against
462/// `GET /accounts/3948dc292e724e71b0deefde0ea95999/tokens/permission_groups`
463/// (R912-F1) — Cloudflare's own catalog calls the group "Cloudflare Tunnel
464/// Write", which is what the dashboard/docs render as "Cloudflare Tunnel:
465/// Edit". Deliberately account-scoped only, no zone grants: a tunnel's
466/// ingress config is an account-level resource, not a per-zone one.
467pub const TUNNEL_EDIT_GRANTS: &[TokenGrant] = &[TokenGrant {
468    group_name: "Cloudflare Tunnel Write",
469    scope: GrantScope::Account,
470    fallback_id: "c07321b023e944ff818fec44d8203567",
471}];
472
473/// Result of minting an account-owned API token. `value` is the secret and is
474/// returned by Cloudflare exactly once — store it immediately.
475#[derive(Debug, Clone, Serialize, Deserialize)]
476#[serde(rename_all = "camelCase")]
477pub struct CreateTokenResult {
478    pub id: String,
479    pub name: String,
480    /// The token secret. Shown once by Cloudflare; never retrievable again.
481    pub value: String,
482}
483
484// ---------- client ----------
485
486/// Cloudflare management API client.
487///
488/// Construct with [`CloudflareClient::new`] passing a pre-resolved API token.
489/// The token scope required per method is noted on each method.
490pub struct CloudflareClient {
491    token: String,
492    http: reqwest::Client,
493}
494
495/// @yah:relay(R907, "A cloudflare-tunnel edge cannot be brought up by apply alone when its tunnel has no configuration yet")
496/// @yah:at(2026-09-14T19:14:34Z)
497/// @yah:status(open)
498/// @yah:assignee(agent:bundle-anthropic-ashguard)
499/// @yah:gotcha("FILED FROM THE noisetable CAMP 2026-09-14 (its R704-T4), where it blocks the staging API door end to end. Filed here rather than there because the fix is entirely in this repo. NOT parented to R845 deliberately: R845 is the right neighbourhood — it landed the edge-first `use` resolution this camp's own fix relies on — but it sits in `review`, and a defect filed onto a review-column relay reaches nobody.")
500/// @yah:gotcha("THE CALLER'S ABORT-BEFORE-PUT BEHAVIOUR AT reconciler/ingress.rs:1272-1276 IS CORRECT AND DELIBERATE, and is not what this relay is about. Its doc says a failed API call is never read as \"delete every hostname rule\", and one tunnel multiplexes every service on a node, so a GET that failed for a real reason must keep aborting before the PUT.")
501///
502/// @yah:ticket(R907-B1, "ensure_tunnel_ingress cannot bootstrap a tunnel that has no configuration yet")
503/// @yah:status(review)
504/// @yah:at(2026-09-14T19:34:22Z)
505/// @yah:assignee(agent:bundle-anthropic-miravel)
506/// @yah:parent(R907)
507/// @yah:severity(high)
508/// @yah:next("Tier: Cleric — a small, well-located error-classification fix in one function, but it sits on a live credentialed write path where the empty-vs-failed distinction has already been reasoned about once and must not be blurred.")
509/// @yah:next("THE FIX BELONGS IN THE GET'S CLASSIFICATION, NOT IN THE CALLER: `tunnel_configuration` should return `Ok(json!({}))` for the specific not-found error code and `Err` for everything else. Do NOT loosen the abort-before-PUT guard at reconciler/ingress.rs:1272-1276 — that guard is correct and is what stops a transient API failure being read as \"delete every hostname rule\" on a tunnel that multiplexes every service on a node.")
510/// @yah:next("Worth a sweep while in here: any other `cf_get` whose \"resource has never been created\" answer arrives as an error object rather than an empty result has the same latent shape.")
511/// @yah:verify("REPRO, MEASURED 2026-09-14 from the noisetable camp: tunnel 1b41587a-7298-4f29-97bf-40d51ed79397, created by `yah cloud cf tunnel ensure` WITHOUT `--hostname` so it has never had a configuration PUT, then `yah cloud apply --env staging --service noisetable-api`. Result: `FAILED: reading ingress configuration of tunnel 1b41587a-…: Configuration for tunnel not found`, apply stops at first failure, no connector workload is deployed.")
512/// @yah:verify("THE CREDENTIAL WAS RULED OUT SEPARATELY, and that is what isolates the defect: with a token lacking Tunnel grants the same call fails with `Not authorized` instead, and moving the edge onto a Tunnel-capable provider changed the error from `Not authorized` to `Configuration for tunnel not found`. Two different failures at the same line means the second one is not an auth problem.")
513/// @yah:verify("FIXED when a tunnel with no configuration takes its first ingress rule through `yah cloud apply` alone, AND a GET that fails for a genuine reason (revoked token) still aborts before any PUT. Both halves — the first alone re-introduces the wholesale-delete risk ingress.rs:1272-1276 exists to prevent.")
514/// @yah:gotcha("MECHANISM, TRACED AND READ RATHER THAN INFERRED. `CloudflareClient::tunnel_configuration` (oss/yubaba/crates/cloud/src/provider/cloudflare.rs:774) already handles an EMPTY configuration correctly further down — `.and_then(|r| r.get(\"config\")).cloned().unwrap_or_else(|| json!({}))` at :787, plus the non-object guard at :790 whose own comment covers a tunnel whose config reads back as an explicit JSON null. But it checks `self.ok(&resp.success, &resp.errors)?` FIRST, at :781. Cloudflare answers `success: false` / `Configuration for tunnel not found` for a token-form tunnel that has never had a configuration PUT, so that guard fires and the tolerant path at :787 is unreachable for exactly the case it looks like it already covers.")
515/// @yah:gotcha("WHY IT IS A BUG AND NOT A MISSING FEATURE: \"no configuration exists\" is an EMPTY read, not a FAILED one. The function already draws that distinction correctly for a null config body; it simply does not draw it for the API's own not-found answer. Same fact, two doors — one arrives as `result.config = null`, the other as an error object.")
516/// @yah:gotcha("CONSEQUENCE: the only code path that would CREATE a first configuration is gated behind READING a configuration that cannot exist until it has been created. Any cloudflare-tunnel edge whose tunnel was minted without `--hostname` is permanently un-bootstrappable through `yah cloud apply`. The apparent workaround is not one: `yah cloud cf tunnel ensure --hostname` also upserts a DNS CNAME, so it takes a live hostname dark if the connector is not up yet.")
517/// @yah:assumes("NOT VERIFIED FROM THE REPORTING CAMP, because settling it would have meant a write to live Cloudflare: whether the GET returns a 404 status specifically, or HTTP 200 with `success: false` in the body. `Configuration for tunnel not found` is the text that surfaced through the `with_context` at reconciler/ingress.rs:1291. Classify on the Cloudflare error CODE in `resp.errors`, not on the message string and not on an assumed status — a message match would layer a second lexical guess on the first.")
518/// @yah:verify("THE NON-WORKAROUND IS PINNED BY A NEGATIVE RESULT, so nobody re-tries it: run `cf tunnel ensure --hostname` against a tunnel with no configuration, then `yah cloud apply` — apply must still fail `Configuration for tunnel not found`. If a future change makes `cf tunnel ensure` write config, that is a DIFFERENT fix and this assertion should flip deliberately rather than quietly.")
519/// @yah:gotcha("SEVERITY IS HIGHER THAN THIS TICKET FIRST STATED, AND THE OBVIOUS ESCAPE HATCH IS NOT ONE. `yah cloud cf tunnel ensure --name <n> --hostname <h> --api-token-slot cloudflare-legacy-yah` DOES NOT WRITE INGRESS CONFIGURATION — measured 2026-09-14 against tunnel 1b41587a-7298-4f29-97bf-40d51ed79397, not reasoned from the source. It printed `reusing existing 'noisetable-account-staging'`, then `existing connector token preserved in keystore slot`, then went STRAIGHT to the DNS upsert. `yah cloud apply --env staging --service noisetable-api` immediately afterwards failed with the byte-identical `Configuration for tunnel not found`. Its own `--help` confirms the scope in one sentence: ensure the tunnel exists, fetch/store the connector token, OPTIONALLY upsert a CNAME — nothing about configuration.")
520/// @yah:gotcha("THEREFORE THERE IS NO VERB IN THE CAMP THAT CAN WRITE A TUNNEL'S FIRST INGRESS CONFIGURATION. This is not 'apply takes an awkward path to a reachable state' — it is a complete dead end for any tunnel minted without `--hostname`, with no CLI escape hatch at all. `apply` is gated behind the GET this ticket describes; `cf tunnel ensure` never writes config; `yah cloud cf` exposes only `token` and `tunnel`. The only two routes out are (a) fixing the GET's error classification, i.e. this ticket, or (b) a human creating the configuration by hand in the Cloudflare dashboard. A consumer camp cannot unblock itself.")
521/// @yah:gotcha("READ THIS BESIDE THE SIBLING R907-B2 BEFORE DESIGNING THE FIX. The reporting camp hit BOTH of R907's children on the same door in the same hour, and they compose badly: with no config (this ticket) and no connector, a DNS cutover (B2) has no end date for its dark window. Fixing this one first is what makes the other one safe to sequence.")
522/// @yah:handoff("Fixed the root cause: tunnel_configuration (cloudflare.rs) now classifies a failed GET on /cfd_tunnel/{id}/configurations as an EMPTY read (Ok(json!({}))) when the Cloudflare error carries code 1003 (the generic 'not found' family), and still returns Err for every other error (auth, rate-limit, transient). Classification is on the numeric code in resp.errors, not the message text, via a new private CloudflareClient::is_resource_not_found(&Option<Vec<Value>>) helper.")
523/// @yah:handoff("Grounding for code 1003: NOT in Cloudflare's own published API reference (checked developers.cloudflare.com's error-response schema for this endpoint 2026-09-14 — it documents only the generic {code, message} shape, no per-condition table). Grounded instead from convergent third-party Cloudflare API clients that hardcode it for this exact family: vana-com/vana-connect's gcp.ts comment '1003: tunnel not found', mandar-karhade/dockflare's error table '1003 | Zone not found', plus independent test fixtures (ratazzi/coulson, MauroDruwel/TunnelDashDesktop) using 1003 for 'Invalid or missing account id' / 'Account not found'. Full sourcing is in the doc comment on is_resource_not_found (cloudflare.rs, just above it). Confirming against a live token was out of scope (ticket forbids live CF calls) and would be worth a follow-up if the sourcing is judged insufficient.")
524/// @yah:handoff("Left reconciler/ingress.rs:1272-1276 (the abort-before-PUT guard) untouched, per the ticket's explicit constraint — only read it, did not edit it.")
525/// @yah:handoff("Sweep for the same latent shape found two more sites with the identical bug: tunnel_dns_records and tunnel_dns_drift (both in cloudflare.rs) each GET a tunnel's configuration and previously called self.ok(...)? unconditionally, so ANY unconfigured tunnel among many would abort the whole DNS-records listing / drift report with an Err instead of contributing zero ingress hostnames for that tunnel. Fixed both the same way, reusing is_resource_not_found. No other cf_get call site in the file shares this shape — the rest are LIST endpoints (naturally return an empty array when nothing exists, not an error) or are already documented as tolerant (upsert_index_rewrite already treats a missing ruleset entrypoint as empty).")
526/// @yah:handoff("Added 3 unit tests pinning is_resource_not_found in both directions plus the None case: resource_not_found_classifies_missing_tunnel_config_as_empty (code 1003 -> true), resource_not_found_does_not_classify_auth_error_as_empty (code 10000 -> false), resource_not_found_false_when_no_errors_present. Tests exercise the classifier directly rather than the full async tunnel_configuration end-to-end, because the crate has no HTTP-mocking dev-dependency (checked Cargo.toml dev-dependencies: tempfile, serde_yaml, tokio — no mockito/wiremock) and adding one is a bigger change than this leaf fix warrants; the classifier is the entire decision this fix adds, so it's what's under test.")
527/// @yah:verify("cargo test -p yah-cloud --lib -- cloudflare: 29 pass / 0 fail (includes the 3 new tests), no skew reported.")
528/// @yah:verify("cargo test -p yah-cloud --lib (whole crate): 1200 pass / 2 fail / 4 ignored. The 2 failures (reconciler::mesofact_static::tests::a_mount_extends_the_prefix_and_is_slash_insensitive, bundled_worker_walks_the_route_table) are in mesofact_static.rs, which was already modified and uncommitted by a peer at session start (per git status) and is untouched by this ticket's diff -- confirmed unrelated via git diff --stat on that file.")
529/// @yah:verify("cargo check -p yah-cloud --lib: clean (1 pre-existing unrelated warning in mesofact_static.rs).")
530/// @yah:verify("CloudflareClient::ok now surfaces the numeric Cloudflare error code in its returned message (e.g. \"Cloudflare error 1003: Configuration for tunnel not found\") instead of dropping it, so the FIRST live run of this fix from the reporting camp will either confirm 1003 or name the real code plainly, closing the sourcing gap.")
531/// @yah:assumes("The not-found classification code (1003) is UNCONFIRMED against a live Cloudflare account — sourced only from convergent third-party API clients (see is_resource_not_found's doc comment), not from Cloudflare's own docs or a live call. Disproved by: a live `yah cloud apply` against an unconfigured tunnel either (a) succeeding, confirming 1003, or (b) still failing with a message now showing a different numeric code, in which case update CF_ERR_NOT_FOUND in is_resource_not_found to that code.")
532/// @yah:handoff("Addendum landed: CloudflareClient::ok now renders every error entry as \\\"Cloudflare error <code>: <message>\\\" (joined with \\\"; \\\" for multiple), instead of dropping the code -- so the reporting camp's next live run will show the real code plainly. Added 3 more unit tests (ok_error_message_carries_the_cloudflare_code, ok_error_message_joins_multiple_errors, ok_error_message_falls_back_when_no_errors_present) via CloudflareClient::new(\\\"test-token\\\"). Grepped the crate for callers/tests asserting the old bare-message string -- none exist outside cloudflare.rs itself; reconciler/ingress.rs:1291 only wraps the error with with_context, doesn't match on it.")
533/// @yah:verify("cargo check -p yah-cloud --lib: clean, no errors, after the ok() change. cargo test -p yah-cloud --lib currently fails to COMPILE (19 errors, all `missing field build in initializer of MirrorConfig` across mesofact_static.rs/pond.rs/local_process.rs/static_asset.rs/mod.rs/mesofact_bundle.rs/cloudflare_worker.rs/ingress.rs) -- confirmed zero of those 19 reference provider/cloudflare.rs; this is a live peer's in-flight MirrorConfig migration (git status showed mesofact_static.rs already modified at session start), not this ticket's diff. The cloudflare-module test run that passed 29/29 (prior handoff entry) was taken before that peer's edit progressed to this state; could not re-run cloudflare-scoped tests after the ok() addendum because the whole test binary now fails to link for an unrelated reason. Re-run `cargo test -p yah-cloud --lib -- cloudflare` once that peer's migration lands.")
534/// @yah:handoff("Addendum complete: CloudflareClient::ok surfaces the numeric Cloudflare error code in its message text instead of discarding it, so the 1003 guess is now self-diagnosing on the first live run. See appended handoff/verify entries for detail and the peer-breakage caveat on re-running the test suite.")
535/// @yah:verify("RE-VERIFIED after the peer's MirrorConfig migration landed: cargo test -p yah-cloud --lib (whole crate) = 1213 pass / 0 fail / 4 ignored, no skew (baseline was 1200/2/4 with the 2 failures being that same peer's in-flight breakage, now resolved). cargo test -p yah-cloud --lib -- cloudflare = 32 pass / 0 fail, includes all 6 new tests (3 for is_resource_not_found, 3 for the ok() addendum).")
536///
537/// @yah:ticket(R907-B2, "yah cloud cf has no DNS-record verb, so a tunnel hostname whose name already holds A records cannot be moved")
538/// @yah:status(review)
539/// @yah:at(2026-09-14T20:23:53Z)
540/// @yah:assignee(agent:bundle-anthropic-miravel)
541/// @yah:parent(R907)
542/// @yah:severity(medium)
543/// @yah:next("Tier: Cleric — small surface to add, but the sequencing semantics are the whole point and a naive delete-then-create verb would be worse than none.")
544/// @yah:next("THE ORDERING HAZARD IS THE DESIGN CONSTRAINT, NOT A CAVEAT. Deleting the incumbent A records BEFORE a connector is up takes a live hostname dark, and when the tunnel also has no configuration (sibling R907-B1) that window has NO END DATE — nothing is standing by to answer. So a DNS verb here should not be a thin delete+create wrapper. Prefer a cutover that is atomic, or gated on connector readiness (tunnel has >=1 healthy connector AND an ingress rule for the hostname) before it touches the incumbent record, with the unsafe ordering available only behind an explicit flag that names the outage it buys.")
545/// @yah:next("Worth deciding deliberately: whether the verb belongs on `yah cloud cf` as a general record CRUD, or whether the tunnel cutover should stay a single higher-level operation (`cf tunnel cutover --hostname`) that owns the safe ordering internally. The second is harder to misuse and matches how `cf tunnel ensure` already hides the record shape from the caller.")
546/// @yah:next("Fix R907-B1 FIRST or in the same change. Sequencing this one alone still leaves a consumer unable to get a configuration onto the tunnel, so the connector-readiness gate above could never be satisfied.")
547/// @yah:verify("REPRO: point `cf tunnel ensure --name <tunnel> --hostname <h> --api-token-slot cloudflare-legacy-yah` at a hostname that currently resolves via A records. Observed: the tunnel-reuse and token-preserve lines print normally, then `Error: upserting CNAME <h> → <id>.cfargotunnel.com: An A, AAAA, or CNAME record with that host already exists.` The API token is not the issue — cloudflare-legacy-yah carries both Tunnel:Edit and DNS:Edit, and the same invocation's tunnel half succeeded in the same run.")
548/// @yah:verify("FIXED when a hostname on A records can be moved onto a tunnel through the CLI alone, AND the safe-ordering property holds — i.e. an attempt to cut over to a tunnel with no healthy connector either refuses or is explicitly opted into. Both halves: the first alone just automates the outage.")
549/// @yah:gotcha("THE SURFACE IS TWO SUBCOMMANDS WIDE. `yah cloud cf --help` lists exactly `token` (mint scoped account-owned tokens) and `tunnel` (idempotent ensure). There is no verb that reads, creates, updates or deletes a DNS record. The only DNS write anywhere in `yah cloud cf` is the CNAME upsert bolted onto `cf tunnel ensure --hostname`.")
550/// @yah:gotcha("AND THAT UPSERT IS NOT A GENERAL UPSERT. Cloudflare forbids a CNAME coexisting with an A/AAAA record at the same name, so the call fails with `An A, AAAA, or CNAME record with that host already exists` whenever the incumbent record is not ALREADY a CNAME. The `--help` text's promise that it 'upserts a CNAME ... so the public hostname is wired even before the origin daemon comes up' therefore holds only for a name that is unused or already CNAME'd — precisely NOT the migration case, where a hostname is being moved onto a tunnel from something else.")
551/// @yah:gotcha("CONCRETE INSTANCE, MEASURED 2026-09-14 from the noisetable camp: api-staging.noisetable.com holds A records 51.81.85.145 and 45.32.194.254 (its production passway edge, currently serving 200), and cannot be moved to tunnel 1b41587a-7298-4f29-97bf-40d51ed79397 because those two records must be DELETED first. With no DNS verb in the CLI, the only route is hand-editing the Cloudflare dashboard or hand-calling the API — which is exactly the class of action a camp's own rules put out of bounds for an agent, so the cutover stalls with no in-tool path.")
552/// @yah:handoff("Built the higher-level operation as directed: `yah cloud cf tunnel cutover --hostname <h> --name <tunnel> [--force-dark-window]` (app/yah/cli/src/cloud.rs). No general DNS CRUD exposed on the CLI -- the DNS methods it needed (list_dns_records, delete_dns_records_matching, upsert_dns_record) already existed as CloudflareClient methods from prior work (R859-F1); only one new client method was added: CloudflareClient::tunnel_conn_state(account_id, tunnel_id) -> Result<TunnelConnState>, a thin wrapper around the existing (private) list_tunnels_meta that surfaces just the connector state for one tunnel by id.")
553/// @yah:handoff("Ordering: lists existing DNS records at the hostname first. If already a correct CNAME with no A/AAAA left beside it, no-op success (idempotent, matches `ensure`'s contract). Otherwise, if A/AAAA records are present, gates on readiness (tunnel_conn_state == Active AND tunnel_configuration's ingress array contains a rule for the hostname) before deleting them -- refuses with a message naming both missing preconditions and the outage it would cause, unless --force-dark-window is passed (which prints an explicit warning naming the outage before proceeding). Deletes only the incumbent A/AAAA records (delete_dns_records_matching, scoped by type+content so a round-robin sibling isn't touched), then upserts the CNAME.")
554/// @yah:handoff("Composes R907-B1 directly: the readiness check's ingress-rule half calls tunnel_configuration, which is what B1 made tolerant of a never-configured tunnel -- so a tunnel with a connector up but no ingress yet correctly reads as not-ready (has_ingress_rule=false) instead of erroring out of the whole cutover attempt.")
555/// @yah:handoff("Pure decision logic extracted and unit-tested directly (not through the async handler, matching how the R907-B1 classifier was tested): cutover_already_complete(records, cname_target), tunnel_ready_for_cutover(conn_state, has_ingress_rule), tunnel_ingress_has_hostname(ingress_json, hostname). 7 new tests in a cf_tunnel_cutover_tests module pin: no-op when already correct CNAME alone; not-complete when an A record remains beside the CNAME; not-complete when only A records present; not-complete when CNAME points elsewhere; readiness requires BOTH Active connector state AND an ingress rule (every other TunnelConnState variant fails the gate even with the rule present); ingress-hostname matching is exact and returns false on an empty/malformed config.")
556/// @yah:handoff("Plumbing: DnsRecordDetail and TunnelConnState were not previously re-exported through cloud::provider::mod.rs / cloud::lib.rs (only used internally); added both to the existing pub-use lists so app/yah/cli can name them. No new pub surface beyond that plus tunnel_conn_state.")
557/// @yah:verify("cargo check -p yah --lib: clean (0 errors), 26 pre-existing warnings none of which touch the new code (spot-checked: the only two cloud.rs warnings, a pre-existing unused `mut` at line ~16965 and an unused `stub` fn at line ~3625, predate this change).")
558/// @yah:verify("cargo test -p yah --lib -- cf_tunnel_cutover_tests: 7 pass / 0 fail, no skew on the final run (an earlier run in this same session hit a transient 92-error compile while a peer's kg-store/blake3 and slot_table.rs edits were mid-flight -- confirmed unrelated by re-running clean after they landed; see gotcha).")
559/// @yah:verify("cargo check -p yah-cloud --lib and cargo test -p yah-cloud --lib both re-verified green after this ticket's additions (1213 pass / 0 fail / 4 ignored) -- the new tunnel_conn_state method didn't regress anything.")
560/// @yah:verify("NOT verified live (explicitly out of scope): no live Cloudflare account was hit. The connector-readiness gate's real-world shape (does list_tunnels_meta's status field actually read 'active' the way TunnelConnState::from_cf_status expects for a freshly-up connector) is asserted only by the existing TunnelConnState tests from prior work, not newly re-verified here.")
561/// @yah:gotcha("Mid-session this crate hit two DIFFERENT transient shared-tree compile breaks, both from live peers, both now resolved and neither touching this ticket's files: (1) oss/yubaba/crates/cloud/src/reconciler/mesofact_static.rs's MirrorConfig gained a `build` field mid-session, breaking `cargo test -p yah-cloud --lib` (test-only, `cargo check --lib` stayed clean throughout) -- landed and reverified 1213/0/4. (2) crates/yah/kg-store/src/camp_config.rs referenced `blake3::hash` before its Cargo.toml dependency landed, breaking `cargo check -p yah --lib` entirely for a window -- also since resolved. Neither was touched by this session; noted here only so a reviewer re-running the same commands mid-flight doesn't misattribute a stale failure to R907-B2.")
562/// @yah:verify("RE-RUN 2026-09-14 ~13:30 per @Ashguard:rose's R897 heads-up: an unattributed git stash wiped the tree at 13:09:46, restored ~13:14 (faba56a4 + pop), and this ticket's two flagged commands (cargo check -p yah --lib, cf_tunnel_cutover_tests) had run inside that window on a prior pass. Confirmed by content first that nothing was lost (git show HEAD still had handle_cf_tunnel_cutover / Cutover variant / cf_tunnel_cutover_tests module, all committed pre-wipe in 30c2c02c per the operator's own diff verification) -- nothing to re-author. Re-ran fresh anyway: cargo check -p yah --lib clean (no skew); cf_tunnel_cutover_tests 7/7 pass (no skew); cargo test -p yah-cloud --lib 1213/0/4 (no skew). All three identical to the pre-wipe results.")
563impl CloudflareClient {
564    /// Create a client for the given API token.
565    pub fn new(token: String) -> Self {
566        Self {
567            token,
568            http: reqwest::Client::new(),
569        }
570    }
571
572    /// List accounts the token can access.
573    /// Requires: `Account: Read`.
574    pub async fn list_accounts(&self) -> Result<Vec<CfAccountInfo>> {
575        #[derive(Deserialize)]
576        struct Entry {
577            id: String,
578            name: String,
579        }
580        let resp: CfPage<Entry> = self.cf_get("/accounts").await?;
581        self.ok(&resp.success, &resp.errors)?;
582        Ok(resp
583            .result
584            .unwrap_or_default()
585            .into_iter()
586            .map(|a| CfAccountInfo {
587                id: a.id,
588                name: a.name,
589            })
590            .collect())
591    }
592
593    /// List non-deleted Cloudflare Tunnels in `account_id` with connection state.
594    /// Requires: `Cloudflare Tunnel: Read`.
595    async fn list_tunnels_meta(&self, account_id: &str) -> Result<Vec<TunnelMeta>> {
596        let resp: CfPage<CfTunnel> = self
597            .cf_get(&format!(
598                "/accounts/{account_id}/cfd_tunnel?is_deleted=false"
599            ))
600            .await?;
601        self.ok(&resp.success, &resp.errors)?;
602        Ok(resp
603            .result
604            .unwrap_or_default()
605            .into_iter()
606            .map(|t| TunnelMeta {
607                conn_state: t
608                    .status
609                    .as_deref()
610                    .map(TunnelConnState::from_cf_status)
611                    .unwrap_or(TunnelConnState::Unknown),
612                conn_since: t.conns_active_at,
613                id: t.id,
614                name: t.name,
615            })
616            .collect())
617    }
618
619    /// List non-deleted Cloudflare Tunnels in `account_id` as `(id, name)` pairs.
620    /// Requires: `Cloudflare Tunnel: Read`.
621    pub async fn list_tunnels(&self, account_id: &str) -> Result<Vec<(String, String)>> {
622        Ok(self
623            .list_tunnels_meta(account_id)
624            .await?
625            .into_iter()
626            .map(|m| (m.id, m.name))
627            .collect())
628    }
629
630    /// Live connection state for one tunnel, by id — the connector-readiness
631    /// half of a safe DNS cutover (R907-B2). `Unknown` when the tunnel
632    /// doesn't appear in the account's tunnel list at all (deleted, wrong
633    /// account) rather than erroring, since a caller gating on `Active` — the
634    /// only value that clears the gate — treats every other variant the
635    /// same way.
636    ///
637    /// Requires: `Cloudflare Tunnel: Read`.
638    pub async fn tunnel_conn_state(
639        &self,
640        account_id: &str,
641        tunnel_id: &str,
642    ) -> Result<TunnelConnState> {
643        Ok(self
644            .list_tunnels_meta(account_id)
645            .await?
646            .into_iter()
647            .find(|t| t.id == tunnel_id)
648            .map(|t| t.conn_state)
649            .unwrap_or(TunnelConnState::Unknown))
650    }
651
652    /// Collect CNAME records for all tunnels across all accessible accounts.
653    ///
654    /// Walks accounts → tunnels → ingress configurations. Returns an empty
655    /// vec when the token has no tunnels or no configured ingress hostnames.
656    pub async fn tunnel_dns_records(&self) -> Result<Vec<TunnelDnsRecord>> {
657        let mut records = Vec::new();
658        let accounts = self.list_accounts().await?;
659
660        for account in &accounts {
661            let tunnels = self.list_tunnels(&account.id).await?;
662            for (tunnel_id, tunnel_name) in &tunnels {
663                let cname_target = format!("{tunnel_id}.cfargotunnel.com");
664                let path = format!(
665                    "/accounts/{}/cfd_tunnel/{tunnel_id}/configurations",
666                    account.id
667                );
668                let config_resp: CfSingle<CfTunnelConfig> = self.cf_get(&path).await?;
669                let ingress = if !config_resp.success
670                    && Self::is_resource_not_found(&config_resp.errors)
671                {
672                    Vec::new()
673                } else {
674                    self.ok(&config_resp.success, &config_resp.errors)?;
675                    config_resp
676                        .result
677                        .and_then(|c| c.config)
678                        .and_then(|c| c.ingress)
679                        .unwrap_or_default()
680                };
681                for rule in ingress {
682                    if let Some(hostname) = rule.hostname.filter(|h| !h.is_empty()) {
683                        records.push(TunnelDnsRecord {
684                            tunnel_name: tunnel_name.clone(),
685                            hostname,
686                            cname_target: cname_target.clone(),
687                        });
688                    }
689                }
690            }
691        }
692        Ok(records)
693    }
694
695    /// Compute DNS drift for every tunnel ingress hostname, enriched with live
696    /// connector connection state.
697    ///
698    /// The declared side is the tunnel ingress config; the live side is the
699    /// zone's DNS records. Each ingress hostname is classified:
700    /// [`TunnelDriftState::Synced`] when a record points at the tunnel's CNAME
701    /// target, `Missing` when none exists, `Mismatch` when one points elsewhere.
702    ///
703    /// Connection state (`conn_state` / `conn_since`) comes from the `status`
704    /// and `conns_active_at` fields on the tunnel list response — fetched in the
705    /// same pass as the ingress configs to avoid an extra `list_accounts` round-trip.
706    ///
707    /// Degrades gracefully — a hostname whose zone can't be resolved or read is
708    /// reported `ZoneUnknown` rather than failing the whole report. Returns an
709    /// empty vec when the token has no tunnels or no ingress hostnames.
710    ///
711    /// Requires: `Cloudflare Tunnel: Read`, `Zone: Read`, `DNS: Read`.
712    pub async fn tunnel_dns_drift(&self) -> Result<Vec<TunnelDriftRow>> {
713        // Walk accounts → tunnels (with live connection state) → ingress configs
714        // in one pass, collecting both declared DNS records and conn-state in a
715        // single list_accounts round-trip.
716        let accounts = self.list_accounts().await?;
717        let mut declared: Vec<TunnelDnsRecord> = Vec::new();
718        let mut conn_by_tunnel: std::collections::HashMap<
719            String,
720            (TunnelConnState, Option<String>),
721        > = Default::default();
722
723        for account in &accounts {
724            let metas = self.list_tunnels_meta(&account.id).await?;
725            for meta in &metas {
726                conn_by_tunnel.insert(
727                    meta.name.clone(),
728                    (meta.conn_state, meta.conn_since.clone()),
729                );
730                let cname_target = format!("{}.cfargotunnel.com", meta.id);
731                let path = format!(
732                    "/accounts/{}/cfd_tunnel/{}/configurations",
733                    account.id, meta.id
734                );
735                let config_resp: CfSingle<CfTunnelConfig> = self.cf_get(&path).await?;
736                let ingress = if !config_resp.success
737                    && Self::is_resource_not_found(&config_resp.errors)
738                {
739                    Vec::new()
740                } else {
741                    self.ok(&config_resp.success, &config_resp.errors)?;
742                    config_resp
743                        .result
744                        .and_then(|c| c.config)
745                        .and_then(|c| c.ingress)
746                        .unwrap_or_default()
747                };
748                for rule in ingress {
749                    if let Some(hostname) = rule.hostname.filter(|h| !h.is_empty()) {
750                        declared.push(TunnelDnsRecord {
751                            tunnel_name: meta.name.clone(),
752                            hostname,
753                            cname_target: cname_target.clone(),
754                        });
755                    }
756                }
757            }
758        }
759
760        if declared.is_empty() {
761            return Ok(Vec::new());
762        }
763
764        // Resolve zones once. A permission failure leaves the set empty, so
765        // every hostname falls through to `ZoneUnknown` instead of erroring.
766        let zones = self.list_zones().await.unwrap_or_default();
767
768        let mut rows = Vec::with_capacity(declared.len());
769        for rec in declared {
770            let (conn_state, conn_since) = conn_by_tunnel
771                .remove(&rec.tunnel_name)
772                .unwrap_or((TunnelConnState::Unknown, None));
773            let (state, live_target) = match best_zone_for(&rec.hostname, &zones) {
774                None => (TunnelDriftState::ZoneUnknown, None),
775                Some(zone_id) => match self.dns_records_named(zone_id, &rec.hostname).await {
776                    Ok(live) => classify_tunnel_drift(&rec.hostname, &rec.cname_target, &live),
777                    Err(_) => (TunnelDriftState::ZoneUnknown, None),
778                },
779            };
780            rows.push(TunnelDriftRow {
781                tunnel_name: rec.tunnel_name,
782                hostname: rec.hostname,
783                expected_target: rec.cname_target,
784                state,
785                live_target,
786                conn_state,
787                conn_since,
788            });
789        }
790        Ok(rows)
791    }
792
793    /// List zones the token can read, as `(zone_id, zone_name)` pairs.
794    /// Requires: `Zone: Read`.
795    pub async fn list_zones(&self) -> Result<Vec<(String, String)>> {
796        #[derive(Deserialize)]
797        struct ZoneEntry {
798            id: String,
799            name: String,
800        }
801        let resp: CfPage<ZoneEntry> = self.cf_get("/zones?per_page=50").await?;
802        self.ok(&resp.success, &resp.errors)?;
803        Ok(resp
804            .result
805            .unwrap_or_default()
806            .into_iter()
807            .map(|z| (z.id, z.name))
808            .collect())
809    }
810
811    /// Fetch DNS records in `zone_id` whose name exactly matches `name`.
812    /// Requires: `DNS: Read`.
813    async fn dns_records_named(&self, zone_id: &str, name: &str) -> Result<Vec<CfDnsRecord>> {
814        let resp: CfPage<CfDnsRecord> = self
815            .cf_get(&format!("/zones/{zone_id}/dns_records?name={name}"))
816            .await?;
817        self.ok(&resp.success, &resp.errors)?;
818        Ok(resp.result.unwrap_or_default())
819    }
820
821    /// Create a new Named Tunnel under `account_id` and return the connector token.
822    /// Requires: `Cloudflare Tunnel: Edit`.
823    pub async fn create_tunnel(&self, account_id: &str, name: &str) -> Result<CreateTunnelResult> {
824        use base64::Engine as _;
825
826        let mut secret_bytes = [0u8; 32];
827        getrandom::getrandom(&mut secret_bytes)
828            .map_err(|e| anyhow!("generate tunnel secret: {e}"))?;
829        let tunnel_secret = base64::engine::general_purpose::STANDARD.encode(secret_bytes);
830
831        #[derive(Serialize)]
832        struct CreateBody<'a> {
833            name: &'a str,
834            tunnel_secret: String,
835        }
836        #[derive(Deserialize)]
837        struct CreatedTunnel {
838            id: String,
839            name: String,
840        }
841        let create_resp: CfSingle<CreatedTunnel> = self
842            .cf_post(
843                &format!("/accounts/{account_id}/cfd_tunnel"),
844                &CreateBody {
845                    name,
846                    tunnel_secret,
847                },
848            )
849            .await?;
850        self.ok(&create_resp.success, &create_resp.errors)?;
851        let created = create_resp
852            .result
853            .ok_or_else(|| anyhow!("tunnel create: no result in response"))?;
854
855        // Fetch the connector JWT.
856        #[derive(Deserialize)]
857        struct TokenResp {
858            success: bool,
859            result: Option<String>,
860            errors: Option<Vec<serde_json::Value>>,
861        }
862        let token_resp: TokenResp = self
863            .cf_get(&format!(
864                "/accounts/{account_id}/cfd_tunnel/{}/token",
865                created.id
866            ))
867            .await?;
868        self.ok(&token_resp.success, &token_resp.errors)?;
869        let connector_token = token_resp
870            .result
871            .ok_or_else(|| anyhow!("no connector token in response"))?;
872
873        Ok(CreateTunnelResult {
874            cname_target: format!("{}.cfargotunnel.com", created.id),
875            tunnel_id: created.id,
876            tunnel_name: created.name,
877            connector_token,
878        })
879    }
880
881    /// Read a tunnel's remotely-managed configuration body as raw JSON
882    /// (R594-F11).
883    ///
884    /// Returns the `result.config` object — the thing a PUT round-trips — or an
885    /// empty object when the tunnel has never been configured. Deliberately
886    /// untyped: the `ingress` list is the only key this crate owns, and every
887    /// sibling (`warp-routing`, `originRequest`, …) must survive a
888    /// read-modify-write untouched.
889    ///
890    /// Requires: `Cloudflare Tunnel: Read`.
891    pub async fn tunnel_configuration(
892        &self,
893        account_id: &str,
894        tunnel_id: &str,
895    ) -> Result<serde_json::Value> {
896        let path = format!("/accounts/{account_id}/cfd_tunnel/{tunnel_id}/configurations");
897        let resp: CfSingle<serde_json::Value> = self.cf_get(&path).await?;
898        if !resp.success && Self::is_resource_not_found(&resp.errors) {
899            return Ok(serde_json::json!({}));
900        }
901        self.ok(&resp.success, &resp.errors)?;
902        let config = resp
903            .result
904            .as_ref()
905            .and_then(|r| r.get("config"))
906            .cloned()
907            .unwrap_or_else(|| serde_json::json!({}));
908        // A tunnel configured with an explicit JSON `null` config reads back as
909        // Value::Null, which has no object to insert `ingress` into.
910        Ok(if config.is_object() {
911            config
912        } else {
913            serde_json::json!({})
914        })
915    }
916
917    /// Replace a tunnel's remotely-managed configuration (R594-F11).
918    ///
919    /// `config` is the whole config body, not a patch — Cloudflare replaces it
920    /// wholesale, which is why callers must GET-merge-PUT rather than PUT a
921    /// freshly-built list. See
922    /// [`reconciler::ingress::ensure_tunnel_ingress`](crate::reconciler::ingress::ensure_tunnel_ingress).
923    ///
924    /// Requires: `Cloudflare Tunnel: Edit`.
925    pub async fn put_tunnel_configuration(
926        &self,
927        account_id: &str,
928        tunnel_id: &str,
929        config: &serde_json::Value,
930    ) -> Result<()> {
931        #[derive(Serialize)]
932        struct ConfigBody<'a> {
933            config: &'a serde_json::Value,
934        }
935        let path = format!("/accounts/{account_id}/cfd_tunnel/{tunnel_id}/configurations");
936        let resp: CfSingle<serde_json::Value> = self.cf_put(&path, &ConfigBody { config }).await?;
937        self.ok(&resp.success, &resp.errors)?;
938        Ok(())
939    }
940
941    /// Create a new R2 bucket under `account_id`.
942    /// Requires: `Account: Cloudflare R2: Edit`.
943    pub async fn create_r2_bucket(
944        &self,
945        account_id: &str,
946        bucket_name: &str,
947    ) -> Result<CreateR2BucketResult> {
948        #[derive(Serialize)]
949        struct CreateBody<'a> {
950            name: &'a str,
951        }
952        let resp: CfSingle<serde_json::Value> = self
953            .cf_post(
954                &format!("/accounts/{account_id}/r2/buckets"),
955                &CreateBody { name: bucket_name },
956            )
957            .await?;
958        self.ok(&resp.success, &resp.errors)?;
959
960        Ok(CreateR2BucketResult {
961            endpoint: format!("https://{account_id}.r2.cloudflarestorage.com"),
962            name: bucket_name.to_string(),
963        })
964    }
965
966    /// Resolve a zone name (e.g. `"yah.dev"`) to its Cloudflare zone ID.
967    /// Requires: `Zone: Read`.
968    pub async fn zone_id_for_name(&self, zone_name: &str) -> Result<String> {
969        #[derive(Deserialize)]
970        struct ZoneEntry {
971            id: String,
972            name: String,
973        }
974        let resp: CfPage<ZoneEntry> = self.cf_get(&format!("/zones?name={zone_name}")).await?;
975        self.ok(&resp.success, &resp.errors)?;
976        resp.result
977            .unwrap_or_default()
978            .into_iter()
979            .find(|z| z.name == zone_name)
980            .map(|z| z.id)
981            .ok_or_else(|| anyhow!("no Cloudflare zone found for name {zone_name:?}"))
982    }
983
984    /// Purge content by cache tags from a zone.
985    ///
986    /// Cache tags must be applied to responses via the `Cache-Tag` header or
987    /// Cloudflare page rules. Returns `Ok(())` when all tags are queued for
988    /// purge. Requires: `Zone: Cache Purge`.
989    pub async fn purge_cache_tags(&self, zone_id: &str, tags: &[String]) -> Result<()> {
990        if tags.is_empty() {
991            return Ok(());
992        }
993        #[derive(Serialize)]
994        struct PurgeBody<'a> {
995            tags: &'a [String],
996        }
997        let resp: CfSingle<serde_json::Value> = self
998            .cf_post(
999                &format!("/zones/{zone_id}/purge_cache"),
1000                &PurgeBody { tags },
1001            )
1002            .await?;
1003        self.ok(&resp.success, &resp.errors)
1004    }
1005
1006    /// Upsert the Transform Rule that rewrites `GET /` → `/index.html` on the
1007    /// zone, identified by the stable description tag `"yah:static-index"`.
1008    ///
1009    /// Idempotent: fetches the existing `http_request_transform` entrypoint,
1010    /// drops any prior `"yah:static-index"` rule, appends the current one,
1011    /// and PUTs the merged list back. Treats a missing entrypoint (no rules
1012    /// yet) as an empty list.
1013    ///
1014    /// Requires: `Zone: Transform Rules: Edit`.
1015    pub async fn upsert_index_rewrite(&self, zone_id: &str) -> Result<()> {
1016        const RULE_DESC: &str = "yah:static-index";
1017        let path = format!("/zones/{zone_id}/rulesets/phases/http_request_transform/entrypoint");
1018
1019        // Fetch existing rules; a missing entrypoint is not an error.
1020        let existing: Vec<serde_json::Value> = {
1021            #[derive(Deserialize)]
1022            struct Rs {
1023                rules: Option<Vec<serde_json::Value>>,
1024            }
1025            match self.cf_get::<CfSingle<Rs>>(&path).await {
1026                Ok(resp) if resp.success => resp.result.and_then(|r| r.rules).unwrap_or_default(),
1027                _ => Vec::new(),
1028            }
1029        };
1030
1031        // Keep every rule except the one we manage, then append ours.
1032        let mut rules: Vec<serde_json::Value> = existing
1033            .into_iter()
1034            .filter(|r| r.get("description").and_then(|v| v.as_str()) != Some(RULE_DESC))
1035            .collect();
1036        rules.push(serde_json::json!({
1037            "action": "rewrite",
1038            "description": RULE_DESC,
1039            "expression": "(http.request.uri.path eq \"/\")",
1040            "action_parameters": {
1041                "uri": { "path": { "value": "/index.html" } }
1042            },
1043            "enabled": true
1044        }));
1045
1046        let resp: CfSingle<serde_json::Value> = self
1047            .cf_put(&path, &serde_json::json!({ "rules": rules }))
1048            .await?;
1049        self.ok(&resp.success, &resp.errors)
1050    }
1051
1052    /// Deploy an ES-module Worker script with typed bindings for runtime config.
1053    ///
1054    /// Each entry in `bindings` becomes one `metadata.bindings[…]` declaration in
1055    /// the upload payload — see [`WorkerBinding`] for the supported variants
1056    /// (plain_text config, R2 bucket references).
1057    ///
1058    /// Uses a manual multipart/form-data upload (CF Workers API requires multipart
1059    /// when metadata/bindings are attached). Idempotent: re-uploading the same script
1060    /// is safe but costs one CF API round-trip — callers should hash-guard this.
1061    ///
1062    /// Requires: `Workers Scripts: Edit` (account-scoped).
1063    pub async fn deploy_worker_script(
1064        &self,
1065        account_id: &str,
1066        script_name: &str,
1067        script_js: &str,
1068        bindings: &[WorkerBinding<'_>],
1069    ) -> Result<WorkerDeployResult> {
1070        let url = format!("{CF_API}/accounts/{account_id}/workers/scripts/{script_name}");
1071        let (content_type, body) = build_worker_multipart(script_js, bindings);
1072        let resp = self
1073            .http
1074            .put(&url)
1075            .header("Authorization", format!("Bearer {}", self.token))
1076            .header("Content-Type", content_type)
1077            .body(body)
1078            .send()
1079            .await
1080            .map_err(|e| anyhow!("PUT {url}: {e}"))?;
1081        let result: CfSingle<WorkerDeployResult> = resp
1082            .json()
1083            .await
1084            .map_err(|e| anyhow!("PUT {url} parse: {e}"))?;
1085        self.ok(&result.success, &result.errors)?;
1086        result
1087            .result
1088            .ok_or_else(|| anyhow!("deploy worker: no result in response"))
1089    }
1090
1091    /// Upsert a Worker route for `pattern` on `zone_id`, pointing at `script_name`.
1092    ///
1093    /// Idempotent: fetches existing routes, skips PUT/POST when the pattern already
1094    /// points at the right script, updates an existing pattern pointing elsewhere,
1095    /// or creates a new route entry.
1096    ///
1097    /// Requires: `Zone: Workers Routes: Edit` (zone-scoped).
1098    pub async fn upsert_worker_route(
1099        &self,
1100        zone_id: &str,
1101        pattern: &str,
1102        script_name: &str,
1103    ) -> Result<()> {
1104        let list_path = format!("/zones/{zone_id}/workers/routes");
1105
1106        #[derive(Deserialize)]
1107        struct RouteEntry {
1108            id: String,
1109            pattern: String,
1110            #[serde(default)]
1111            script: Option<String>,
1112        }
1113        let list: CfPage<RouteEntry> = self.cf_get(&list_path).await?;
1114        self.ok(&list.success, &list.errors)?;
1115        let routes = list.result.unwrap_or_default();
1116
1117        #[derive(Serialize)]
1118        struct RouteBody<'a> {
1119            pattern: &'a str,
1120            script: &'a str,
1121        }
1122
1123        if let Some(existing) = routes.iter().find(|r| r.pattern == pattern) {
1124            if existing.script.as_deref() == Some(script_name) {
1125                return Ok(());
1126            }
1127            let resp: CfSingle<serde_json::Value> = self
1128                .cf_put(
1129                    &format!("/zones/{zone_id}/workers/routes/{}", existing.id),
1130                    &RouteBody {
1131                        pattern,
1132                        script: script_name,
1133                    },
1134                )
1135                .await?;
1136            self.ok(&resp.success, &resp.errors)
1137        } else {
1138            let resp: CfSingle<serde_json::Value> = self
1139                .cf_post(
1140                    &list_path,
1141                    &RouteBody {
1142                        pattern,
1143                        script: script_name,
1144                    },
1145                )
1146                .await?;
1147            self.ok(&resp.success, &resp.errors)
1148        }
1149    }
1150
1151    /// Idempotently attach `hostname` (e.g. `cr.yah.dev`) as a Workers Custom
1152    /// Domain on `script_name`. Custom Domains route every request for the
1153    /// hostname into the Worker — distinct from a Worker Route, which only
1154    /// matches a URL pattern within an already-proxied zone.
1155    ///
1156    /// Walks the existing Custom Domains list first; if `hostname` is already
1157    /// bound to `script_name` on `zone_id`, returns Ok without an extra PUT.
1158    /// Otherwise PUTs `/accounts/{account_id}/workers/domains`, which CF treats
1159    /// as an upsert keyed on `(hostname, environment)`.
1160    ///
1161    /// Requires: `Workers Scripts: Edit` (account-scoped).
1162    pub async fn upsert_worker_custom_domain(
1163        &self,
1164        account_id: &str,
1165        zone_id: &str,
1166        hostname: &str,
1167        script_name: &str,
1168    ) -> Result<()> {
1169        let list_path = format!("/accounts/{account_id}/workers/domains");
1170
1171        #[derive(Deserialize)]
1172        struct DomainEntry {
1173            #[serde(default)]
1174            hostname: Option<String>,
1175            #[serde(default)]
1176            service: Option<String>,
1177            #[serde(default, rename = "zone_id")]
1178            zone_id: Option<String>,
1179        }
1180        let list: CfPage<DomainEntry> = self.cf_get(&list_path).await?;
1181        self.ok(&list.success, &list.errors)?;
1182        let domains = list.result.unwrap_or_default();
1183        if domains.iter().any(|d| {
1184            d.hostname.as_deref() == Some(hostname)
1185                && d.service.as_deref() == Some(script_name)
1186                && d.zone_id.as_deref() == Some(zone_id)
1187        }) {
1188            return Ok(());
1189        }
1190
1191        #[derive(Serialize)]
1192        struct DomainBody<'a> {
1193            environment: &'a str,
1194            hostname: &'a str,
1195            service: &'a str,
1196            zone_id: &'a str,
1197        }
1198        let resp: CfSingle<serde_json::Value> = self
1199            .cf_put(
1200                &list_path,
1201                &DomainBody {
1202                    environment: "production",
1203                    hostname,
1204                    service: script_name,
1205                    zone_id,
1206                },
1207            )
1208            .await?;
1209        self.ok(&resp.success, &resp.errors)
1210    }
1211
1212    /// Delete an R2 bucket under `account_id`.
1213    ///
1214    /// Cloudflare's management API handles non-empty buckets — objects do not
1215    /// need to be drained first. Returns `Ok(())` on success, `Err` if the
1216    /// API returns a failure (including "bucket not found" — callers that
1217    /// need idempotency should probe [`Self::list_r2_buckets`] first).
1218    ///
1219    /// Requires: `Account: Cloudflare R2: Edit`.
1220    pub async fn delete_r2_bucket(&self, account_id: &str, bucket_name: &str) -> Result<()> {
1221        let resp: CfSingle<serde_json::Value> = self
1222            .cf_delete(&format!("/accounts/{account_id}/r2/buckets/{bucket_name}"))
1223            .await?;
1224        self.ok(&resp.success, &resp.errors)
1225    }
1226
1227    /// Idempotently upsert a DNS record in `zone_id`. Fetches existing records
1228    /// with the same name and type: updates the first match if found, creates
1229    /// a new record otherwise. Returns the provider-issued record ID.
1230    ///
1231    /// Requires: `DNS: Edit` (zone-scoped).
1232    pub async fn upsert_dns_record(
1233        &self,
1234        zone_id: &str,
1235        name: &str,
1236        record_type: &str,
1237        content: &str,
1238        ttl: u32,
1239        proxied: bool,
1240    ) -> Result<String> {
1241        self.upsert_dns_record_matching(zone_id, name, record_type, content, ttl, proxied, false)
1242            .await
1243    }
1244
1245    /// [`upsert_dns_record`](Self::upsert_dns_record) with control over what
1246    /// counts as "the existing record" — R859-F1.
1247    ///
1248    /// `match_content = false` reproduces the original behaviour: the first
1249    /// record sharing `name` + `record_type` is updated in place. That is right
1250    /// for a single-valued name (one CNAME at `cdn.yah.dev`) and **wrong for a
1251    /// multi-valued RRset**: adding the second A record of a round-robin apex
1252    /// would rewrite the first one's content, silently halving the origin set
1253    /// to one box.
1254    ///
1255    /// `match_content = true` keys the lookup on `(name, type, content)`, so
1256    /// the call means "ensure exactly this record exists" — a no-op update when
1257    /// it already does, a create when it does not, and never a mutation of a
1258    /// sibling record at the same name.
1259    #[allow(clippy::too_many_arguments)]
1260    pub async fn upsert_dns_record_matching(
1261        &self,
1262        zone_id: &str,
1263        name: &str,
1264        record_type: &str,
1265        content: &str,
1266        ttl: u32,
1267        proxied: bool,
1268        match_content: bool,
1269    ) -> Result<String> {
1270        let existing_id = if match_content {
1271            self.list_dns_records(zone_id, Some(name), Some(record_type))
1272                .await?
1273                .into_iter()
1274                .find(|r| r.content == content)
1275                .map(|r| r.id)
1276        } else {
1277            self.find_dns_record(zone_id, name, record_type).await?
1278        };
1279
1280        #[derive(Serialize)]
1281        struct RecordBody<'a> {
1282            name: &'a str,
1283            #[serde(rename = "type")]
1284            record_type: &'a str,
1285            content: &'a str,
1286            ttl: u32,
1287            proxied: bool,
1288        }
1289        let body = RecordBody {
1290            name,
1291            record_type,
1292            content,
1293            ttl,
1294            proxied,
1295        };
1296
1297        #[derive(Deserialize)]
1298        struct RecordResult {
1299            id: String,
1300        }
1301
1302        if let Some(id) = existing_id {
1303            let resp: CfSingle<RecordResult> = self
1304                .cf_put(&format!("/zones/{zone_id}/dns_records/{id}"), &body)
1305                .await?;
1306            self.ok(&resp.success, &resp.errors)?;
1307            resp.result
1308                .map(|r| r.id)
1309                .ok_or_else(|| anyhow!("dns record update: no id in response"))
1310        } else {
1311            let resp: CfSingle<RecordResult> = self
1312                .cf_post(&format!("/zones/{zone_id}/dns_records"), &body)
1313                .await?;
1314            self.ok(&resp.success, &resp.errors)?;
1315            resp.result
1316                .map(|r| r.id)
1317                .ok_or_else(|| anyhow!("dns record create: no id in response"))
1318        }
1319    }
1320
1321    /// Read the DNS records in `zone_id`, optionally narrowed to one `name`
1322    /// and/or one `record_type` — R859-F1, the read half the `dns.*` catalog
1323    /// was missing.
1324    ///
1325    /// Unlike [`dns_records_named`](Self::dns_records_named) (drift-detection
1326    /// only, name + content) this returns the record **id**, type, ttl and
1327    /// proxy flag, which is what a reconciler needs to decide what to change.
1328    ///
1329    /// Requires: `DNS: Read` (zone-scoped).
1330    pub async fn list_dns_records(
1331        &self,
1332        zone_id: &str,
1333        name: Option<&str>,
1334        record_type: Option<&str>,
1335    ) -> Result<Vec<DnsRecordDetail>> {
1336        let mut query = format!("/zones/{zone_id}/dns_records?per_page=100");
1337        if let Some(n) = name {
1338            query.push_str(&format!("&name={n}"));
1339        }
1340        if let Some(t) = record_type {
1341            query.push_str(&format!("&type={t}"));
1342        }
1343        let resp: CfPage<DnsRecordDetail> = self.cf_get(&query).await?;
1344        self.ok(&resp.success, &resp.errors)?;
1345        Ok(resp.result.unwrap_or_default())
1346    }
1347
1348    /// Delete all DNS records in `zone_id` whose name matches `name` (and
1349    /// optionally `record_type`). Returns the count of records deleted.
1350    /// A count of 0 is not an error — the records may already have been absent.
1351    ///
1352    /// Requires: `DNS: Edit` (zone-scoped).
1353    pub async fn delete_dns_records(
1354        &self,
1355        zone_id: &str,
1356        name: &str,
1357        record_type: Option<&str>,
1358    ) -> Result<u32> {
1359        self.delete_dns_records_matching(zone_id, name, record_type, None)
1360            .await
1361    }
1362
1363    /// [`delete_dns_records`](Self::delete_dns_records) narrowed to records
1364    /// carrying one exact value — R859-F1.
1365    ///
1366    /// A round-robin apex holds several A records under one name, so
1367    /// "delete the A records at `yah.dev`" is not a way to withdraw *one*
1368    /// origin: it takes the live ones with it. `content = Some(ip)` deletes
1369    /// only the withdrawn member.
1370    ///
1371    /// Requires: `DNS: Edit` (zone-scoped).
1372    pub async fn delete_dns_records_matching(
1373        &self,
1374        zone_id: &str,
1375        name: &str,
1376        record_type: Option<&str>,
1377        content: Option<&str>,
1378    ) -> Result<u32> {
1379        let ids: Vec<String> = self
1380            .list_dns_records(zone_id, Some(name), record_type)
1381            .await?
1382            .into_iter()
1383            .filter(|r| content.is_none_or(|c| r.content == c))
1384            .map(|r| r.id)
1385            .collect();
1386        let mut deleted = 0u32;
1387        for id in &ids {
1388            let del: CfSingle<serde_json::Value> = self
1389                .cf_delete(&format!("/zones/{zone_id}/dns_records/{id}"))
1390                .await?;
1391            self.ok(&del.success, &del.errors)?;
1392            deleted += 1;
1393        }
1394        Ok(deleted)
1395    }
1396
1397    /// Fetch the ID of the first DNS record matching `name` and `record_type`.
1398    /// Returns `None` when no matching record exists.
1399    async fn find_dns_record(
1400        &self,
1401        zone_id: &str,
1402        name: &str,
1403        record_type: &str,
1404    ) -> Result<Option<String>> {
1405        #[derive(Deserialize)]
1406        struct RecordEntry {
1407            id: String,
1408        }
1409        let resp: CfPage<RecordEntry> = self
1410            .cf_get(&format!(
1411                "/zones/{zone_id}/dns_records?name={name}&type={record_type}"
1412            ))
1413            .await?;
1414        self.ok(&resp.success, &resp.errors)?;
1415        Ok(resp
1416            .result
1417            .unwrap_or_default()
1418            .into_iter()
1419            .next()
1420            .map(|r| r.id))
1421    }
1422
1423    /// List R2 buckets in `account_id`.
1424    /// Requires: `Account: Cloudflare R2: Read`.
1425    ///
1426    /// Unlike `/accounts` and `/zones`, the R2 list endpoint nests the array
1427    /// under `result.buckets` rather than returning `result` as a bare array,
1428    /// so it needs `CfSingle<R2ListResult>` and not `CfPage<BucketEntry>`.
1429    pub async fn list_r2_buckets(&self, account_id: &str) -> Result<Vec<R2BucketInfo>> {
1430        #[derive(Deserialize)]
1431        struct BucketEntry {
1432            name: String,
1433            #[serde(default)]
1434            location: Option<String>,
1435            #[serde(default)]
1436            creation_date: Option<String>,
1437        }
1438        #[derive(Deserialize)]
1439        struct R2ListResult {
1440            #[serde(default)]
1441            buckets: Vec<BucketEntry>,
1442        }
1443        let resp: CfSingle<R2ListResult> = self
1444            .cf_get(&format!("/accounts/{account_id}/r2/buckets"))
1445            .await?;
1446        self.ok(&resp.success, &resp.errors)?;
1447        Ok(resp
1448            .result
1449            .map(|r| r.buckets)
1450            .unwrap_or_default()
1451            .into_iter()
1452            .map(|b| R2BucketInfo {
1453                name: b.name,
1454                location: b.location,
1455                creation_date: b.creation_date,
1456            })
1457            .collect())
1458    }
1459
1460    /// List R2 custom-domain bindings on `bucket_name`.
1461    ///
1462    /// Requires: `Workers R2 Storage: Read` (or Write, which implies Read).
1463    /// The response nests the array under `result.domains`, mirroring
1464    /// `list_r2_buckets`'s `result.buckets` shape.
1465    pub async fn list_r2_custom_domains(
1466        &self,
1467        account_id: &str,
1468        bucket_name: &str,
1469    ) -> Result<Vec<R2CustomDomain>> {
1470        #[derive(Deserialize)]
1471        struct DomainEntry {
1472            domain: String,
1473            #[serde(default)]
1474            enabled: bool,
1475        }
1476        #[derive(Deserialize)]
1477        struct R2DomainListResult {
1478            #[serde(default)]
1479            domains: Vec<DomainEntry>,
1480        }
1481        let resp: CfSingle<R2DomainListResult> = self
1482            .cf_get(&format!(
1483                "/accounts/{account_id}/r2/buckets/{bucket_name}/domains/custom"
1484            ))
1485            .await?;
1486        self.ok(&resp.success, &resp.errors)?;
1487        Ok(resp
1488            .result
1489            .map(|r| r.domains)
1490            .unwrap_or_default()
1491            .into_iter()
1492            .map(|d| R2CustomDomain {
1493                domain: d.domain,
1494                enabled: d.enabled,
1495            })
1496            .collect())
1497    }
1498
1499    /// Bind a custom domain to an R2 bucket.
1500    ///
1501    /// `zone_id` names the zone that owns `domain` (resolve via
1502    /// [`Self::zone_id_for_name`]). CF requires it so the CNAME write into
1503    /// that zone is authorized — even though the caller is the bucket-side
1504    /// API. CF creates the CNAME automatically; no separate DNS-side call.
1505    /// Requires: `Workers R2 Storage: Edit` (account-scoped).
1506    ///
1507    /// `enabled: true` activates the binding immediately. CF still has to
1508    /// validate ownership + provision TLS in the background — the binding
1509    /// returns success the moment the record is queued, not when the
1510    /// hostname is fully resolvable. First-time DNS propagation is on the
1511    /// order of seconds to a minute.
1512    pub async fn add_r2_custom_domain(
1513        &self,
1514        account_id: &str,
1515        bucket_name: &str,
1516        domain: &str,
1517        zone_id: &str,
1518    ) -> Result<()> {
1519        #[derive(Serialize)]
1520        #[serde(rename_all = "camelCase")]
1521        struct AddBody<'a> {
1522            domain: &'a str,
1523            enabled: bool,
1524            zone_id: &'a str,
1525        }
1526        let resp: CfSingle<serde_json::Value> = self
1527            .cf_post(
1528                &format!("/accounts/{account_id}/r2/buckets/{bucket_name}/domains/custom"),
1529                &AddBody {
1530                    domain,
1531                    enabled: true,
1532                    zone_id,
1533                },
1534            )
1535            .await?;
1536        self.ok(&resp.success, &resp.errors)
1537    }
1538
1539    /// Fetch the account's permission-group catalog as a name → id map, used
1540    /// to resolve [`TokenGrant`] names before minting a token. Requires the
1541    /// calling token to carry `API Tokens: Read` (implied by `Write`).
1542    pub async fn list_permission_group_ids(
1543        &self,
1544        account_id: &str,
1545    ) -> Result<std::collections::BTreeMap<String, String>> {
1546        #[derive(Deserialize)]
1547        struct PgEntry {
1548            id: String,
1549            name: String,
1550        }
1551        let resp: CfPage<PgEntry> = self
1552            .cf_get(&format!(
1553                "/accounts/{account_id}/tokens/permission_groups?per_page=500"
1554            ))
1555            .await?;
1556        self.ok(&resp.success, &resp.errors)?;
1557        Ok(resp
1558            .result
1559            .unwrap_or_default()
1560            .into_iter()
1561            .map(|p| (p.name, p.id))
1562            .collect())
1563    }
1564
1565    /// Mint an **account-owned** API token under `account_id` from `grants`
1566    /// scoped to `account_id` + `zone_id`.
1567    ///
1568    /// Resolves each grant's permission-group name against the live catalog
1569    /// (falling back to its baked-in ID), groups the IDs into account- and
1570    /// zone-scoped policy blocks, and POSTs to `/accounts/{id}/tokens`.
1571    /// Requires the calling token to carry `API Tokens: Write` — but the
1572    /// minted token is bounded by the *account's* access, not the calling
1573    /// token's, so the caller may hold only `API Tokens: Write`.
1574    ///
1575    /// The returned [`CreateTokenResult::value`] is the secret — Cloudflare
1576    /// reveals it only here.
1577    pub async fn create_account_token(
1578        &self,
1579        account_id: &str,
1580        zone_id: &str,
1581        token_name: &str,
1582        grants: &[TokenGrant],
1583    ) -> Result<CreateTokenResult> {
1584        // A catalog lookup failure is non-fatal: build_token_body falls back
1585        // to the baked-in IDs for any name the (empty) catalog can't resolve.
1586        let catalog = self
1587            .list_permission_group_ids(account_id)
1588            .await
1589            .unwrap_or_default();
1590        let body = build_token_body(token_name, account_id, zone_id, grants, &catalog);
1591
1592        let resp: CfSingle<CreateTokenResult> = self
1593            .cf_post(&format!("/accounts/{account_id}/tokens"), &body)
1594            .await?;
1595        self.ok(&resp.success, &resp.errors)?;
1596        resp.result
1597            .ok_or_else(|| anyhow!("token create: no result in response"))
1598    }
1599
1600    // ---------- HTTP helpers ----------
1601
1602    async fn cf_get<T: serde::de::DeserializeOwned>(&self, path: &str) -> Result<T> {
1603        let url = format!("{CF_API}{path}");
1604        let resp = self
1605            .http
1606            .get(&url)
1607            .header("Authorization", format!("Bearer {}", self.token))
1608            .header("Content-Type", "application/json")
1609            .send()
1610            .await
1611            .map_err(|e| anyhow!("GET {url}: {e}"))?;
1612        resp.json::<T>()
1613            .await
1614            .map_err(|e| anyhow!("GET {url} parse: {e}"))
1615    }
1616
1617    async fn cf_post<B: Serialize, T: serde::de::DeserializeOwned>(
1618        &self,
1619        path: &str,
1620        body: &B,
1621    ) -> Result<T> {
1622        let url = format!("{CF_API}{path}");
1623        let resp = self
1624            .http
1625            .post(&url)
1626            .header("Authorization", format!("Bearer {}", self.token))
1627            .header("Content-Type", "application/json")
1628            .json(body)
1629            .send()
1630            .await
1631            .map_err(|e| anyhow!("POST {url}: {e}"))?;
1632        resp.json::<T>()
1633            .await
1634            .map_err(|e| anyhow!("POST {url} parse: {e}"))
1635    }
1636
1637    async fn cf_put<B: Serialize, T: serde::de::DeserializeOwned>(
1638        &self,
1639        path: &str,
1640        body: &B,
1641    ) -> Result<T> {
1642        let url = format!("{CF_API}{path}");
1643        let resp = self
1644            .http
1645            .put(&url)
1646            .header("Authorization", format!("Bearer {}", self.token))
1647            .header("Content-Type", "application/json")
1648            .json(body)
1649            .send()
1650            .await
1651            .map_err(|e| anyhow!("PUT {url}: {e}"))?;
1652        resp.json::<T>()
1653            .await
1654            .map_err(|e| anyhow!("PUT {url} parse: {e}"))
1655    }
1656
1657    async fn cf_delete<T: serde::de::DeserializeOwned>(&self, path: &str) -> Result<T> {
1658        let url = format!("{CF_API}{path}");
1659        let resp = self
1660            .http
1661            .delete(&url)
1662            .header("Authorization", format!("Bearer {}", self.token))
1663            .header("Content-Type", "application/json")
1664            .send()
1665            .await
1666            .map_err(|e| anyhow!("DELETE {url}: {e}"))?;
1667        resp.json::<T>()
1668            .await
1669            .map_err(|e| anyhow!("DELETE {url} parse: {e}"))
1670    }
1671
1672    /// R907-B1 addendum: the numeric error CODE is carried into the returned
1673    /// message alongside the text, not dropped. Before this, a caller (and
1674    /// anyone debugging from outside the process) could only see the
1675    /// message string — which is exactly why `is_resource_not_found`'s 1003
1676    /// had to be sourced from third-party reverse-engineering instead of
1677    /// read off a live failure. The next live "Configuration for tunnel not
1678    /// found" this produces will show its real code plainly.
1679    fn ok(&self, success: &bool, errors: &Option<Vec<serde_json::Value>>) -> Result<()> {
1680        if *success {
1681            return Ok(());
1682        }
1683        let entries = errors.as_deref().unwrap_or(&[]);
1684        if entries.is_empty() {
1685            return Err(anyhow!("Cloudflare returned an error"));
1686        }
1687        let rendered: Vec<String> = entries
1688            .iter()
1689            .map(|e| {
1690                let msg = e
1691                    .get("message")
1692                    .and_then(|m| m.as_str())
1693                    .unwrap_or("(no message)");
1694                match e.get("code").and_then(|c| c.as_i64()) {
1695                    Some(code) => format!("Cloudflare error {code}: {msg}"),
1696                    None => format!("Cloudflare error: {msg}"),
1697                }
1698            })
1699            .collect();
1700        Err(anyhow!("{}", rendered.join("; ")))
1701    }
1702
1703    /// Classifies a failed response as "the resource has never been
1704    /// created" (an EMPTY read, safe to treat as absent) rather than a
1705    /// genuine failure (auth, rate limit, transient API error). Callers
1706    /// must still call [`Self::ok`] for anything this returns `false` for —
1707    /// it only tells you when NOT to.
1708    ///
1709    /// Classified on the error CODE, not the message text: Cloudflare
1710    /// answers a not-yet-created tunnel configuration with code 1003, the
1711    /// same generic "not found" family it uses for a missing account or
1712    /// zone lookup. That code is not in Cloudflare's own published API
1713    /// reference as of 2026-09-14 (checked: the docs site lists only the
1714    /// generic `{code, message}` error shape, no per-condition table) —
1715    /// grounded instead from convergent third-party Cloudflare API clients
1716    /// that hardcode it: `vana-com/vana-connect`'s
1717    /// `connect/src/lib/server-provider/gcp.ts` comment "1003: tunnel not
1718    /// found", `mandar-karhade/dockflare`'s
1719    /// `docs/design/03-cloudflare-integration.md` error table "1003 | Zone
1720    /// not found", and several independent test fixtures (e.g.
1721    /// `ratazzi/coulson`, `MauroDruwel/TunnelDashDesktop`) using 1003 for
1722    /// "Invalid or missing account id" / "Account not found". Confirming
1723    /// against a live Cloudflare token is out of scope for R907-B1.
1724    fn is_resource_not_found(errors: &Option<Vec<serde_json::Value>>) -> bool {
1725        const CF_ERR_NOT_FOUND: i64 = 1003;
1726        errors
1727            .iter()
1728            .flatten()
1729            .filter_map(|e| e.get("code").and_then(|c| c.as_i64()))
1730            .any(|code| code == CF_ERR_NOT_FOUND)
1731    }
1732}
1733
1734// ---------- worker-upload helpers (pure, network-free) ----------
1735
1736/// Build a `multipart/form-data` body for uploading a CF Worker script with
1737/// typed bindings for runtime config.
1738///
1739/// Each [`WorkerBinding`] becomes one entry in the Worker metadata's
1740/// `bindings` array — plain_text strings, R2 bucket references, etc.
1741///
1742/// Returns `(content_type_header_value, raw_body_bytes)`.
1743fn build_worker_multipart(script_js: &str, bindings: &[WorkerBinding<'_>]) -> (String, Vec<u8>) {
1744    const BOUNDARY: &str = "yahWorkerUpload0";
1745    const SCRIPT_FILENAME: &str = "worker.js";
1746    let binding_json: Vec<serde_json::Value> = bindings
1747        .iter()
1748        .map(|b| match *b {
1749            WorkerBinding::PlainText { name, text } => serde_json::json!({
1750                "type": "plain_text",
1751                "name": name,
1752                "text": text,
1753            }),
1754            WorkerBinding::R2Bucket { name, bucket_name } => serde_json::json!({
1755                "type": "r2_bucket",
1756                "name": name,
1757                "bucket_name": bucket_name,
1758            }),
1759        })
1760        .collect();
1761    let metadata = serde_json::json!({
1762        "main_module": SCRIPT_FILENAME,
1763        "bindings": binding_json,
1764    });
1765    let mut body: Vec<u8> = Vec::new();
1766    let push = |v: &mut Vec<u8>, s: &str| v.extend_from_slice(s.as_bytes());
1767    // metadata part
1768    push(&mut body, &format!("--{BOUNDARY}\r\n"));
1769    push(
1770        &mut body,
1771        "Content-Disposition: form-data; name=\"metadata\"\r\n",
1772    );
1773    push(&mut body, "Content-Type: application/json\r\n\r\n");
1774    push(&mut body, &metadata.to_string());
1775    push(&mut body, "\r\n");
1776    // script part
1777    push(&mut body, &format!("--{BOUNDARY}\r\n"));
1778    push(&mut body, &format!(
1779        "Content-Disposition: form-data; name=\"{SCRIPT_FILENAME}\"; filename=\"{SCRIPT_FILENAME}\"\r\n"
1780    ));
1781    push(
1782        &mut body,
1783        "Content-Type: application/javascript+module\r\n\r\n",
1784    );
1785    push(&mut body, script_js);
1786    push(&mut body, &format!("\r\n--{BOUNDARY}--\r\n"));
1787    (format!("multipart/form-data; boundary={BOUNDARY}"), body)
1788}
1789
1790// ---------- token-create helpers (pure, network-free) ----------
1791
1792/// Resolve a grant's permission-group ID: prefer the live catalog entry for its
1793/// name, fall back to the baked-in constant.
1794fn resolve_grant_id(
1795    grant: &TokenGrant,
1796    catalog: &std::collections::BTreeMap<String, String>,
1797) -> String {
1798    catalog
1799        .get(grant.group_name)
1800        .cloned()
1801        .unwrap_or_else(|| grant.fallback_id.to_string())
1802}
1803
1804/// Build the `POST /accounts/{id}/tokens` request body for `grants`, grouping
1805/// account- and zone-scoped permission groups into separate policy blocks
1806/// (Cloudflare rejects a single block mixing the two scopes). A scope with no
1807/// grants produces no block.
1808fn build_token_body(
1809    token_name: &str,
1810    account_id: &str,
1811    zone_id: &str,
1812    grants: &[TokenGrant],
1813    catalog: &std::collections::BTreeMap<String, String>,
1814) -> serde_json::Value {
1815    let ids_for = |scope: GrantScope| -> Vec<serde_json::Value> {
1816        grants
1817            .iter()
1818            .filter(|g| g.scope == scope)
1819            .map(|g| serde_json::json!({ "id": resolve_grant_id(g, catalog) }))
1820            .collect()
1821    };
1822    let block = |resource: String, groups: Vec<serde_json::Value>| -> Option<serde_json::Value> {
1823        if groups.is_empty() {
1824            return None;
1825        }
1826        let mut resources = serde_json::Map::new();
1827        resources.insert(resource, serde_json::Value::String("*".into()));
1828        Some(serde_json::json!({
1829            "effect": "allow",
1830            "resources": serde_json::Value::Object(resources),
1831            "permission_groups": groups,
1832        }))
1833    };
1834
1835    let policies: Vec<serde_json::Value> = [
1836        block(
1837            format!("com.cloudflare.api.account.{account_id}"),
1838            ids_for(GrantScope::Account),
1839        ),
1840        block(
1841            format!("com.cloudflare.api.account.zone.{zone_id}"),
1842            ids_for(GrantScope::Zone),
1843        ),
1844    ]
1845    .into_iter()
1846    .flatten()
1847    .collect();
1848
1849    serde_json::json!({ "name": token_name, "policies": policies })
1850}
1851
1852// ---------- drift helpers (pure, network-free) ----------
1853
1854/// Normalise a DNS name/target for comparison: trim, drop a single trailing
1855/// dot, lowercase. So `yah.dev.` and `YAH.DEV` both compare equal to `yah.dev`.
1856fn norm_dns(s: &str) -> String {
1857    s.trim().trim_end_matches('.').to_ascii_lowercase()
1858}
1859
1860/// Pick the most specific accessible zone for `hostname` — the longest
1861/// zone-name suffix that is the apex of, or a parent of, the hostname.
1862/// Returns the matching zone id, or `None` when no accessible zone covers it.
1863fn best_zone_for<'a>(hostname: &str, zones: &'a [(String, String)]) -> Option<&'a str> {
1864    let h = norm_dns(hostname);
1865    zones
1866        .iter()
1867        .filter(|(_, name)| {
1868            let z = norm_dns(name);
1869            h == z || h.ends_with(&format!(".{z}"))
1870        })
1871        .max_by_key(|(_, name)| name.len())
1872        .map(|(id, _)| id.as_str())
1873}
1874
1875/// Classify drift for one ingress hostname against the live records fetched
1876/// for it. `live` is the zone's records (filtered by name upstream, but we
1877/// re-filter defensively so this stays a self-contained pure function).
1878fn classify_tunnel_drift(
1879    hostname: &str,
1880    expected_target: &str,
1881    live: &[CfDnsRecord],
1882) -> (TunnelDriftState, Option<String>) {
1883    let h = norm_dns(hostname);
1884    let matching: Vec<&CfDnsRecord> = live.iter().filter(|r| norm_dns(&r.name) == h).collect();
1885    if matching.is_empty() {
1886        return (TunnelDriftState::Missing, None);
1887    }
1888    let want = norm_dns(expected_target);
1889    if matching.iter().any(|r| norm_dns(&r.content) == want) {
1890        return (TunnelDriftState::Synced, None);
1891    }
1892    (
1893        TunnelDriftState::Mismatch,
1894        Some(matching[0].content.clone()),
1895    )
1896}
1897
1898#[cfg(test)]
1899mod tests {
1900    use super::*;
1901
1902    fn rec(name: &str, content: &str) -> CfDnsRecord {
1903        CfDnsRecord {
1904            name: name.into(),
1905            content: content.into(),
1906        }
1907    }
1908
1909    /// R907-B1: a tunnel that has never had a configuration PUT is an EMPTY
1910    /// read, not a failed one — pins the not-found code on the "treat as
1911    /// empty" side of `is_resource_not_found`.
1912    #[test]
1913    fn resource_not_found_classifies_missing_tunnel_config_as_empty() {
1914        let errors = Some(vec![serde_json::json!({
1915            "code": 1003,
1916            "message": "Configuration for tunnel not found",
1917        })]);
1918        assert!(CloudflareClient::is_resource_not_found(&errors));
1919    }
1920
1921    /// R907-B1: a genuine failure (here, an auth error) must still surface
1922    /// as `Err` from `tunnel_configuration` — this is what keeps the
1923    /// abort-before-PUT guard at reconciler/ingress.rs:1272-1276 correct.
1924    #[test]
1925    fn resource_not_found_does_not_classify_auth_error_as_empty() {
1926        let errors = Some(vec![serde_json::json!({
1927            "code": 10000,
1928            "message": "Authentication error",
1929        })]);
1930        assert!(!CloudflareClient::is_resource_not_found(&errors));
1931    }
1932
1933    #[test]
1934    fn resource_not_found_false_when_no_errors_present() {
1935        assert!(!CloudflareClient::is_resource_not_found(&None));
1936    }
1937
1938    /// R907-B1 addendum: `ok`'s error text must carry the numeric code, not
1939    /// just the message — that's what lets the next live failure name its
1940    /// real Cloudflare code instead of forcing another reverse-engineering
1941    /// pass like the one that produced `is_resource_not_found`'s 1003.
1942    #[test]
1943    fn ok_error_message_carries_the_cloudflare_code() {
1944        let client = CloudflareClient::new("test-token".to_string());
1945        let errors = Some(vec![serde_json::json!({
1946            "code": 1003,
1947            "message": "Configuration for tunnel not found",
1948        })]);
1949        let err = client.ok(&false, &errors).unwrap_err();
1950        assert_eq!(
1951            err.to_string(),
1952            "Cloudflare error 1003: Configuration for tunnel not found"
1953        );
1954    }
1955
1956    #[test]
1957    fn ok_error_message_joins_multiple_errors() {
1958        let client = CloudflareClient::new("test-token".to_string());
1959        let errors = Some(vec![
1960            serde_json::json!({"code": 1003, "message": "first"}),
1961            serde_json::json!({"code": 9109, "message": "second"}),
1962        ]);
1963        let err = client.ok(&false, &errors).unwrap_err();
1964        assert_eq!(
1965            err.to_string(),
1966            "Cloudflare error 1003: first; Cloudflare error 9109: second"
1967        );
1968    }
1969
1970    #[test]
1971    fn ok_error_message_falls_back_when_no_errors_present() {
1972        let client = CloudflareClient::new("test-token".to_string());
1973        let err = client.ok(&false, &None).unwrap_err();
1974        assert_eq!(err.to_string(), "Cloudflare returned an error");
1975    }
1976
1977    #[test]
1978    fn drift_synced_when_live_matches() {
1979        let live = vec![rec("yubaba.yah.dev", "9e4d.cfargotunnel.com")];
1980        let (state, target) =
1981            classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
1982        assert_eq!(state, TunnelDriftState::Synced);
1983        assert!(target.is_none());
1984    }
1985
1986    #[test]
1987    fn drift_missing_when_no_matching_record() {
1988        let live = vec![rec("other.yah.dev", "x.cfargotunnel.com")];
1989        let (state, target) =
1990            classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
1991        assert_eq!(state, TunnelDriftState::Missing);
1992        assert!(target.is_none());
1993    }
1994
1995    #[test]
1996    fn drift_missing_when_zone_empty() {
1997        let (state, _) = classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &[]);
1998        assert_eq!(state, TunnelDriftState::Missing);
1999    }
2000
2001    #[test]
2002    fn drift_mismatch_surfaces_live_target() {
2003        let live = vec![rec("yubaba.yah.dev", "stale.cfargotunnel.com")];
2004        let (state, target) =
2005            classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
2006        assert_eq!(state, TunnelDriftState::Mismatch);
2007        assert_eq!(target.as_deref(), Some("stale.cfargotunnel.com"));
2008    }
2009
2010    #[test]
2011    fn drift_normalises_trailing_dot_and_case() {
2012        // Cloudflare returns FQDNs and CNAME content with/without trailing dots
2013        // and arbitrary case; normalisation must treat these as synced.
2014        let live = vec![rec("Yubaba.YAH.dev.", "9E4D.cfargotunnel.com.")];
2015        let (state, _) = classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
2016        assert_eq!(state, TunnelDriftState::Synced);
2017    }
2018
2019    #[test]
2020    fn best_zone_picks_longest_suffix() {
2021        let zones = vec![
2022            ("z_apex".to_string(), "dev".to_string()),
2023            ("z_zone".to_string(), "yah.dev".to_string()),
2024        ];
2025        assert_eq!(best_zone_for("yubaba.yah.dev", &zones), Some("z_zone"));
2026        assert_eq!(best_zone_for("yah.dev", &zones), Some("z_zone"));
2027    }
2028
2029    #[test]
2030    fn best_zone_none_when_no_suffix_covers() {
2031        let zones = vec![("z1".to_string(), "yah.dev".to_string())];
2032        assert_eq!(best_zone_for("example.com", &zones), None);
2033        // A label that merely ends in the zone string but isn't a subdomain
2034        // must NOT match: `notyah.dev` is not under `yah.dev`.
2035        assert_eq!(best_zone_for("notyah.dev", &zones), None);
2036    }
2037
2038    #[test]
2039    fn best_zone_apex_self_match() {
2040        let zones = vec![("z1".to_string(), "yah.dev".to_string())];
2041        assert_eq!(best_zone_for("yah.dev", &zones), Some("z1"));
2042    }
2043
2044    #[test]
2045    fn token_body_splits_scopes_and_resolves_ids() {
2046        use std::collections::BTreeMap;
2047        // Catalog overrides Zone Read's id; everything else falls back to the
2048        // baked-in constant.
2049        let mut catalog = BTreeMap::new();
2050        catalog.insert("Zone Read".to_string(), "CATALOG_ZONE_READ".to_string());
2051
2052        let body = build_token_body("t", "ACCT", "ZONE", MESOFACT_STATIC_GRANTS, &catalog);
2053        let policies = body["policies"].as_array().unwrap();
2054        assert_eq!(policies.len(), 2, "one account block + one zone block");
2055
2056        let acct = &policies[0];
2057        assert_eq!(acct["resources"]["com.cloudflare.api.account.ACCT"], "*");
2058        let acct_ids: Vec<&str> = acct["permission_groups"]
2059            .as_array()
2060            .unwrap()
2061            .iter()
2062            .map(|g| g["id"].as_str().unwrap())
2063            .collect();
2064        assert!(acct_ids.contains(&"c1fde68c7bcc44588cbb6ddbc16d6480")); // Account Settings Read
2065        assert!(acct_ids.contains(&"bf7481a1826f439697cb59a20b22293e")); // R2 Storage Write
2066        assert!(acct_ids.contains(&"e086da7e2179491d91ee5f35b3ca210a")); // Workers Scripts Write
2067
2068        let zone = &policies[1];
2069        assert_eq!(
2070            zone["resources"]["com.cloudflare.api.account.zone.ZONE"],
2071            "*"
2072        );
2073        let zone_ids: Vec<&str> = zone["permission_groups"]
2074            .as_array()
2075            .unwrap()
2076            .iter()
2077            .map(|g| g["id"].as_str().unwrap())
2078            .collect();
2079        assert!(
2080            zone_ids.contains(&"CATALOG_ZONE_READ"),
2081            "catalog id wins over fallback"
2082        );
2083        assert!(zone_ids.contains(&"0ac90a90249747bca6b047d97f0803e9")); // Zone Transform Rules Write
2084        assert!(zone_ids.contains(&"28f4b596e7d643029c524985477ae49a")); // Workers Routes Write
2085        assert!(zone_ids.contains(&"e17beae8b8cb423a99b1730f21238bed")); // Cache Purge
2086    }
2087
2088    #[test]
2089    fn token_body_omits_empty_scope_block() {
2090        use std::collections::BTreeMap;
2091        let only_zone = &[TokenGrant {
2092            group_name: "Zone Read",
2093            scope: GrantScope::Zone,
2094            fallback_id: "ZR",
2095        }];
2096        let body = build_token_body("t", "A", "Z", only_zone, &BTreeMap::new());
2097        let policies = body["policies"].as_array().unwrap();
2098        assert_eq!(policies.len(), 1, "no account block when no account grants");
2099        assert_eq!(
2100            policies[0]["resources"]["com.cloudflare.api.account.zone.Z"],
2101            "*"
2102        );
2103    }
2104
2105    /// R912-F1: TUNNEL_EDIT_GRANTS is account-scoped only, so a token built
2106    /// from it must carry exactly one policy block (no zone block at all,
2107    /// not an empty one), with the Tunnel Write permission group.
2108    #[test]
2109    fn tunnel_edit_grants_build_account_only_policy() {
2110        use std::collections::BTreeMap;
2111        let body = build_token_body("t", "ACCT", "ZONE", TUNNEL_EDIT_GRANTS, &BTreeMap::new());
2112        let policies = body["policies"].as_array().unwrap();
2113        assert_eq!(policies.len(), 1, "no zone block for an account-only profile");
2114        assert_eq!(
2115            policies[0]["resources"]["com.cloudflare.api.account.ACCT"],
2116            "*"
2117        );
2118        let ids: Vec<&str> = policies[0]["permission_groups"]
2119            .as_array()
2120            .unwrap()
2121            .iter()
2122            .map(|g| g["id"].as_str().unwrap())
2123            .collect();
2124        assert_eq!(ids, vec!["c07321b023e944ff818fec44d8203567"]);
2125    }
2126
2127    /// Decode the multipart `metadata` part and return its parsed JSON.
2128    fn extract_metadata_json(body: &[u8]) -> serde_json::Value {
2129        let s = std::str::from_utf8(body).expect("multipart body is utf-8 for these tests");
2130        let (_, after) = s
2131            .split_once("name=\"metadata\"")
2132            .expect("metadata part present");
2133        let (_, after) = after
2134            .split_once("\r\n\r\n")
2135            .expect("metadata body delimited");
2136        let (json, _) = after
2137            .split_once("\r\n--")
2138            .expect("metadata terminated by boundary");
2139        serde_json::from_str(json).expect("metadata JSON parses")
2140    }
2141
2142    #[test]
2143    fn multipart_includes_r2_bucket_binding_metadata() {
2144        let bindings = [WorkerBinding::R2Bucket {
2145            name: "CACHE",
2146            bucket_name: "yah-cr-cache",
2147        }];
2148        let (content_type, body) = build_worker_multipart("export default {}", &bindings);
2149
2150        assert!(
2151            content_type.starts_with("multipart/form-data; boundary="),
2152            "content-type advertises multipart with boundary: got {content_type}",
2153        );
2154
2155        let metadata = extract_metadata_json(&body);
2156        assert_eq!(metadata["main_module"], "worker.js");
2157        let bindings = metadata["bindings"].as_array().expect("bindings array");
2158        assert_eq!(bindings.len(), 1);
2159        assert_eq!(bindings[0]["type"], "r2_bucket");
2160        assert_eq!(bindings[0]["name"], "CACHE");
2161        assert_eq!(bindings[0]["bucket_name"], "yah-cr-cache");
2162    }
2163
2164    #[test]
2165    fn multipart_mixes_plain_text_and_r2_bindings() {
2166        let bindings = [
2167            WorkerBinding::PlainText {
2168                name: "MODE",
2169                text: "cache",
2170            },
2171            WorkerBinding::R2Bucket {
2172                name: "CACHE",
2173                bucket_name: "yah-cr-cache",
2174            },
2175        ];
2176        let (_, body) = build_worker_multipart("export default {}", &bindings);
2177
2178        let metadata = extract_metadata_json(&body);
2179        let entries = metadata["bindings"].as_array().unwrap();
2180        assert_eq!(entries.len(), 2);
2181        assert_eq!(entries[0]["type"], "plain_text");
2182        assert_eq!(entries[0]["text"], "cache");
2183        assert_eq!(entries[1]["type"], "r2_bucket");
2184        assert_eq!(entries[1]["bucket_name"], "yah-cr-cache");
2185    }
2186}