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/// Zone-scoped DNS Read + DNS Write and nothing else (R936-B7): the token the
474/// three front doors hold so the raft leader's `yubaba::door_dns` reconciler
475/// can withdraw a dead door's A record and re-add it once the door is live,
476/// across every zone the doors front (minted with one `--zone` per zone). Same
477/// group names and fallback ids as the DNS pair in [`MESOFACT_STATIC_GRANTS`];
478/// deliberately none of that profile's Workers/R2/Transform grants, because
479/// this one lives on internet-facing boxes.
480pub const DOOR_DNS_GRANTS: &[TokenGrant] = &[
481 TokenGrant {
482 group_name: "DNS Read",
483 scope: GrantScope::Zone,
484 fallback_id: "82e64a83756745bbbb1c9c2701bf816b",
485 },
486 TokenGrant {
487 group_name: "DNS Write",
488 scope: GrantScope::Zone,
489 fallback_id: "4755a26eedb94da69e1066d98aa820be",
490 },
491];
492
493/// Result of minting an account-owned API token. `value` is the secret and is
494/// returned by Cloudflare exactly once — store it immediately.
495#[derive(Debug, Clone, Serialize, Deserialize)]
496#[serde(rename_all = "camelCase")]
497pub struct CreateTokenResult {
498 pub id: String,
499 pub name: String,
500 /// The token secret. Shown once by Cloudflare; never retrievable again.
501 pub value: String,
502}
503
504// ---------- client ----------
505
506/// Cloudflare management API client.
507///
508/// Construct with [`CloudflareClient::new`] passing a pre-resolved API token.
509/// The token scope required per method is noted on each method.
510pub struct CloudflareClient {
511 token: String,
512 http: reqwest::Client,
513}
514
515/// @yah:relay(R907, "A cloudflare-tunnel edge cannot be brought up by apply alone when its tunnel has no configuration yet")
516/// @yah:at(2026-09-14T19:14:34Z)
517/// @yah:status(open)
518/// @yah:assignee(agent:bundle-anthropic-ashguard)
519/// @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.")
520/// @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.")
521///
522/// @yah:ticket(R907-B1, "ensure_tunnel_ingress cannot bootstrap a tunnel that has no configuration yet")
523/// @yah:status(review)
524/// @yah:at(2026-09-14T19:34:22Z)
525/// @yah:assignee(agent:bundle-anthropic-miravel)
526/// @yah:parent(R907)
527/// @yah:severity(high)
528/// @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.")
529/// @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.")
530/// @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.")
531/// @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.")
532/// @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.")
533/// @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.")
534/// @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.")
535/// @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.")
536/// @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.")
537/// @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.")
538/// @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.")
539/// @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.")
540/// @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.")
541/// @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.")
542/// @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.")
543/// @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.")
544/// @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.")
545/// @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).")
546/// @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.")
547/// @yah:verify("cargo test -p yah-cloud --lib -- cloudflare: 29 pass / 0 fail (includes the 3 new tests), no skew reported.")
548/// @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.")
549/// @yah:verify("cargo check -p yah-cloud --lib: clean (1 pre-existing unrelated warning in mesofact_static.rs).")
550/// @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.")
551/// @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.")
552/// @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.")
553/// @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.")
554/// @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.")
555/// @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).")
556///
557/// @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")
558/// @yah:status(review)
559/// @yah:at(2026-09-14T20:23:53Z)
560/// @yah:assignee(agent:bundle-anthropic-miravel)
561/// @yah:parent(R907)
562/// @yah:severity(medium)
563/// @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.")
564/// @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.")
565/// @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.")
566/// @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.")
567/// @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.")
568/// @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.")
569/// @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`.")
570/// @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.")
571/// @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.")
572/// @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.")
573/// @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.")
574/// @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.")
575/// @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.")
576/// @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.")
577/// @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).")
578/// @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).")
579/// @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.")
580/// @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.")
581/// @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.")
582/// @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.")
583impl CloudflareClient {
584 /// Create a client for the given API token.
585 pub fn new(token: String) -> Self {
586 Self {
587 token,
588 http: reqwest::Client::new(),
589 }
590 }
591
592 /// List accounts the token can access.
593 /// Requires: `Account: Read`.
594 pub async fn list_accounts(&self) -> Result<Vec<CfAccountInfo>> {
595 #[derive(Deserialize)]
596 struct Entry {
597 id: String,
598 name: String,
599 }
600 let resp: CfPage<Entry> = self.cf_get("/accounts").await?;
601 self.ok(&resp.success, &resp.errors)?;
602 Ok(resp
603 .result
604 .unwrap_or_default()
605 .into_iter()
606 .map(|a| CfAccountInfo {
607 id: a.id,
608 name: a.name,
609 })
610 .collect())
611 }
612
613 /// List non-deleted Cloudflare Tunnels in `account_id` with connection state.
614 /// Requires: `Cloudflare Tunnel: Read`.
615 async fn list_tunnels_meta(&self, account_id: &str) -> Result<Vec<TunnelMeta>> {
616 let resp: CfPage<CfTunnel> = self
617 .cf_get(&format!(
618 "/accounts/{account_id}/cfd_tunnel?is_deleted=false"
619 ))
620 .await?;
621 self.ok(&resp.success, &resp.errors)?;
622 Ok(resp
623 .result
624 .unwrap_or_default()
625 .into_iter()
626 .map(|t| TunnelMeta {
627 conn_state: t
628 .status
629 .as_deref()
630 .map(TunnelConnState::from_cf_status)
631 .unwrap_or(TunnelConnState::Unknown),
632 conn_since: t.conns_active_at,
633 id: t.id,
634 name: t.name,
635 })
636 .collect())
637 }
638
639 /// List non-deleted Cloudflare Tunnels in `account_id` as `(id, name)` pairs.
640 /// Requires: `Cloudflare Tunnel: Read`.
641 pub async fn list_tunnels(&self, account_id: &str) -> Result<Vec<(String, String)>> {
642 Ok(self
643 .list_tunnels_meta(account_id)
644 .await?
645 .into_iter()
646 .map(|m| (m.id, m.name))
647 .collect())
648 }
649
650 /// Live connection state for one tunnel, by id — the connector-readiness
651 /// half of a safe DNS cutover (R907-B2). `Unknown` when the tunnel
652 /// doesn't appear in the account's tunnel list at all (deleted, wrong
653 /// account) rather than erroring, since a caller gating on `Active` — the
654 /// only value that clears the gate — treats every other variant the
655 /// same way.
656 ///
657 /// Requires: `Cloudflare Tunnel: Read`.
658 pub async fn tunnel_conn_state(
659 &self,
660 account_id: &str,
661 tunnel_id: &str,
662 ) -> Result<TunnelConnState> {
663 Ok(self
664 .list_tunnels_meta(account_id)
665 .await?
666 .into_iter()
667 .find(|t| t.id == tunnel_id)
668 .map(|t| t.conn_state)
669 .unwrap_or(TunnelConnState::Unknown))
670 }
671
672 /// Collect CNAME records for all tunnels across all accessible accounts.
673 ///
674 /// Walks accounts → tunnels → ingress configurations. Returns an empty
675 /// vec when the token has no tunnels or no configured ingress hostnames.
676 pub async fn tunnel_dns_records(&self) -> Result<Vec<TunnelDnsRecord>> {
677 let mut records = Vec::new();
678 let accounts = self.list_accounts().await?;
679
680 for account in &accounts {
681 let tunnels = self.list_tunnels(&account.id).await?;
682 for (tunnel_id, tunnel_name) in &tunnels {
683 let cname_target = format!("{tunnel_id}.cfargotunnel.com");
684 let path = format!(
685 "/accounts/{}/cfd_tunnel/{tunnel_id}/configurations",
686 account.id
687 );
688 let config_resp: CfSingle<CfTunnelConfig> = self.cf_get(&path).await?;
689 let ingress = if !config_resp.success
690 && Self::is_resource_not_found(&config_resp.errors)
691 {
692 Vec::new()
693 } else {
694 self.ok(&config_resp.success, &config_resp.errors)?;
695 config_resp
696 .result
697 .and_then(|c| c.config)
698 .and_then(|c| c.ingress)
699 .unwrap_or_default()
700 };
701 for rule in ingress {
702 if let Some(hostname) = rule.hostname.filter(|h| !h.is_empty()) {
703 records.push(TunnelDnsRecord {
704 tunnel_name: tunnel_name.clone(),
705 hostname,
706 cname_target: cname_target.clone(),
707 });
708 }
709 }
710 }
711 }
712 Ok(records)
713 }
714
715 /// Compute DNS drift for every tunnel ingress hostname, enriched with live
716 /// connector connection state.
717 ///
718 /// The declared side is the tunnel ingress config; the live side is the
719 /// zone's DNS records. Each ingress hostname is classified:
720 /// [`TunnelDriftState::Synced`] when a record points at the tunnel's CNAME
721 /// target, `Missing` when none exists, `Mismatch` when one points elsewhere.
722 ///
723 /// Connection state (`conn_state` / `conn_since`) comes from the `status`
724 /// and `conns_active_at` fields on the tunnel list response — fetched in the
725 /// same pass as the ingress configs to avoid an extra `list_accounts` round-trip.
726 ///
727 /// Degrades gracefully — a hostname whose zone can't be resolved or read is
728 /// reported `ZoneUnknown` rather than failing the whole report. Returns an
729 /// empty vec when the token has no tunnels or no ingress hostnames.
730 ///
731 /// Requires: `Cloudflare Tunnel: Read`, `Zone: Read`, `DNS: Read`.
732 pub async fn tunnel_dns_drift(&self) -> Result<Vec<TunnelDriftRow>> {
733 // Walk accounts → tunnels (with live connection state) → ingress configs
734 // in one pass, collecting both declared DNS records and conn-state in a
735 // single list_accounts round-trip.
736 let accounts = self.list_accounts().await?;
737 let mut declared: Vec<TunnelDnsRecord> = Vec::new();
738 let mut conn_by_tunnel: std::collections::HashMap<
739 String,
740 (TunnelConnState, Option<String>),
741 > = Default::default();
742
743 for account in &accounts {
744 let metas = self.list_tunnels_meta(&account.id).await?;
745 for meta in &metas {
746 conn_by_tunnel.insert(
747 meta.name.clone(),
748 (meta.conn_state, meta.conn_since.clone()),
749 );
750 let cname_target = format!("{}.cfargotunnel.com", meta.id);
751 let path = format!(
752 "/accounts/{}/cfd_tunnel/{}/configurations",
753 account.id, meta.id
754 );
755 let config_resp: CfSingle<CfTunnelConfig> = self.cf_get(&path).await?;
756 let ingress = if !config_resp.success
757 && Self::is_resource_not_found(&config_resp.errors)
758 {
759 Vec::new()
760 } else {
761 self.ok(&config_resp.success, &config_resp.errors)?;
762 config_resp
763 .result
764 .and_then(|c| c.config)
765 .and_then(|c| c.ingress)
766 .unwrap_or_default()
767 };
768 for rule in ingress {
769 if let Some(hostname) = rule.hostname.filter(|h| !h.is_empty()) {
770 declared.push(TunnelDnsRecord {
771 tunnel_name: meta.name.clone(),
772 hostname,
773 cname_target: cname_target.clone(),
774 });
775 }
776 }
777 }
778 }
779
780 if declared.is_empty() {
781 return Ok(Vec::new());
782 }
783
784 // Resolve zones once. A permission failure leaves the set empty, so
785 // every hostname falls through to `ZoneUnknown` instead of erroring.
786 let zones = self.list_zones().await.unwrap_or_default();
787
788 let mut rows = Vec::with_capacity(declared.len());
789 for rec in declared {
790 let (conn_state, conn_since) = conn_by_tunnel
791 .remove(&rec.tunnel_name)
792 .unwrap_or((TunnelConnState::Unknown, None));
793 let (state, live_target) = match best_zone_for(&rec.hostname, &zones) {
794 None => (TunnelDriftState::ZoneUnknown, None),
795 Some(zone_id) => match self.dns_records_named(zone_id, &rec.hostname).await {
796 Ok(live) => classify_tunnel_drift(&rec.hostname, &rec.cname_target, &live),
797 Err(_) => (TunnelDriftState::ZoneUnknown, None),
798 },
799 };
800 rows.push(TunnelDriftRow {
801 tunnel_name: rec.tunnel_name,
802 hostname: rec.hostname,
803 expected_target: rec.cname_target,
804 state,
805 live_target,
806 conn_state,
807 conn_since,
808 });
809 }
810 Ok(rows)
811 }
812
813 /// List zones the token can read, as `(zone_id, zone_name)` pairs.
814 /// Requires: `Zone: Read`.
815 pub async fn list_zones(&self) -> Result<Vec<(String, String)>> {
816 #[derive(Deserialize)]
817 struct ZoneEntry {
818 id: String,
819 name: String,
820 }
821 let resp: CfPage<ZoneEntry> = self.cf_get("/zones?per_page=50").await?;
822 self.ok(&resp.success, &resp.errors)?;
823 Ok(resp
824 .result
825 .unwrap_or_default()
826 .into_iter()
827 .map(|z| (z.id, z.name))
828 .collect())
829 }
830
831 /// Fetch DNS records in `zone_id` whose name exactly matches `name`.
832 /// Requires: `DNS: Read`.
833 async fn dns_records_named(&self, zone_id: &str, name: &str) -> Result<Vec<CfDnsRecord>> {
834 let resp: CfPage<CfDnsRecord> = self
835 .cf_get(&format!("/zones/{zone_id}/dns_records?name={name}"))
836 .await?;
837 self.ok(&resp.success, &resp.errors)?;
838 Ok(resp.result.unwrap_or_default())
839 }
840
841 /// Create a new Named Tunnel under `account_id` and return the connector token.
842 /// Requires: `Cloudflare Tunnel: Edit`.
843 pub async fn create_tunnel(&self, account_id: &str, name: &str) -> Result<CreateTunnelResult> {
844 use base64::Engine as _;
845
846 let mut secret_bytes = [0u8; 32];
847 getrandom::getrandom(&mut secret_bytes)
848 .map_err(|e| anyhow!("generate tunnel secret: {e}"))?;
849 let tunnel_secret = base64::engine::general_purpose::STANDARD.encode(secret_bytes);
850
851 #[derive(Serialize)]
852 struct CreateBody<'a> {
853 name: &'a str,
854 tunnel_secret: String,
855 }
856 #[derive(Deserialize)]
857 struct CreatedTunnel {
858 id: String,
859 name: String,
860 }
861 let create_resp: CfSingle<CreatedTunnel> = self
862 .cf_post(
863 &format!("/accounts/{account_id}/cfd_tunnel"),
864 &CreateBody {
865 name,
866 tunnel_secret,
867 },
868 )
869 .await?;
870 self.ok(&create_resp.success, &create_resp.errors)?;
871 let created = create_resp
872 .result
873 .ok_or_else(|| anyhow!("tunnel create: no result in response"))?;
874
875 // Fetch the connector JWT.
876 #[derive(Deserialize)]
877 struct TokenResp {
878 success: bool,
879 result: Option<String>,
880 errors: Option<Vec<serde_json::Value>>,
881 }
882 let token_resp: TokenResp = self
883 .cf_get(&format!(
884 "/accounts/{account_id}/cfd_tunnel/{}/token",
885 created.id
886 ))
887 .await?;
888 self.ok(&token_resp.success, &token_resp.errors)?;
889 let connector_token = token_resp
890 .result
891 .ok_or_else(|| anyhow!("no connector token in response"))?;
892
893 Ok(CreateTunnelResult {
894 cname_target: format!("{}.cfargotunnel.com", created.id),
895 tunnel_id: created.id,
896 tunnel_name: created.name,
897 connector_token,
898 })
899 }
900
901 /// Read a tunnel's remotely-managed configuration body as raw JSON
902 /// (R594-F11).
903 ///
904 /// Returns the `result.config` object — the thing a PUT round-trips — or an
905 /// empty object when the tunnel has never been configured. Deliberately
906 /// untyped: the `ingress` list is the only key this crate owns, and every
907 /// sibling (`warp-routing`, `originRequest`, …) must survive a
908 /// read-modify-write untouched.
909 ///
910 /// Requires: `Cloudflare Tunnel: Read`.
911 pub async fn tunnel_configuration(
912 &self,
913 account_id: &str,
914 tunnel_id: &str,
915 ) -> Result<serde_json::Value> {
916 let path = format!("/accounts/{account_id}/cfd_tunnel/{tunnel_id}/configurations");
917 let resp: CfSingle<serde_json::Value> = self.cf_get(&path).await?;
918 if !resp.success && Self::is_resource_not_found(&resp.errors) {
919 return Ok(serde_json::json!({}));
920 }
921 self.ok(&resp.success, &resp.errors)?;
922 let config = resp
923 .result
924 .as_ref()
925 .and_then(|r| r.get("config"))
926 .cloned()
927 .unwrap_or_else(|| serde_json::json!({}));
928 // A tunnel configured with an explicit JSON `null` config reads back as
929 // Value::Null, which has no object to insert `ingress` into.
930 Ok(if config.is_object() {
931 config
932 } else {
933 serde_json::json!({})
934 })
935 }
936
937 /// Replace a tunnel's remotely-managed configuration (R594-F11).
938 ///
939 /// `config` is the whole config body, not a patch — Cloudflare replaces it
940 /// wholesale, which is why callers must GET-merge-PUT rather than PUT a
941 /// freshly-built list. See
942 /// [`reconciler::ingress::ensure_tunnel_ingress`](crate::reconciler::ingress::ensure_tunnel_ingress).
943 ///
944 /// Requires: `Cloudflare Tunnel: Edit`.
945 pub async fn put_tunnel_configuration(
946 &self,
947 account_id: &str,
948 tunnel_id: &str,
949 config: &serde_json::Value,
950 ) -> Result<()> {
951 #[derive(Serialize)]
952 struct ConfigBody<'a> {
953 config: &'a serde_json::Value,
954 }
955 let path = format!("/accounts/{account_id}/cfd_tunnel/{tunnel_id}/configurations");
956 let resp: CfSingle<serde_json::Value> = self.cf_put(&path, &ConfigBody { config }).await?;
957 self.ok(&resp.success, &resp.errors)?;
958 Ok(())
959 }
960
961 /// Create a new R2 bucket under `account_id`.
962 /// Requires: `Account: Cloudflare R2: Edit`.
963 pub async fn create_r2_bucket(
964 &self,
965 account_id: &str,
966 bucket_name: &str,
967 ) -> Result<CreateR2BucketResult> {
968 #[derive(Serialize)]
969 struct CreateBody<'a> {
970 name: &'a str,
971 }
972 let resp: CfSingle<serde_json::Value> = self
973 .cf_post(
974 &format!("/accounts/{account_id}/r2/buckets"),
975 &CreateBody { name: bucket_name },
976 )
977 .await?;
978 self.ok(&resp.success, &resp.errors)?;
979
980 Ok(CreateR2BucketResult {
981 endpoint: format!("https://{account_id}.r2.cloudflarestorage.com"),
982 name: bucket_name.to_string(),
983 })
984 }
985
986 /// Resolve a zone name (e.g. `"yah.dev"`) to its Cloudflare zone ID.
987 /// Requires: `Zone: Read`.
988 pub async fn zone_id_for_name(&self, zone_name: &str) -> Result<String> {
989 #[derive(Deserialize)]
990 struct ZoneEntry {
991 id: String,
992 name: String,
993 }
994 let resp: CfPage<ZoneEntry> = self.cf_get(&format!("/zones?name={zone_name}")).await?;
995 self.ok(&resp.success, &resp.errors)?;
996 resp.result
997 .unwrap_or_default()
998 .into_iter()
999 .find(|z| z.name == zone_name)
1000 .map(|z| z.id)
1001 .ok_or_else(|| anyhow!("no Cloudflare zone found for name {zone_name:?}"))
1002 }
1003
1004 /// Purge content by cache tags from a zone.
1005 ///
1006 /// Cache tags must be applied to responses via the `Cache-Tag` header or
1007 /// Cloudflare page rules. Returns `Ok(())` when all tags are queued for
1008 /// purge. Requires: `Zone: Cache Purge`.
1009 pub async fn purge_cache_tags(&self, zone_id: &str, tags: &[String]) -> Result<()> {
1010 if tags.is_empty() {
1011 return Ok(());
1012 }
1013 #[derive(Serialize)]
1014 struct PurgeBody<'a> {
1015 tags: &'a [String],
1016 }
1017 let resp: CfSingle<serde_json::Value> = self
1018 .cf_post(
1019 &format!("/zones/{zone_id}/purge_cache"),
1020 &PurgeBody { tags },
1021 )
1022 .await?;
1023 self.ok(&resp.success, &resp.errors)
1024 }
1025
1026 /// Upsert the Transform Rule that rewrites `GET /` → `/index.html` on the
1027 /// zone, identified by the stable description tag `"yah:static-index"`.
1028 ///
1029 /// Idempotent: fetches the existing `http_request_transform` entrypoint,
1030 /// drops any prior `"yah:static-index"` rule, appends the current one,
1031 /// and PUTs the merged list back. Treats a missing entrypoint (no rules
1032 /// yet) as an empty list.
1033 ///
1034 /// Requires: `Zone: Transform Rules: Edit`.
1035 pub async fn upsert_index_rewrite(&self, zone_id: &str) -> Result<()> {
1036 const RULE_DESC: &str = "yah:static-index";
1037 let path = format!("/zones/{zone_id}/rulesets/phases/http_request_transform/entrypoint");
1038
1039 // Fetch existing rules; a missing entrypoint is not an error.
1040 let existing: Vec<serde_json::Value> = {
1041 #[derive(Deserialize)]
1042 struct Rs {
1043 rules: Option<Vec<serde_json::Value>>,
1044 }
1045 match self.cf_get::<CfSingle<Rs>>(&path).await {
1046 Ok(resp) if resp.success => resp.result.and_then(|r| r.rules).unwrap_or_default(),
1047 _ => Vec::new(),
1048 }
1049 };
1050
1051 // Keep every rule except the one we manage, then append ours.
1052 let mut rules: Vec<serde_json::Value> = existing
1053 .into_iter()
1054 .filter(|r| r.get("description").and_then(|v| v.as_str()) != Some(RULE_DESC))
1055 .collect();
1056 rules.push(serde_json::json!({
1057 "action": "rewrite",
1058 "description": RULE_DESC,
1059 "expression": "(http.request.uri.path eq \"/\")",
1060 "action_parameters": {
1061 "uri": { "path": { "value": "/index.html" } }
1062 },
1063 "enabled": true
1064 }));
1065
1066 let resp: CfSingle<serde_json::Value> = self
1067 .cf_put(&path, &serde_json::json!({ "rules": rules }))
1068 .await?;
1069 self.ok(&resp.success, &resp.errors)
1070 }
1071
1072 /// Deploy an ES-module Worker script with typed bindings for runtime config.
1073 ///
1074 /// Each entry in `bindings` becomes one `metadata.bindings[…]` declaration in
1075 /// the upload payload — see [`WorkerBinding`] for the supported variants
1076 /// (plain_text config, R2 bucket references).
1077 ///
1078 /// Uses a manual multipart/form-data upload (CF Workers API requires multipart
1079 /// when metadata/bindings are attached). Idempotent: re-uploading the same script
1080 /// is safe but costs one CF API round-trip — callers should hash-guard this.
1081 ///
1082 /// Requires: `Workers Scripts: Edit` (account-scoped).
1083 pub async fn deploy_worker_script(
1084 &self,
1085 account_id: &str,
1086 script_name: &str,
1087 script_js: &str,
1088 bindings: &[WorkerBinding<'_>],
1089 ) -> Result<WorkerDeployResult> {
1090 let url = format!("{CF_API}/accounts/{account_id}/workers/scripts/{script_name}");
1091 let (content_type, body) = build_worker_multipart(script_js, bindings);
1092 let resp = self
1093 .http
1094 .put(&url)
1095 .header("Authorization", format!("Bearer {}", self.token))
1096 .header("Content-Type", content_type)
1097 .body(body)
1098 .send()
1099 .await
1100 .map_err(|e| anyhow!("PUT {url}: {e}"))?;
1101 let result: CfSingle<WorkerDeployResult> = resp
1102 .json()
1103 .await
1104 .map_err(|e| anyhow!("PUT {url} parse: {e}"))?;
1105 self.ok(&result.success, &result.errors)?;
1106 result
1107 .result
1108 .ok_or_else(|| anyhow!("deploy worker: no result in response"))
1109 }
1110
1111 /// Upsert a Worker route for `pattern` on `zone_id`, pointing at `script_name`.
1112 ///
1113 /// Idempotent: fetches existing routes, skips PUT/POST when the pattern already
1114 /// points at the right script, updates an existing pattern pointing elsewhere,
1115 /// or creates a new route entry.
1116 ///
1117 /// Requires: `Zone: Workers Routes: Edit` (zone-scoped).
1118 pub async fn upsert_worker_route(
1119 &self,
1120 zone_id: &str,
1121 pattern: &str,
1122 script_name: &str,
1123 ) -> Result<()> {
1124 let list_path = format!("/zones/{zone_id}/workers/routes");
1125
1126 #[derive(Deserialize)]
1127 struct RouteEntry {
1128 id: String,
1129 pattern: String,
1130 #[serde(default)]
1131 script: Option<String>,
1132 }
1133 let list: CfPage<RouteEntry> = self.cf_get(&list_path).await?;
1134 self.ok(&list.success, &list.errors)?;
1135 let routes = list.result.unwrap_or_default();
1136
1137 #[derive(Serialize)]
1138 struct RouteBody<'a> {
1139 pattern: &'a str,
1140 script: &'a str,
1141 }
1142
1143 if let Some(existing) = routes.iter().find(|r| r.pattern == pattern) {
1144 if existing.script.as_deref() == Some(script_name) {
1145 return Ok(());
1146 }
1147 let resp: CfSingle<serde_json::Value> = self
1148 .cf_put(
1149 &format!("/zones/{zone_id}/workers/routes/{}", existing.id),
1150 &RouteBody {
1151 pattern,
1152 script: script_name,
1153 },
1154 )
1155 .await?;
1156 self.ok(&resp.success, &resp.errors)
1157 } else {
1158 let resp: CfSingle<serde_json::Value> = self
1159 .cf_post(
1160 &list_path,
1161 &RouteBody {
1162 pattern,
1163 script: script_name,
1164 },
1165 )
1166 .await?;
1167 self.ok(&resp.success, &resp.errors)
1168 }
1169 }
1170
1171 /// Idempotently attach `hostname` (e.g. `cr.yah.dev`) as a Workers Custom
1172 /// Domain on `script_name`. Custom Domains route every request for the
1173 /// hostname into the Worker — distinct from a Worker Route, which only
1174 /// matches a URL pattern within an already-proxied zone.
1175 ///
1176 /// Walks the existing Custom Domains list first; if `hostname` is already
1177 /// bound to `script_name` on `zone_id`, returns Ok without an extra PUT.
1178 /// Otherwise PUTs `/accounts/{account_id}/workers/domains`, which CF treats
1179 /// as an upsert keyed on `(hostname, environment)`.
1180 ///
1181 /// Requires: `Workers Scripts: Edit` (account-scoped).
1182 pub async fn upsert_worker_custom_domain(
1183 &self,
1184 account_id: &str,
1185 zone_id: &str,
1186 hostname: &str,
1187 script_name: &str,
1188 ) -> Result<()> {
1189 let list_path = format!("/accounts/{account_id}/workers/domains");
1190
1191 #[derive(Deserialize)]
1192 struct DomainEntry {
1193 #[serde(default)]
1194 hostname: Option<String>,
1195 #[serde(default)]
1196 service: Option<String>,
1197 #[serde(default, rename = "zone_id")]
1198 zone_id: Option<String>,
1199 }
1200 let list: CfPage<DomainEntry> = self.cf_get(&list_path).await?;
1201 self.ok(&list.success, &list.errors)?;
1202 let domains = list.result.unwrap_or_default();
1203 if domains.iter().any(|d| {
1204 d.hostname.as_deref() == Some(hostname)
1205 && d.service.as_deref() == Some(script_name)
1206 && d.zone_id.as_deref() == Some(zone_id)
1207 }) {
1208 return Ok(());
1209 }
1210
1211 #[derive(Serialize)]
1212 struct DomainBody<'a> {
1213 environment: &'a str,
1214 hostname: &'a str,
1215 service: &'a str,
1216 zone_id: &'a str,
1217 }
1218 let resp: CfSingle<serde_json::Value> = self
1219 .cf_put(
1220 &list_path,
1221 &DomainBody {
1222 environment: "production",
1223 hostname,
1224 service: script_name,
1225 zone_id,
1226 },
1227 )
1228 .await?;
1229 self.ok(&resp.success, &resp.errors)
1230 }
1231
1232 /// Delete an R2 bucket under `account_id`.
1233 ///
1234 /// Cloudflare's management API handles non-empty buckets — objects do not
1235 /// need to be drained first. Returns `Ok(())` on success, `Err` if the
1236 /// API returns a failure (including "bucket not found" — callers that
1237 /// need idempotency should probe [`Self::list_r2_buckets`] first).
1238 ///
1239 /// Requires: `Account: Cloudflare R2: Edit`.
1240 pub async fn delete_r2_bucket(&self, account_id: &str, bucket_name: &str) -> Result<()> {
1241 let resp: CfSingle<serde_json::Value> = self
1242 .cf_delete(&format!("/accounts/{account_id}/r2/buckets/{bucket_name}"))
1243 .await?;
1244 self.ok(&resp.success, &resp.errors)
1245 }
1246
1247 /// Idempotently upsert a DNS record in `zone_id`. Fetches existing records
1248 /// with the same name and type: updates the first match if found, creates
1249 /// a new record otherwise. Returns the provider-issued record ID.
1250 ///
1251 /// Requires: `DNS: Edit` (zone-scoped).
1252 pub async fn upsert_dns_record(
1253 &self,
1254 zone_id: &str,
1255 name: &str,
1256 record_type: &str,
1257 content: &str,
1258 ttl: u32,
1259 proxied: bool,
1260 ) -> Result<String> {
1261 self.upsert_dns_record_matching(zone_id, name, record_type, content, ttl, proxied, false)
1262 .await
1263 }
1264
1265 /// [`upsert_dns_record`](Self::upsert_dns_record) with control over what
1266 /// counts as "the existing record" — R859-F1.
1267 ///
1268 /// `match_content = false` reproduces the original behaviour: the first
1269 /// record sharing `name` + `record_type` is updated in place. That is right
1270 /// for a single-valued name (one CNAME at `cdn.yah.dev`) and **wrong for a
1271 /// multi-valued RRset**: adding the second A record of a round-robin apex
1272 /// would rewrite the first one's content, silently halving the origin set
1273 /// to one box.
1274 ///
1275 /// `match_content = true` keys the lookup on `(name, type, content)`, so
1276 /// the call means "ensure exactly this record exists" — a no-op update when
1277 /// it already does, a create when it does not, and never a mutation of a
1278 /// sibling record at the same name.
1279 #[allow(clippy::too_many_arguments)]
1280 pub async fn upsert_dns_record_matching(
1281 &self,
1282 zone_id: &str,
1283 name: &str,
1284 record_type: &str,
1285 content: &str,
1286 ttl: u32,
1287 proxied: bool,
1288 match_content: bool,
1289 ) -> Result<String> {
1290 let existing_id = if match_content {
1291 self.list_dns_records(zone_id, Some(name), Some(record_type))
1292 .await?
1293 .into_iter()
1294 .find(|r| r.content == content)
1295 .map(|r| r.id)
1296 } else {
1297 self.find_dns_record(zone_id, name, record_type).await?
1298 };
1299
1300 #[derive(Serialize)]
1301 struct RecordBody<'a> {
1302 name: &'a str,
1303 #[serde(rename = "type")]
1304 record_type: &'a str,
1305 content: &'a str,
1306 ttl: u32,
1307 proxied: bool,
1308 }
1309 let body = RecordBody {
1310 name,
1311 record_type,
1312 content,
1313 ttl,
1314 proxied,
1315 };
1316
1317 #[derive(Deserialize)]
1318 struct RecordResult {
1319 id: String,
1320 }
1321
1322 if let Some(id) = existing_id {
1323 let resp: CfSingle<RecordResult> = self
1324 .cf_put(&format!("/zones/{zone_id}/dns_records/{id}"), &body)
1325 .await?;
1326 self.ok(&resp.success, &resp.errors)?;
1327 resp.result
1328 .map(|r| r.id)
1329 .ok_or_else(|| anyhow!("dns record update: no id in response"))
1330 } else {
1331 let resp: CfSingle<RecordResult> = self
1332 .cf_post(&format!("/zones/{zone_id}/dns_records"), &body)
1333 .await?;
1334 self.ok(&resp.success, &resp.errors)?;
1335 resp.result
1336 .map(|r| r.id)
1337 .ok_or_else(|| anyhow!("dns record create: no id in response"))
1338 }
1339 }
1340
1341 /// Read the DNS records in `zone_id`, optionally narrowed to one `name`
1342 /// and/or one `record_type` — R859-F1, the read half the `dns.*` catalog
1343 /// was missing.
1344 ///
1345 /// Unlike [`dns_records_named`](Self::dns_records_named) (drift-detection
1346 /// only, name + content) this returns the record **id**, type, ttl and
1347 /// proxy flag, which is what a reconciler needs to decide what to change.
1348 ///
1349 /// Requires: `DNS: Read` (zone-scoped).
1350 pub async fn list_dns_records(
1351 &self,
1352 zone_id: &str,
1353 name: Option<&str>,
1354 record_type: Option<&str>,
1355 ) -> Result<Vec<DnsRecordDetail>> {
1356 let mut query = format!("/zones/{zone_id}/dns_records?per_page=100");
1357 if let Some(n) = name {
1358 query.push_str(&format!("&name={n}"));
1359 }
1360 if let Some(t) = record_type {
1361 query.push_str(&format!("&type={t}"));
1362 }
1363 let resp: CfPage<DnsRecordDetail> = self.cf_get(&query).await?;
1364 self.ok(&resp.success, &resp.errors)?;
1365 Ok(resp.result.unwrap_or_default())
1366 }
1367
1368 /// Delete all DNS records in `zone_id` whose name matches `name` (and
1369 /// optionally `record_type`). Returns the count of records deleted.
1370 /// A count of 0 is not an error — the records may already have been absent.
1371 ///
1372 /// Requires: `DNS: Edit` (zone-scoped).
1373 pub async fn delete_dns_records(
1374 &self,
1375 zone_id: &str,
1376 name: &str,
1377 record_type: Option<&str>,
1378 ) -> Result<u32> {
1379 self.delete_dns_records_matching(zone_id, name, record_type, None)
1380 .await
1381 }
1382
1383 /// [`delete_dns_records`](Self::delete_dns_records) narrowed to records
1384 /// carrying one exact value — R859-F1.
1385 ///
1386 /// A round-robin apex holds several A records under one name, so
1387 /// "delete the A records at `yah.dev`" is not a way to withdraw *one*
1388 /// origin: it takes the live ones with it. `content = Some(ip)` deletes
1389 /// only the withdrawn member.
1390 ///
1391 /// Requires: `DNS: Edit` (zone-scoped).
1392 pub async fn delete_dns_records_matching(
1393 &self,
1394 zone_id: &str,
1395 name: &str,
1396 record_type: Option<&str>,
1397 content: Option<&str>,
1398 ) -> Result<u32> {
1399 let ids: Vec<String> = self
1400 .list_dns_records(zone_id, Some(name), record_type)
1401 .await?
1402 .into_iter()
1403 .filter(|r| content.is_none_or(|c| r.content == c))
1404 .map(|r| r.id)
1405 .collect();
1406 let mut deleted = 0u32;
1407 for id in &ids {
1408 let del: CfSingle<serde_json::Value> = self
1409 .cf_delete(&format!("/zones/{zone_id}/dns_records/{id}"))
1410 .await?;
1411 self.ok(&del.success, &del.errors)?;
1412 deleted += 1;
1413 }
1414 Ok(deleted)
1415 }
1416
1417 /// Fetch the ID of the first DNS record matching `name` and `record_type`.
1418 /// Returns `None` when no matching record exists.
1419 async fn find_dns_record(
1420 &self,
1421 zone_id: &str,
1422 name: &str,
1423 record_type: &str,
1424 ) -> Result<Option<String>> {
1425 #[derive(Deserialize)]
1426 struct RecordEntry {
1427 id: String,
1428 }
1429 let resp: CfPage<RecordEntry> = self
1430 .cf_get(&format!(
1431 "/zones/{zone_id}/dns_records?name={name}&type={record_type}"
1432 ))
1433 .await?;
1434 self.ok(&resp.success, &resp.errors)?;
1435 Ok(resp
1436 .result
1437 .unwrap_or_default()
1438 .into_iter()
1439 .next()
1440 .map(|r| r.id))
1441 }
1442
1443 /// List R2 buckets in `account_id`.
1444 /// Requires: `Account: Cloudflare R2: Read`.
1445 ///
1446 /// Unlike `/accounts` and `/zones`, the R2 list endpoint nests the array
1447 /// under `result.buckets` rather than returning `result` as a bare array,
1448 /// so it needs `CfSingle<R2ListResult>` and not `CfPage<BucketEntry>`.
1449 pub async fn list_r2_buckets(&self, account_id: &str) -> Result<Vec<R2BucketInfo>> {
1450 #[derive(Deserialize)]
1451 struct BucketEntry {
1452 name: String,
1453 #[serde(default)]
1454 location: Option<String>,
1455 #[serde(default)]
1456 creation_date: Option<String>,
1457 }
1458 #[derive(Deserialize)]
1459 struct R2ListResult {
1460 #[serde(default)]
1461 buckets: Vec<BucketEntry>,
1462 }
1463 let resp: CfSingle<R2ListResult> = self
1464 .cf_get(&format!("/accounts/{account_id}/r2/buckets"))
1465 .await?;
1466 self.ok(&resp.success, &resp.errors)?;
1467 Ok(resp
1468 .result
1469 .map(|r| r.buckets)
1470 .unwrap_or_default()
1471 .into_iter()
1472 .map(|b| R2BucketInfo {
1473 name: b.name,
1474 location: b.location,
1475 creation_date: b.creation_date,
1476 })
1477 .collect())
1478 }
1479
1480 /// List R2 custom-domain bindings on `bucket_name`.
1481 ///
1482 /// Requires: `Workers R2 Storage: Read` (or Write, which implies Read).
1483 /// The response nests the array under `result.domains`, mirroring
1484 /// `list_r2_buckets`'s `result.buckets` shape.
1485 pub async fn list_r2_custom_domains(
1486 &self,
1487 account_id: &str,
1488 bucket_name: &str,
1489 ) -> Result<Vec<R2CustomDomain>> {
1490 #[derive(Deserialize)]
1491 struct DomainEntry {
1492 domain: String,
1493 #[serde(default)]
1494 enabled: bool,
1495 }
1496 #[derive(Deserialize)]
1497 struct R2DomainListResult {
1498 #[serde(default)]
1499 domains: Vec<DomainEntry>,
1500 }
1501 let resp: CfSingle<R2DomainListResult> = self
1502 .cf_get(&format!(
1503 "/accounts/{account_id}/r2/buckets/{bucket_name}/domains/custom"
1504 ))
1505 .await?;
1506 self.ok(&resp.success, &resp.errors)?;
1507 Ok(resp
1508 .result
1509 .map(|r| r.domains)
1510 .unwrap_or_default()
1511 .into_iter()
1512 .map(|d| R2CustomDomain {
1513 domain: d.domain,
1514 enabled: d.enabled,
1515 })
1516 .collect())
1517 }
1518
1519 /// Bind a custom domain to an R2 bucket.
1520 ///
1521 /// `zone_id` names the zone that owns `domain` (resolve via
1522 /// [`Self::zone_id_for_name`]). CF requires it so the CNAME write into
1523 /// that zone is authorized — even though the caller is the bucket-side
1524 /// API. CF creates the CNAME automatically; no separate DNS-side call.
1525 /// Requires: `Workers R2 Storage: Edit` (account-scoped).
1526 ///
1527 /// `enabled: true` activates the binding immediately. CF still has to
1528 /// validate ownership + provision TLS in the background — the binding
1529 /// returns success the moment the record is queued, not when the
1530 /// hostname is fully resolvable. First-time DNS propagation is on the
1531 /// order of seconds to a minute.
1532 pub async fn add_r2_custom_domain(
1533 &self,
1534 account_id: &str,
1535 bucket_name: &str,
1536 domain: &str,
1537 zone_id: &str,
1538 ) -> Result<()> {
1539 #[derive(Serialize)]
1540 #[serde(rename_all = "camelCase")]
1541 struct AddBody<'a> {
1542 domain: &'a str,
1543 enabled: bool,
1544 zone_id: &'a str,
1545 }
1546 let resp: CfSingle<serde_json::Value> = self
1547 .cf_post(
1548 &format!("/accounts/{account_id}/r2/buckets/{bucket_name}/domains/custom"),
1549 &AddBody {
1550 domain,
1551 enabled: true,
1552 zone_id,
1553 },
1554 )
1555 .await?;
1556 self.ok(&resp.success, &resp.errors)
1557 }
1558
1559 /// Fetch the account's permission-group catalog as a name → id map, used
1560 /// to resolve [`TokenGrant`] names before minting a token. Requires the
1561 /// calling token to carry `API Tokens: Read` (implied by `Write`).
1562 pub async fn list_permission_group_ids(
1563 &self,
1564 account_id: &str,
1565 ) -> Result<std::collections::BTreeMap<String, String>> {
1566 #[derive(Deserialize)]
1567 struct PgEntry {
1568 id: String,
1569 name: String,
1570 }
1571 let resp: CfPage<PgEntry> = self
1572 .cf_get(&format!(
1573 "/accounts/{account_id}/tokens/permission_groups?per_page=500"
1574 ))
1575 .await?;
1576 self.ok(&resp.success, &resp.errors)?;
1577 Ok(resp
1578 .result
1579 .unwrap_or_default()
1580 .into_iter()
1581 .map(|p| (p.name, p.id))
1582 .collect())
1583 }
1584
1585 /// Mint an **account-owned** API token under `account_id` from `grants`
1586 /// scoped to `account_id` + every zone in `zone_ids` (zone-scoped grants
1587 /// apply to all of them; at least one is required).
1588 ///
1589 /// Resolves each grant's permission-group name against the live catalog
1590 /// (falling back to its baked-in ID), groups the IDs into account- and
1591 /// zone-scoped policy blocks, and POSTs to `/accounts/{id}/tokens`.
1592 /// Requires the calling token to carry `API Tokens: Write` — but the
1593 /// minted token is bounded by the *account's* access, not the calling
1594 /// token's, so the caller may hold only `API Tokens: Write`.
1595 ///
1596 /// The returned [`CreateTokenResult::value`] is the secret — Cloudflare
1597 /// reveals it only here.
1598 pub async fn create_account_token(
1599 &self,
1600 account_id: &str,
1601 zone_ids: &[String],
1602 token_name: &str,
1603 grants: &[TokenGrant],
1604 ) -> Result<CreateTokenResult> {
1605 anyhow::ensure!(!zone_ids.is_empty(), "token create: no zone given");
1606 // A catalog lookup failure is non-fatal: build_token_body falls back
1607 // to the baked-in IDs for any name the (empty) catalog can't resolve.
1608 let catalog = self
1609 .list_permission_group_ids(account_id)
1610 .await
1611 .unwrap_or_default();
1612 let body = build_token_body(token_name, account_id, zone_ids, grants, &catalog);
1613
1614 let resp: CfSingle<CreateTokenResult> = self
1615 .cf_post(&format!("/accounts/{account_id}/tokens"), &body)
1616 .await?;
1617 self.ok(&resp.success, &resp.errors)?;
1618 resp.result
1619 .ok_or_else(|| anyhow!("token create: no result in response"))
1620 }
1621
1622 // ---------- HTTP helpers ----------
1623
1624 async fn cf_get<T: serde::de::DeserializeOwned>(&self, path: &str) -> Result<T> {
1625 let url = format!("{CF_API}{path}");
1626 let resp = self
1627 .http
1628 .get(&url)
1629 .header("Authorization", format!("Bearer {}", self.token))
1630 .header("Content-Type", "application/json")
1631 .send()
1632 .await
1633 .map_err(|e| anyhow!("GET {url}: {e}"))?;
1634 resp.json::<T>()
1635 .await
1636 .map_err(|e| anyhow!("GET {url} parse: {e}"))
1637 }
1638
1639 async fn cf_post<B: Serialize, T: serde::de::DeserializeOwned>(
1640 &self,
1641 path: &str,
1642 body: &B,
1643 ) -> Result<T> {
1644 let url = format!("{CF_API}{path}");
1645 let resp = self
1646 .http
1647 .post(&url)
1648 .header("Authorization", format!("Bearer {}", self.token))
1649 .header("Content-Type", "application/json")
1650 .json(body)
1651 .send()
1652 .await
1653 .map_err(|e| anyhow!("POST {url}: {e}"))?;
1654 resp.json::<T>()
1655 .await
1656 .map_err(|e| anyhow!("POST {url} parse: {e}"))
1657 }
1658
1659 async fn cf_put<B: Serialize, T: serde::de::DeserializeOwned>(
1660 &self,
1661 path: &str,
1662 body: &B,
1663 ) -> Result<T> {
1664 let url = format!("{CF_API}{path}");
1665 let resp = self
1666 .http
1667 .put(&url)
1668 .header("Authorization", format!("Bearer {}", self.token))
1669 .header("Content-Type", "application/json")
1670 .json(body)
1671 .send()
1672 .await
1673 .map_err(|e| anyhow!("PUT {url}: {e}"))?;
1674 resp.json::<T>()
1675 .await
1676 .map_err(|e| anyhow!("PUT {url} parse: {e}"))
1677 }
1678
1679 async fn cf_delete<T: serde::de::DeserializeOwned>(&self, path: &str) -> Result<T> {
1680 let url = format!("{CF_API}{path}");
1681 let resp = self
1682 .http
1683 .delete(&url)
1684 .header("Authorization", format!("Bearer {}", self.token))
1685 .header("Content-Type", "application/json")
1686 .send()
1687 .await
1688 .map_err(|e| anyhow!("DELETE {url}: {e}"))?;
1689 resp.json::<T>()
1690 .await
1691 .map_err(|e| anyhow!("DELETE {url} parse: {e}"))
1692 }
1693
1694 /// R907-B1 addendum: the numeric error CODE is carried into the returned
1695 /// message alongside the text, not dropped. Before this, a caller (and
1696 /// anyone debugging from outside the process) could only see the
1697 /// message string — which is exactly why `is_resource_not_found`'s 1003
1698 /// had to be sourced from third-party reverse-engineering instead of
1699 /// read off a live failure. The next live "Configuration for tunnel not
1700 /// found" this produces will show its real code plainly.
1701 fn ok(&self, success: &bool, errors: &Option<Vec<serde_json::Value>>) -> Result<()> {
1702 if *success {
1703 return Ok(());
1704 }
1705 let entries = errors.as_deref().unwrap_or(&[]);
1706 if entries.is_empty() {
1707 return Err(anyhow!("Cloudflare returned an error"));
1708 }
1709 let rendered: Vec<String> = entries
1710 .iter()
1711 .map(|e| {
1712 let msg = e
1713 .get("message")
1714 .and_then(|m| m.as_str())
1715 .unwrap_or("(no message)");
1716 match e.get("code").and_then(|c| c.as_i64()) {
1717 Some(code) => format!("Cloudflare error {code}: {msg}"),
1718 None => format!("Cloudflare error: {msg}"),
1719 }
1720 })
1721 .collect();
1722 Err(anyhow!("{}", rendered.join("; ")))
1723 }
1724
1725 /// Classifies a failed response as "the resource has never been
1726 /// created" (an EMPTY read, safe to treat as absent) rather than a
1727 /// genuine failure (auth, rate limit, transient API error). Callers
1728 /// must still call [`Self::ok`] for anything this returns `false` for —
1729 /// it only tells you when NOT to.
1730 ///
1731 /// Classified on the error CODE, not the message text: Cloudflare
1732 /// answers a not-yet-created tunnel configuration with code 1003, the
1733 /// same generic "not found" family it uses for a missing account or
1734 /// zone lookup. That code is not in Cloudflare's own published API
1735 /// reference as of 2026-09-14 (checked: the docs site lists only the
1736 /// generic `{code, message}` error shape, no per-condition table) —
1737 /// grounded instead from convergent third-party Cloudflare API clients
1738 /// that hardcode it: `vana-com/vana-connect`'s
1739 /// `connect/src/lib/server-provider/gcp.ts` comment "1003: tunnel not
1740 /// found", `mandar-karhade/dockflare`'s
1741 /// `docs/design/03-cloudflare-integration.md` error table "1003 | Zone
1742 /// not found", and several independent test fixtures (e.g.
1743 /// `ratazzi/coulson`, `MauroDruwel/TunnelDashDesktop`) using 1003 for
1744 /// "Invalid or missing account id" / "Account not found". Confirming
1745 /// against a live Cloudflare token is out of scope for R907-B1.
1746 fn is_resource_not_found(errors: &Option<Vec<serde_json::Value>>) -> bool {
1747 const CF_ERR_NOT_FOUND: i64 = 1003;
1748 errors
1749 .iter()
1750 .flatten()
1751 .filter_map(|e| e.get("code").and_then(|c| c.as_i64()))
1752 .any(|code| code == CF_ERR_NOT_FOUND)
1753 }
1754}
1755
1756// ---------- worker-upload helpers (pure, network-free) ----------
1757
1758/// Build a `multipart/form-data` body for uploading a CF Worker script with
1759/// typed bindings for runtime config.
1760///
1761/// Each [`WorkerBinding`] becomes one entry in the Worker metadata's
1762/// `bindings` array — plain_text strings, R2 bucket references, etc.
1763///
1764/// Returns `(content_type_header_value, raw_body_bytes)`.
1765fn build_worker_multipart(script_js: &str, bindings: &[WorkerBinding<'_>]) -> (String, Vec<u8>) {
1766 const BOUNDARY: &str = "yahWorkerUpload0";
1767 const SCRIPT_FILENAME: &str = "worker.js";
1768 let binding_json: Vec<serde_json::Value> = bindings
1769 .iter()
1770 .map(|b| match *b {
1771 WorkerBinding::PlainText { name, text } => serde_json::json!({
1772 "type": "plain_text",
1773 "name": name,
1774 "text": text,
1775 }),
1776 WorkerBinding::R2Bucket { name, bucket_name } => serde_json::json!({
1777 "type": "r2_bucket",
1778 "name": name,
1779 "bucket_name": bucket_name,
1780 }),
1781 })
1782 .collect();
1783 let metadata = serde_json::json!({
1784 "main_module": SCRIPT_FILENAME,
1785 "bindings": binding_json,
1786 });
1787 let mut body: Vec<u8> = Vec::new();
1788 let push = |v: &mut Vec<u8>, s: &str| v.extend_from_slice(s.as_bytes());
1789 // metadata part
1790 push(&mut body, &format!("--{BOUNDARY}\r\n"));
1791 push(
1792 &mut body,
1793 "Content-Disposition: form-data; name=\"metadata\"\r\n",
1794 );
1795 push(&mut body, "Content-Type: application/json\r\n\r\n");
1796 push(&mut body, &metadata.to_string());
1797 push(&mut body, "\r\n");
1798 // script part
1799 push(&mut body, &format!("--{BOUNDARY}\r\n"));
1800 push(&mut body, &format!(
1801 "Content-Disposition: form-data; name=\"{SCRIPT_FILENAME}\"; filename=\"{SCRIPT_FILENAME}\"\r\n"
1802 ));
1803 push(
1804 &mut body,
1805 "Content-Type: application/javascript+module\r\n\r\n",
1806 );
1807 push(&mut body, script_js);
1808 push(&mut body, &format!("\r\n--{BOUNDARY}--\r\n"));
1809 (format!("multipart/form-data; boundary={BOUNDARY}"), body)
1810}
1811
1812// ---------- token-create helpers (pure, network-free) ----------
1813
1814/// Resolve a grant's permission-group ID: prefer the live catalog entry for its
1815/// name, fall back to the baked-in constant.
1816fn resolve_grant_id(
1817 grant: &TokenGrant,
1818 catalog: &std::collections::BTreeMap<String, String>,
1819) -> String {
1820 catalog
1821 .get(grant.group_name)
1822 .cloned()
1823 .unwrap_or_else(|| grant.fallback_id.to_string())
1824}
1825
1826/// Build the `POST /accounts/{id}/tokens` request body for `grants`, grouping
1827/// account- and zone-scoped permission groups into separate policy blocks
1828/// (Cloudflare rejects a single block mixing the two scopes). A scope with no
1829/// grants produces no block.
1830fn build_token_body(
1831 token_name: &str,
1832 account_id: &str,
1833 zone_ids: &[String],
1834 grants: &[TokenGrant],
1835 catalog: &std::collections::BTreeMap<String, String>,
1836) -> serde_json::Value {
1837 let ids_for = |scope: GrantScope| -> Vec<serde_json::Value> {
1838 grants
1839 .iter()
1840 .filter(|g| g.scope == scope)
1841 .map(|g| serde_json::json!({ "id": resolve_grant_id(g, catalog) }))
1842 .collect()
1843 };
1844 let block = |resource_keys: Vec<String>,
1845 groups: Vec<serde_json::Value>|
1846 -> Option<serde_json::Value> {
1847 if groups.is_empty() || resource_keys.is_empty() {
1848 return None;
1849 }
1850 let mut resources = serde_json::Map::new();
1851 for key in resource_keys {
1852 resources.insert(key, serde_json::Value::String("*".into()));
1853 }
1854 Some(serde_json::json!({
1855 "effect": "allow",
1856 "resources": serde_json::Value::Object(resources),
1857 "permission_groups": groups,
1858 }))
1859 };
1860
1861 let policies: Vec<serde_json::Value> = [
1862 block(
1863 vec![format!("com.cloudflare.api.account.{account_id}")],
1864 ids_for(GrantScope::Account),
1865 ),
1866 block(
1867 zone_ids
1868 .iter()
1869 .map(|z| format!("com.cloudflare.api.account.zone.{z}"))
1870 .collect(),
1871 ids_for(GrantScope::Zone),
1872 ),
1873 ]
1874 .into_iter()
1875 .flatten()
1876 .collect();
1877
1878 serde_json::json!({ "name": token_name, "policies": policies })
1879}
1880
1881// ---------- drift helpers (pure, network-free) ----------
1882
1883/// Normalise a DNS name/target for comparison: trim, drop a single trailing
1884/// dot, lowercase. So `yah.dev.` and `YAH.DEV` both compare equal to `yah.dev`.
1885fn norm_dns(s: &str) -> String {
1886 s.trim().trim_end_matches('.').to_ascii_lowercase()
1887}
1888
1889/// Pick the most specific accessible zone for `hostname` — the longest
1890/// zone-name suffix that is the apex of, or a parent of, the hostname.
1891/// Returns the matching zone id, or `None` when no accessible zone covers it.
1892fn best_zone_for<'a>(hostname: &str, zones: &'a [(String, String)]) -> Option<&'a str> {
1893 let h = norm_dns(hostname);
1894 zones
1895 .iter()
1896 .filter(|(_, name)| {
1897 let z = norm_dns(name);
1898 h == z || h.ends_with(&format!(".{z}"))
1899 })
1900 .max_by_key(|(_, name)| name.len())
1901 .map(|(id, _)| id.as_str())
1902}
1903
1904/// Classify drift for one ingress hostname against the live records fetched
1905/// for it. `live` is the zone's records (filtered by name upstream, but we
1906/// re-filter defensively so this stays a self-contained pure function).
1907fn classify_tunnel_drift(
1908 hostname: &str,
1909 expected_target: &str,
1910 live: &[CfDnsRecord],
1911) -> (TunnelDriftState, Option<String>) {
1912 let h = norm_dns(hostname);
1913 let matching: Vec<&CfDnsRecord> = live.iter().filter(|r| norm_dns(&r.name) == h).collect();
1914 if matching.is_empty() {
1915 return (TunnelDriftState::Missing, None);
1916 }
1917 let want = norm_dns(expected_target);
1918 if matching.iter().any(|r| norm_dns(&r.content) == want) {
1919 return (TunnelDriftState::Synced, None);
1920 }
1921 (
1922 TunnelDriftState::Mismatch,
1923 Some(matching[0].content.clone()),
1924 )
1925}
1926
1927#[cfg(test)]
1928mod tests {
1929 use super::*;
1930
1931 fn rec(name: &str, content: &str) -> CfDnsRecord {
1932 CfDnsRecord {
1933 name: name.into(),
1934 content: content.into(),
1935 }
1936 }
1937
1938 /// R907-B1: a tunnel that has never had a configuration PUT is an EMPTY
1939 /// read, not a failed one — pins the not-found code on the "treat as
1940 /// empty" side of `is_resource_not_found`.
1941 #[test]
1942 fn resource_not_found_classifies_missing_tunnel_config_as_empty() {
1943 let errors = Some(vec![serde_json::json!({
1944 "code": 1003,
1945 "message": "Configuration for tunnel not found",
1946 })]);
1947 assert!(CloudflareClient::is_resource_not_found(&errors));
1948 }
1949
1950 /// R907-B1: a genuine failure (here, an auth error) must still surface
1951 /// as `Err` from `tunnel_configuration` — this is what keeps the
1952 /// abort-before-PUT guard at reconciler/ingress.rs:1272-1276 correct.
1953 #[test]
1954 fn resource_not_found_does_not_classify_auth_error_as_empty() {
1955 let errors = Some(vec![serde_json::json!({
1956 "code": 10000,
1957 "message": "Authentication error",
1958 })]);
1959 assert!(!CloudflareClient::is_resource_not_found(&errors));
1960 }
1961
1962 #[test]
1963 fn resource_not_found_false_when_no_errors_present() {
1964 assert!(!CloudflareClient::is_resource_not_found(&None));
1965 }
1966
1967 /// R907-B1 addendum: `ok`'s error text must carry the numeric code, not
1968 /// just the message — that's what lets the next live failure name its
1969 /// real Cloudflare code instead of forcing another reverse-engineering
1970 /// pass like the one that produced `is_resource_not_found`'s 1003.
1971 #[test]
1972 fn ok_error_message_carries_the_cloudflare_code() {
1973 let client = CloudflareClient::new("test-token".to_string());
1974 let errors = Some(vec![serde_json::json!({
1975 "code": 1003,
1976 "message": "Configuration for tunnel not found",
1977 })]);
1978 let err = client.ok(&false, &errors).unwrap_err();
1979 assert_eq!(
1980 err.to_string(),
1981 "Cloudflare error 1003: Configuration for tunnel not found"
1982 );
1983 }
1984
1985 #[test]
1986 fn ok_error_message_joins_multiple_errors() {
1987 let client = CloudflareClient::new("test-token".to_string());
1988 let errors = Some(vec![
1989 serde_json::json!({"code": 1003, "message": "first"}),
1990 serde_json::json!({"code": 9109, "message": "second"}),
1991 ]);
1992 let err = client.ok(&false, &errors).unwrap_err();
1993 assert_eq!(
1994 err.to_string(),
1995 "Cloudflare error 1003: first; Cloudflare error 9109: second"
1996 );
1997 }
1998
1999 #[test]
2000 fn ok_error_message_falls_back_when_no_errors_present() {
2001 let client = CloudflareClient::new("test-token".to_string());
2002 let err = client.ok(&false, &None).unwrap_err();
2003 assert_eq!(err.to_string(), "Cloudflare returned an error");
2004 }
2005
2006 #[test]
2007 fn drift_synced_when_live_matches() {
2008 let live = vec![rec("yubaba.yah.dev", "9e4d.cfargotunnel.com")];
2009 let (state, target) =
2010 classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
2011 assert_eq!(state, TunnelDriftState::Synced);
2012 assert!(target.is_none());
2013 }
2014
2015 #[test]
2016 fn drift_missing_when_no_matching_record() {
2017 let live = vec![rec("other.yah.dev", "x.cfargotunnel.com")];
2018 let (state, target) =
2019 classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
2020 assert_eq!(state, TunnelDriftState::Missing);
2021 assert!(target.is_none());
2022 }
2023
2024 #[test]
2025 fn drift_missing_when_zone_empty() {
2026 let (state, _) = classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &[]);
2027 assert_eq!(state, TunnelDriftState::Missing);
2028 }
2029
2030 #[test]
2031 fn drift_mismatch_surfaces_live_target() {
2032 let live = vec![rec("yubaba.yah.dev", "stale.cfargotunnel.com")];
2033 let (state, target) =
2034 classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
2035 assert_eq!(state, TunnelDriftState::Mismatch);
2036 assert_eq!(target.as_deref(), Some("stale.cfargotunnel.com"));
2037 }
2038
2039 #[test]
2040 fn drift_normalises_trailing_dot_and_case() {
2041 // Cloudflare returns FQDNs and CNAME content with/without trailing dots
2042 // and arbitrary case; normalisation must treat these as synced.
2043 let live = vec![rec("Yubaba.YAH.dev.", "9E4D.cfargotunnel.com.")];
2044 let (state, _) = classify_tunnel_drift("yubaba.yah.dev", "9e4d.cfargotunnel.com", &live);
2045 assert_eq!(state, TunnelDriftState::Synced);
2046 }
2047
2048 #[test]
2049 fn best_zone_picks_longest_suffix() {
2050 let zones = vec![
2051 ("z_apex".to_string(), "dev".to_string()),
2052 ("z_zone".to_string(), "yah.dev".to_string()),
2053 ];
2054 assert_eq!(best_zone_for("yubaba.yah.dev", &zones), Some("z_zone"));
2055 assert_eq!(best_zone_for("yah.dev", &zones), Some("z_zone"));
2056 }
2057
2058 #[test]
2059 fn best_zone_none_when_no_suffix_covers() {
2060 let zones = vec![("z1".to_string(), "yah.dev".to_string())];
2061 assert_eq!(best_zone_for("example.com", &zones), None);
2062 // A label that merely ends in the zone string but isn't a subdomain
2063 // must NOT match: `notyah.dev` is not under `yah.dev`.
2064 assert_eq!(best_zone_for("notyah.dev", &zones), None);
2065 }
2066
2067 #[test]
2068 fn best_zone_apex_self_match() {
2069 let zones = vec![("z1".to_string(), "yah.dev".to_string())];
2070 assert_eq!(best_zone_for("yah.dev", &zones), Some("z1"));
2071 }
2072
2073 #[test]
2074 fn token_body_splits_scopes_and_resolves_ids() {
2075 use std::collections::BTreeMap;
2076 // Catalog overrides Zone Read's id; everything else falls back to the
2077 // baked-in constant.
2078 let mut catalog = BTreeMap::new();
2079 catalog.insert("Zone Read".to_string(), "CATALOG_ZONE_READ".to_string());
2080
2081 let body = build_token_body("t", "ACCT", &["ZONE".into()], MESOFACT_STATIC_GRANTS, &catalog);
2082 let policies = body["policies"].as_array().unwrap();
2083 assert_eq!(policies.len(), 2, "one account block + one zone block");
2084
2085 let acct = &policies[0];
2086 assert_eq!(acct["resources"]["com.cloudflare.api.account.ACCT"], "*");
2087 let acct_ids: Vec<&str> = acct["permission_groups"]
2088 .as_array()
2089 .unwrap()
2090 .iter()
2091 .map(|g| g["id"].as_str().unwrap())
2092 .collect();
2093 assert!(acct_ids.contains(&"c1fde68c7bcc44588cbb6ddbc16d6480")); // Account Settings Read
2094 assert!(acct_ids.contains(&"bf7481a1826f439697cb59a20b22293e")); // R2 Storage Write
2095 assert!(acct_ids.contains(&"e086da7e2179491d91ee5f35b3ca210a")); // Workers Scripts Write
2096
2097 let zone = &policies[1];
2098 assert_eq!(
2099 zone["resources"]["com.cloudflare.api.account.zone.ZONE"],
2100 "*"
2101 );
2102 let zone_ids: Vec<&str> = zone["permission_groups"]
2103 .as_array()
2104 .unwrap()
2105 .iter()
2106 .map(|g| g["id"].as_str().unwrap())
2107 .collect();
2108 assert!(
2109 zone_ids.contains(&"CATALOG_ZONE_READ"),
2110 "catalog id wins over fallback"
2111 );
2112 assert!(zone_ids.contains(&"0ac90a90249747bca6b047d97f0803e9")); // Zone Transform Rules Write
2113 assert!(zone_ids.contains(&"28f4b596e7d643029c524985477ae49a")); // Workers Routes Write
2114 assert!(zone_ids.contains(&"e17beae8b8cb423a99b1730f21238bed")); // Cache Purge
2115 }
2116
2117 #[test]
2118 fn token_body_omits_empty_scope_block() {
2119 use std::collections::BTreeMap;
2120 let only_zone = &[TokenGrant {
2121 group_name: "Zone Read",
2122 scope: GrantScope::Zone,
2123 fallback_id: "ZR",
2124 }];
2125 let body = build_token_body("t", "A", &["Z".into()], only_zone, &BTreeMap::new());
2126 let policies = body["policies"].as_array().unwrap();
2127 assert_eq!(policies.len(), 1, "no account block when no account grants");
2128 assert_eq!(
2129 policies[0]["resources"]["com.cloudflare.api.account.zone.Z"],
2130 "*"
2131 );
2132 }
2133
2134 /// R936-B7: DOOR_DNS_GRANTS minted across several zones yields ONE zone
2135 /// block naming every zone, carrying exactly DNS Read + DNS Write, and no
2136 /// account block at all.
2137 #[test]
2138 fn door_dns_grants_cover_every_zone_and_nothing_else() {
2139 use std::collections::BTreeMap;
2140 let zones: Vec<String> = vec!["Z1".into(), "Z2".into(), "Z3".into()];
2141 let body = build_token_body("t", "ACCT", &zones, DOOR_DNS_GRANTS, &BTreeMap::new());
2142 let policies = body["policies"].as_array().unwrap();
2143 assert_eq!(policies.len(), 1, "zone block only");
2144 let res = policies[0]["resources"].as_object().unwrap();
2145 assert_eq!(res.len(), 3);
2146 for z in &zones {
2147 assert_eq!(res[&format!("com.cloudflare.api.account.zone.{z}")], "*");
2148 }
2149 let mut ids: Vec<&str> = policies[0]["permission_groups"]
2150 .as_array()
2151 .unwrap()
2152 .iter()
2153 .map(|g| g["id"].as_str().unwrap())
2154 .collect();
2155 ids.sort();
2156 assert_eq!(
2157 ids,
2158 vec!["4755a26eedb94da69e1066d98aa820be", "82e64a83756745bbbb1c9c2701bf816b"]
2159 );
2160 }
2161
2162 /// R912-F1: TUNNEL_EDIT_GRANTS is account-scoped only, so a token built
2163 /// from it must carry exactly one policy block (no zone block at all,
2164 /// not an empty one), with the Tunnel Write permission group.
2165 #[test]
2166 fn tunnel_edit_grants_build_account_only_policy() {
2167 use std::collections::BTreeMap;
2168 let body = build_token_body("t", "ACCT", &["ZONE".into()], TUNNEL_EDIT_GRANTS, &BTreeMap::new());
2169 let policies = body["policies"].as_array().unwrap();
2170 assert_eq!(policies.len(), 1, "no zone block for an account-only profile");
2171 assert_eq!(
2172 policies[0]["resources"]["com.cloudflare.api.account.ACCT"],
2173 "*"
2174 );
2175 let ids: Vec<&str> = policies[0]["permission_groups"]
2176 .as_array()
2177 .unwrap()
2178 .iter()
2179 .map(|g| g["id"].as_str().unwrap())
2180 .collect();
2181 assert_eq!(ids, vec!["c07321b023e944ff818fec44d8203567"]);
2182 }
2183
2184 /// Decode the multipart `metadata` part and return its parsed JSON.
2185 fn extract_metadata_json(body: &[u8]) -> serde_json::Value {
2186 let s = std::str::from_utf8(body).expect("multipart body is utf-8 for these tests");
2187 let (_, after) = s
2188 .split_once("name=\"metadata\"")
2189 .expect("metadata part present");
2190 let (_, after) = after
2191 .split_once("\r\n\r\n")
2192 .expect("metadata body delimited");
2193 let (json, _) = after
2194 .split_once("\r\n--")
2195 .expect("metadata terminated by boundary");
2196 serde_json::from_str(json).expect("metadata JSON parses")
2197 }
2198
2199 #[test]
2200 fn multipart_includes_r2_bucket_binding_metadata() {
2201 let bindings = [WorkerBinding::R2Bucket {
2202 name: "CACHE",
2203 bucket_name: "yah-cr-cache",
2204 }];
2205 let (content_type, body) = build_worker_multipart("export default {}", &bindings);
2206
2207 assert!(
2208 content_type.starts_with("multipart/form-data; boundary="),
2209 "content-type advertises multipart with boundary: got {content_type}",
2210 );
2211
2212 let metadata = extract_metadata_json(&body);
2213 assert_eq!(metadata["main_module"], "worker.js");
2214 let bindings = metadata["bindings"].as_array().expect("bindings array");
2215 assert_eq!(bindings.len(), 1);
2216 assert_eq!(bindings[0]["type"], "r2_bucket");
2217 assert_eq!(bindings[0]["name"], "CACHE");
2218 assert_eq!(bindings[0]["bucket_name"], "yah-cr-cache");
2219 }
2220
2221 #[test]
2222 fn multipart_mixes_plain_text_and_r2_bindings() {
2223 let bindings = [
2224 WorkerBinding::PlainText {
2225 name: "MODE",
2226 text: "cache",
2227 },
2228 WorkerBinding::R2Bucket {
2229 name: "CACHE",
2230 bucket_name: "yah-cr-cache",
2231 },
2232 ];
2233 let (_, body) = build_worker_multipart("export default {}", &bindings);
2234
2235 let metadata = extract_metadata_json(&body);
2236 let entries = metadata["bindings"].as_array().unwrap();
2237 assert_eq!(entries.len(), 2);
2238 assert_eq!(entries[0]["type"], "plain_text");
2239 assert_eq!(entries[0]["text"], "cache");
2240 assert_eq!(entries[1]["type"], "r2_bucket");
2241 assert_eq!(entries[1]["bucket_name"], "yah-cr-cache");
2242 }
2243}