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