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