Skip to main content

cloud/reconciler/
static_asset.rs

1//! Reconciler for `kind = "static-asset"` components (R429-F2).
2//!
3//! Syncs a content-addressed file catalog declared in `workload.toml` to the
4//! mirror's `providers.object_store` slot (Cloudflare R2 or pond MinIO).
5//!
6//! ## Per-asset pipeline
7//!
8//! For each `[[asset]]` row:
9//! 1. Read the source file from disk.
10//! 2. Compute its BLAKE3 hash; compare to `entry.blake3`. Mismatch → surface
11//!    as a hard error and skip upload (don't push bad data).
12//! 3. HEAD the object key in the bucket — if already present, skip PUT.
13//! 4. PUT the file (single-part; R2 accepts up to 5 GB per PUT).
14//!
15//! ## Drift detection (v1)
16//!
17//! A `_yah-asset-catalog.json` sidecar stored in the bucket records what was
18//! uploaded. On the next run, the manifest∖current-catalog delta surfaces as
19//! prune candidates (logged as warnings; nothing is deleted). Exhaustive bucket
20//! listing via S3 ListObjects (which would find bucket∖catalog stragglers added
21//! outside of yah) is deferred to v2 when an XML list response parser exists.
22//!
23//! ## Auto-delete is OFF
24//!
25//! The reconciler NEVER deletes from the bucket. Prune candidates are reported
26//! for operator review; deletion requires `yah service prune`.
27//!
28//! @arch:see(.yah/docs/working/W164-derived-static-assets.md)
29//!
30//! @yah:ticket(R438-T8, "Worked examples: whisper-derive e2e + mesofact in-container build")
31//! @yah:assignee(agent:claude)
32//! @yah:at(2026-06-04T21:07:41Z)
33//! @yah:status(review)
34//! @yah:phase(P3)
35//! @yah:parent(R438)
36//! @yah:next("Whisper-derive e2e against fake R2: cargo run --example whisper_derive_e2e — two consecutive runs upload-then-skip")
37//! @yah:next("Post-prune --derive run re-materializes from scratch and produces bit-identical bytes (reproducibility check)")
38//! @yah:next("In-tree mesofact-static workload with build_mode=in_container builds green in CI")
39//! @yah:verify("Examples run green in CI matrix")
40//! @yah:verify("Reproducibility: independent operator/machine yields identical output blake3 (pinned container guarantees this)")
41//! @arch:see(.yah/docs/working/W164-derived-static-assets.md)
42//! @arch:see(.yah/docs/working/W165-mesofact-build-mode-lowering.md)
43//! @yah:depends_on(R438-T7)
44//! @yah:handoff("T8 landed. Three artifacts: (1) crates/yah/cloud/testdata/mesofact-in-container/workload.toml — in-tree fixture with build_mode = { mode = \"in_container\", image = \"oven/bun:1.2.13@sha256:aaa...\" }; (2) new test in_container_fixture_roundtrips_as_container_build_mode in mesofact_static.rs tests: loads fixture via read_mesofact_build, asserts BuildMode::InContainer returned, runs in CI without docker — passes; (3) crates/yah/cloud/tests/whisper_derive_e2e.rs — integration test derive_pipeline_upload_skip_prune_reproducibility. Uses in-process axum fake-S3 (HEAD/GET/PUT, state map) + fake upstream HTTP server + MockCopyExecutor (copies argv[1]→argv[2], counts invocations). Three-run scenario: run1=cold→executor invoked 1x, model.bin PUT; run2=warm→executor NOT invoked again, HEAD 200 skip; run3=prune cache+evict S3 object→executor invoked again, re-uploaded bytes have identical BLAKE3. cargo test -p cloud --lib: 304 pass, 5 pre-existing failures (port collision + cloud_init drift, same as T6 gotcha). cargo test -p cloud --test whisper_derive_e2e: 1 pass. cargo check --workspace --locked: clean.")
45//! @yah:verify("cargo test -p cloud --lib reconciler::mesofact_static::tests::in_container_fixture_roundtrips_as_container_build_mode  # 1 pass")
46//! @yah:verify("cargo test -p cloud --test whisper_derive_e2e  # 1 pass")
47//! @yah:verify("cargo check --workspace --locked  # clean")
48//!
49//! @yah:ticket(R438-F11, "HTTP fetch retry/resume policy for materialize step (W164 OQ#4)")
50//! @yah:assignee(agent:claude)
51//! @yah:at(2026-06-04T21:08:05Z)
52//! @yah:status(review)
53//! @yah:parent(R438)
54//! @yah:next("Exponential backoff with bounded attempts on transient network failures")
55//! @yah:next("Range: header for resumable downloads on multi-GB blobs (whisper-large is 1.5GB)")
56//! @yah:next("Surface progress to yah QED/task-pane per long-running→yah-surface rule (not Bash bg)")
57//! @yah:verify("Simulated 50%-mid-download disconnect resumes via Range and completes")
58//! @arch:see(.yah/docs/working/W164-derived-static-assets.md)
59//! @yah:depends_on(R438-T5)
60//! @yah:handoff("F11 landed. HTTP fetch retry + Range resume for materialize_fetch (W164 OQ#4). Changes: (1) Added stream feature to reqwest in cloud/Cargo.toml. (2) Replaced the single-shot reqwest::get+bytes() call in materialize_fetch with: FETCH_MAX_ATTEMPTS=5 (1 in test) retry loop with exponential backoff (FETCH_BASE_DELAY_MS=1000ms, capped at 30s; 1ms/5ms in tests so suite stays fast); fetch_once helper streams via resp.chunk() to a <hash>.partial file; Range: bytes=<offset>- header sent on retry when partial file exists; 206 response appended, 200 response truncates-and-restarts, 416 clears partial and retries fresh, 4xx is fatal (no retry), 5xx/429 is transient (retry). (3) Progress logged via info!() every 100MiB (FETCH_LOG_INTERVAL_BYTES) to surface through task-pane. (4) FetchOnceFail enum (Fatal/Transient) separates non-retriable from transient outcomes. After successful fetch_once, blake3 is verified on the partial file then atomic renamed to <hash>.bin. (5) Two new tests: materialize_fetch_retries_on_server_error (500 first, 200 second, asserts 2 requests made) and materialize_fetch_resumes_via_range_header (pre-writes half to .partial, server verifies Range header and returns 206+second-half, asserts full body in cache). cargo test -p cloud --lib reconciler::static_asset: 34 pass. cargo check -p cloud: clean.")
61//! @yah:verify("cargo test -p cloud --lib reconciler::static_asset::tests::materialize_fetch_retries_on_server_error  # passes")
62//! @yah:verify("cargo test -p cloud --lib reconciler::static_asset::tests::materialize_fetch_resumes_via_range_header  # passes, Range header asserted")
63//! @yah:verify("cargo check -p cloud  # clean")
64//!
65//! @yah:ticket(R438-T15, "cloud reconciler materialize step: HTTP fetch + recipe lowering + content-addressed cache")
66//! @yah:at(2026-06-05T00:03:49Z)
67//! @yah:status(review)
68//! @yah:phase(P2)
69//! @yah:parent(R438)
70//! @yah:next("Inject executor: Arc<dyn ForgeExecutor> via StaticAssetReconciler::with_executor(...) setter (defaults to Arc::new(LocalForgeDriver::default())).")
71//! @yah:verify("Two consecutive runs against a fake R2 / MinIO mock: upload-then-skip; derive-mode assets resolve through cache and round-trip.")
72//! @arch:see(.yah/docs/working/W164-derived-static-assets.md)
73//! @yah:depends_on(R438-T13)
74//! @yah:handoff("T15 landed. Cloud reconciler materialize step (W164) wired end-to-end. Changes: (1) StaticAssetReconciler gains `executor: Arc<dyn ForgeExecutor>` field + `with_executor(...)` setter; default Arc::new(LocalForgeDriver::default()). Threaded through up→sync_to_r2/sync_to_minio→sync_assets. (2) New materialize_asset() in crates/yah/cloud/src/reconciler/static_asset.rs lifts legacy `source=...` to disk path verbatim; derive-mode lowers to materialize_fetch() + optional materialize_transform(). (3) materialize_fetch: HTTP GET cached by upstream-blake3 to .yah/cache/derive/fetch/<hex>.bin; cache HIT verifies hash to catch bit-rot; cache MISS does atomic tmp+rename. (4) materialize_transform: TransformRecipeLoader.load → substitute_argv binds YAH_TRANSFORM_IN_0/_OUT + recipe.params at argv-element granularity → lowers each step to ForgeSpec{Subprocess{argv, image}, TaskPlacement{Local, recipe.runtime}, timeout, label, initiator=Gnome{static-asset-reconciler}}, hands to executor.execute(); cached by output-blake3. Unresolved {{placeholder}} after substitute_argv is a hard error. (5) Materialized path replaces entry.source for the existing BLAKE3-verify + S3 PUT loop — derive-mode assets ride the same upload path as legacy ones.")
75//! @yah:next("Sign off → archive R438-T15.")
76//! @yah:next("R422-F11 unblocked. Picker can resume: workload.toml gets [[asset]] + [asset.derive.fetch] + [asset.derive.transform], and the recipe TOML at .yah/qed/transforms/whisper-quantize.toml needs to be authored (R438-T8 worked-example covers the e2e verification).")
77//! @yah:next("Follow-up R438-F11 (HTTP fetch retry/resume policy) is the natural next-step for production whisper-large fetches; today's path single-shots reqwest::get.")
78//! @yah:verify("cargo test -p cloud --lib reconciler::static_asset  # 32 pass including 7 new W164 tests: materialize_legacy_source_returns_workload_dir_path + materialize_fetch_cache_miss_then_hit + materialize_fetch_blake3_mismatch_is_hard_error + materialize_transform_cache_miss_then_hit + materialize_transform_blake3_mismatch_on_output_is_hard_error + materialize_transform_recipe_failure_surfaces_stderr + materialize_derive_fetch_only_uses_fetch_path_for_upload")
79//! @yah:verify("cargo check --workspace  # clean (only pre-existing warnings in desktop)")
80//! @yah:verify("no new cloud→qed dep edge — cloud's new dep is task only (verified via cloud/Cargo.toml diff)")
81//! @yah:gotcha("Architectural decision during T15: moved qed::transforms → task::transforms because cloud cannot dep on qed (per the original T5 verify clause). transforms.rs only deps on task::TaskRuntime + workload_spec::ImageRef; task is the right home. qed re-exports dropped (no external callers existed). Added `toml = '0.8'` + `pub mod transforms;` to task/Cargo.toml + task/src/lib.rs. Cloud's Cargo.toml now has `task = { path = '../task' }`. The W164 doc's piece-placement table that says 'Recipe TOML loader | qed (existing)' is now stale; transforms lives in task.")
82//! @yah:gotcha("Stale test fixed: task::transforms::tests::rejects_recipe_with_struct_image_missing_digest now expects RecipeError::Parse (was ImageNotPinned). After R438-T3 tightened ImageRef.digest to non-Optional String, struct-form bare-tag fails at serde-deserialize, not at the post-parse ImageNotPinned check (which is now belt-and-braces against an empty-string digest).")
83//! @yah:gotcha("Reconciler is per-asset sequential — W164 calls for bounded semaphore (default 4) cross-asset concurrency. Filed as R438-F11-style follow-up rather than added in this ticket to keep the diff focused on correctness. Not blocking for R422-F11.")
84//!
85//! @yah:ticket(R546-B6, "yah cloud cache seed computes the fetched-input blake3 then discards it — leaves [asset.derive.fetch].blake3 sentinel unfilled, silently disarming the shared lock fast-path")
86//! @yah:phase(P1)
87//! @yah:status(review)
88//! @yah:assignee(agent:bundle-anthropic-ashguard)
89//! @yah:at(2026-08-02T23:38:04Z)
90//! @yah:parent(R546)
91//! @yah:next("Make `yah cloud cache seed` print the fetched-input blake3 alongside the other three, e.g. `[asset.derive.fetch].blake3 = \"<fetched_hash>\"`, so the operator can paste all four and arm BOTH skip layers. The value is already in scope as `fetched_hash` in seed_derivation_for_target — this is a print-line change plus threading it out through SeededDerivation.")
92//! @yah:next("Verified fix shape by hand for the x86_64 row: fetch.blake3 = 3568931fb074a4e1b4d43098db5810e683cdadd06887a8bdd964a5719bca6481 (b3sum of the rusty_v8-149.4.0 archive; the fetch cache is content-addressed so the blob FILENAME is the blake3). After pasting, a re-seed reproduces the identical derive_key 6238d80e..., i.e. lock_skip_hash's recomputation now matches the lock.")
93//! @yah:next("Consider also asserting it: if seed writes a lock whose fast-path cannot engage (fetch pin still sentinel), emit a warning rather than printing a success message that implies the job is done.")
94//! @yah:gotcha("Surfaced 2026-07-20 by leif while closing R546-T3: 'why wasn't fetch.blake3 filled in already?'. `seed_derivation_for_target` ALREADY has the value — it calls materialize_fetch(&derive.fetch, ...) and binds `fetched_hash`, uses it to compute the derive_key, then prints ONLY [[asset]].blake3 + lock.input_hash + lock.output_blake3. The fetched-input hash is dropped on the floor.")
95//! @yah:gotcha("CONSEQUENCE (non-obvious, cost real time to find): with [asset.derive.fetch].blake3 left as the zero sentinel, lock_skip_hash() bails at its FIRST guard (`is_bootstrap_sentinel(&derive.fetch.blake3.0) -> return None`, static_asset.rs ~line 770). So the W212 substituter fast-path never engages. The local action cache still skips the build, which MASKS the problem on the seeding machine — but a clean checkout / CI / another operator gets no skip at all. The seed command's whole purpose is 'prepare a no-rebuild apply', so leaving this unfilled defeats the shared half of it.")
96//! @yah:handoff("DONE. `yah cloud cache seed` now prints all FOUR paste-back values and a pin verdict. (1) oss/yubaba/crates/cloud/src/reconciler/static_asset.rs: SeededDerivation gains fetch_blake3 (the fetched-input hash the derive_key was keyed on, previously computed then dropped) and fetch_pin: Option<FetchPinState>; new pub enum FetchPinState {Pinned, Sentinel, Mismatch{committed}} plus pure fn classify_fetch_pin(committed, fetched). seed_transform_derivation fills fetch_blake3 and logs it; seed_derivation_for_target classifies the COMMITTED [asset.derive.fetch].blake3 against what it actually fetched. (2) app/yah/cli/src/cloud.rs handle_cache_seed: prints [asset.derive.fetch].blake3 first (columns realigned), then a per-verdict block -- Pinned = fast-path arms as soon as the lock lands; Sentinel = WARNING that lock_skip_hash declines and the skip is local-action-cache-only so CI / a clean checkout / another operator still rebuilds; Mismatch = WARNING printing both hashes (moved version anchor or wrong pin).")
97//! @yah:verify("LIVE-VERIFIED offline against the real workload (fetch anchor + artifact both already cached): `./target/debug/yah cloud cache seed --workload .yah/services/yah-cloud/components/rusty-v8-musl/workload.toml --target x86_64-unknown-linux-musl --artifact .yah/cache/artifacts/567e8f9c...` printed [asset.derive.fetch].blake3 = 3568931fb074a4e1b4d43098db5810e683cdadd06887a8bdd964a5719bca6481 -- byte-identical to the value leif verified by hand in this ticket -- and input_hash f188580181d5ae26004ecd3267b48ca7703b68969e8ad1e1db68e6c4b94d5207, reproducing the committed lock exactly. Verdict line printed: 'fetch pin already committed - the W212 fast-path arms as soon as the lock lands' (both rows of that workload were pinned by hand in an earlier session, so Pinned is correct there).")
98//! @yah:verify("Tests: 2 new in static_asset.rs. seed_surfaces_fetched_input_hash_and_pin_state drives seed_derivation_for_target end-to-end with a pre-warmed fetch cache (materialize_fetch HIT path = hermetic, no network) and asserts (a) fetch_blake3 equals the fetched hash, (b) fetch_pin == Pinned, (c) derivation_key(fetch_blake3, recipe_bk, params) reproduces seeded.derive_key -- i.e. the surfaced value really IS the key input. classify_fetch_pin_covers_all_three_verdicts covers Sentinel / Pinned (case-insensitive) / Mismatch without a network fetch. `cargo test -p yah-cloud --lib` from oss/yubaba: 595 passed, 1 pre-existing flake (reconciler::pond::tests::ensure_sim_port_free_ok_when_unbound -- binds a real port, passes in isolation, unrelated). `cargo check -p yah --bin yah` clean.")
99//! @yah:gotcha("RESOLVED-IN-PLACE, not a leftover: the two [asset.derive.fetch].blake3 rows in .yah/services/yah-cloud/components/rusty-v8-musl/workload.toml were ALREADY pinned to 3568931f... by hand in an earlier session, so this workload's fast-path is armed today. This ticket fixed the TOOLING that made them get missed; nothing in the workload needed editing. Consequence for review: the Sentinel and Mismatch warning branches are covered by unit test only, not by a live run -- exercising Sentinel live would require an unpinned workload and a real network fetch.")
100//!
101//! @yah:ticket(R546-B8, "static-asset transform passes a RELATIVE YAH_TRANSFORM_OUT into the container — any recipe that chdirs writes its artifact inside the container and silently loses it (cost a ~2h rusty-v8 arm64 build)")
102//! @yah:phase(P1)
103//! @yah:status(review)
104//! @yah:assignee(agent:bundle-anthropic-ashguard)
105//! @yah:at(2026-08-03T00:53:38Z)
106//! @yah:parent(R546)
107//! @yah:next("FIX LANDED: canonicalize `cache_dir` and join the filename (the tmp file itself does not exist yet so it cannot be canonicalized directly); also canonicalize `input_path` for the YAH_TRANSFORM_IN_0 binding for the same reason. static_asset tests 46/46 green.")
108//! @yah:next("ADD A REGRESSION TEST that asserts both YAH_TRANSFORM_IN_0 and YAH_TRANSFORM_OUT bindings are absolute after substitution. Not added yet — the existing tests exercise recipes that never chdir, so they cannot catch this class. A cheap version: assert Path::new(&params[ENV_TRANSFORM_OUT]).is_absolute().")
109//! @yah:next("HARDENING: the reconciler should fail LOUDLY when a transform step exits 0 but produces no output at YAH_TRANSFORM_OUT — today the missing-file error surfaces as a confusing BLAKE3 read failure that reads like cache corruption rather than 'your recipe never wrote its output'.")
110//! @yah:gotcha("FAILURE MODE IS SILENT AND EXPENSIVE: the step exits 0, no 'step failed' is reported, and the reconciler dies afterwards on `reading transform output ./.yah/cache/derive/transform/<derive_key>.tmp for BLAKE3: No such file or directory`. Hit 2026-07-20 after a ~2h native arm64 V8 build that had actually SUCCEEDED — the finished tar was written inside the container and thrown away with it.")
111//! @yah:gotcha("ROOT CAUSE: materialize_transform built `tmp_output = cache_dir.join(\"<derive_key>.tmp\")` from `workspace_root`, which is \".\" by CLI default, then bound that RELATIVE string into the recipe argv as {{YAH_TRANSFORM_OUT}}. Note the code immediately below it already canonicalized `workspace_abs` for the container cwd (docker rejects `-w .`) — the OUT binding just never got the same treatment. A relative OUT only survives if the recipe never leaves its cwd: true for the whisper recipes (which is why this went unnoticed for months), FALSE for rusty-v8 because build-v8.sh chdirs into /tmp/tmp.XXXX/v8src to build V8.")
112//! @yah:gotcha("CONTRAST: the QED/P018 offload path was unaffected because it passes an ABSOLUTE container path (/yah/produced/...). Only the LOCAL static-asset transform path had the relative binding.")
113//! @yah:handoff("DONE — both remaining next-steps landed (the canonicalize fix itself was already in the tree when I picked this up). (1) REGRESSION TEST: reconciler::static_asset::tests::transform_bindings_are_absolute_from_a_relative_workspace_root asserts BOTH YAH_TRANSFORM_IN_0 and YAH_TRANSFORM_OUT are absolute in the substituted argv. New ArgvCapture test executor records the argv (and can decline to write the output while still reporting exit 0 — the exact shape of the bug). (2) HARDENING: materialize_transform now checks tmp_output exists before reading it and bails with 'recipe X completed successfully but produced no output at <path> (YAH_TRANSFORM_OUT)' plus the three things to check. Covered by transform_that_writes_no_output_names_the_recipe_contract, which also asserts the raw 'No such file or directory' does NOT leak — that string is what sent the last reader into the cache subsystem instead of the recipe.")
114//! @yah:handoff("THE FIXTURE IS THE LOAD-BEARING PART OF THE TEST, and it is not obvious: it uses tempfile::TempDir::new_in(\".\") and rebuilds the relative form by hand (new_in returns an ABSOLUTE path), because `yah cloud apply` defaults --path to \".\". The repo's standard Fixture uses an absolute tempdir, and against an absolute root cache_dir is absolute either way — so an absolute-tempdir version of this test passes against the BUGGY code and proves nothing. Commented at the test.")
115//! @yah:handoff("ALSO recorded the invariant at the definition site so a future edit cannot re-break it silently: oss/qed/crates/velveteen-exec/src/transforms.rs — ENV_TRANSFORM_OUT / ENV_TRANSFORM_IN_0 doc comments now state 'Bind an ABSOLUTE path', why (a recipe is free to cd), and what it cost. Doc-only, no behaviour change; cargo check -p velveteen-exec clean.")
116//! @yah:handoff("SWEEP: grepped every ENV_TRANSFORM_IN_0/ENV_TRANSFORM_OUT binding site across oss/yubaba, oss/qed and app/yah/cli. static_asset.rs is the only production binder; velveteen-exec only defines the constants and its own test already uses absolute paths. The QED/P018 offload path was never affected (it passes /yah/produced/... absolute), confirming the ticket's contrast note. No other instance of this bug class exists.")
117//! @yah:verify("THE TEST WAS PROVEN TO FAIL AGAINST THE BUG, not just to pass against the fix. I reverted both canonicalizations in place (tmp_output back to cache_dir.join(...), input_abs back to input_path.to_path_buf()), re-ran, and got: 'YAH_TRANSFORM_IN_0 binding must be absolute, got \"./.tmpX7vsrP/.yah/cache/derive/fetch/in.bin\"'. Then restored both and re-ran green. Without this step the test would have been indistinguishable from one that can never fire.")
118//! @yah:verify("GREEN: cargo test -p yah-cloud --lib — 617 passed, 0 failed, 4 ignored (the pond real-port flakes did not fire this run; they are unrelated and pass in isolation). cargo test -p yah-cloud --test whisper_derive_e2e — 1/1. cargo check -p velveteen-exec clean.")
119//!
120//! @yah:ticket(R546-B10, "static-asset publish skips PUT on key EXISTENCE (HEAD), not content — bucket keeps stale bytes while the manifest records the new hash")
121//! @yah:status(review)
122//! @yah:assignee(agent:claude)
123//! @yah:at(2026-07-21T21:00:26Z)
124//! @yah:parent(R546)
125//! @yah:severity(high)
126//! @yah:next("Repro (real, cost this relay a wrong pin): static_asset.rs:670 does `if object_exists(...) { report.already_synced.push(key); new_manifest.insert(key, actual_hash); continue; }`. object_exists issues an existence-only probe (HTTP-verb HEAD) with no content comparison. So when the locally-materialized artifact differs from the object already at that key, the reconciler skips the PUT (bucket keeps the OLD bytes) yet still records the NEW locally-computed hash in the manifest. The published _yah-asset-catalog.json then advertises a hash that describes bytes which were never uploaded.")
127//! @yah:next("Live evidence: cdn.yah.dev/yah-cloud/rusty-v8/v149.4.0/rusty-v8-aarch64-unknown-linux-musl.tar.gz serves blake3 d322b4a1… (last-modified 2026-06-20, hand-uploaded), while _yah-asset-catalog.json claims ebdd842d… — the artifact the 2026-07-20 apply built and then silently declined to upload. Both hashes are real; they describe different tarballs.")
128//! @yah:next("Fix shape: when the object already exists, compare its hash against actual_hash before skipping. Cheapest probe is the ETag when it is a plain MD5 (single-part upload), but ETag is NOT a content hash for multipart uploads — do not key correctness on it alone. Preferred: store the blake3 as object metadata (x-amz-meta-blake3) at PUT time and compare that on the existence probe; fall back to re-PUT when the metadata is absent (older objects). Re-PUT on mismatch instead of skipping.")
129//! @yah:next("Also decide the drift POSTURE, not just the mechanism: a content mismatch at a stable key is exactly the DriftBucket condition the journal already models (AssetState::DriftBucket, used at static_asset.rs:645 for source-vs-declared mismatch). Emitting the same event here, rather than silently overwriting, is probably the right default for a published CDN object — make it explicit rather than incidental.")
130//! @yah:verify("Regression test: seed a bucket object at key K with bytes A, run the reconciler with a workload whose materialized artifact for K is bytes B (B != A), assert the object at K afterwards is B (or that a DriftBucket event was emitted) AND that the manifest never records B while the bucket still holds A.")
131//! @yah:gotcha("Latent for most assets, which is why it went unnoticed for so long: W164 catalogs put the VERSION IN THE PATH (yah-cloud/rusty-v8/v149.4.0/…), so bytes at a given key normally never change and the skip is correct. It only bites when an object is placed at a key out-of-band (hand-uploaded) or when a rebuild is non-deterministic — both true for rusty-v8 arm64.")
132//! @yah:gotcha("The manifest is NOT load-bearing for consumer correctness, so do not overstate the blast radius: prior_manifest is read only for PRUNE-candidate detection (static_asset.rs:534-549), and consumers verify pulled bytes against [[asset]].blake3 in workload.toml, not against the catalog JSON. The damage is a lying audit record + a silently-skipped publish, not corrupted downloads.")
133//! @yah:assumes("Tier: Cleric — mechanism is a few lines, but the drift posture (overwrite vs refuse-and-report) is a real design call, and the ETag/multipart trap makes the naive fix wrong.")
134
135use std::collections::{BTreeMap, HashMap};
136use std::path::{Path, PathBuf};
137use std::sync::Arc;
138
139use anyhow::{Context, Result};
140use async_trait::async_trait;
141use chrono::Utc;
142use sha2::{Digest, Sha256};
143use tracing::{debug, info, warn};
144
145use super::{ReconcileCtx, Reconciler, RunningWorkload};
146use crate::asset_journal::{AssetState, AssetStatusEvent, AssetStatusJournal};
147use crate::provider::cloudflare::CloudflareClient;
148use crate::reconciler::pond::DEFAULT_MINIO_PASSWORD;
149use crate::reconciler::pond::DEFAULT_MINIO_USER;
150use crate::{MirrorProviderSlot, Provider};
151
152use local_driver::s3_sign::{sign_s3_empty_body, sign_s3_put_object, uri_encode_key};
153use velveteen::{ForgeCommand, ForgeSpec, Initiator, MeshAccess, TaskLocation, TaskPlacement};
154use velveteen_exec::transforms::{
155    substitute_argv, RecipeStep, TransformRecipe, TransformRecipeLoader, ENV_TRANSFORM_IN_0,
156    ENV_TRANSFORM_OUT,
157};
158use velveteen_exec::{ExecContext, ForgeExecutor, LocalForgeDriver, PlacementRouter};
159use workload_spec::validate::shape_static_asset;
160use workload_spec::{AssetEntry, FetchSource, Millis, StaticAssetWorkload, TransformSpec};
161
162/// Workload kind handled by this reconciler.
163pub const WORKLOAD_KIND: &str = "static-asset";
164
165/// S3 region string for Cloudflare R2.
166const R2_REGION: &str = "auto";
167/// S3 region string for MinIO (SigV4 requires a non-empty value).
168const MINIO_REGION: &str = "us-east-1";
169
170/// Bucket key for the per-run catalog manifest sidecar.
171const CATALOG_MANIFEST_KEY: &str = "_yah-asset-catalog.json";
172
173/// 64 zero hex digits — the "not yet pinned" sentinel for `BlakeHash`.
174///
175/// Bootstrap mode: a derive-mode asset can ship with this value in any
176/// `blake3` field (`[[asset]].blake3` or `[asset.derive.fetch].blake3`)
177/// before its first apply. The reconciler computes the actual hash, accepts
178/// the bytes, uploads, and surfaces the discovered hash in the report for
179/// paste-back. Once pinned, subsequent runs verify normally. Treats the
180/// `blake3` field as the lockfile's *output*, not its precondition — first
181/// publish or disaster-recovery hydration just works.
182#[allow(dead_code)] // referenced by tests + serves as documentation of the sentinel literal
183const ZERO_SENTINEL_HEX: &str = "0000000000000000000000000000000000000000000000000000000000000000";
184
185/// True when `hex` is the 64-zero "not pinned yet" sentinel for `BlakeHash`.
186fn is_bootstrap_sentinel(hex: &str) -> bool {
187    hex.len() == 64 && hex.bytes().all(|b| b == b'0')
188}
189
190/// W209/F4: per-bootstrap output key on `$YAH_OUTPUTS`. The filename is
191/// included so multi-asset workloads round-trip — each asset's bind in the
192/// publish-assets pipeline TOML references its own key. Key shape:
193///   - `discovered_asset_blake3:<filename>` for post-transform output
194///   - `discovered_fetch_blake3:<filename>` for upstream content pin
195fn bootstrap_output_key(b: &BootstrappedHash) -> String {
196    let prefix = match b.kind {
197        BootstrapHashKind::Output => "discovered_asset_blake3",
198        BootstrapHashKind::Fetch => "discovered_fetch_blake3",
199        // W212/R518: the derivation key for the in-tree `[asset.derive.lock]`.
200        BootstrapHashKind::Input => "discovered_input_hash",
201    };
202    format!("{prefix}:{}", b.filename)
203}
204
205/// Append discovered BLAKE3 values to `$YAH_OUTPUTS` so the QED runner's
206/// per-step output collector picks them up after `yah cloud apply` returns
207/// (W209/F4). No-op outside a QED pipeline (env var unset).
208fn write_bootstrap_outputs(bootstrapped: &[BootstrappedHash]) -> std::io::Result<()> {
209    let Some(path) = std::env::var_os("YAH_OUTPUTS") else {
210        return Ok(());
211    };
212    append_bootstrap_outputs(Path::new(&path), bootstrapped)
213}
214
215/// Inner write — extracted so unit tests can exercise the format without
216/// racing on the `YAH_OUTPUTS` env var (cargo's parallel test runner makes
217/// process-wide env mutation unsafe).
218fn append_bootstrap_outputs(path: &Path, bootstrapped: &[BootstrappedHash]) -> std::io::Result<()> {
219    use std::io::Write;
220    let mut file = std::fs::OpenOptions::new()
221        .create(true)
222        .append(true)
223        .open(path)?;
224    for b in bootstrapped {
225        writeln!(file, "{}={}", bootstrap_output_key(b), b.hash)?;
226    }
227    Ok(())
228}
229
230/// Where a discovered BLAKE3 belongs in `workload.toml` for paste-back.
231#[derive(Debug, Clone, Copy, PartialEq, Eq)]
232pub enum BootstrapHashKind {
233    /// `[[asset]].blake3` — post-transform (or post-fetch when no transform) output.
234    Output,
235    /// `[asset.derive.fetch].blake3` — upstream content pin.
236    Fetch,
237    /// W212/R518: `[asset.derive.lock].input_hash` — the input-addressed
238    /// derivation key. Emitted on every successful transform build so the bind
239    /// path refreshes the lock whenever the inputs change.
240    Input,
241}
242
243/// One BLAKE3 value discovered during a bootstrap-mode apply.
244///
245/// Operator pastes `hash` into the field named by `kind` for the row matching
246/// `filename`. Subsequent runs verify against the pinned value.
247#[derive(Debug, Clone)]
248pub struct BootstrappedHash {
249    pub filename: String,
250    pub kind: BootstrapHashKind,
251    pub hash: String,
252}
253
254/// In-bucket sidecar: maps object key → blake3_hex for what was last uploaded.
255/// Used to detect prune candidates across runs without listing the bucket.
256type CatalogManifest = HashMap<String, String>;
257
258/// Summary of a completed static-asset sync.
259#[derive(Debug, Default)]
260pub struct StaticAssetSyncReport {
261    /// Asset filenames that were already present in the bucket with a matching
262    /// BLAKE3 stamp (genuinely skipped — the bytes are known to be current).
263    pub already_synced: Vec<String>,
264    /// Asset filenames uploaded this run.
265    pub uploaded: Vec<String>,
266    /// R546-B10: filenames whose bucket object existed but did NOT match the
267    /// declared artifact — either drifted content or a legacy object with no
268    /// `x-amz-meta-blake3` stamp — and was therefore re-uploaded. A key showing
269    /// up here repeatedly across runs means something outside this reconciler
270    /// keeps rewriting it; a key appearing exactly once is the expected
271    /// one-time heal of a pre-metadata object.
272    pub republished: Vec<String>,
273    /// Asset filenames whose source BLAKE3 hash didn't match the manifest.
274    /// These are NOT uploaded — operator must rebuild or fix the declaration.
275    pub hash_mismatch: Vec<String>,
276    /// Asset filenames that were in the stored catalog manifest but are no
277    /// longer in the current `workload.toml`. Prune candidates — not deleted.
278    pub prune_candidates: Vec<String>,
279    /// BLAKE3 values discovered during a bootstrap-mode apply (a `blake3` field
280    /// shipped as the zero sentinel). Operator pastes these into the catalog
281    /// to pin it; the assets were nonetheless uploaded this run.
282    pub bootstrapped: Vec<BootstrappedHash>,
283}
284
285/// Reconciles `kind = "static-asset"` components.
286///
287/// The `executor` field handles W164 transform recipes for derive-mode assets
288/// (R438-T15). Default is a [`PlacementRouter`] over `LocalForgeDriver` with no
289/// remote side: local recipes run on the host that owns the cache, and a
290/// remotely-placed one is refused with a message naming the missing wiring.
291///
292/// It is a router rather than a bare `LocalForgeDriver` because this crate
293/// *cannot* build the remote half — a `RemoteForgeDriver` needs a
294/// `WardenClient` over the camp's machine inventory, which lives one layer up
295/// in the CLI (R555-F3). So the default has to be the honest "no cloud from
296/// here" case, and the caller that does have one injects it through
297/// [`Self::with_executor`] (see `yah cloud apply`).
298pub struct StaticAssetReconciler {
299    executor: Arc<dyn ForgeExecutor>,
300}
301
302impl StaticAssetReconciler {
303    pub fn new() -> Self {
304        Self {
305            executor: Arc::new(PlacementRouter::local_only(Arc::new(
306                LocalForgeDriver::default(),
307            ))),
308        }
309    }
310
311    /// Swap the [`ForgeExecutor`] used to materialize derive-mode transforms.
312    /// Production callers inject a router that can also reach the camp's fleet;
313    /// tests inject a mock.
314    pub fn with_executor(mut self, executor: Arc<dyn ForgeExecutor>) -> Self {
315        self.executor = executor;
316        self
317    }
318}
319
320impl Default for StaticAssetReconciler {
321    fn default() -> Self {
322        Self::new()
323    }
324}
325
326#[async_trait]
327impl Reconciler for StaticAssetReconciler {
328    fn kind(&self) -> &'static str {
329        WORKLOAD_KIND
330    }
331
332    async fn up(&self, ctx: ReconcileCtx<'_>) -> Result<RunningWorkload> {
333        let workload_dir = ctx.workload_dir();
334
335        // Load and shape-validate the workload.
336        let workload = load_workload(&workload_dir)?;
337        shape_static_asset(&workload).with_context(|| {
338            format!(
339                "{}/workload.toml: closed-catalog invariant violated",
340                workload_dir.display()
341            )
342        })?;
343
344        // Resolve the object_store provider slot.
345        let slot = ctx.slot("object_store").with_context(|| {
346            format!(
347                "mirror has no `providers.object_store` slot — required for \
348                 kind=\"static-asset\" (service={}, env={})",
349                ctx.service.name, ctx.env,
350            )
351        })?;
352
353        let journal = AssetStatusJournal::at_workspace(ctx.workspace_root);
354        let service_name = ctx.service.name.as_str();
355
356        let report = match slot {
357            MirrorProviderSlot::Reference {
358                provider_id,
359                fields,
360            } => {
361                // Any cloudflare-*kind* provider routes here (dispatch is on
362                // the resolved kind, not the literal name) so a workspace can
363                // declare several — e.g. `cloudflare` + `cloudflare-scrabcake`.
364                let cf = super::cf_creds::CfProvider::resolve_scoped(
365                    ctx.workspace_root,
366                    provider_id,
367                    &ctx.scope.tenant,
368                    &ctx.scope.namespace,
369                )?;
370                anyhow::ensure!(
371                    matches!(cf.cfg.kind, Provider::Cloudflare),
372                    "providers.object_store.use = {provider_id:?} (kind={:?}) not supported \
373                     for static-asset — only cloudflare-kind reference providers are",
374                    cf.cfg.kind,
375                );
376                sync_to_r2(
377                    &ctx,
378                    cf,
379                    &workload,
380                    fields,
381                    self.executor.clone(),
382                    service_name,
383                    &journal,
384                )
385                .await?
386            }
387            MirrorProviderSlot::Inline {
388                kind: Provider::MinioContainer,
389                fields,
390            } => {
391                sync_to_minio(
392                    &ctx,
393                    &workload,
394                    fields,
395                    self.executor.clone(),
396                    service_name,
397                    &journal,
398                )
399                .await?
400            }
401            MirrorProviderSlot::Inline { kind, .. } => {
402                anyhow::bail!(
403                    "providers.object_store.kind = {kind:?} not supported for static-asset \
404                     (expected minio-container or a cloudflare reference)"
405                );
406            }
407        };
408
409        if !report.hash_mismatch.is_empty() {
410            warn!(
411                files = ?report.hash_mismatch,
412                "BLAKE3 mismatch for {} asset(s) — rebuild source files before syncing",
413                report.hash_mismatch.len(),
414            );
415            anyhow::bail!(
416                "static-asset sync failed: {} source file(s) have BLAKE3 mismatches: {:?}",
417                report.hash_mismatch.len(),
418                report.hash_mismatch,
419            );
420        }
421
422        if !report.prune_candidates.is_empty() {
423            warn!(
424                files = ?report.prune_candidates,
425                "{} file(s) no longer in catalog — run `yah service prune` to remove",
426                report.prune_candidates.len(),
427            );
428        }
429
430        // Surface bootstrap-mode discoveries as pipeline outputs (W209 §
431        // Migration #1). When this reconciler runs inside a QED step, the
432        // runner sets `$YAH_OUTPUTS` to a per-step KEY=VALUE sidechannel;
433        // discovered BLAKE3s flow into the run's OutputMap and feed any
434        // `[[bind]]` declarations in the publish-assets pipeline TOML — the
435        // applier writes them back into workload.toml mid-pipeline. Key
436        // shape encodes the asset filename so multi-asset workloads round-
437        // trip: `discovered_asset_blake3:<filename>=<hex>`.
438        if !report.bootstrapped.is_empty() {
439            if let Err(err) = write_bootstrap_outputs(&report.bootstrapped) {
440                // A write failure shouldn't poison the apply — the operator
441                // still has the summary info! line below for triage.
442                warn!(error = %err, "failed to write discovered BLAKE3s to $YAH_OUTPUTS");
443            }
444            info!(
445                count = report.bootstrapped.len(),
446                discoveries = ?report
447                    .bootstrapped
448                    .iter()
449                    .map(|b| format!("{}={}", bootstrap_output_key(b), b.hash))
450                    .collect::<Vec<_>>(),
451                "discovered {} BLAKE3 value(s) — surfaced via $YAH_OUTPUTS for pipeline bind",
452                report.bootstrapped.len(),
453            );
454        }
455
456        info!(
457            uploaded = report.uploaded.len(),
458            already_synced = report.already_synced.len(),
459            prune_candidates = report.prune_candidates.len(),
460            bootstrapped = report.bootstrapped.len(),
461            "static-asset sync complete",
462        );
463
464        Ok(
465            RunningWorkload::adopted(WORKLOAD_KIND, "object_store", None)
466                .with_notes(render_sync_notes(&report)),
467        )
468    }
469}
470
471/// R546-B12: what this run actually DID, one line per outcome class, for the
472/// apply console.
473///
474/// Everything above goes to `info!`, which the CLI does not render — so a clean
475/// static-asset reconcile printed the component header and then nothing at all.
476/// A successful publish and a successful no-op looked identical, and the view
477/// that would have told them apart (`yah cloud status`) was blind for the same
478/// underlying reason. The no-op case is stated explicitly rather than omitted:
479/// "nothing to do" is a result, and silence is not.
480fn render_sync_notes(report: &StaticAssetSyncReport) -> Vec<String> {
481    let mut notes = Vec::new();
482    if !report.uploaded.is_empty() {
483        notes.push(format!("published {} asset(s):", report.uploaded.len()));
484        notes.extend(report.uploaded.iter().map(|k| format!("  + {k}")));
485    }
486    if !report.republished.is_empty() {
487        notes.push(format!(
488            "re-published {} asset(s) whose bucket bytes did not match the declared hash:",
489            report.republished.len()
490        ));
491        notes.extend(report.republished.iter().map(|k| format!("  ~ {k}")));
492    }
493    if !report.already_synced.is_empty() {
494        notes.push(format!(
495            "{} asset(s) already current (bucket hash matches) — no upload",
496            report.already_synced.len()
497        ));
498    }
499    if !report.bootstrapped.is_empty() {
500        notes.push(format!(
501            "discovered {} BLAKE3 value(s) for paste-back into workload.toml:",
502            report.bootstrapped.len()
503        ));
504        notes.extend(
505            report
506                .bootstrapped
507                .iter()
508                .map(|b| format!("  {} = {}", bootstrap_output_key(b), b.hash)),
509        );
510    }
511    if !report.prune_candidates.is_empty() {
512        notes.push(format!(
513            "{} bucket object(s) no longer declared — `yah cloud cache prune` to review",
514            report.prune_candidates.len()
515        ));
516    }
517    if notes.is_empty() {
518        notes.push("no assets declared — nothing to reconcile".to_string());
519    }
520    notes
521}
522
523// ── Backend dispatch ──────────────────────────────────────────────────────────
524
525async fn sync_to_r2(
526    ctx: &ReconcileCtx<'_>,
527    cf_provider: super::cf_creds::CfProvider,
528    workload: &StaticAssetWorkload,
529    slot_fields: &std::collections::BTreeMap<String, toml::Value>,
530    executor: Arc<dyn ForgeExecutor>,
531    service_name: &str,
532    journal: &AssetStatusJournal,
533) -> Result<StaticAssetSyncReport> {
534    // `cf_provider` resolved from the mirror slot's `use = "<id>"` — supplies
535    // account_id + management token + R2 S3 keys, all per-provider.
536    let account_id = cf_provider.account_id.clone();
537
538    let bucket = slot_fields
539        .get("bucket")
540        .and_then(|v| v.as_str())
541        .context("providers.object_store missing `bucket` field for cloudflare static-asset sync")?
542        .to_string();
543
544    let (access_key, secret_key) = cf_provider.r2_keys()?;
545    let api_token = cf_provider.api_token()?;
546
547    ensure_r2_bucket(&api_token, &account_id, &bucket).await?;
548
549    let endpoint = format!("https://{account_id}.r2.cloudflarestorage.com");
550    let client = reqwest::Client::new();
551
552    sync_assets(
553        workload,
554        ctx.workspace_root,
555        &ctx.workload_dir(),
556        &client,
557        &endpoint,
558        &bucket,
559        R2_REGION,
560        &access_key,
561        &secret_key,
562        executor,
563        service_name,
564        journal,
565    )
566    .await
567}
568
569async fn sync_to_minio(
570    ctx: &ReconcileCtx<'_>,
571    workload: &StaticAssetWorkload,
572    slot_fields: &std::collections::BTreeMap<String, toml::Value>,
573    executor: Arc<dyn ForgeExecutor>,
574    service_name: &str,
575    journal: &AssetStatusJournal,
576) -> Result<StaticAssetSyncReport> {
577    use super::slot_field_u16;
578    use crate::reconciler::pond::DEFAULT_MINIO_API_PORT;
579
580    let api_port = slot_field_u16(slot_fields, "api_port").unwrap_or(DEFAULT_MINIO_API_PORT);
581    let bucket = slot_fields
582        .get("bucket")
583        .and_then(|v| v.as_str())
584        .context(
585            "providers.object_store missing `bucket` field for minio-container static-asset sync",
586        )?
587        .to_string();
588
589    let endpoint = format!("http://127.0.0.1:{api_port}");
590    let client = reqwest::Client::new();
591
592    sync_assets(
593        workload,
594        ctx.workspace_root,
595        &ctx.workload_dir(),
596        &client,
597        &endpoint,
598        &bucket,
599        MINIO_REGION,
600        DEFAULT_MINIO_USER,
601        DEFAULT_MINIO_PASSWORD,
602        executor,
603        service_name,
604        journal,
605    )
606    .await
607}
608
609// ── Core sync logic ───────────────────────────────────────────────────────────
610
611async fn sync_assets(
612    workload: &StaticAssetWorkload,
613    workspace_root: &Path,
614    workload_dir: &Path,
615    client: &reqwest::Client,
616    endpoint: &str,
617    bucket: &str,
618    region: &str,
619    access_key: &str,
620    secret_key: &str,
621    executor: Arc<dyn ForgeExecutor>,
622    service_name: &str,
623    journal: &AssetStatusJournal,
624) -> Result<StaticAssetSyncReport> {
625    let mut report = StaticAssetSyncReport::default();
626    let endpoint = endpoint.trim_end_matches('/');
627
628    // Load the stored catalog manifest to detect prune candidates.
629    let prior_manifest =
630        load_catalog_manifest(client, endpoint, bucket, region, access_key, secret_key).await;
631
632    // Catalog set for prune detection: current catalog filenames.
633    let current_filenames: std::collections::HashSet<&str> = workload
634        .assets
635        .iter()
636        .map(|a| a.filename.as_str())
637        .collect();
638
639    // Prune candidates: filenames in the stored manifest but not in current catalog.
640    for prior_key in prior_manifest.keys() {
641        if !current_filenames.contains(prior_key.as_str()) {
642            report.prune_candidates.push(prior_key.clone());
643        }
644    }
645
646    // Updated manifest built from this run.
647    let mut new_manifest = CatalogManifest::new();
648
649    for entry in &workload.assets {
650        // W212/R518: substituter fast-path. When the committed `[asset.derive.lock]`
651        // still matches the inputs recomputed from the current pins (purely
652        // local — no fetch) AND the bucket already holds the output, skip the
653        // entire build: no model download, no Docker transform, no PUT. This is
654        // the Nix-substituter / Bazel-remote-cache behaviour, backed by the
655        // checked-in lock as the action cache and R2 as the CAS.
656        if let Some(out_hash) = lock_skip_hash(entry, workspace_root).await {
657            // R630-B1: encode at URL-construction time (see `uri_encode_key`).
658            let object_url = format!("{endpoint}/{bucket}/{}", uri_encode_key(&entry.filename));
659            // R546-B10: the lock says WHAT the output should be; the bucket must
660            // actually hold it. An unstamped legacy object is not evidence of a
661            // match, so it falls through and rebuilds rather than blessing bytes
662            // we cannot identify — the expensive-but-correct direction.
663            let remote = head_object(client, &object_url, region, access_key, secret_key).await?;
664            if remote.matches(&out_hash) {
665                debug!(
666                    filename = %entry.filename,
667                    "derivation lock in sync + object content matches — skipping build",
668                );
669                report.already_synced.push(entry.filename.clone());
670                new_manifest.insert(entry.filename.clone(), out_hash);
671                continue;
672            }
673            if remote.exists() {
674                debug!(
675                    filename = %entry.filename,
676                    "derivation lock in sync but bucket object is unstamped or drifted \
677                     — not trusting it, falling through to build/publish",
678                );
679            }
680        }
681
682        // R438-T15: derive-mode assets materialize to a content-addressed cache
683        // path (W164); legacy `source = "..."` assets read straight from disk.
684        // Both arms return a real on-disk path that the existing BLAKE3 verify +
685        // S3 PUT loop below treats uniformly.
686        let materialized =
687            materialize_asset(entry, workspace_root, workload_dir, executor.as_ref())
688                .await
689                .with_context(|| format!("materializing asset {:?}", entry.filename))?;
690        let source_path = materialized.path;
691        // Capture before the discovered_fetch_hash is moved out below.
692        let fetch_was_bootstrap = materialized.discovered_fetch_hash.is_some();
693        // W212/R518: derivation key for this build — emitted below so the bind
694        // path refreshes `[asset.derive.lock].input_hash`.
695        let build_derive_key = materialized.derive_key.clone();
696
697        // Surface any discovered upstream BLAKE3 (bootstrap mode for
698        // `[asset.derive.fetch].blake3`). Collected into `report.bootstrapped`
699        // and emitted to `$YAH_OUTPUTS` by the summary block in `up` (W209/F4).
700        if let Some(hash) = materialized.discovered_fetch_hash {
701            debug!(
702                filename = %entry.filename,
703                discovered = %hash,
704                "bootstrap discovered upstream blake3",
705            );
706            report.bootstrapped.push(BootstrappedHash {
707                filename: entry.filename.clone(),
708                kind: BootstrapHashKind::Fetch,
709                hash,
710            });
711        }
712
713        // Read source file.
714        let body = tokio::fs::read(&source_path).await.with_context(|| {
715            format!(
716                "reading source file {} for asset {:?}",
717                source_path.display(),
718                entry.filename
719            )
720        })?;
721
722        // BLAKE3 verification — strict mode rejects mismatch and skips upload;
723        // bootstrap mode (entry.blake3 == zero sentinel) accepts the computed
724        // hash and surfaces it in the report for paste-back. The bytes are
725        // uploaded either way — the difference is whether the pin is being
726        // *verified against* (strict) or *discovered* (bootstrap).
727        let actual_hash = blake3_hex(&body);
728        if is_bootstrap_sentinel(&entry.blake3.0) {
729            debug!(
730                filename = %entry.filename,
731                discovered = %actual_hash,
732                "bootstrap discovered asset blake3",
733            );
734            report.bootstrapped.push(BootstrappedHash {
735                filename: entry.filename.clone(),
736                kind: BootstrapHashKind::Output,
737                hash: actual_hash.clone(),
738            });
739        } else if !hashes_equal(&actual_hash, &entry.blake3.0) {
740            warn!(
741                filename = %entry.filename,
742                declared = %entry.blake3.0,
743                actual = %actual_hash,
744                "BLAKE3 mismatch — source file doesn't match declared hash",
745            );
746            report.hash_mismatch.push(entry.filename.clone());
747            journal
748                .append(&AssetStatusEvent {
749                    at: Utc::now(),
750                    asset: format!("{service_name}:{}", entry.filename),
751                    from: None,
752                    to: AssetState::DriftBucket,
753                    bytes: None,
754                    blake3: None,
755                })
756                .await;
757            continue;
758        }
759
760        // W212/R518: a verified build emits its derivation key so the bind path
761        // refreshes `[asset.derive.lock].input_hash`. Paired with the output
762        // hash above, this records the action-cache entry the next run's
763        // substituter fast-path consults. Only emitted when a build actually
764        // ran (the lock-skip path `continue`d before reaching here).
765        if let Some(dk) = &build_derive_key {
766            report.bootstrapped.push(BootstrappedHash {
767                filename: entry.filename.clone(),
768                kind: BootstrapHashKind::Input,
769                hash: dk.clone(),
770            });
771        }
772
773        let key = &entry.filename;
774        // R630-B1: encode at URL-construction time (see `uri_encode_key`).
775        let object_url = format!("{endpoint}/{bucket}/{}", uri_encode_key(key));
776
777        // R546-B10: skip the PUT only when the bucket demonstrably holds THESE
778        // bytes. Existence alone is not enough — an object placed at this key
779        // out-of-band, or a rebuild that isn't byte-reproducible, leaves content
780        // that differs from the artifact we just verified. Skipping on mere
781        // existence let the bucket keep stale bytes while `new_manifest` below
782        // recorded the new hash, so the published catalog described bytes that
783        // were never uploaded.
784        //
785        // Note `actual_hash` is already verified equal to the declared
786        // `entry.blake3` at this point (the mismatch arm above `continue`s), so
787        // re-PUTting on drift converges the bucket toward the committed catalog
788        // rather than clobbering it with something unvetted.
789        let remote = head_object(client, &object_url, region, access_key, secret_key).await?;
790        if remote.matches(&actual_hash) {
791            debug!(key, "already present with matching blake3 — skipping PUT");
792            report.already_synced.push(key.clone());
793            new_manifest.insert(key.clone(), actual_hash);
794            continue;
795        }
796        if remote.exists() {
797            // Drifted, or a legacy object predating the metadata stamp. Either
798            // way re-PUT: identical bytes make it a cheap idempotent overwrite
799            // that heals the missing stamp, and differing bytes are exactly the
800            // case that must not be skipped. Legacy objects therefore cost ONE
801            // re-upload, after which the stamp makes the skip work again.
802            warn!(
803                key,
804                expected = %actual_hash,
805                remote = ?match &remote {
806                    RemoteObject::Present { blake3 } => blake3.as_deref(),
807                    RemoteObject::Absent => None,
808                },
809                "bucket object does not match the declared artifact — re-publishing",
810            );
811            report.republished.push(key.clone());
812        }
813
814        // PUT the file.
815        let body_sha256 = sha256_hex(&body);
816        let content_length = body.len();
817        let content_type = content_type_for(&source_path);
818        let headers = sign_s3_put_object(
819            &object_url,
820            &body_sha256,
821            content_type,
822            content_length,
823            region,
824            access_key,
825            secret_key,
826            // R546-B10: stamp the content hash so a later run can tell whether
827            // the bucket already holds these exact bytes.
828            Some(actual_hash.as_str()),
829        )
830        .with_context(|| format!("signing PUT {object_url}"))?;
831
832        let resp = client
833            .put(&object_url)
834            .headers(headers)
835            .body(body)
836            .send()
837            .await
838            .with_context(|| format!("PUT {object_url}"))?;
839
840        if !resp.status().is_success() {
841            let status = resp.status();
842            let body_text = resp.text().await.unwrap_or_default();
843            anyhow::bail!("PUT {object_url} → {status}: {}", body_text.trim());
844        }
845
846        info!(key, "uploaded");
847        report.uploaded.push(key.clone());
848        new_manifest.insert(key.clone(), actual_hash.clone());
849        let from_state = if fetch_was_bootstrap {
850            AssetState::PlaceholderFetch
851        } else if is_bootstrap_sentinel(&entry.blake3.0) {
852            AssetState::PlaceholderOutput
853        } else {
854            AssetState::PinnedNotPublished
855        };
856        journal
857            .append(&AssetStatusEvent {
858                at: Utc::now(),
859                asset: format!("{service_name}:{key}"),
860                from: Some(from_state),
861                to: AssetState::Published,
862                bytes: Some(content_length as u64),
863                blake3: Some(actual_hash),
864            })
865            .await;
866    }
867
868    // Save the updated manifest (non-fatal: data is already in the bucket).
869    if let Err(e) = save_catalog_manifest(
870        client,
871        endpoint,
872        bucket,
873        &new_manifest,
874        region,
875        access_key,
876        secret_key,
877    )
878    .await
879    {
880        warn!(
881            error = %e,
882            "failed to save asset catalog manifest (non-fatal) — prune detection may miss \
883             candidates on the next run"
884        );
885    }
886
887    Ok(report)
888}
889
890// ── W164 materialize step (R438-T15) ──────────────────────────────────────────
891
892/// Output of [`materialize_asset`]: the on-disk path the upload loop reads,
893/// plus any BLAKE3 discovered during bootstrap-mode fetch.
894struct MaterializedAsset {
895    path: PathBuf,
896    /// Set when `derive.fetch.blake3` was the zero sentinel and the reconciler
897    /// computed the actual upstream hash. Operator pastes this back into
898    /// `[asset.derive.fetch].blake3`. `None` in strict mode and for legacy
899    /// `source = "..."` assets.
900    discovered_fetch_hash: Option<String>,
901    /// W212/R518: the input-addressed derivation key for this build, emitted to
902    /// `$YAH_OUTPUTS` so the bind path refreshes `[asset.derive.lock].input_hash`.
903    /// `None` for `source` / no-transform assets.
904    derive_key: Option<String>,
905}
906
907/// BLAKE3 of a transform recipe's TOML bytes — an input to the derivation key
908/// (W212). The file carries the pinned container digest, steps, and placement,
909/// so hashing it captures every recipe-side input in one value. `None` when the
910/// recipe file is unreadable (the caller treats that as "can't memoize").
911async fn recipe_blake3(workspace_root: &Path, recipe: &str) -> Option<String> {
912    let path = workspace_root
913        .join(".yah/qed/transforms")
914        .join(format!("{recipe}.toml"));
915    let bytes = tokio::fs::read(&path).await.ok()?;
916    Some(blake3_hex(&bytes))
917}
918
919/// W212/R518: the substituter fast-path. Returns `Some(output_hash)` when the
920/// asset's committed `[asset.derive.lock]` is *current* — i.e. the derivation
921/// key recomputed from the COMMITTED pins (no network) equals `lock.input_hash`
922/// and the lock describes the currently-pinned output. The caller still
923/// confirms the bytes exist in the bucket before skipping the build. Returns
924/// `None` (→ build normally) for source assets, no-transform assets, bootstrap
925/// rows, a missing lock, or any input drift.
926async fn lock_skip_hash(entry: &AssetEntry, workspace_root: &Path) -> Option<String> {
927    let derive = entry.derive.as_ref()?;
928    let transform = derive.transform.as_ref()?;
929    let lock = derive.lock.as_ref()?;
930    // A sentinel pin means nothing has been pinned yet → must build.
931    if is_bootstrap_sentinel(&entry.blake3.0) || is_bootstrap_sentinel(&derive.fetch.blake3.0) {
932        return None;
933    }
934    // The lock must describe the currently-pinned output, else it's stale.
935    if !hashes_equal(&lock.output_blake3, &entry.blake3.0) {
936        return None;
937    }
938    // Recompute the derivation key from the committed pins — purely local, no
939    // fetch. The fetch pin is the declared input identity (a fixed-output
940    // derivation), so in strict mode it equals the bytes that fed the build.
941    let recipe_bk = recipe_blake3(workspace_root, &transform.recipe).await?;
942    let key = derivation_key(&derive.fetch.blake3.0, &recipe_bk, &transform.params);
943    (key == lock.input_hash).then(|| entry.blake3.0.clone())
944}
945
946/// Resolve an `[[asset]]` row to an on-disk path the upload loop can read.
947///
948/// Two modes:
949///
950/// - **Legacy `source`**: returns `workload_dir.join(source_rel)` — bytes are
951///   already on disk, the rest of the pipeline is unchanged.
952/// - **Derive (`derive.fetch[+transform]`)**: fetches the upstream blob into
953///   `.yah/cache/derive/fetch/<upstream-blake3>.bin`, optionally runs the
954///   transform recipe into `.yah/cache/derive/transform/<output-blake3>.bin`,
955///   and returns whichever cache path holds the final bytes. The shape
956///   validator (`shape_static_asset`) guarantees exactly one of the two is
957///   set, so the `else` branch is total.
958///
959/// **Bootstrap mode.** When `derive.fetch.blake3` or `entry.blake3` is the
960/// zero sentinel, the reconciler discovers the actual hash instead of
961/// verifying against the pinned value. The discovered upstream hash is
962/// returned in `discovered_fetch_hash`; the discovered output hash is
963/// computed at the upload site (the bytes are read there anyway) and
964/// surfaced separately.
965async fn materialize_asset(
966    entry: &AssetEntry,
967    workspace_root: &Path,
968    workload_dir: &Path,
969    executor: &dyn ForgeExecutor,
970) -> Result<MaterializedAsset> {
971    if let Some(source_rel) = entry.source.as_ref() {
972        return Ok(MaterializedAsset {
973            path: workload_dir.join(source_rel),
974            discovered_fetch_hash: None,
975            derive_key: None,
976        });
977    }
978
979    let derive = entry
980        .derive
981        .as_ref()
982        .expect("AssetEntry shape: exactly one of source/derive must be set (shape_static_asset)");
983
984    let cache_root = workspace_root.join(".yah/cache/derive");
985    let fetch_bootstrap = is_bootstrap_sentinel(&derive.fetch.blake3.0);
986    let (fetched_path, fetched_hash) =
987        materialize_fetch(&derive.fetch, &cache_root.join("fetch")).await?;
988    let discovered_fetch_hash = fetch_bootstrap.then(|| fetched_hash.clone());
989
990    let path = if let Some(transform) = &derive.transform {
991        materialize_transform(
992            transform,
993            &fetched_path,
994            &fetched_hash,
995            &entry.blake3.0,
996            &cache_root.join("transform"),
997            workspace_root,
998            executor,
999        )
1000        .await?
1001    } else {
1002        // No transform — the fetched bytes ARE the output. In strict mode,
1003        // fetch.blake3 must equal entry.blake3; let the upload loop's BLAKE3
1004        // verify surface any mismatch. In bootstrap mode the upload loop
1005        // discovers entry.blake3 directly from the file bytes.
1006        fetched_path
1007    };
1008
1009    // W212/R518: derivation key for this build — the same value the local action
1010    // cache used (fetched-input hash ⊕ recipe-file hash ⊕ params). Surfaced so
1011    // the bind path can refresh `[asset.derive.lock].input_hash`.
1012    let derive_key = match &derive.transform {
1013        Some(transform) => recipe_blake3(workspace_root, &transform.recipe)
1014            .await
1015            .map(|rb| derivation_key(&fetched_hash, &rb, &transform.params)),
1016        None => None,
1017    };
1018
1019    Ok(MaterializedAsset {
1020        path,
1021        discovered_fetch_hash,
1022        derive_key,
1023    })
1024}
1025
1026/// Maximum number of fetch attempts (initial try + retries).
1027#[cfg(not(test))]
1028const FETCH_MAX_ATTEMPTS: u32 = 5;
1029#[cfg(test)]
1030const FETCH_MAX_ATTEMPTS: u32 = 3; // fewer retries in tests
1031
1032/// Initial backoff between retries in milliseconds (doubles each attempt).
1033#[cfg(not(test))]
1034const FETCH_BASE_DELAY_MS: u64 = 1_000;
1035#[cfg(test)]
1036const FETCH_BASE_DELAY_MS: u64 = 1; // near-instant in tests
1037
1038/// Ceiling on retry backoff in milliseconds.
1039#[cfg(not(test))]
1040const FETCH_MAX_DELAY_MS: u64 = 30_000;
1041#[cfg(test)]
1042const FETCH_MAX_DELAY_MS: u64 = 5; // near-instant in tests
1043
1044/// Log download progress every N bytes.
1045const FETCH_LOG_INTERVAL_BYTES: u64 = 100 * 1024 * 1024; // 100 MiB
1046
1047/// Outcome of a single download attempt that did not fully succeed. Used by
1048/// [`fetch_once`] to distinguish a transient failure (retry + Range resume)
1049/// from a non-retriable 4xx (bail immediately).
1050enum FetchOnceFail {
1051    /// Server returned a non-retriable 4xx. Caller should surface and stop.
1052    Fatal(reqwest::StatusCode),
1053    /// Network error or retriable server response (5xx / 429). Partial file
1054    /// is preserved on disk for Range-header resumption on the next attempt.
1055    Transient(anyhow::Error),
1056}
1057
1058/// Fetch `fetch.url` into `<cache_dir>/<fetch.blake3>.bin`, verifying BLAKE3.
1059///
1060/// On cache HIT the network is skipped entirely; a hash mismatch surfaces as a
1061/// hard error (bit-rot / hand-edited cache, not silently re-fetched).
1062///
1063/// Cache MISS downloads the blob with:
1064/// - **Streaming to disk** — `Response::chunk()` rather than buffering the full
1065///   body in RAM (required for multi-GB blobs like whisper-large at 1.5 GB).
1066/// - **Exponential-backoff retry** — up to [`FETCH_MAX_ATTEMPTS`] attempts on
1067///   transient network errors and 5xx / 429 HTTP responses.
1068/// - **Range-header resume** — if a `.partial` file exists from a prior attempt,
1069///   subsequent tries send `Range: bytes=<offset>-` to avoid re-downloading
1070///   already-received bytes.
1071/// - **Progress logging** — `info!()` every 100 MiB surfaces download progress
1072///   through the task-pane / QED log surface (W164 OQ#4).
1073async fn materialize_fetch(fetch: &FetchSource, cache_dir: &Path) -> Result<(PathBuf, String)> {
1074    let bootstrap = is_bootstrap_sentinel(&fetch.blake3.0);
1075
1076    tokio::fs::create_dir_all(cache_dir)
1077        .await
1078        .with_context(|| format!("creating fetch cache dir {}", cache_dir.display()))?;
1079
1080    // Cache HIT — only meaningful when we know the expected hash. Bootstrap
1081    // mode has no stable cache key to look up; it always re-downloads on the
1082    // first run, then strict mode (after the operator pins the discovered
1083    // hash) takes the HIT path on subsequent runs.
1084    if !bootstrap {
1085        let cache_path = cache_dir.join(format!("{}.bin", fetch.blake3.0));
1086        if tokio::fs::try_exists(&cache_path).await.unwrap_or(false) {
1087            verify_blake3_path(&cache_path, &fetch.blake3.0)
1088                .await
1089                .with_context(|| {
1090                    format!(
1091                        "fetch cache HIT for {} but bytes don't match pinned BLAKE3 — \
1092                     rm the cache entry or fix the pin",
1093                        cache_path.display()
1094                    )
1095                })?;
1096            debug!(url = %fetch.url, "fetch cache HIT");
1097            return Ok((cache_path, fetch.blake3.0.clone()));
1098        }
1099    }
1100
1101    // Partial filename: stable per (URL, mode) so Range resume works across
1102    // attempts. In bootstrap mode every fetch shares the zero "expected" hash,
1103    // so partials are namespaced by URL hash to avoid cross-asset collisions.
1104    let partial_path = if bootstrap {
1105        cache_dir.join(format!(
1106            "bootstrap-{}.partial",
1107            blake3_hex(fetch.url.as_bytes())
1108        ))
1109    } else {
1110        cache_dir.join(format!("{}.partial", fetch.blake3.0))
1111    };
1112
1113    let client = reqwest::Client::new();
1114    let mut last_err = anyhow::anyhow!("no attempt made");
1115
1116    for attempt in 0..FETCH_MAX_ATTEMPTS {
1117        if attempt > 0 {
1118            let delay_ms = (FETCH_BASE_DELAY_MS << (attempt - 1)).min(FETCH_MAX_DELAY_MS);
1119            warn!(
1120                url = %fetch.url,
1121                attempt,
1122                delay_ms,
1123                err = %last_err,
1124                "fetch transient failure; retrying with backoff",
1125            );
1126            tokio::time::sleep(tokio::time::Duration::from_millis(delay_ms)).await;
1127        }
1128
1129        match fetch_once(&client, &fetch.url, &partial_path).await {
1130            Ok(()) => {
1131                // Compute hash once — used for verify (strict) and cache
1132                // naming (both modes; bootstrap mode names by the discovered
1133                // value rather than the zero placeholder).
1134                let bytes = tokio::fs::read(&partial_path)
1135                    .await
1136                    .with_context(|| format!("reading {} for BLAKE3", partial_path.display()))?;
1137                let actual = blake3_hex(&bytes);
1138                drop(bytes);
1139
1140                if !bootstrap && !hashes_equal(&actual, &fetch.blake3.0) {
1141                    anyhow::bail!(
1142                        "downloaded {} but BLAKE3 doesn't match pin — \
1143                         check the pin in workload.toml (actual {actual} != expected {})",
1144                        fetch.url,
1145                        fetch.blake3.0,
1146                    );
1147                }
1148
1149                let cache_path = cache_dir.join(format!("{actual}.bin"));
1150                tokio::fs::rename(&partial_path, &cache_path)
1151                    .await
1152                    .with_context(|| {
1153                        format!(
1154                            "promoting {} → {}",
1155                            partial_path.display(),
1156                            cache_path.display()
1157                        )
1158                    })?;
1159                if bootstrap {
1160                    info!(url = %fetch.url, discovered = %actual, "fetch complete (bootstrap)");
1161                } else {
1162                    info!(url = %fetch.url, "fetch complete");
1163                }
1164                return Ok((cache_path, actual));
1165            }
1166            Err(FetchOnceFail::Fatal(status)) => {
1167                anyhow::bail!("GET {} → HTTP {status} (not retriable)", fetch.url);
1168            }
1169            Err(FetchOnceFail::Transient(e)) => {
1170                last_err = e;
1171            }
1172        }
1173    }
1174
1175    Err(last_err).with_context(|| {
1176        format!(
1177            "GET {} failed after {FETCH_MAX_ATTEMPTS} attempts",
1178            fetch.url
1179        )
1180    })
1181}
1182
1183/// Single download attempt: GET `url` (with `Range` if `partial_path` is non-empty),
1184/// stream the response body to `partial_path`, and return once the body is exhausted.
1185///
1186/// Returns `Ok(())` when all bytes were received and flushed to disk.
1187/// Returns `Err(FetchOnceFail::Fatal)` for non-retriable 4xx responses.
1188/// Returns `Err(FetchOnceFail::Transient)` for connection errors, 5xx, or 429;
1189/// the partial file is left intact so the next attempt can resume via `Range`.
1190async fn fetch_once(
1191    client: &reqwest::Client,
1192    url: &str,
1193    partial_path: &Path,
1194) -> Result<(), FetchOnceFail> {
1195    use tokio::io::AsyncWriteExt;
1196
1197    let resume_offset = tokio::fs::metadata(partial_path)
1198        .await
1199        .ok()
1200        .map(|m| m.len())
1201        .filter(|&n| n > 0);
1202
1203    let mut req = client.get(url);
1204    if let Some(offset) = resume_offset {
1205        req = req.header(reqwest::header::RANGE, format!("bytes={offset}-"));
1206        info!(url, offset, "resuming partial download");
1207    } else {
1208        info!(url, "downloading");
1209    }
1210
1211    let resp = req
1212        .send()
1213        .await
1214        .map_err(|e| FetchOnceFail::Transient(anyhow::anyhow!("connecting to {url}: {e}")))?;
1215
1216    let status = resp.status();
1217    let is_partial_response = status == reqwest::StatusCode::PARTIAL_CONTENT;
1218
1219    if status == reqwest::StatusCode::RANGE_NOT_SATISFIABLE {
1220        // Partial file may already hold all bytes (prior complete-but-not-verified
1221        // attempt). Delete it and let the next retry start fresh.
1222        let _ = tokio::fs::remove_file(partial_path).await;
1223        return Err(FetchOnceFail::Transient(anyhow::anyhow!(
1224            "GET {url} → 416 Range Not Satisfiable; partial cleared"
1225        )));
1226    }
1227    if status.is_client_error() {
1228        return Err(FetchOnceFail::Fatal(status));
1229    }
1230    if status.is_server_error() || status == reqwest::StatusCode::TOO_MANY_REQUESTS {
1231        return Err(FetchOnceFail::Transient(anyhow::anyhow!(
1232            "GET {url} → HTTP {status}"
1233        )));
1234    }
1235    if !status.is_success() && !is_partial_response {
1236        return Err(FetchOnceFail::Fatal(status));
1237    }
1238
1239    // Server returned 200 instead of 206 — it ignored our Range request.
1240    // Discard any partial bytes and stream the full body fresh.
1241    if !is_partial_response && resume_offset.is_some() {
1242        let _ = tokio::fs::remove_file(partial_path).await;
1243    }
1244
1245    let mut file = tokio::fs::OpenOptions::new()
1246        .write(true)
1247        .create(true)
1248        .append(is_partial_response)
1249        .truncate(!is_partial_response)
1250        .open(partial_path)
1251        .await
1252        .map_err(|e| {
1253            FetchOnceFail::Transient(anyhow::anyhow!("opening {}: {e}", partial_path.display()))
1254        })?;
1255
1256    let mut downloaded = if is_partial_response {
1257        resume_offset.unwrap_or(0)
1258    } else {
1259        0
1260    };
1261    let mut resp = resp;
1262    loop {
1263        match resp.chunk().await {
1264            Ok(Some(chunk)) => {
1265                file.write_all(&chunk).await.map_err(|e| {
1266                    FetchOnceFail::Transient(anyhow::anyhow!(
1267                        "writing chunk to {}: {e}",
1268                        partial_path.display()
1269                    ))
1270                })?;
1271                let prev = downloaded;
1272                downloaded += chunk.len() as u64;
1273                if downloaded / FETCH_LOG_INTERVAL_BYTES > prev / FETCH_LOG_INTERVAL_BYTES {
1274                    info!(
1275                        url,
1276                        downloaded_mib = downloaded / (1024 * 1024),
1277                        "fetch progress",
1278                    );
1279                }
1280            }
1281            Ok(None) => break,
1282            Err(e) => {
1283                return Err(FetchOnceFail::Transient(anyhow::anyhow!(
1284                    "reading body from {url}: {e}"
1285                )));
1286            }
1287        }
1288    }
1289
1290    file.flush().await.map_err(|e| {
1291        FetchOnceFail::Transient(anyhow::anyhow!("flushing {}: {e}", partial_path.display()))
1292    })?;
1293
1294    Ok(())
1295}
1296
1297/// Run the transform recipe against `input_path`, writing the output to
1298/// `<cache_dir>/<actual_blake3>.bin`. In strict mode the actual hash is
1299/// verified against `output_blake3`; in bootstrap mode (`output_blake3` ==
1300/// zero sentinel) the actual hash is accepted as-is and named accordingly.
1301/// Cache HIT (file present + hash matches) skips recipe execution; bootstrap
1302/// mode has no stable cache key and always re-runs the recipe.
1303/// W212 schema version for the derivation key. Bump when the *way* we hash or
1304/// lower inputs changes (not when an individual recipe changes — that's
1305/// captured by the recipe-file hash). A bump invalidates every action-cache
1306/// entry, forcing a clean rebuild under the new keying.
1307const DERIVE_KEY_SCHEMA: u32 = 1;
1308
1309/// Compute the input-addressed action-cache key for a derive-mode transform
1310/// (W212). A pure function of the complete declared input set:
1311///
1312/// - `input_hash` — BLAKE3 of the fetched transform input (the fixed-output /
1313///   version anchor; for whisper this is `config.json`, which the model's repo
1314///   version tracks).
1315/// - `recipe_blake3` — BLAKE3 of the recipe TOML bytes. The file carries the
1316///   pinned container digest, the steps, and the placement, so hashing it
1317///   covers all three: a digest bump or a step edit flips the key.
1318/// - `params` — the workload-side invocation params (`BTreeMap` is already
1319///   sorted, so the encoding is deterministic).
1320///
1321/// **Hermeticity invariant:** every input that can change the output bytes MUST
1322/// feed this key. A missing dimension reintroduces the W164/R438 silent-stale
1323/// cache bug. The per-dimension unit test guards this.
1324fn derivation_key(
1325    input_hash: &str,
1326    recipe_blake3: &str,
1327    params: &BTreeMap<String, String>,
1328) -> String {
1329    let mut buf = String::new();
1330    buf.push_str(&format!("derive-key-v{DERIVE_KEY_SCHEMA}\n"));
1331    buf.push_str(&format!("input={input_hash}\n"));
1332    buf.push_str(&format!("recipe={recipe_blake3}\n"));
1333    for (k, v) in params {
1334        buf.push_str(&format!("param.{k}={v}\n"));
1335    }
1336    blake3_hex(buf.as_bytes())
1337}
1338
1339/// Read a recorded action-cache output hash, if present. A missing or
1340/// unreadable / empty entry is a cache miss, never an error.
1341async fn read_action_cache(ac_path: &Path) -> Option<String> {
1342    let raw = tokio::fs::read_to_string(ac_path).await.ok()?;
1343    let hash = raw.trim();
1344    (!hash.is_empty()).then(|| hash.to_string())
1345}
1346
1347/// Record an action-cache entry (`derive_key → output hash`). Written
1348/// atomically via tmp + rename so a crashed run never leaves a half-written
1349/// entry a later run would trust.
1350///
1351/// R925: the staging path used to be `<derive_key>.out.tmp`, fixed and therefore
1352/// shared. The cache root is `<workspace_root>/.yah/cache/derive/transform/`, a
1353/// per-*workspace* directory that every `yah` process in that tree writes, so two
1354/// concurrent runs of the same `derive_key` interleaved into one staging file.
1355/// The damage is silent in both directions: [`read_action_cache`] reads an empty
1356/// or truncated entry as a cache MISS rather than an error, and a torn hash that
1357/// still looks well-formed points at a CAS entry that does not exist, poisoning
1358/// the entry so every later run falls through to a full re-materialisation.
1359async fn write_action_cache(ac_path: &Path, output_hash: &str) -> Result<()> {
1360    crate::atomic_write::write_atomic_async(ac_path, format!("{output_hash}\n").as_bytes())
1361        .await
1362        .with_context(|| format!("recording action-cache entry {}", ac_path.display()))
1363}
1364
1365/// Removes [`materialize_transform`]'s staging output when the call unwinds
1366/// without publishing it (R925).
1367///
1368/// `Drop` rather than a cleanup call because the path between creating the file
1369/// and the publishing rename has a dozen exits — the executor's error, the
1370/// "recipe produced no output" bail, the BLAKE3 mismatch bail, every `?` on a
1371/// read — and a cleanup that is only on the ones someone remembered is the same
1372/// leak with more code. There is no async `Drop`; a blocking `unlink` of one
1373/// file in a local cache directory is the right trade rather than a reason to
1374/// keep a collidable fixed name.
1375struct StagingOutput(PathBuf);
1376
1377impl Drop for StagingOutput {
1378    fn drop(&mut self) {
1379        // ENOENT on the success path: the rename already consumed it.
1380        let _ = std::fs::remove_file(&self.0);
1381    }
1382}
1383
1384async fn materialize_transform(
1385    transform: &TransformSpec,
1386    input_path: &Path,
1387    input_hash: &str,
1388    output_blake3: &str,
1389    cache_dir: &Path,
1390    workspace_root: &Path,
1391    executor: &dyn ForgeExecutor,
1392) -> Result<PathBuf> {
1393    let bootstrap = is_bootstrap_sentinel(output_blake3);
1394
1395    tokio::fs::create_dir_all(cache_dir)
1396        .await
1397        .with_context(|| format!("creating transform cache dir {}", cache_dir.display()))?;
1398
1399    // Load + hash the recipe up front: the recipe TOML's bytes are an INPUT to
1400    // the derivation key (W212), so the key — and therefore the cache decision —
1401    // depends on them. The file carries the pinned container digest, steps, and
1402    // placement, so a change to any of those flips the key and forces a rebuild.
1403    let transforms_dir = workspace_root.join(".yah/qed/transforms");
1404    let loader = TransformRecipeLoader::new(&transforms_dir);
1405    let recipe_path = loader.recipe_path(&transform.recipe);
1406    let recipe = loader.load(&transform.recipe).with_context(|| {
1407        format!(
1408            "loading recipe {:?} from {}",
1409            transform.recipe,
1410            transforms_dir.display()
1411        )
1412    })?;
1413    let recipe_bytes = tokio::fs::read(&recipe_path).await.with_context(|| {
1414        format!(
1415            "reading recipe {} for derivation key",
1416            recipe_path.display()
1417        )
1418    })?;
1419    let recipe_blake3 = blake3_hex(&recipe_bytes);
1420
1421    // W212: input-addressed action cache. The skip decision is a pure function
1422    // of the declared INPUTS (fetched input content ⊕ recipe definition ⊕
1423    // invocation params ⊕ lowering-schema version), not of the expected output
1424    // hash. A recipe / digest / param change flips `derive_key` → guaranteed
1425    // miss → rebuild, eliminating the warm-cache silent-stale skip (the
1426    // W164/R438 defect this fixes). Applies in bootstrap mode too: identical
1427    // inputs surface the recorded output without re-running.
1428    let derive_key = derivation_key(input_hash, &recipe_blake3, &transform.params);
1429    let ac_path = cache_dir.join("ac").join(format!("{derive_key}.out"));
1430
1431    if let Some(recorded) = read_action_cache(&ac_path).await {
1432        let cas_path = cache_dir.join(format!("{recorded}.bin"));
1433        if tokio::fs::try_exists(&cas_path).await.unwrap_or(false)
1434            && verify_blake3_path(&cas_path, &recorded).await.is_ok()
1435        {
1436            // The output pin (`entry.blake3`) is still enforced downstream by
1437            // the upload loop's strict BLAKE3 verify, so a pin that has drifted
1438            // from the inputs surfaces there even though we skip the transform.
1439            debug!(recipe = %recipe.name, derive_key, "derivation cache HIT — skipping transform");
1440            return Ok(cas_path);
1441        }
1442        // AC entry present but the CAS bytes are gone / corrupt → fall through
1443        // and re-materialise (the AC records what the output should be; the CAS
1444        // lost the bytes).
1445    }
1446
1447    // Absolute workspace_root for the container bind-mount — docker rejects
1448    // `-w .` ("the working directory '.' is invalid"). Callers may pass a
1449    // relative path (the CLI `--path` default is "."), so canonicalize once.
1450    let workspace_abs = workspace_root
1451        .canonicalize()
1452        .with_context(|| format!("canonicalizing workspace root {}", workspace_root.display()))?;
1453
1454    // R546-B8: the IN/OUT bindings handed to the recipe MUST be absolute.
1455    // `cache_dir` is derived from workspace_root, which is "." by CLI default,
1456    // so these used to be relative (`./.yah/cache/derive/transform/<key>.tmp`).
1457    // A relative path only survives if the recipe never leaves its cwd — true
1458    // for whisper's recipes, FALSE for rusty-v8: build-v8.sh chdirs into a
1459    // scratch dir (/tmp/tmp.XXXX/v8src) to build V8, so it wrote its finished
1460    // tar to "<scratch>/./.yah/cache/..." INSIDE the container, invisible to the
1461    // host. The step then exited 0 and the reconciler failed reading a file that
1462    // was never going to be there — after a ~2h build. Canonicalize the parent
1463    // (it exists; created above) and join the filename, since the tmp output
1464    // itself does not exist yet and cannot be canonicalized directly.
1465    let cache_dir_abs = cache_dir
1466        .canonicalize()
1467        .with_context(|| format!("canonicalizing transform cache dir {}", cache_dir.display()))?;
1468
1469    // The derivation key is unique per INPUT SET, which stops two different rows
1470    // (and zero-sentinel bootstrap rows) colliding — but R925: it does NOT stop
1471    // two concurrent runs of the SAME row. This cache dir hangs off
1472    // `workspace_root`, so every `yah` process in the tree shares it, and
1473    // `<derive_key>.tmp` handed the same path to two recipes at once as
1474    // {{YAH_TRANSFORM_OUT}}. In strict mode the spliced bytes fail the BLAKE3
1475    // check below, i.e. a hard failure on a run that did nothing wrong, possibly
1476    // after hours of build. In BOOTSTRAP mode there is nothing to check against:
1477    // the torn output is hashed, published to the CAS under that hash, and
1478    // surfaced as the discovered pin. So: pid + sequence disambiguated.
1479    let tmp_output =
1480        crate::atomic_write::staging_path(&cache_dir_abs.join(format!("{derive_key}.out")));
1481    // A unique staging name is never reused, so a leaked one accumulates forever
1482    // in a long-lived cache dir. This unlinks it on every exit below — the `?`s,
1483    // the `bail!`s, and a cancelled future alike — and no-ops (ENOENT) on the
1484    // success path, where the publishing rename has already consumed it.
1485    let _staging = StagingOutput(tmp_output.clone());
1486
1487    let input_abs = input_path
1488        .canonicalize()
1489        .unwrap_or_else(|_| input_path.to_path_buf());
1490
1491    // R555-F3: a remotely-placed recipe runs on a worker with its own
1492    // filesystem, so the IN/OUT contract has to be re-read there. `remote_out`
1493    // is the worker-side path bound as {{YAH_TRANSFORM_OUT}}; the bytes are
1494    // pulled back to `tmp_output` after the step, and everything downstream
1495    // (exists-check, BLAKE3, rename, action cache) is untouched.
1496    let remote_out = match &recipe.placement.location {
1497        TaskLocation::Local => {
1498            refuse_secrets_on_a_local_recipe(&recipe)?;
1499            None
1500        }
1501        _ => Some(remote_transform_out(&recipe, &derive_key)?),
1502    };
1503
1504    let mut params: BTreeMap<String, String> = BTreeMap::new();
1505    params.insert(
1506        ENV_TRANSFORM_IN_0.to_string(),
1507        input_abs.to_string_lossy().into_owned(),
1508    );
1509    params.insert(
1510        ENV_TRANSFORM_OUT.to_string(),
1511        remote_out
1512            .as_ref()
1513            .unwrap_or(&tmp_output)
1514            .to_string_lossy()
1515            .into_owned(),
1516    );
1517    for (k, v) in &transform.params {
1518        params.insert(k.clone(), v.clone());
1519    }
1520
1521    info!(recipe = %recipe.name, "transform cache MISS, running");
1522    for step in &recipe.steps {
1523        let argv = substitute_argv(&step.argv, &params);
1524        // substitute_argv preserves `{{key}}` for unknown keys (no shell, no
1525        // string concat). Surface a missing-binding bug as a hard error rather
1526        // than letting the subprocess receive a literal placeholder.
1527        if let Some(unresolved) = argv.iter().find(|a| a.contains("{{")) {
1528            anyhow::bail!(
1529                "recipe {:?} step {:?}: unresolved placeholder in argv element {:?}",
1530                recipe.name,
1531                step.name,
1532                unresolved,
1533            );
1534        }
1535        // R555-T2 opened the recipe surface to `remote` / `remote_any`, so the
1536        // lowered spec carries whatever the recipe declared and reaches the
1537        // matching driver through the injected router (R555-F3). The two
1538        // placements need different execution context, and the difference is
1539        // invisible at the type level — hence the split below rather than one
1540        // shared `ctx`.
1541        let spec = lower_recipe_step_to_forge_spec(&recipe, step, argv);
1542        let mut ctx = ExecContext::default();
1543        if let Some(remote_out) = &remote_out {
1544            // NO cwd: `workspace_abs` is a path on THIS box. The local driver
1545            // bind-mounts it; a remote workdir is just a string interpreted on
1546            // the worker, where that path doesn't exist. A recipe needing a
1547            // source tree on the worker must bring it (the rusty-v8 builder
1548            // image clones its own).
1549            ctx = ctx.with_produced(remote_out.clone(), tmp_output.clone());
1550            // R555-F4 / W235 §(c): a signed recipe carries its admission grant
1551            // to kamaji. The grant is DERIVED from the recipe here rather than
1552            // stored in the TOML, so editing what a recipe runs un-signs it;
1553            // `None` for an unsigned recipe, which a node running the default
1554            // `permissive` policy still admits.
1555            if let Some(envelope) =
1556                velveteen_exec::envelope_for_recipe_step(&recipe, step)
1557                    .with_context(|| format!("deriving admission grant for recipe {:?}", recipe.name))?
1558            {
1559                ctx = ctx.with_admission(envelope);
1560            }
1561            // R555-F5: the credentials the recipe declared, mounted per run.
1562            // Only ever on the remote leg — the local driver refuses them,
1563            // because a dev box has no cluster secret store to resolve them
1564            // from and a silently-dropped credential fails deep inside the
1565            // build instead of here.
1566            if !recipe.secrets.is_empty() {
1567                ctx = ctx
1568                    .with_secrets(recipe.secrets.iter().map(|s| s.to_mount()).collect());
1569            }
1570        } else {
1571            ctx = ctx.with_cwd(workspace_abs.clone());
1572            // `platform` asks a HOST container runtime for foreign-arch
1573            // emulation. It has no remote referent — a remote run picks
1574            // architecture by scheduling (mesh_tags) — and RemoteForgeDriver
1575            // refuses it rather than handing back a wrong-arch artifact, so it
1576            // is only ever applied on the local leg.
1577            if let Some(platform) = &recipe.placement.platform {
1578                ctx = ctx.with_platform(platform.clone());
1579            }
1580        }
1581        let outcome = executor
1582            .execute(spec, ctx, None)
1583            .await
1584            .with_context(|| format!("executing recipe {:?} step {:?}", recipe.name, step.name))?;
1585        if !outcome.succeeded() {
1586            anyhow::bail!(
1587                "recipe {:?} step {:?} failed ({}): {}",
1588                recipe.name,
1589                step.name,
1590                outcome.status.discriminant(),
1591                outcome.stderr_tail
1592            );
1593        }
1594    }
1595
1596    // R546-B8: every step exited 0, so if the output is missing the recipe
1597    // simply never wrote it. Say THAT, rather than letting the `read` below
1598    // surface `No such file or directory` on a cache path — which reads like
1599    // cache corruption and sent the last person debugging the wrong subsystem
1600    // after a ~2h build. The relative-OUT bug that caused it is fixed above,
1601    // but a recipe can still write to the wrong place on its own, and this is
1602    // the only moment we can tell the operator exactly which contract broke.
1603    if !tokio::fs::try_exists(&tmp_output).await.unwrap_or(false) {
1604        anyhow::bail!(
1605            "recipe {:?} completed successfully but produced no output at {} \
1606             ({ENV_TRANSFORM_OUT}). The recipe's last step must write its artifact \
1607             to that exact path — check that it doesn't chdir and then write to a \
1608             relative location, and that it isn't writing to a directory instead \
1609             of a file.",
1610            recipe.name,
1611            // A remote run was told a worker-side path; naming `tmp_output` (a
1612            // local path it never saw) would send the reader hunting for a
1613            // binding bug that isn't there.
1614            remote_out.as_ref().unwrap_or(&tmp_output).display(),
1615        );
1616    }
1617
1618    // Compute the actual hash once — used for verify (strict) and cache
1619    // naming (both modes). Bootstrap mode names by the discovered value;
1620    // strict mode names by the expected value (which equals actual after
1621    // the verify below).
1622    let output_bytes = tokio::fs::read(&tmp_output).await.with_context(|| {
1623        format!(
1624            "reading transform output {} for BLAKE3",
1625            tmp_output.display()
1626        )
1627    })?;
1628    let actual = blake3_hex(&output_bytes);
1629    drop(output_bytes);
1630
1631    if !bootstrap && !hashes_equal(&actual, output_blake3) {
1632        anyhow::bail!(
1633            "transform output from recipe {:?} doesn't match pinned BLAKE3 \
1634             (actual {actual} != expected {output_blake3})",
1635            recipe.name
1636        );
1637    }
1638
1639    // tmp_output already lives in cache_dir — atomic publish is one rename.
1640    let cache_path = cache_dir.join(format!("{actual}.bin"));
1641    tokio::fs::rename(&tmp_output, &cache_path)
1642        .await
1643        .with_context(|| {
1644            format!(
1645                "renaming {} → {}",
1646                tmp_output.display(),
1647                cache_path.display()
1648            )
1649        })?;
1650
1651    // W212: record the action-cache entry (derive_key → output hash) so the
1652    // next run with identical inputs skips this transform entirely.
1653    write_action_cache(&ac_path, &actual)
1654        .await
1655        .with_context(|| format!("recording action-cache entry {}", ac_path.display()))?;
1656
1657    if bootstrap {
1658        info!(
1659            recipe = %recipe.name,
1660            discovered = %actual,
1661            "transform complete (bootstrap)",
1662        );
1663    }
1664
1665    Ok(cache_path)
1666}
1667
1668/// State of the committed `[asset.derive.fetch].blake3` pin, relative to the
1669/// fetched-input hash a seed run actually keyed its derivation on (R546-B6).
1670///
1671/// The W212 substituter fast-path ([`lock_skip_hash`]) recomputes the derivation
1672/// key from the COMMITTED pins with no network access, so it bails immediately
1673/// when the fetch pin is still the zero sentinel. Seeding a lock without also
1674/// pinning the fetch input therefore leaves that fast-path disarmed for every
1675/// machine except the one that seeded (whose local action cache masks it).
1676#[derive(Debug, Clone, PartialEq, Eq)]
1677pub enum FetchPinState {
1678    /// Committed pin equals the fetched-input hash — the fast-path can engage.
1679    Pinned,
1680    /// Committed pin is the 64-zero sentinel — the fast-path is disarmed until
1681    /// the operator pastes the fetched-input hash back into the workload.
1682    Sentinel,
1683    /// Committed pin names a DIFFERENT blob than the one this seed fetched.
1684    /// Either the version anchor moved or the pin is wrong; the fast-path will
1685    /// silently decline (input drift) until they agree.
1686    Mismatch { committed: String },
1687}
1688
1689/// Classify the committed `[asset.derive.fetch].blake3` against the hash the
1690/// fetch actually resolved to (R546-B6). Pure so the three arms are testable
1691/// without a network fetch.
1692fn classify_fetch_pin(committed: &str, fetched: &str) -> FetchPinState {
1693    if is_bootstrap_sentinel(committed) {
1694        FetchPinState::Sentinel
1695    } else if hashes_equal(committed, fetched) {
1696        FetchPinState::Pinned
1697    } else {
1698        FetchPinState::Mismatch {
1699            committed: committed.to_string(),
1700        }
1701    }
1702}
1703
1704/// Outcome of seeding the transform derivation cache from a pre-built artifact
1705/// (the qed→W164 bridge, R546-T3).
1706#[derive(Debug, Clone)]
1707pub struct SeededDerivation {
1708    /// The W212 derivation key the seeded action-cache entry is filed under —
1709    /// byte-identical to the one [`materialize_transform`] will look up, so the
1710    /// next `yah cloud apply` HITs it. Equals the value to paste into
1711    /// `[asset.derive.lock].input_hash`.
1712    pub derive_key: String,
1713    /// BLAKE3 of the pre-built artifact — the transform output hash. Equals the
1714    /// value to paste into `[[asset]].blake3` and `[asset.derive.lock].output_blake3`.
1715    pub output_blake3: String,
1716    /// BLAKE3 of the fetched *input* the derivation key was keyed on — the
1717    /// fourth paste-back value, `[asset.derive.fetch].blake3` (R546-B6). Without
1718    /// it the seeded lock is unusable off the seeding machine.
1719    pub fetch_blake3: String,
1720    /// What the workload currently commits for that pin, when the caller went
1721    /// through [`seed_derivation_for_target`] (which has the `[[asset]]` row in
1722    /// hand). `None` for direct [`seed_transform_derivation`] callers.
1723    pub fetch_pin: Option<FetchPinState>,
1724    /// The content-addressed store path the artifact was landed at.
1725    pub cas_path: PathBuf,
1726}
1727
1728/// Seed the transform derivation cache (AC + CAS) from an artifact built OUT OF
1729/// BAND — the qed→W164 bridge (R546-T3).
1730///
1731/// # Why this exists
1732///
1733/// The static-asset reconciler lowers every transform to
1734/// [`TaskLocation::Local`](task::TaskLocation::Local) (see
1735/// [`lower_recipe_step_to_forge_spec`]) — it has no fleet-offload path, so on a
1736/// foreign-arch host it can only build the recipe under emulation (which OOMs for
1737/// the rusty_v8 musl build, the reason R546 exists). `yah qed run <pipeline>`
1738/// (e.g. rusty-v8-musl) DOES offload to an arch-matched build-worker and, via R590-F6,
1739/// retrieves the produced tar content-addressed onto the caller. This function
1740/// bridges that pre-built tar into the W164 substituter: it writes the exact
1741/// action-cache + CAS entries [`materialize_transform`] would have written, so
1742/// the next `yah cloud apply` finds a HIT and PUBLISHES the pre-built bytes to
1743/// the consumer key + records the lock, instead of re-running the transform.
1744///
1745/// # Hermeticity
1746///
1747/// The derivation key is computed with the reconciler's OWN [`derivation_key`] +
1748/// [`recipe_blake3`], keyed on the SAME `fetched_hash` the reconciler derives
1749/// from the version anchor. Reusing those (not reimplementing) is what makes the
1750/// seeded key byte-identical to the lookup key — a drift would reintroduce the
1751/// W164/R438 silent-stale-cache bug. `fetched_hash` MUST be the value
1752/// [`materialize_fetch`] returns for `derive.fetch` (see [`seed_derivation_for_target`]).
1753pub async fn seed_transform_derivation(
1754    workspace_root: &Path,
1755    transform: &TransformSpec,
1756    fetched_hash: &str,
1757    artifact_path: &Path,
1758) -> Result<SeededDerivation> {
1759    let recipe_bk = recipe_blake3(workspace_root, &transform.recipe)
1760        .await
1761        .ok_or_else(|| {
1762            anyhow::anyhow!(
1763                "recipe {:?} unreadable under .yah/qed/transforms — cannot compute derivation key",
1764                transform.recipe,
1765            )
1766        })?;
1767    let derive_key = derivation_key(fetched_hash, &recipe_bk, &transform.params);
1768
1769    let bytes = tokio::fs::read(artifact_path)
1770        .await
1771        .with_context(|| format!("reading pre-built artifact {}", artifact_path.display()))?;
1772    let output_blake3 = blake3_hex(&bytes);
1773
1774    let cache_dir = workspace_root.join(".yah/cache/derive/transform");
1775    tokio::fs::create_dir_all(&cache_dir)
1776        .await
1777        .with_context(|| format!("creating transform cache dir {}", cache_dir.display()))?;
1778
1779    // CAS: <cache_dir>/<output_blake3>.bin — the exact name materialize_transform
1780    // renames its output to (so its HIT-path verify_blake3_path passes). Atomic
1781    // tmp+rename; idempotent when the entry already exists.
1782    let cas_path = cache_dir.join(format!("{output_blake3}.bin"));
1783    if !tokio::fs::try_exists(&cas_path).await.unwrap_or(false) {
1784        // R925: `<output_blake3>.seed.tmp` was fixed, and this cache dir is
1785        // shared by every process in the workspace — two concurrent seeds of the
1786        // same artifact spliced into one staging file and published the result.
1787        // A CAS entry is trusted by hash, so a torn one is only caught later by
1788        // the HIT-path verify, as a hard failure on a run that did nothing wrong.
1789        crate::atomic_write::write_atomic_async(&cas_path, &bytes)
1790            .await
1791            .with_context(|| format!("seeding CAS entry {}", cas_path.display()))?;
1792    }
1793
1794    // Action cache: <cache_dir>/ac/<derive_key>.out = <output_blake3>. Reuses the
1795    // reconciler's writer so the on-disk format can't drift from the reader.
1796    let ac_path = cache_dir.join("ac").join(format!("{derive_key}.out"));
1797    write_action_cache(&ac_path, &output_blake3)
1798        .await
1799        .with_context(|| format!("recording seeded action-cache entry {}", ac_path.display()))?;
1800
1801    info!(
1802        recipe = %transform.recipe,
1803        derive_key,
1804        output_blake3,
1805        fetch_blake3 = fetched_hash,
1806        "seeded transform derivation cache from pre-built artifact (qed→W164 bridge)",
1807    );
1808
1809    Ok(SeededDerivation {
1810        derive_key,
1811        output_blake3,
1812        fetch_blake3: fetched_hash.to_string(),
1813        fetch_pin: None,
1814        cas_path,
1815    })
1816}
1817
1818/// CLI-facing entry for the qed→W164 bridge (R546-T3): seed the derivation cache
1819/// for the `[[asset]]` row whose `derive.transform.params["target"]` equals
1820/// `target`, from a pre-built `artifact_path`.
1821///
1822/// Loads the static-asset workload, resolves the matching asset row, computes the
1823/// fetched-input hash via the reconciler's own [`materialize_fetch`] (downloading
1824/// the version anchor once; cached thereafter), then delegates to
1825/// [`seed_transform_derivation`]. After this returns, `yah cloud apply` on the
1826/// same service publishes the pre-built bytes without re-running the transform.
1827pub async fn seed_derivation_for_target(
1828    workspace_root: &Path,
1829    workload_path: &Path,
1830    target: &str,
1831    artifact_path: &Path,
1832) -> Result<SeededDerivation> {
1833    let raw = tokio::fs::read_to_string(workload_path)
1834        .await
1835        .with_context(|| format!("reading workload {}", workload_path.display()))?;
1836    let workload: StaticAssetWorkload = toml::from_str(&raw)
1837        .with_context(|| format!("parsing static-asset workload {}", workload_path.display()))?;
1838
1839    let entry = workload
1840        .assets
1841        .iter()
1842        .find(|a| {
1843            a.derive
1844                .as_ref()
1845                .and_then(|d| d.transform.as_ref())
1846                .map(|t| t.params.get("target").map(String::as_str) == Some(target))
1847                .unwrap_or(false)
1848        })
1849        .ok_or_else(|| {
1850            anyhow::anyhow!(
1851                "no [[asset]] with derive.transform.params.target = {target:?} in {}",
1852                workload_path.display(),
1853            )
1854        })?;
1855
1856    let derive = entry
1857        .derive
1858        .as_ref()
1859        .expect("asset matched by derive.transform above");
1860    let transform = derive
1861        .transform
1862        .as_ref()
1863        .expect("asset matched by derive.transform above");
1864
1865    // Recompute the fetched-input hash exactly as the reconciler will — same
1866    // fetch source, same cache dir.
1867    let cache_root = workspace_root.join(".yah/cache/derive");
1868    let (_fetched_path, fetched_hash) =
1869        materialize_fetch(&derive.fetch, &cache_root.join("fetch")).await?;
1870
1871    // R546-B6: classify the COMMITTED fetch pin against what we just fetched, so
1872    // the caller can tell the operator whether the W212 fast-path is armed. A
1873    // seeded lock with a sentinel fetch pin skips builds only on this machine
1874    // (via the local action cache) — everywhere else `lock_skip_hash` declines.
1875    let fetch_pin = classify_fetch_pin(&derive.fetch.blake3.0, &fetched_hash);
1876
1877    let mut seeded =
1878        seed_transform_derivation(workspace_root, transform, &fetched_hash, artifact_path).await?;
1879    seeded.fetch_pin = Some(fetch_pin);
1880    Ok(seeded)
1881}
1882
1883/// Lower a single recipe step to a [`ForgeSpec`] (W164).
1884///
1885/// - `image` is always `Some(recipe.image)` — recipes always run inside the
1886///   pinned container.
1887/// - `where_` mirrors `recipe.placement` straight through — both the
1888///   recipe-declared location and runtime. Until W235 this hard-coded
1889///   `TaskLocation::Local` because `RecipeLocation` had no other variant;
1890///   R555-T2 opened the recipe surface to `remote` / `remote_any`, so pinning
1891///   here would have quietly demoted every remote recipe back to the dev box.
1892/// - `timeout=0` in the recipe means "no timeout" (omitted from the spec).
1893/// - `label = "transform:<recipe>:<step>"`; initiator carries the reconciler
1894///   identity in the Gnome variant so audit traces attribute the run.
1895///
1896/// Pure function: callers feed it the already-substituted argv (or the raw
1897/// one in tests) and decide what to do with the resulting spec. Exposed at
1898/// `pub(crate)` for golden-test parity with the BuildMode lowering helper
1899/// (R438-T7).
1900pub(crate) fn lower_recipe_step_to_forge_spec(
1901    recipe: &TransformRecipe,
1902    step: &RecipeStep,
1903    substituted_argv: Vec<String>,
1904) -> ForgeSpec {
1905    ForgeSpec {
1906        command: ForgeCommand::Subprocess {
1907            argv: substituted_argv,
1908            image: Some(recipe.image.clone()),
1909        },
1910        where_: TaskPlacement::new(recipe.placement.location.clone(), recipe.placement.runtime),
1911        timeout: if step.timeout == 0 {
1912            None
1913        } else {
1914            Some(Millis::from_secs(step.timeout))
1915        },
1916        label: Some(format!("transform:{}:{}", recipe.name, step.name)),
1917        initiator: Initiator::Gnome {
1918            camp: "static-asset-reconciler".into(),
1919            shift: format!("derive-{}", recipe.name),
1920        },
1921        mesh_access: MeshAccess::default(),
1922        cache_key: None,
1923    }
1924}
1925
1926/// Validate a remotely-placed transform recipe and return the worker-side path
1927/// its output must be written to (R555-F3).
1928///
1929/// Three things a local recipe may do that a remote one may not, each refused
1930/// here rather than at the point it would produce a wrong artifact:
1931///
1932/// 1. **Reference `{{YAH_TRANSFORM_IN_0}}`.** The fetched input lives on this
1933///    box. [`velveteen_exec::WardenClient`] has a retrieval leg
1934///    (`fetch_produced_file`) and no upload leg, so there is no transport that
1935///    puts those bytes on the worker. Substituting the host path would hand
1936///    the recipe a path that doesn't resolve there — and the recipes that
1937///    ignore IN_0 entirely (rusty-v8-musl, whisper-bundle-tar drive their own
1938///    source) are exactly the ones worth dispatching, so this is a real
1939///    boundary rather than a stopgap.
1940/// 2. **Declare more than one step.** Each step is a separate one-shot
1941///    workload with its own container and its own produced dir — locally the
1942///    steps share a filesystem, remotely they share nothing. A two-step remote
1943///    recipe would silently lose whatever step 1 wrote.
1944/// 3. **Declare `[placement] platform`.** That asks a *host* container runtime
1945///    for foreign-arch emulation; a yubaba node has no such knob and picks
1946///    architecture by scheduling instead. `RemoteForgeDriver` refuses it too —
1947///    this refusal exists so the message names the recipe (R555-T7 is the
1948///    ticket that collapses the per-arch recipe fork onto `mesh_tags`).
1949fn remote_transform_out(recipe: &TransformRecipe, derive_key: &str) -> Result<PathBuf> {
1950    if recipe.steps.len() > 1 {
1951        anyhow::bail!(
1952            "recipe {:?} declares {} steps and a remote placement. Each step dispatches as \
1953             its own one-shot workload with its own container and produced dir, so steps \
1954             cannot hand files to each other the way they do locally — collapse them into \
1955             one step (the builder image is the right place for the sequencing), or run \
1956             the recipe locally",
1957            recipe.name,
1958            recipe.steps.len(),
1959        );
1960    }
1961    if let Some(platform) = &recipe.placement.platform {
1962        anyhow::bail!(
1963            "recipe {:?} declares [placement] platform = {:?} together with a remote \
1964             location. `platform` asks a host container runtime for foreign-arch \
1965             emulation; a remote node has no such knob and selects architecture by \
1966             scheduling instead. Drop it and express the arch as mesh tags, e.g. \
1967             location = {{ kind = \"remote_any\", tier = \"infra\", mesh_tags = \
1968             [\"tag:build-worker\", \"arch:x86\", \"os:linux\"] }}",
1969            recipe.name,
1970            platform,
1971        );
1972    }
1973    let in_0 = format!("{{{{{ENV_TRANSFORM_IN_0}}}}}");
1974    for step in &recipe.steps {
1975        if step.argv.iter().any(|a| a.contains(&in_0)) {
1976            anyhow::bail!(
1977                "recipe {:?} step {:?} references {} but is placed remotely. The fetched \
1978                 input lives on this machine and there is no upload leg to the worker \
1979                 (WardenClient can retrieve produced files, not send inputs), so the \
1980                 binding would name a path that does not exist on the node. Either have \
1981                 the recipe fetch its own input, or run it locally",
1982                recipe.name,
1983                step.name,
1984                in_0,
1985            );
1986        }
1987    }
1988    Ok(PathBuf::from(workload_spec::forge_produced::CONTAINER_DIR).join(format!("{derive_key}.out")))
1989}
1990
1991/// The mirror of [`remote_transform_out`]'s refusals, for the local leg: one
1992/// thing a remote recipe may do that a local one may not (R555-F5).
1993///
1994/// Vault credentials are resolved by yubaba on the node that runs the workload,
1995/// against the cluster secret store and that node's KEK. A local run has no such
1996/// node, so `LocalForgeDriver` refuses `ExecContext::secrets` outright. Refused
1997/// here as well so the message names the *recipe* — a `location = "local"`
1998/// recipe with a `[[secrets]]` block is a mistake in the recipe, and a driver
1999/// error names only the dispatch.
2000fn refuse_secrets_on_a_local_recipe(recipe: &TransformRecipe) -> Result<()> {
2001    if !recipe.secrets.is_empty() {
2002        anyhow::bail!(
2003            "recipe {:?} declares {} [[secrets]] entr{} but is placed locally. Vault \
2004             credentials are resolved by yubaba on the node that runs the workload, \
2005             and a local run has no such node — give the recipe a remote placement \
2006             (location = {{ kind = \"remote_any\", … }}), or drop the secret and have \
2007             the step read its credential from the environment it already runs in \
2008             (W235 §(c) / R555-F5)",
2009            recipe.name,
2010            recipe.secrets.len(),
2011            if recipe.secrets.len() == 1 { "y" } else { "ies" },
2012        );
2013    }
2014    Ok(())
2015}
2016
2017/// Read `path` and assert its BLAKE3 hex matches `expected_hex` (case-insensitive).
2018async fn verify_blake3_path(path: &Path, expected_hex: &str) -> Result<()> {
2019    let bytes = tokio::fs::read(path)
2020        .await
2021        .with_context(|| format!("reading {} for BLAKE3 verify", path.display()))?;
2022    let actual = blake3_hex(&bytes);
2023    if !hashes_equal(&actual, expected_hex) {
2024        anyhow::bail!(
2025            "BLAKE3 mismatch for {}: actual {} != expected {}",
2026            path.display(),
2027            actual,
2028            expected_hex,
2029        );
2030    }
2031    Ok(())
2032}
2033
2034// ── Manifest helpers ──────────────────────────────────────────────────────────
2035
2036async fn load_catalog_manifest(
2037    client: &reqwest::Client,
2038    endpoint: &str,
2039    bucket: &str,
2040    region: &str,
2041    access_key: &str,
2042    secret_key: &str,
2043) -> CatalogManifest {
2044    let url = format!("{endpoint}/{bucket}/{CATALOG_MANIFEST_KEY}");
2045    let Ok(headers) = sign_s3_empty_body("GET", &url, region, access_key, secret_key) else {
2046        return HashMap::new();
2047    };
2048    let Ok(resp) = client.get(&url).headers(headers).send().await else {
2049        return HashMap::new();
2050    };
2051    if !resp.status().is_success() {
2052        return HashMap::new();
2053    }
2054    let Ok(bytes) = resp.bytes().await else {
2055        return HashMap::new();
2056    };
2057    serde_json::from_slice(&bytes).unwrap_or_default()
2058}
2059
2060async fn save_catalog_manifest(
2061    client: &reqwest::Client,
2062    endpoint: &str,
2063    bucket: &str,
2064    manifest: &CatalogManifest,
2065    region: &str,
2066    access_key: &str,
2067    secret_key: &str,
2068) -> Result<()> {
2069    let body = serde_json::to_vec(manifest).context("serializing catalog manifest")?;
2070    let body_sha256 = sha256_hex(&body);
2071    let url = format!("{endpoint}/{bucket}/{CATALOG_MANIFEST_KEY}");
2072    let headers = sign_s3_put_object(
2073        &url,
2074        &body_sha256,
2075        "application/json",
2076        body.len(),
2077        region,
2078        access_key,
2079        secret_key,
2080        // The manifest sidecar is rewritten every run by definition — stamping it
2081        // would be noise, and nothing skips a PUT on it (R546-B10).
2082        None,
2083    )
2084    .context("signing manifest PUT")?;
2085    let resp = client
2086        .put(&url)
2087        .headers(headers)
2088        .body(body)
2089        .send()
2090        .await
2091        .context("PUT catalog manifest")?;
2092    if !resp.status().is_success() {
2093        let status = resp.status();
2094        let body_text = resp.text().await.unwrap_or_default();
2095        anyhow::bail!("PUT {url} → {status}: {}", body_text.trim());
2096    }
2097    Ok(())
2098}
2099
2100// ── S3 helpers ────────────────────────────────────────────────────────────────
2101
2102/// What a HEAD probe found at an object key (R546-B10).
2103///
2104/// The distinction that matters is *not* present-vs-absent but whether the
2105/// present bytes are the ones we are about to publish. `Present { blake3: None }`
2106/// is an object written before we started stamping `x-amz-meta-blake3` — its
2107/// content is unknowable from a HEAD, so callers must treat it as "might be
2108/// stale" rather than "matches".
2109#[derive(Debug, Clone, PartialEq, Eq)]
2110enum RemoteObject {
2111    Absent,
2112    Present { blake3: Option<String> },
2113}
2114
2115impl RemoteObject {
2116    /// True only when the remote object is known to hold exactly `expected`.
2117    /// Unknown provenance (no metadata) is deliberately NOT a match.
2118    fn matches(&self, expected: &str) -> bool {
2119        matches!(self, RemoteObject::Present { blake3: Some(b3) } if hashes_equal(b3, expected))
2120    }
2121
2122    fn exists(&self) -> bool {
2123        matches!(self, RemoteObject::Present { .. })
2124    }
2125}
2126
2127/// HEAD an object key and report existence *plus* the recorded BLAKE3.
2128///
2129/// R546-B10: the previous `object_exists` returned a bare bool, so the publish
2130/// loop could not tell "same bytes already there" from "different bytes already
2131/// there" and skipped the PUT for both — leaving stale content in the bucket
2132/// while the catalog manifest advertised the new hash. We read our own
2133/// `x-amz-meta-blake3` stamp rather than the ETag, because ETag is only an MD5
2134/// of the content for single-part uploads; for multipart it is a digest-of-
2135/// digests and comparing it to a content hash is simply wrong.
2136async fn head_object(
2137    client: &reqwest::Client,
2138    url: &str,
2139    region: &str,
2140    access_key: &str,
2141    secret_key: &str,
2142) -> Result<RemoteObject> {
2143    let headers = sign_s3_empty_body("HEAD", url, region, access_key, secret_key)
2144        .with_context(|| format!("signing HEAD {url}"))?;
2145    let resp = client
2146        .head(url)
2147        .headers(headers)
2148        .send()
2149        .await
2150        .with_context(|| format!("HEAD {url}"))?;
2151    if !resp.status().is_success() {
2152        return Ok(RemoteObject::Absent);
2153    }
2154    let blake3 = resp
2155        .headers()
2156        .get("x-amz-meta-blake3")
2157        .and_then(|v| v.to_str().ok())
2158        .map(str::to_string);
2159    Ok(RemoteObject::Present { blake3 })
2160}
2161
2162// ── Crypto helpers ────────────────────────────────────────────────────────────
2163
2164fn blake3_hex(body: &[u8]) -> String {
2165    hex::encode(blake3::hash(body).as_bytes())
2166}
2167
2168fn sha256_hex(body: &[u8]) -> String {
2169    hex::encode(Sha256::digest(body))
2170}
2171
2172/// Case-insensitive hex comparison (BLAKE3 crate outputs lowercase; stored
2173/// hashes may have been authored in uppercase).
2174fn hashes_equal(a: &str, b: &str) -> bool {
2175    a.eq_ignore_ascii_case(b)
2176}
2177
2178// ── Content-type ──────────────────────────────────────────────────────────────
2179
2180fn content_type_for(path: &Path) -> &'static str {
2181    match path.extension().and_then(|e| e.to_str()).unwrap_or("") {
2182        "bin" => "application/octet-stream",
2183        "json" => "application/json",
2184        "txt" => "text/plain; charset=utf-8",
2185        "wasm" => "application/wasm",
2186        _ => "application/octet-stream",
2187    }
2188}
2189
2190// ── Bucket auto-create ────────────────────────────────────────────────────────
2191
2192/// List R2 buckets under `account_id`; create `bucket_name` only when absent.
2193///
2194/// Takes the management `api_token` already resolved from the service's
2195/// provider (see [`super::cf_creds::CfProvider`]). Mirrors the idempotent pattern in
2196/// `cloudflare_worker.rs::ensure_r2_bucket` — list-first keeps the reconcile
2197/// loop safe to re-run without depending on CF's 4xx error shape on duplicate
2198/// create. R422-T12: lets the first `yah cloud apply --env cloud --service <s>`
2199/// against a static-asset mirror provision its bucket without an out-of-band
2200/// dashboard step.
2201async fn ensure_r2_bucket(api_token: &str, account_id: &str, bucket_name: &str) -> Result<()> {
2202    let cf = CloudflareClient::new(api_token.to_string());
2203    let existing = cf
2204        .list_r2_buckets(account_id)
2205        .await
2206        .context("listing R2 buckets")?;
2207    if existing.iter().any(|b| b.name == bucket_name) {
2208        debug!(bucket_name, "R2 bucket already exists — skipping create");
2209        return Ok(());
2210    }
2211    cf.create_r2_bucket(account_id, bucket_name)
2212        .await
2213        .with_context(|| format!("creating R2 bucket {bucket_name}"))?;
2214    info!(bucket_name, "R2 bucket created");
2215    Ok(())
2216}
2217
2218// ── Workload loading ──────────────────────────────────────────────────────────
2219
2220fn load_workload(workload_dir: &Path) -> Result<StaticAssetWorkload> {
2221    let path = workload_dir.join("workload.toml");
2222    let src =
2223        std::fs::read_to_string(&path).with_context(|| format!("reading {}", path.display()))?;
2224    // R546-B7 (RESOLVED — this comment used to say the envelope was unusable).
2225    // `workload_spec::Workload` is no longer externally tagged for text
2226    // formats: its hand-written Deserialize branches on `is_human_readable`, so
2227    // TOML/JSON get the flat `kind`-tagged shape real files use while postcard
2228    // keeps the variant-index encoding the kamaji UDS needs.
2229    //
2230    // This path still probes `kind` and deserializes `StaticAssetWorkload`
2231    // directly rather than matching on the envelope, because it needs to reject
2232    // a non-static-asset workload with a precise message naming the kind it
2233    // found — an envelope match would only say "expected variant". The bypass
2234    // is now a choice, not a workaround.
2235    #[derive(serde::Deserialize)]
2236    struct KindProbe {
2237        kind: String,
2238    }
2239    let probe: KindProbe =
2240        toml::from_str(&src).with_context(|| format!("parsing {}", path.display()))?;
2241    if probe.kind != "static-asset" {
2242        anyhow::bail!(
2243            "{}: expected kind=\"static-asset\" but found kind={:?}",
2244            path.display(),
2245            probe.kind
2246        );
2247    }
2248    let workload: StaticAssetWorkload =
2249        toml::from_str(&src).with_context(|| format!("parsing {}", path.display()))?;
2250
2251    // R546-B7: enforce the closed-catalog invariant the type's docs promise
2252    // ("every value in [aliases] must be a filename that exists in [[asset]],
2253    // enforced by validate::shape_static_asset"). It was NOT actually being run
2254    // here — `up_bails_when_catalog_alias_orphaned` only passed because the old
2255    // fixture failed to PARSE, so the expected error came from the wrong place.
2256    // With parsing fixed, an orphaned alias would have silently reconciled.
2257    workload_spec::validate::shape_static_asset(&workload)
2258        .map_err(|e| anyhow::anyhow!("{}: {e}", path.display()))?;
2259
2260    Ok(workload)
2261}
2262
2263// ── Tests ─────────────────────────────────────────────────────────────────────
2264
2265#[cfg(test)]
2266mod tests {
2267    use super::*;
2268    use crate::asset_journal::AssetStatusJournal;
2269    use crate::{MirrorConfig, MirrorShape, ServiceComponent, ServiceConfig};
2270    use std::collections::BTreeMap;
2271    use std::path::PathBuf;
2272    use tempfile::tempdir;
2273    use workload_spec::{AssetEntry, BlakeHash};
2274
2275    const HASH_64: &str = "abcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890";
2276
2277    /// Minimal test fixture that mirrors the mesofact_static test pattern.
2278    struct Fixture {
2279        _workspace: tempfile::TempDir,
2280        workspace_root: PathBuf,
2281        service: ServiceConfig,
2282        component: ServiceComponent,
2283        mirror: MirrorConfig,
2284        env: String,
2285    }
2286
2287    impl Fixture {
2288        fn new(slot: MirrorProviderSlot) -> Self {
2289            let workspace = tempdir().unwrap();
2290            let workspace_root = workspace.path().to_path_buf();
2291            let workload_dir = workspace_root.join("app/assets/whisper");
2292            std::fs::create_dir_all(&workload_dir).unwrap();
2293
2294            let mut providers = BTreeMap::new();
2295            providers.insert("object_store".to_string(), slot);
2296            let mirror = MirrorConfig {
2297                schema_version: 1,
2298                shape: MirrorShape::Local,
2299                providers,
2300                ingress: Default::default(),
2301                ingress_machines: Vec::new(),
2302                drivers: Default::default(),
2303                asset_aliases: BTreeMap::new(),
2304                build: Default::default(),
2305            };
2306            let service = ServiceConfig {
2307                schema_version: 1,
2308                name: "yah-desktop".to_string(),
2309                address: crate::config::ServiceAddress::front_door("releases.yah.dev".to_string()),
2310                description: None,
2311                components: vec![],
2312                db: crate::DbCatalog::default(),
2313            };
2314            let component = ServiceComponent {
2315                mount: None,
2316                id: "whisper-models".to_string(),
2317                kind: "static-asset".to_string(),
2318                path: "app/assets/whisper".to_string(),
2319                role: "assets".to_string(),
2320                publishes: None,
2321                wave: 0,
2322                git: None,
2323                deploy: Default::default(),
2324            };
2325            Self {
2326                _workspace: workspace,
2327                workspace_root,
2328                service,
2329                component,
2330                mirror,
2331                env: "pond".to_string(),
2332            }
2333        }
2334
2335        fn ctx(&self) -> ReconcileCtx<'_> {
2336            ReconcileCtx {
2337                workspace_root: &self.workspace_root,
2338                service: &self.service,
2339                component: &self.component,
2340                mirror: &self.mirror,
2341                env: &self.env,
2342                scope: crate::reconciler::ProviderScope::singleton(),
2343            }
2344        }
2345
2346        fn workload_dir(&self) -> PathBuf {
2347            self.workspace_root.join("app/assets/whisper")
2348        }
2349
2350        fn write_workload(&self, extra: &str) {
2351            // R546-B7: emit the FLAT shape every real on-disk workload.toml uses
2352            // (`kind = "static-asset"` beside `schema_version`), not the
2353            // externally-tagged `[static-asset]` wrapper. The old fixture encoded
2354            // the `workload_spec::Workload` enum's (untagged) representation,
2355            // which NO real file has ever used — so these tests were asserting a
2356            // shape that could not occur in production and masked the fact that
2357            // `yah cloud apply` could not parse any actual catalog.
2358            let toml = format!(
2359                r#"schema_version = "V1"
2360kind = "static-asset"
2361{extra}
2362"#
2363            );
2364            std::fs::write(self.workload_dir().join("workload.toml"), toml).unwrap();
2365        }
2366
2367        fn write_source_file(&self, name: &str, content: &[u8]) -> PathBuf {
2368            let path = self.workload_dir().join(name);
2369            if let Some(parent) = path.parent() {
2370                std::fs::create_dir_all(parent).unwrap();
2371            }
2372            std::fs::write(&path, content).unwrap();
2373            path
2374        }
2375    }
2376
2377    fn minio_slot() -> MirrorProviderSlot {
2378        let mut fields = BTreeMap::new();
2379        fields.insert("api_port".to_string(), toml::Value::Integer(9000));
2380        fields.insert(
2381            "bucket".to_string(),
2382            toml::Value::String("yah-dev".to_string()),
2383        );
2384        MirrorProviderSlot::Inline {
2385            kind: Provider::MinioContainer,
2386            fields,
2387        }
2388    }
2389
2390    /// Default executor for tests that don't exercise transforms — legacy
2391    /// `source = "..."` assets never touch the executor, so any impl works.
2392    fn test_executor() -> Arc<dyn ForgeExecutor> {
2393        Arc::new(LocalForgeDriver::default())
2394    }
2395
2396    // ── Mock executor for W164 materialize-transform tests ────────────────────
2397
2398    use tokio::sync::mpsc::UnboundedSender;
2399    use tokio::sync::Mutex;
2400    use velveteen::ForgeStatus;
2401    use velveteen_exec::executor::{ExecEvent, ExecOutcome, ForgeExecutorError};
2402
2403    /// Executor that writes a caller-supplied byte string to whatever path the
2404    /// recipe sets `YAH_TRANSFORM_OUT` to (via the substituted argv). Returns
2405    /// success unless `fail_with` is set. Tracks invocation count for HIT/MISS
2406    /// assertions.
2407    struct MockExecutor {
2408        out_bytes: Vec<u8>,
2409        invocations: Arc<Mutex<u32>>,
2410        fail_with: Option<String>,
2411    }
2412
2413    impl MockExecutor {
2414        fn new(out_bytes: Vec<u8>) -> (Arc<Self>, Arc<Mutex<u32>>) {
2415            let counter = Arc::new(Mutex::new(0));
2416            let me = Arc::new(Self {
2417                out_bytes,
2418                invocations: counter.clone(),
2419                fail_with: None,
2420            });
2421            (me, counter)
2422        }
2423
2424        fn failing(reason: String) -> Arc<Self> {
2425            Arc::new(Self {
2426                out_bytes: Vec::new(),
2427                invocations: Arc::new(Mutex::new(0)),
2428                fail_with: Some(reason),
2429            })
2430        }
2431    }
2432
2433    #[async_trait]
2434    impl ForgeExecutor for MockExecutor {
2435        async fn execute(
2436            &self,
2437            spec: ForgeSpec,
2438            _ctx: ExecContext,
2439            _sink: Option<UnboundedSender<ExecEvent>>,
2440        ) -> Result<ExecOutcome, ForgeExecutorError> {
2441            *self.invocations.lock().await += 1;
2442            if let Some(reason) = &self.fail_with {
2443                return Ok(ExecOutcome {
2444                    status: ForgeStatus::Done {
2445                        exit_code: 1,
2446                        ended_at: 0,
2447                    },
2448                    stderr_tail: reason.clone(),
2449                });
2450            }
2451            // Find the YAH_TRANSFORM_OUT path in the substituted argv. The
2452            // recipe convention is that {{YAH_TRANSFORM_OUT}} resolves to the
2453            // tmp path we want to write to; it appears as a positional arg.
2454            let ForgeCommand::Subprocess { argv, .. } = &spec.command else {
2455                return Err(ForgeExecutorError::Unsupported(
2456                    "mock only handles Subprocess",
2457                ));
2458            };
2459            // The test recipes follow the W164 convention: argv contains the
2460            // tmp output path verbatim. materialize_transform stages at
2461            // `<cache_dir>/<derive_key>.out.<pid>.<seq>.tmp` (R925; then renames
2462            // to `<hash>.bin` on hash match), so the output element is the one
2463            // ending in `.tmp`. Any future change to that staging name must keep
2464            // the `.tmp` suffix or this find() silently stops matching.
2465            let out_path = argv.iter().find(|a| a.ends_with(".tmp")).ok_or(
2466                ForgeExecutorError::Unsupported("mock recipe must pass a .tmp output path"),
2467            )?;
2468            std::fs::write(out_path, &self.out_bytes).map_err(ForgeExecutorError::Io)?;
2469            Ok(ExecOutcome {
2470                status: ForgeStatus::Done {
2471                    exit_code: 0,
2472                    ended_at: 0,
2473                },
2474                stderr_tail: String::new(),
2475            })
2476        }
2477    }
2478
2479    /// R546-B8: records the substituted argv and, optionally, declines to write
2480    /// the output while still reporting exit 0 — the exact shape of a recipe
2481    /// that chdirs and then writes its artifact to a relative path inside the
2482    /// container.
2483    struct ArgvCapture {
2484        argv: Arc<Mutex<Vec<String>>>,
2485        write_output: bool,
2486    }
2487
2488    impl ArgvCapture {
2489        fn new(write_output: bool) -> (Arc<Self>, Arc<Mutex<Vec<String>>>) {
2490            let argv = Arc::new(Mutex::new(Vec::new()));
2491            let me = Arc::new(Self {
2492                argv: argv.clone(),
2493                write_output,
2494            });
2495            (me, argv)
2496        }
2497    }
2498
2499    #[async_trait]
2500    impl ForgeExecutor for ArgvCapture {
2501        async fn execute(
2502            &self,
2503            spec: ForgeSpec,
2504            _ctx: ExecContext,
2505            _sink: Option<UnboundedSender<ExecEvent>>,
2506        ) -> Result<ExecOutcome, ForgeExecutorError> {
2507            let ForgeCommand::Subprocess { argv, .. } = &spec.command else {
2508                return Err(ForgeExecutorError::Unsupported(
2509                    "capture only handles Subprocess",
2510                ));
2511            };
2512            *self.argv.lock().await = argv.clone();
2513            if self.write_output {
2514                let out = argv.iter().find(|a| a.ends_with(".tmp")).ok_or(
2515                    ForgeExecutorError::Unsupported("recipe must pass a .tmp output path"),
2516                )?;
2517                std::fs::write(out, b"produced").map_err(ForgeExecutorError::Io)?;
2518            }
2519            Ok(ExecOutcome {
2520                status: ForgeStatus::Done {
2521                    exit_code: 0,
2522                    ended_at: 0,
2523                },
2524                stderr_tail: String::new(),
2525            })
2526        }
2527    }
2528
2529    /// Write a `whisper-quantize.toml`-style recipe under
2530    /// `<workspace>/.yah/qed/transforms/<name>.toml`. Returns nothing — the
2531    /// recipe loader resolves the path itself.
2532    fn write_recipe(workspace_root: &Path, name: &str) {
2533        let transforms_dir = workspace_root.join(".yah/qed/transforms");
2534        std::fs::create_dir_all(&transforms_dir).unwrap();
2535        let toml = format!(
2536            r#"
2537name  = "{name}"
2538label = "test recipe"
2539image = "ghcr.io/test/tool:v1@sha256:{HASH_64}"
2540
2541[placement]
2542location = "local"
2543runtime  = "container"
2544
2545[[steps]]
2546name = "transform"
2547argv = ["./tool", "{{{{YAH_TRANSFORM_IN_0}}}}", "{{{{YAH_TRANSFORM_OUT}}}}"]
2548"#
2549        );
2550        std::fs::write(transforms_dir.join(format!("{name}.toml")), toml).unwrap();
2551    }
2552
2553    /// Like [`write_recipe`] but with a caller-chosen container digest, so a
2554    /// test can simulate a digest bump (which must invalidate the W212 action
2555    /// cache and force a re-run).
2556    fn write_recipe_with_digest(workspace_root: &Path, name: &str, digest: &str) {
2557        let transforms_dir = workspace_root.join(".yah/qed/transforms");
2558        std::fs::create_dir_all(&transforms_dir).unwrap();
2559        let toml = format!(
2560            r#"
2561name  = "{name}"
2562label = "test recipe"
2563image = "ghcr.io/test/tool:v1@sha256:{digest}"
2564
2565[placement]
2566location = "local"
2567runtime  = "container"
2568
2569[[steps]]
2570name = "transform"
2571argv = ["./tool", "{{{{YAH_TRANSFORM_IN_0}}}}", "{{{{YAH_TRANSFORM_OUT}}}}"]
2572"#
2573        );
2574        std::fs::write(transforms_dir.join(format!("{name}.toml")), toml).unwrap();
2575    }
2576
2577    /// Write a remotely-placed recipe. `placement_extra` appends lines to the
2578    /// `[placement]` table (e.g. a `platform` the remote leg must refuse) and
2579    /// `steps` is the whole `[[steps]]` block, so a test can pass two of them.
2580    fn write_remote_recipe(
2581        workspace_root: &Path,
2582        name: &str,
2583        placement_extra: &str,
2584        steps: &str,
2585    ) {
2586        let transforms_dir = workspace_root.join(".yah/qed/transforms");
2587        std::fs::create_dir_all(&transforms_dir).unwrap();
2588        let toml = format!(
2589            r#"
2590name  = "{name}"
2591label = "test remote recipe"
2592image = "ghcr.io/test/tool:v1@sha256:{HASH_64}"
2593
2594[placement]
2595location = {{ kind = "remote_any", tier = "infra", mesh_tags = ["arch:x86"] }}
2596runtime  = "container"
2597{placement_extra}
2598
2599{steps}
2600"#
2601        );
2602        std::fs::write(transforms_dir.join(format!("{name}.toml")), toml).unwrap();
2603    }
2604
2605    /// Stands in for `RemoteForgeDriver`: records the context it was handed and
2606    /// performs the retrieval half by writing the "built" bytes to
2607    /// `ctx.produced.dest`, exactly as the real driver does after
2608    /// `fetch_produced_file`.
2609    struct RemoteMock {
2610        out_bytes: Vec<u8>,
2611        seen: Arc<Mutex<Option<(Vec<String>, ExecContext)>>>,
2612    }
2613
2614    impl RemoteMock {
2615        fn new(out_bytes: Vec<u8>) -> (Arc<Self>, Arc<Mutex<Option<(Vec<String>, ExecContext)>>>) {
2616            let seen = Arc::new(Mutex::new(None));
2617            let me = Arc::new(Self {
2618                out_bytes,
2619                seen: seen.clone(),
2620            });
2621            (me, seen)
2622        }
2623    }
2624
2625    #[async_trait]
2626    impl ForgeExecutor for RemoteMock {
2627        async fn execute(
2628            &self,
2629            spec: ForgeSpec,
2630            ctx: ExecContext,
2631            _sink: Option<UnboundedSender<ExecEvent>>,
2632        ) -> Result<ExecOutcome, ForgeExecutorError> {
2633            let ForgeCommand::Subprocess { argv, .. } = &spec.command else {
2634                return Err(ForgeExecutorError::Unsupported("remote mock: Subprocess only"));
2635            };
2636            assert!(
2637                !matches!(spec.where_.location, TaskLocation::Local),
2638                "the remote mock must only ever see a remotely-placed spec",
2639            );
2640            if let Some(produced) = &ctx.produced {
2641                std::fs::write(&produced.dest, &self.out_bytes).map_err(ForgeExecutorError::Io)?;
2642            }
2643            *self.seen.lock().await = Some((argv.clone(), ctx));
2644            Ok(ExecOutcome {
2645                status: ForgeStatus::Done {
2646                    exit_code: 0,
2647                    ended_at: 0,
2648                },
2649                stderr_tail: String::new(),
2650            })
2651        }
2652    }
2653
2654    /// R555-F3, the whole point: a remotely-placed recipe binds its output to
2655    /// the worker's durable produced dir, and the bytes come back to the local
2656    /// derivation cache so the publish leg above is untouched.
2657    #[tokio::test]
2658    async fn a_remote_recipe_binds_out_on_the_worker_and_lands_the_bytes_locally() {
2659        let fx = Fixture::new(minio_slot());
2660        write_remote_recipe(
2661            &fx.workspace_root,
2662            "remote-recipe",
2663            "",
2664            "[[steps]]\nname = \"build\"\nargv = [\"build.sh\", \"{{YAH_TRANSFORM_OUT}}\"]",
2665        );
2666        let out_bytes = b"remotely built".to_vec();
2667        let out_hash = blake3_hex(&out_bytes);
2668        let (mock, seen) = RemoteMock::new(out_bytes.clone());
2669        let executor: Arc<dyn ForgeExecutor> = mock;
2670
2671        let fetch_path = fx.workspace_root.join(".yah/cache/derive/fetch/in.bin");
2672        std::fs::create_dir_all(fetch_path.parent().unwrap()).unwrap();
2673        std::fs::write(&fetch_path, b"anchor").unwrap();
2674
2675        let landed = materialize_transform(
2676            &TransformSpec {
2677                recipe: "remote-recipe".to_string(),
2678                params: BTreeMap::new(),
2679            },
2680            &fetch_path,
2681            &blake3_hex(b"anchor"),
2682            &out_hash,
2683            &fx.workspace_root.join(".yah/cache/derive/transform"),
2684            &fx.workspace_root,
2685            executor.as_ref(),
2686        )
2687        .await
2688        .unwrap();
2689
2690        assert_eq!(std::fs::read(&landed).unwrap(), out_bytes);
2691
2692        let seen = seen.lock().await;
2693        let (argv, ctx) = seen.as_ref().expect("the remote step ran");
2694        let bound_out = argv.last().expect("argv carries the OUT binding");
2695        assert!(
2696            bound_out.starts_with(workload_spec::forge_produced::CONTAINER_DIR),
2697            "OUT must be bound under the worker's durable produced dir, got {bound_out}",
2698        );
2699        let produced = ctx.produced.as_ref().expect("retrieval must be requested");
2700        assert_eq!(produced.remote_path, PathBuf::from(bound_out));
2701        // A host path handed to a remote workdir names a directory that does
2702        // not exist on the node.
2703        assert!(ctx.cwd.is_none(), "no host cwd may travel to a worker");
2704    }
2705
2706    /// The fetched input lives on this box and there is no upload leg, so a
2707    /// remote recipe that reads it would be handed an unresolvable path.
2708    #[tokio::test]
2709    async fn a_remote_recipe_reading_the_fetched_input_is_refused() {
2710        let fx = Fixture::new(minio_slot());
2711        write_remote_recipe(
2712            &fx.workspace_root,
2713            "needs-input",
2714            "",
2715            "[[steps]]\nname = \"build\"\nargv = [\"tool\", \"{{YAH_TRANSFORM_IN_0}}\", \"{{YAH_TRANSFORM_OUT}}\"]",
2716        );
2717        let msg = remote_refusal(&fx, "needs-input").await;
2718        assert!(msg.contains("YAH_TRANSFORM_IN_0"), "got: {msg}");
2719        assert!(msg.contains("no upload leg"), "got: {msg}");
2720    }
2721
2722    /// Steps share a filesystem locally and share nothing remotely.
2723    #[tokio::test]
2724    async fn a_multi_step_remote_recipe_is_refused() {
2725        let fx = Fixture::new(minio_slot());
2726        write_remote_recipe(
2727            &fx.workspace_root,
2728            "two-steps",
2729            "",
2730            "[[steps]]\nname = \"one\"\nargv = [\"a\"]\n\n[[steps]]\nname = \"two\"\nargv = [\"b\", \"{{YAH_TRANSFORM_OUT}}\"]",
2731        );
2732        let msg = remote_refusal(&fx, "two-steps").await;
2733        assert!(msg.contains("one-shot workload"), "got: {msg}");
2734    }
2735
2736    /// R555-F5: the recipe's declared vault credentials have to reach the
2737    /// driver as mounts, or the build runs without the credential it asked for
2738    /// and fails as a missing-auth error deep inside itself.
2739    #[tokio::test]
2740    async fn a_remote_recipe_hands_its_declared_secrets_to_the_driver() {
2741        let fx = Fixture::new(minio_slot());
2742        write_remote_recipe(
2743            &fx.workspace_root,
2744            "needs-r2",
2745            "",
2746            "[[steps]]\nname = \"build\"\nargv = [\"build.sh\", \"{{YAH_TRANSFORM_OUT}}\"]\n\n\
2747             [[secrets]]\ncluster = \"r2/write\"\npath = \"/run/yah/r2.json\"",
2748        );
2749        let out_bytes = b"built with a credential".to_vec();
2750        let (mock, seen) = RemoteMock::new(out_bytes.clone());
2751        let executor: Arc<dyn ForgeExecutor> = mock;
2752
2753        let fetch_path = fx.workspace_root.join(".yah/cache/derive/fetch/in.bin");
2754        std::fs::create_dir_all(fetch_path.parent().unwrap()).unwrap();
2755        std::fs::write(&fetch_path, b"anchor").unwrap();
2756
2757        materialize_transform(
2758            &TransformSpec {
2759                recipe: "needs-r2".to_string(),
2760                params: BTreeMap::new(),
2761            },
2762            &fetch_path,
2763            &blake3_hex(b"anchor"),
2764            &blake3_hex(&out_bytes),
2765            &fx.workspace_root.join(".yah/cache/derive/transform"),
2766            &fx.workspace_root,
2767            executor.as_ref(),
2768        )
2769        .await
2770        .unwrap();
2771
2772        let seen = seen.lock().await;
2773        let (_, ctx) = seen.as_ref().expect("the remote step ran");
2774        assert_eq!(
2775            ctx.secrets,
2776            vec![workload_spec::SecretMount {
2777                source: workload_spec::SecretRef::Cluster {
2778                    name: "r2/write".into()
2779                },
2780                target: workload_spec::SecretTarget::File {
2781                    path: PathBuf::from("/run/yah/r2.json"),
2782                    // Owner-only is the default a recipe gets without asking.
2783                    mode: 0o400,
2784                },
2785            }]
2786        );
2787    }
2788
2789    /// The local mirror: yubaba resolves credentials on the node that runs the
2790    /// workload, and a local run has no node. Refused by name so the operator
2791    /// is pointed at the recipe rather than at a driver error.
2792    #[tokio::test]
2793    async fn a_local_recipe_declaring_secrets_is_refused_before_dispatch() {
2794        let fx = Fixture::new(minio_slot());
2795        let transforms_dir = fx.workspace_root.join(".yah/qed/transforms");
2796        std::fs::create_dir_all(&transforms_dir).unwrap();
2797        std::fs::write(
2798            transforms_dir.join("local-with-secret.toml"),
2799            format!(
2800                r#"
2801name  = "local-with-secret"
2802label = "test local recipe"
2803image = "ghcr.io/test/tool:v1@sha256:{HASH_64}"
2804
2805[placement]
2806location = "local"
2807runtime  = "container"
2808
2809[[steps]]
2810name = "build"
2811argv = ["build.sh", "{{{{YAH_TRANSFORM_OUT}}}}"]
2812
2813[[secrets]]
2814cluster = "r2/write"
2815path    = "/run/yah/r2.json"
2816"#
2817            ),
2818        )
2819        .unwrap();
2820        let msg = remote_refusal(&fx, "local-with-secret").await;
2821        assert!(msg.contains("placed locally"), "got: {msg}");
2822        assert!(msg.contains("local-with-secret"), "got: {msg}");
2823    }
2824
2825    /// `platform` is host-emulation; remote picks arch by scheduling. Silently
2826    /// dropping it is how R546 got a wrong-arch artifact in the first place.
2827    #[tokio::test]
2828    async fn a_remote_recipe_declaring_a_platform_is_refused_and_named_mesh_tags() {
2829        let fx = Fixture::new(minio_slot());
2830        write_remote_recipe(
2831            &fx.workspace_root,
2832            "emulated",
2833            "platform = \"linux/amd64\"",
2834            "[[steps]]\nname = \"build\"\nargv = [\"build.sh\", \"{{YAH_TRANSFORM_OUT}}\"]",
2835        );
2836        let msg = remote_refusal(&fx, "emulated").await;
2837        assert!(msg.contains("mesh_tags"), "got: {msg}");
2838    }
2839
2840    /// Run `recipe` through `materialize_transform` with an executor that
2841    /// panics if reached, and return the refusal message. Every caller here
2842    /// asserts the recipe is rejected *before* anything is dispatched.
2843    async fn remote_refusal(fx: &Fixture, recipe: &str) -> String {
2844        struct NeverRuns;
2845        #[async_trait]
2846        impl ForgeExecutor for NeverRuns {
2847            async fn execute(
2848                &self,
2849                _spec: ForgeSpec,
2850                _ctx: ExecContext,
2851                _sink: Option<UnboundedSender<ExecEvent>>,
2852            ) -> Result<ExecOutcome, ForgeExecutorError> {
2853                panic!("the recipe must be refused before dispatch");
2854            }
2855        }
2856        let fetch_path = fx.workspace_root.join(".yah/cache/derive/fetch/in.bin");
2857        std::fs::create_dir_all(fetch_path.parent().unwrap()).unwrap();
2858        std::fs::write(&fetch_path, b"anchor").unwrap();
2859        let executor: Arc<dyn ForgeExecutor> = Arc::new(NeverRuns);
2860        let err = materialize_transform(
2861            &TransformSpec {
2862                recipe: recipe.to_string(),
2863                params: BTreeMap::new(),
2864            },
2865            &fetch_path,
2866            &blake3_hex(b"anchor"),
2867            &blake3_hex(b"whatever"),
2868            &fx.workspace_root.join(".yah/cache/derive/transform"),
2869            &fx.workspace_root,
2870            executor.as_ref(),
2871        )
2872        .await
2873        .expect_err("remote recipe must be refused");
2874        format!("{err:#}")
2875    }
2876
2877    // ── Slot validation ───────────────────────────────────────────────────────
2878
2879    #[tokio::test]
2880    async fn up_bails_when_object_store_slot_missing() {
2881        let mut fx = Fixture::new(minio_slot());
2882        fx.mirror.providers.clear();
2883        fx.write_workload("");
2884        let err = StaticAssetReconciler::new().up(fx.ctx()).await.unwrap_err();
2885        let msg = format!("{err:#}");
2886        assert!(msg.contains("providers.object_store"), "got: {msg}");
2887    }
2888
2889    #[tokio::test]
2890    async fn up_bails_on_unsupported_inline_slot_kind() {
2891        let slot = MirrorProviderSlot::Inline {
2892            kind: Provider::MiniflareNative,
2893            fields: BTreeMap::new(),
2894        };
2895        let fx = Fixture::new(slot);
2896        fx.write_workload("");
2897        let err = StaticAssetReconciler::new().up(fx.ctx()).await.unwrap_err();
2898        let msg = format!("{err:#}");
2899        assert!(msg.contains("minio-container"), "got: {msg}");
2900    }
2901
2902    #[tokio::test]
2903    async fn up_bails_on_unsupported_reference_provider() {
2904        let slot = MirrorProviderSlot::Reference {
2905            provider_id: "hetzner".to_string(),
2906            fields: BTreeMap::new(),
2907        };
2908        let fx = Fixture::new(slot);
2909        fx.write_workload("");
2910        // Declare a real non-cloudflare provider so dispatch resolves it and
2911        // reaches the kind check (the new dispatch keys on kind, not name).
2912        let providers_dir = fx.workspace_root.join(".yah/infra/providers");
2913        std::fs::create_dir_all(&providers_dir).unwrap();
2914        std::fs::write(
2915            providers_dir.join("hetzner.toml"),
2916            "schema_version = 1\nid = \"hetzner\"\nkind = \"hetzner\"\naccount_id = \"x\"\n",
2917        )
2918        .unwrap();
2919        let err = StaticAssetReconciler::new().up(fx.ctx()).await.unwrap_err();
2920        let msg = format!("{err:#}");
2921        assert!(msg.contains("cloudflare-kind"), "got: {msg}");
2922    }
2923
2924    #[tokio::test]
2925    async fn up_bails_when_workload_toml_missing() {
2926        let fx = Fixture::new(minio_slot());
2927        // No workload.toml written.
2928        let err = StaticAssetReconciler::new().up(fx.ctx()).await.unwrap_err();
2929        let msg = format!("{err:#}");
2930        assert!(msg.contains("workload.toml"), "got: {msg}");
2931    }
2932
2933    #[tokio::test]
2934    async fn up_bails_when_catalog_alias_orphaned() {
2935        // Alias points at a filename not in the catalog — shape_static_asset rejects it.
2936        let fx = Fixture::new(minio_slot());
2937        fx.write_workload(
2938            r#"[aliases]
2939"default" = "nonexistent/file.bin"
2940"#,
2941        );
2942        let err = StaticAssetReconciler::new().up(fx.ctx()).await.unwrap_err();
2943        let msg = format!("{err:#}");
2944        assert!(
2945            msg.contains("closed-catalog") || msg.contains("aliases"),
2946            "got: {msg}"
2947        );
2948    }
2949
2950    // ── BLAKE3 verification ───────────────────────────────────────────────────
2951
2952    #[test]
2953    fn blake3_mismatch_detected() {
2954        let body = b"hello world";
2955        let actual = blake3_hex(body);
2956        // Fabricate a hash that's wrong.
2957        let wrong = HASH_64;
2958        assert!(!hashes_equal(&actual, wrong));
2959    }
2960
2961    #[test]
2962    fn blake3_match_is_case_insensitive() {
2963        let body = b"hello world";
2964        let lower = blake3_hex(body);
2965        let upper = lower.to_uppercase();
2966        assert!(hashes_equal(&lower, &upper));
2967    }
2968
2969    #[test]
2970    fn sha256_hex_is_64_chars() {
2971        let h = sha256_hex(b"test");
2972        assert_eq!(h.len(), 64);
2973    }
2974
2975    /// Core scenario: source file whose BLAKE3 matches the declared hash should
2976    /// be identified as upload-ready (no mismatch). Tests the in-memory path
2977    /// without hitting any S3 endpoint.
2978    #[test]
2979    fn matching_blake3_produces_no_mismatch_entry() {
2980        let body = b"real model weights here";
2981        let real_hash = blake3_hex(body);
2982        // Simulate what the reconciler checks.
2983        let stored = BlakeHash(real_hash.clone());
2984        assert!(hashes_equal(&real_hash, &stored.0));
2985    }
2986
2987    #[tokio::test]
2988    async fn blake3_mismatch_on_source_file_is_caught() {
2989        let fx = Fixture::new(minio_slot());
2990        let body = b"real model content";
2991        let real_hash = blake3_hex(body);
2992        // Write a workload that declares the correct hash but we'll give it
2993        // a different file content to verify the check fires before any network.
2994        let wrong_content = b"tampered content";
2995        fx.write_source_file("model.bin", wrong_content);
2996        let wrong_hash_for_real = blake3_hex(wrong_content);
2997        // Sanity: the wrong content should produce a different hash.
2998        assert_ne!(real_hash, wrong_hash_for_real);
2999
3000        let workload = StaticAssetWorkload {
3001            assets: vec![AssetEntry {
3002                filename: "model.bin".to_string(),
3003                source: Some("model.bin".into()),
3004                derive: None,
3005                blake3: BlakeHash(real_hash),
3006            }],
3007            aliases: BTreeMap::new(),
3008        };
3009
3010        // Drive the sync with a non-existent MinIO so it would fail on network
3011        // if it ever got past the BLAKE3 check.
3012        let client = reqwest::Client::new();
3013        let journal = AssetStatusJournal::new(fx.workspace_root.join(".yah/cloud/status.jsonl"));
3014        let report = sync_assets(
3015            &workload,
3016            &fx.workspace_root,
3017            &fx.workload_dir(),
3018            &client,
3019            "http://127.0.0.1:19999", // unreachable
3020            "test-bucket",
3021            MINIO_REGION,
3022            "user",
3023            "pass",
3024            test_executor(),
3025            "test-service",
3026            &journal,
3027        )
3028        .await
3029        .unwrap();
3030
3031        assert_eq!(report.hash_mismatch, vec!["model.bin"]);
3032        assert!(report.uploaded.is_empty());
3033    }
3034
3035    #[tokio::test]
3036    async fn correct_blake3_proceeds_to_s3_head() {
3037        let fx = Fixture::new(minio_slot());
3038        let body = b"correct content";
3039        let hash = blake3_hex(body);
3040        fx.write_source_file("model.bin", body);
3041
3042        let workload = StaticAssetWorkload {
3043            assets: vec![AssetEntry {
3044                filename: "model.bin".to_string(),
3045                source: Some("model.bin".into()),
3046                derive: None,
3047                blake3: BlakeHash(hash),
3048            }],
3049            aliases: BTreeMap::new(),
3050        };
3051
3052        // BLAKE3 passes → reconciler proceeds to HEAD the S3 endpoint.
3053        // The unreachable endpoint causes a network error (not a BLAKE3 error).
3054        let client = reqwest::Client::new();
3055        let journal = AssetStatusJournal::new(fx.workspace_root.join(".yah/cloud/status.jsonl"));
3056        let err = sync_assets(
3057            &workload,
3058            &fx.workspace_root,
3059            &fx.workload_dir(),
3060            &client,
3061            "http://127.0.0.1:19999", // unreachable
3062            "test-bucket",
3063            MINIO_REGION,
3064            "user",
3065            "pass",
3066            test_executor(),
3067            "test-service",
3068            &journal,
3069        )
3070        .await
3071        .unwrap_err();
3072
3073        // Error must be about network, not BLAKE3.
3074        let msg = format!("{err:#}");
3075        assert!(
3076            !msg.contains("BLAKE3") && !msg.contains("mismatch"),
3077            "should fail at network, not BLAKE3; got: {msg}"
3078        );
3079    }
3080
3081    // ── Catalog manifest ──────────────────────────────────────────────────────
3082
3083    #[test]
3084    fn catalog_manifest_roundtrip() {
3085        let mut m = CatalogManifest::new();
3086        m.insert("whisper/model-v1.bin".to_string(), HASH_64.to_string());
3087        let json = serde_json::to_vec(&m).unwrap();
3088        let m2: CatalogManifest = serde_json::from_slice(&json).unwrap();
3089        assert_eq!(m, m2);
3090    }
3091
3092    // ── R546-B10: publish-skip must be content-aware, not existence-aware ─────
3093
3094    /// The defect this encodes: an object that EXISTS is not evidence that the
3095    /// object holds the bytes we are about to publish. Skipping the PUT on mere
3096    /// existence left stale bytes in the bucket while the catalog manifest
3097    /// advertised the new hash — the published catalog described bytes that were
3098    /// never uploaded (real incident: rusty-v8 aarch64, bucket held d322b4a1
3099    /// while the manifest claimed ebdd842d).
3100    #[test]
3101    fn remote_object_matches_only_on_identical_content() {
3102        let expected = "a".repeat(64);
3103        let other = "b".repeat(64);
3104
3105        // Same bytes → genuine no-op, safe to skip the PUT.
3106        assert!(RemoteObject::Present {
3107            blake3: Some(expected.clone())
3108        }
3109        .matches(&expected));
3110
3111        // Different bytes → MUST NOT be treated as synced. This is the bug.
3112        assert!(!RemoteObject::Present {
3113            blake3: Some(other)
3114        }
3115        .matches(&expected));
3116
3117        // Absent → nothing to match.
3118        assert!(!RemoteObject::Absent.matches(&expected));
3119        assert!(!RemoteObject::Absent.exists());
3120    }
3121
3122    /// A legacy object written before we stamped `x-amz-meta-blake3` has
3123    /// unknowable content. It must NOT count as a match — otherwise every
3124    /// pre-existing object keeps its old bytes forever. It re-PUTs once, which
3125    /// installs the stamp and restores the cheap skip on subsequent runs.
3126    #[test]
3127    fn unstamped_remote_object_is_not_a_match_but_does_exist() {
3128        let expected = "c".repeat(64);
3129        let legacy = RemoteObject::Present { blake3: None };
3130        assert!(
3131            !legacy.matches(&expected),
3132            "unstamped object must not be trusted as current"
3133        );
3134        assert!(
3135            legacy.exists(),
3136            "it does exist — caller re-publishes rather than treating as absent"
3137        );
3138    }
3139
3140    /// Hashes may have been authored in either case; comparison goes through
3141    /// `hashes_equal`, so an uppercase stamp must still match.
3142    #[test]
3143    fn remote_object_match_is_case_insensitive() {
3144        let lower = "abcdef".repeat(10) + "abcd";
3145        let upper = lower.to_uppercase();
3146        assert_eq!(lower.len(), 64);
3147        assert!(RemoteObject::Present {
3148            blake3: Some(upper)
3149        }
3150        .matches(&lower));
3151    }
3152
3153    // ── Prune candidate detection ─────────────────────────────────────────────
3154
3155    #[tokio::test]
3156    async fn prune_candidate_from_prior_manifest() {
3157        let fx = Fixture::new(minio_slot());
3158        let body = b"current content";
3159        let hash = blake3_hex(body);
3160        fx.write_source_file("current.bin", body);
3161
3162        // Workload only has "current.bin"; prior manifest also had "old.bin".
3163        let workload = StaticAssetWorkload {
3164            assets: vec![AssetEntry {
3165                filename: "current.bin".to_string(),
3166                source: Some("current.bin".into()),
3167                derive: None,
3168                blake3: BlakeHash(hash),
3169            }],
3170            aliases: BTreeMap::new(),
3171        };
3172
3173        // We can't inject a prior manifest without a real/mock server, but we
3174        // can verify the prune detection logic directly.
3175        let current_filenames: std::collections::HashSet<&str> = workload
3176            .assets
3177            .iter()
3178            .map(|a| a.filename.as_str())
3179            .collect();
3180        let prior: CatalogManifest = [("old.bin".to_string(), HASH_64.to_string())]
3181            .into_iter()
3182            .collect();
3183
3184        let prune: Vec<_> = prior
3185            .keys()
3186            .filter(|k| !current_filenames.contains(k.as_str()))
3187            .cloned()
3188            .collect();
3189        assert_eq!(prune, vec!["old.bin"]);
3190    }
3191
3192    // ── W164 materialize step (R438-T15) ──────────────────────────────────────
3193
3194    use workload_spec::{AssetDerive, FetchSource, License, TransformSpec};
3195
3196    /// Helper: build an `AssetEntry` whose `source` is already on disk —
3197    /// the legacy arm of `materialize_asset` should return the disk path
3198    /// verbatim without ever consulting `derive`/cache/executor.
3199    fn legacy_entry(filename: &str, source_rel: &str, hash: &str) -> AssetEntry {
3200        AssetEntry {
3201            filename: filename.to_string(),
3202            source: Some(source_rel.into()),
3203            derive: None,
3204            blake3: BlakeHash(hash.to_string()),
3205        }
3206    }
3207
3208    #[tokio::test]
3209    async fn materialize_legacy_source_returns_workload_dir_path() {
3210        let fx = Fixture::new(minio_slot());
3211        let body = b"on-disk bytes";
3212        fx.write_source_file("model.bin", body);
3213        let entry = legacy_entry("model.bin", "model.bin", &blake3_hex(body));
3214
3215        let materialized = materialize_asset(
3216            &entry,
3217            &fx.workspace_root,
3218            &fx.workload_dir(),
3219            test_executor().as_ref(),
3220        )
3221        .await
3222        .expect("legacy materialize must succeed");
3223
3224        assert_eq!(materialized.path, fx.workload_dir().join("model.bin"));
3225        assert!(
3226            materialized.discovered_fetch_hash.is_none(),
3227            "legacy source mode has no fetch step → never bootstrap-discovers"
3228        );
3229    }
3230
3231    #[tokio::test]
3232    async fn materialize_fetch_cache_miss_then_hit() {
3233        let fx = Fixture::new(minio_slot());
3234        let body = b"upstream bytes";
3235        let hash = blake3_hex(body);
3236
3237        // Server returns the body once; second call would 500 if hit.
3238        use axum::routing::get;
3239        let counter = Arc::new(std::sync::atomic::AtomicU32::new(0));
3240        let counter_c = counter.clone();
3241        let body_c = body.to_vec();
3242        let app = axum::Router::new().route(
3243            "/blob.bin",
3244            get(move || {
3245                let counter = counter_c.clone();
3246                let body = body_c.clone();
3247                async move {
3248                    counter.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
3249                    body
3250                }
3251            }),
3252        );
3253        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
3254        let addr = listener.local_addr().unwrap();
3255        tokio::spawn(async move {
3256            axum::serve(listener, app).await.unwrap();
3257        });
3258
3259        let fetch = FetchSource {
3260            url: format!("http://{addr}/blob.bin"),
3261            blake3: BlakeHash(hash.clone()),
3262            license: License::Mit,
3263        };
3264        let cache_dir = fx.workspace_root.join(".yah/cache/derive/fetch");
3265
3266        // MISS — fetch + write to cache.
3267        let (p1, h1) = materialize_fetch(&fetch, &cache_dir).await.unwrap();
3268        assert_eq!(counter.load(std::sync::atomic::Ordering::SeqCst), 1);
3269        assert!(p1.exists(), "cache file must exist after MISS");
3270        assert_eq!(h1, hash, "strict mode returns the pinned hash");
3271
3272        // HIT — no second HTTP call.
3273        let (p2, h2) = materialize_fetch(&fetch, &cache_dir).await.unwrap();
3274        assert_eq!(p1, p2);
3275        assert_eq!(h1, h2, "HIT returns the same hash");
3276        assert_eq!(
3277            counter.load(std::sync::atomic::Ordering::SeqCst),
3278            1,
3279            "cache HIT must not re-fetch"
3280        );
3281    }
3282
3283    #[tokio::test]
3284    async fn materialize_fetch_blake3_mismatch_is_hard_error() {
3285        let fx = Fixture::new(minio_slot());
3286        let actual_body = b"real bytes";
3287        let actual_hash = blake3_hex(actual_body);
3288        let pinned_hash = HASH_64; // deliberately wrong
3289        assert_ne!(actual_hash, pinned_hash);
3290
3291        use axum::routing::get;
3292        let body_c = actual_body.to_vec();
3293        let app = axum::Router::new().route(
3294            "/blob.bin",
3295            get(move || {
3296                let body = body_c.clone();
3297                async move { body }
3298            }),
3299        );
3300        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
3301        let addr = listener.local_addr().unwrap();
3302        tokio::spawn(async move {
3303            axum::serve(listener, app).await.unwrap();
3304        });
3305
3306        let fetch = FetchSource {
3307            url: format!("http://{addr}/blob.bin"),
3308            blake3: BlakeHash(pinned_hash.to_string()),
3309            license: License::Mit,
3310        };
3311        let err = materialize_fetch(&fetch, &fx.workspace_root.join(".yah/cache/derive/fetch"))
3312            .await
3313            .expect_err("BLAKE3 mismatch must surface as hard error");
3314        let msg = format!("{err:#}");
3315        assert!(
3316            msg.contains("BLAKE3") || msg.contains(&actual_hash),
3317            "got: {msg}"
3318        );
3319        assert!(msg.contains(pinned_hash), "diff must mention pin: {msg}");
3320    }
3321
3322    #[tokio::test]
3323    async fn materialize_fetch_retries_on_server_error() {
3324        // First request returns 500 (transient); second returns 200 with the body.
3325        // Verifies exponential-backoff retry (W164 OQ#4, R438-F11).
3326        use axum::http::StatusCode;
3327        use axum::routing::get;
3328
3329        let fx = Fixture::new(minio_slot());
3330        let body: Vec<u8> = b"retry-me bytes".to_vec();
3331        let hash = blake3_hex(&body);
3332
3333        let counter = Arc::new(std::sync::atomic::AtomicU32::new(0));
3334        let counter_c = counter.clone();
3335        let body_c = body.clone();
3336        let app = axum::Router::new().route(
3337            "/blob.bin",
3338            get(move || {
3339                let n = counter_c.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
3340                let body = body_c.clone();
3341                async move {
3342                    if n == 0 {
3343                        (StatusCode::INTERNAL_SERVER_ERROR, vec![])
3344                    } else {
3345                        (StatusCode::OK, body)
3346                    }
3347                }
3348            }),
3349        );
3350        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
3351        let addr = listener.local_addr().unwrap();
3352        tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
3353
3354        let fetch = FetchSource {
3355            url: format!("http://{addr}/blob.bin"),
3356            blake3: BlakeHash(hash.clone()),
3357            license: License::Mit,
3358        };
3359        let cache_dir = fx.workspace_root.join(".yah/cache/derive/fetch");
3360        let (path, _hash) = materialize_fetch(&fetch, &cache_dir).await.unwrap();
3361        assert!(
3362            path.exists(),
3363            "cache file must exist after successful retry"
3364        );
3365        assert_eq!(
3366            counter.load(std::sync::atomic::Ordering::SeqCst),
3367            2,
3368            "exactly two attempts: one 500 + one 200"
3369        );
3370        assert_eq!(tokio::fs::read(&path).await.unwrap(), body);
3371    }
3372
3373    #[tokio::test]
3374    async fn materialize_fetch_resumes_via_range_header() {
3375        // Pre-populate the partial file with the first half of the blob. The
3376        // server should receive a Range header and return 206 + the second half.
3377        // Verifies Range-resume logic (W164 OQ#4, R438-F11).
3378        use axum::body::Body;
3379        use axum::http::{HeaderMap, StatusCode};
3380        use axum::routing::get;
3381
3382        let fx = Fixture::new(minio_slot());
3383        let body: Vec<u8> = (0u8..100u8).collect();
3384        let half = body.len() / 2;
3385        let hash = blake3_hex(&body);
3386
3387        // Pre-write first half to the .partial file.
3388        let cache_dir = fx.workspace_root.join(".yah/cache/derive/fetch");
3389        tokio::fs::create_dir_all(&cache_dir).await.unwrap();
3390        let partial_path = cache_dir.join(format!("{hash}.partial"));
3391        tokio::fs::write(&partial_path, &body[..half])
3392            .await
3393            .unwrap();
3394
3395        let body_c = body.clone();
3396        let received_range: Arc<std::sync::Mutex<Option<String>>> =
3397            Arc::new(std::sync::Mutex::new(None));
3398        let received_range_c = received_range.clone();
3399        let app = axum::Router::new().route(
3400            "/blob.bin",
3401            get(move |headers: HeaderMap| {
3402                let body = body_c.clone();
3403                let received = received_range_c.clone();
3404                async move {
3405                    let range_hdr = headers
3406                        .get("range")
3407                        .and_then(|v| v.to_str().ok())
3408                        .map(str::to_string);
3409                    *received.lock().unwrap() = range_hdr.clone();
3410
3411                    let Some(range) = range_hdr else {
3412                        return axum::response::Response::builder()
3413                            .status(StatusCode::BAD_REQUEST)
3414                            .body(Body::empty())
3415                            .unwrap();
3416                    };
3417                    let offset: usize = range
3418                        .strip_prefix("bytes=")
3419                        .and_then(|s| s.strip_suffix('-'))
3420                        .and_then(|s| s.parse().ok())
3421                        .unwrap_or(0);
3422                    let total = body.len();
3423                    axum::response::Response::builder()
3424                        .status(StatusCode::PARTIAL_CONTENT)
3425                        .header(
3426                            "Content-Range",
3427                            format!("bytes {offset}-{}/{total}", total - 1),
3428                        )
3429                        .body(Body::from(body[offset..].to_vec()))
3430                        .unwrap()
3431                }
3432            }),
3433        );
3434        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
3435        let addr = listener.local_addr().unwrap();
3436        tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
3437
3438        let fetch = FetchSource {
3439            url: format!("http://{addr}/blob.bin"),
3440            blake3: BlakeHash(hash.clone()),
3441            license: License::Mit,
3442        };
3443        let (path, _hash) = materialize_fetch(&fetch, &cache_dir).await.unwrap();
3444
3445        // Range header must carry the pre-existing partial size.
3446        let range = received_range.lock().unwrap().clone().unwrap();
3447        assert_eq!(
3448            range,
3449            format!("bytes={half}-"),
3450            "Range header must resume from partial offset"
3451        );
3452        // Final cache file must contain the complete body.
3453        assert_eq!(
3454            tokio::fs::read(&path).await.unwrap(),
3455            body,
3456            "cache file must contain the full blob after Range resume"
3457        );
3458    }
3459
3460    #[tokio::test]
3461    async fn materialize_transform_cache_miss_then_hit() {
3462        let fx = Fixture::new(minio_slot());
3463        write_recipe(&fx.workspace_root, "noop-recipe");
3464        let out_bytes = b"transformed output".to_vec();
3465        let out_hash = blake3_hex(&out_bytes);
3466        let (mock, invocations) = MockExecutor::new(out_bytes.clone());
3467        let executor: Arc<dyn ForgeExecutor> = mock;
3468
3469        // Pre-seed a fetch cache entry the transform reads from.
3470        let fetch_path = fx.workspace_root.join(".yah/cache/derive/fetch/in.bin");
3471        std::fs::create_dir_all(fetch_path.parent().unwrap()).unwrap();
3472        std::fs::write(&fetch_path, b"fetch bytes").unwrap();
3473
3474        let transform = TransformSpec {
3475            recipe: "noop-recipe".to_string(),
3476            params: BTreeMap::new(),
3477        };
3478        let cache_dir = fx.workspace_root.join(".yah/cache/derive/transform");
3479
3480        // MISS — runs the recipe.
3481        let p1 = materialize_transform(
3482            &transform,
3483            &fetch_path,
3484            &blake3_hex(b"fetch bytes"),
3485            &out_hash,
3486            &cache_dir,
3487            &fx.workspace_root,
3488            executor.as_ref(),
3489        )
3490        .await
3491        .unwrap();
3492        assert!(p1.exists());
3493        assert_eq!(*invocations.lock().await, 1);
3494
3495        // HIT — recipe NOT re-run.
3496        let p2 = materialize_transform(
3497            &transform,
3498            &fetch_path,
3499            &blake3_hex(b"fetch bytes"),
3500            &out_hash,
3501            &cache_dir,
3502            &fx.workspace_root,
3503            executor.as_ref(),
3504        )
3505        .await
3506        .unwrap();
3507        assert_eq!(p1, p2);
3508        assert_eq!(*invocations.lock().await, 1, "cache HIT must not re-run");
3509    }
3510
3511    /// R546-T3 qed→W164 bridge: seeding the derivation cache from a pre-built
3512    /// artifact makes the NEXT `materialize_transform` a HIT that returns the
3513    /// seeded bytes without running the recipe. This is the whole point —
3514    /// offload the arch-locked, expensive build to a fleet worker (rusty-v8-musl), retrieve
3515    /// its tar (R590-F6), seed the cache, and let `yah cloud apply` publish the
3516    /// pre-built bytes instead of re-building under emulation.
3517    #[tokio::test]
3518    async fn seed_transform_derivation_makes_next_transform_a_hit() {
3519        let fx = Fixture::new(minio_slot());
3520        write_recipe(&fx.workspace_root, "noop-recipe");
3521
3522        let transform = TransformSpec {
3523            recipe: "noop-recipe".to_string(),
3524            params: BTreeMap::new(),
3525        };
3526        let fetched_hash = blake3_hex(b"fetch bytes");
3527
3528        // A pre-built artifact standing in for the F6-retrieved rusty_v8 tar.
3529        let prebuilt = fx.workspace_root.join("prebuilt.tar.gz");
3530        let prebuilt_bytes = b"prebuilt rusty_v8 tarball".to_vec();
3531        std::fs::write(&prebuilt, &prebuilt_bytes).unwrap();
3532
3533        let seeded =
3534            seed_transform_derivation(&fx.workspace_root, &transform, &fetched_hash, &prebuilt)
3535                .await
3536                .unwrap();
3537        assert_eq!(seeded.output_blake3, blake3_hex(&prebuilt_bytes));
3538        assert!(seeded.cas_path.exists(), "CAS entry must be written");
3539
3540        // A materialize_transform with the SAME inputs must HIT the seeded cache:
3541        // the executor is never invoked and the returned bytes are the seeded
3542        // artifact (NOT what the recipe would have produced).
3543        let fetch_path = fx.workspace_root.join(".yah/cache/derive/fetch/in.bin");
3544        std::fs::create_dir_all(fetch_path.parent().unwrap()).unwrap();
3545        std::fs::write(&fetch_path, b"fetch bytes").unwrap();
3546        let cache_dir = fx.workspace_root.join(".yah/cache/derive/transform");
3547        let (mock, invocations) = MockExecutor::new(b"RECIPE OUTPUT (must not appear)".to_vec());
3548
3549        let out = materialize_transform(
3550            &transform,
3551            &fetch_path,
3552            &fetched_hash,
3553            &seeded.output_blake3,
3554            &cache_dir,
3555            &fx.workspace_root,
3556            mock.as_ref(),
3557        )
3558        .await
3559        .unwrap();
3560
3561        assert_eq!(
3562            *invocations.lock().await,
3563            0,
3564            "seeded derivation cache must make the transform a HIT (no recipe run)"
3565        );
3566        assert_eq!(
3567            std::fs::read(&out).unwrap(),
3568            prebuilt_bytes,
3569            "HIT must return the seeded pre-built bytes",
3570        );
3571    }
3572
3573    /// W212/R518: the derivation key is hermetic — every declared input
3574    /// dimension flips it, and identical inputs reproduce it. This is the
3575    /// correctness invariant that makes the silent-stale-cache bug impossible:
3576    /// a missing dimension would let two distinct inputs collide on one key.
3577    #[test]
3578    fn derivation_key_is_hermetic_over_every_input_dimension() {
3579        let input = "a".repeat(64);
3580        let recipe = "b".repeat(64);
3581        let mut params = BTreeMap::new();
3582        params.insert("quant".to_string(), "q5_1".to_string());
3583
3584        let base = derivation_key(&input, &recipe, &params);
3585
3586        // Fetched-input content (the model/version anchor) flips it.
3587        assert_ne!(base, derivation_key(&"c".repeat(64), &recipe, &params));
3588        // Recipe-file bytes (steps / pinned container digest / placement) flip it.
3589        assert_ne!(base, derivation_key(&input, &"d".repeat(64), &params));
3590        // A changed param value flips it.
3591        let mut p2 = params.clone();
3592        p2.insert("quant".to_string(), "q4_0".to_string());
3593        assert_ne!(base, derivation_key(&input, &recipe, &p2));
3594        // An added param flips it.
3595        let mut p3 = params.clone();
3596        p3.insert("extra".to_string(), "x".to_string());
3597        assert_ne!(base, derivation_key(&input, &recipe, &p3));
3598        // Identical inputs reproduce the key (deterministic, order-stable).
3599        assert_eq!(base, derivation_key(&input, &recipe, &params));
3600    }
3601
3602    /// W212/R518: the action cache is keyed on INPUTS, so a recipe change (here
3603    /// a container-digest bump) forces a re-run even on a warm cache. This is
3604    /// the regression guard against the W164/R438 output-keyed silent-stale
3605    /// skip — the exact defect this relay fixes.
3606    #[tokio::test]
3607    async fn transform_recipe_change_busts_cache_and_reruns() {
3608        let fx = Fixture::new(minio_slot());
3609        let digest_a = HASH_64;
3610        let digest_b = "1".repeat(64);
3611        write_recipe_with_digest(&fx.workspace_root, "pinned", digest_a);
3612
3613        let out_bytes = b"transformed output".to_vec();
3614        let out_hash = blake3_hex(&out_bytes);
3615        let (mock, invocations) = MockExecutor::new(out_bytes);
3616        let executor: Arc<dyn ForgeExecutor> = mock;
3617
3618        let fetch_path = fx.workspace_root.join(".yah/cache/derive/fetch/in.bin");
3619        std::fs::create_dir_all(fetch_path.parent().unwrap()).unwrap();
3620        std::fs::write(&fetch_path, b"fetch bytes").unwrap();
3621        let in_hash = blake3_hex(b"fetch bytes");
3622
3623        let transform = TransformSpec {
3624            recipe: "pinned".to_string(),
3625            params: BTreeMap::new(),
3626        };
3627        let cache_dir = fx.workspace_root.join(".yah/cache/derive/transform");
3628        let wr = &fx.workspace_root;
3629
3630        // Run 1 — cold, digest A.
3631        materialize_transform(
3632            &transform,
3633            &fetch_path,
3634            &in_hash,
3635            &out_hash,
3636            &cache_dir,
3637            wr,
3638            executor.as_ref(),
3639        )
3640        .await
3641        .unwrap();
3642        assert_eq!(*invocations.lock().await, 1);
3643
3644        // Re-run with the SAME recipe → warm HIT, no re-run.
3645        materialize_transform(
3646            &transform,
3647            &fetch_path,
3648            &in_hash,
3649            &out_hash,
3650            &cache_dir,
3651            wr,
3652            executor.as_ref(),
3653        )
3654        .await
3655        .unwrap();
3656        assert_eq!(*invocations.lock().await, 1, "identical inputs must skip");
3657
3658        // Bump the container digest A → B. Output-keyed caching silently
3659        // skipped this; input-addressed caching MUST re-run.
3660        write_recipe_with_digest(&fx.workspace_root, "pinned", &digest_b);
3661        materialize_transform(
3662            &transform,
3663            &fetch_path,
3664            &in_hash,
3665            &out_hash,
3666            &cache_dir,
3667            wr,
3668            executor.as_ref(),
3669        )
3670        .await
3671        .unwrap();
3672        assert_eq!(
3673            *invocations.lock().await,
3674            2,
3675            "a digest bump must invalidate the cache and re-run",
3676        );
3677    }
3678
3679    /// W212/R518-P2b: the substituter fast-path honours the committed lock —
3680    /// `Some(output)` when the key recomputed from the pins matches, `None` the
3681    /// moment any input drifts (here a recipe-digest bump).
3682    #[tokio::test]
3683    async fn lock_skip_hash_honours_committed_lock() {
3684        let fx = Fixture::new(minio_slot());
3685        write_recipe_with_digest(&fx.workspace_root, "pinned", HASH_64);
3686
3687        let fetch_pin = "a".repeat(64);
3688        let out_pin = "b".repeat(64);
3689        // The lock's input_hash must equal the key the reconciler recomputes
3690        // from the committed pins.
3691        let recipe_bk = recipe_blake3(&fx.workspace_root, "pinned").await.unwrap();
3692        let key = derivation_key(&fetch_pin, &recipe_bk, &BTreeMap::new());
3693
3694        let toml = format!(
3695            r#"kind = "static-asset"
3696schema_version = "V1"
3697
3698[[asset]]
3699filename = "whisper/coreml.tar.gz"
3700blake3   = "{out_pin}"
3701
3702[asset.derive.fetch]
3703url     = "https://example/config.json"
3704blake3  = "{fetch_pin}"
3705license = "mit"
3706
3707[asset.derive.transform]
3708recipe = "pinned"
3709
3710[asset.derive.lock]
3711input_hash    = "{key}"
3712output_blake3 = "{out_pin}"
3713"#
3714        );
3715        let workload: StaticAssetWorkload = toml::from_str(&toml).unwrap();
3716        let entry = &workload.assets[0];
3717
3718        // Lock current → Some(output hash) (caller then HEADs the bucket).
3719        assert_eq!(
3720            lock_skip_hash(entry, &fx.workspace_root).await.as_deref(),
3721            Some(out_pin.as_str()),
3722        );
3723
3724        // Bump the recipe digest → recipe-file hash changes → recomputed key no
3725        // longer matches the lock → None (must rebuild). The exact drift the
3726        // output-keyed cache missed.
3727        write_recipe_with_digest(&fx.workspace_root, "pinned", &"1".repeat(64));
3728        assert!(
3729            lock_skip_hash(entry, &fx.workspace_root).await.is_none(),
3730            "a stale recipe must not skip the build",
3731        );
3732    }
3733
3734    /// W212/R518-P2b: a bootstrap (zero-sentinel) row never short-circuits via
3735    /// the lock — there is no pinned output to trust yet.
3736    #[tokio::test]
3737    async fn lock_skip_hash_declines_bootstrap_rows() {
3738        let fx = Fixture::new(minio_slot());
3739        write_recipe_with_digest(&fx.workspace_root, "pinned", HASH_64);
3740        let sentinel = "0".repeat(64);
3741        let toml = format!(
3742            r#"kind = "static-asset"
3743schema_version = "V1"
3744
3745[[asset]]
3746filename = "whisper/coreml.tar.gz"
3747blake3   = "{sentinel}"
3748
3749[asset.derive.fetch]
3750url     = "https://example/config.json"
3751blake3  = "{sentinel}"
3752license = "mit"
3753
3754[asset.derive.transform]
3755recipe = "pinned"
3756
3757[asset.derive.lock]
3758input_hash    = "whatever"
3759output_blake3 = "{sentinel}"
3760"#
3761        );
3762        let workload: StaticAssetWorkload = toml::from_str(&toml).unwrap();
3763        assert!(
3764            lock_skip_hash(&workload.assets[0], &fx.workspace_root)
3765                .await
3766                .is_none(),
3767            "bootstrap rows must always build",
3768        );
3769    }
3770
3771    /// R546-B6: the fetched-input hash is the FOURTH paste-back value. Seeding
3772    /// used to compute it, key the derivation on it, and drop it — leaving
3773    /// `[asset.derive.fetch].blake3` at the sentinel, which makes
3774    /// `lock_skip_hash` decline at its first guard for everyone but the seeding
3775    /// machine (whose local action cache masks it).
3776    #[tokio::test]
3777    async fn seed_surfaces_fetched_input_hash_and_pin_state() {
3778        let fx = Fixture::new(minio_slot());
3779        write_recipe(&fx.workspace_root, "pinned");
3780
3781        // Warm the fetch cache so materialize_fetch takes its HIT path — no
3782        // network, and the committed pin is already correct (Pinned arm).
3783        let fetch_bytes = b"upstream source tarball".to_vec();
3784        let fetch_hash = blake3_hex(&fetch_bytes);
3785        let fetch_cache = fx.workspace_root.join(".yah/cache/derive/fetch");
3786        std::fs::create_dir_all(&fetch_cache).unwrap();
3787        std::fs::write(fetch_cache.join(format!("{fetch_hash}.bin")), &fetch_bytes).unwrap();
3788
3789        let sentinel = "0".repeat(64);
3790        let workload_path = fx.workspace_root.join("workload.toml");
3791        std::fs::write(
3792            &workload_path,
3793            format!(
3794                r#"kind = "static-asset"
3795schema_version = "V1"
3796
3797[[asset]]
3798filename = "rusty-v8/x86_64.tar.gz"
3799blake3   = "{sentinel}"
3800
3801[asset.derive.fetch]
3802url     = "https://example/v8.tar.gz"
3803blake3  = "{fetch_hash}"
3804license = "mit"
3805
3806[asset.derive.transform]
3807recipe = "pinned"
3808params = {{ target = "x86_64-unknown-linux-musl" }}
3809"#
3810            ),
3811        )
3812        .unwrap();
3813
3814        let prebuilt = fx.workspace_root.join("prebuilt.tar.gz");
3815        std::fs::write(&prebuilt, b"prebuilt rusty_v8 tarball").unwrap();
3816
3817        let seeded = seed_derivation_for_target(
3818            &fx.workspace_root,
3819            &workload_path,
3820            "x86_64-unknown-linux-musl",
3821            &prebuilt,
3822        )
3823        .await
3824        .unwrap();
3825
3826        assert_eq!(
3827            seeded.fetch_blake3, fetch_hash,
3828            "the hash the derivation key was keyed on must be surfaced for paste-back",
3829        );
3830        assert_eq!(seeded.fetch_pin, Some(FetchPinState::Pinned));
3831
3832        // And it really is the key input: recomputing from the surfaced value
3833        // reproduces the seeded lock's input_hash.
3834        let recipe_bk = recipe_blake3(&fx.workspace_root, "pinned").await.unwrap();
3835        let mut params = BTreeMap::new();
3836        params.insert(
3837            "target".to_string(),
3838            "x86_64-unknown-linux-musl".to_string(),
3839        );
3840        assert_eq!(
3841            derivation_key(&seeded.fetch_blake3, &recipe_bk, &params),
3842            seeded.derive_key,
3843        );
3844    }
3845
3846    /// R546-B6: the three pin verdicts the seed command reports on. `Sentinel`
3847    /// is the disarmed case the ticket was filed for; `Mismatch` catches a
3848    /// moved version anchor before the operator pastes a lock that can only
3849    /// ever decline on input drift.
3850    #[test]
3851    fn classify_fetch_pin_covers_all_three_verdicts() {
3852        let fetched = blake3_hex(b"upstream");
3853        assert_eq!(
3854            classify_fetch_pin(&"0".repeat(64), &fetched),
3855            FetchPinState::Sentinel,
3856        );
3857        assert_eq!(
3858            classify_fetch_pin(&fetched.to_uppercase(), &fetched),
3859            FetchPinState::Pinned,
3860            "hash comparison is case-insensitive like the rest of the reconciler",
3861        );
3862        let stale = blake3_hex(b"a different upstream");
3863        assert_eq!(
3864            classify_fetch_pin(&stale, &fetched),
3865            FetchPinState::Mismatch {
3866                committed: stale.clone()
3867            },
3868        );
3869    }
3870
3871    /// R546-B12: a clean reconcile must SAY something. The failure this guards
3872    /// is not a wrong line, it is an absent one — a successful publish and a
3873    /// successful no-op printing identically (i.e. nothing), which is what made
3874    /// the blind `yah cloud status` view so costly to notice.
3875    #[test]
3876    fn sync_notes_distinguish_publish_from_no_op_and_are_never_empty() {
3877        // Published.
3878        let mut report = StaticAssetSyncReport::default();
3879        report.uploaded.push("a/one.bin".to_string());
3880        let notes = render_sync_notes(&report);
3881        assert!(
3882            notes.iter().any(|n| n.contains("published 1 asset")),
3883            "{notes:?}"
3884        );
3885        assert!(notes.iter().any(|n| n.contains("a/one.bin")), "{notes:?}");
3886
3887        // No-op: everything already current. Must NOT look like a publish, and
3888        // must not be silent.
3889        let mut report = StaticAssetSyncReport::default();
3890        report.already_synced.push("a/one.bin".to_string());
3891        let notes = render_sync_notes(&report);
3892        assert!(
3893            notes.iter().any(|n| n.contains("already current")),
3894            "{notes:?}"
3895        );
3896        assert!(
3897            !notes.iter().any(|n| n.contains("published")),
3898            "a no-op must not read as a publish: {notes:?}"
3899        );
3900
3901        // Bootstrap discoveries are paste-back values — the operator needs the
3902        // hashes themselves, not a count.
3903        let mut report = StaticAssetSyncReport::default();
3904        report.bootstrapped.push(BootstrappedHash {
3905            filename: "a/one.bin".into(),
3906            kind: BootstrapHashKind::Output,
3907            hash: HASH_64.to_string(),
3908        });
3909        let notes = render_sync_notes(&report);
3910        assert!(notes.iter().any(|n| n.contains(HASH_64)), "{notes:?}");
3911
3912        // The degenerate case still speaks.
3913        let notes = render_sync_notes(&StaticAssetSyncReport::default());
3914        assert_eq!(notes.len(), 1, "{notes:?}");
3915        assert!(notes[0].contains("nothing to reconcile"), "{notes:?}");
3916    }
3917
3918    #[test]
3919    fn bootstrap_output_key_input_variant() {
3920        let b = BootstrappedHash {
3921            filename: "whisper/coreml.tar.gz".into(),
3922            kind: BootstrapHashKind::Input,
3923            hash: "deadbeef".into(),
3924        };
3925        assert_eq!(
3926            bootstrap_output_key(&b),
3927            "discovered_input_hash:whisper/coreml.tar.gz",
3928        );
3929    }
3930
3931    /// R546-B8: both recipe bindings must be ABSOLUTE, even when the caller
3932    /// passes a relative workspace root.
3933    ///
3934    /// The fixture deliberately uses `TempDir::new_in(".")` rather than the
3935    /// usual absolute tempdir — that is the whole point. `yah cloud apply`
3936    /// defaults `--path` to `"."`, so `cache_dir` came out relative
3937    /// (`./.yah/cache/derive/transform/<key>.tmp`) and got handed to the recipe
3938    /// verbatim. A relative OUT survives only while the recipe stays in its
3939    /// cwd, which is why the whisper recipes never tripped it; rusty-v8's
3940    /// build-v8.sh chdirs into a scratch dir to build V8, so it wrote the
3941    /// finished tar to `<scratch>/./.yah/cache/...` inside the container and the
3942    /// bytes died with it — after a ~2h build, with the step reporting exit 0.
3943    ///
3944    /// An absolute-tempdir fixture cannot catch this: `cache_dir` is then
3945    /// absolute either way and the assertion passes against the buggy code too.
3946    #[tokio::test]
3947    async fn transform_bindings_are_absolute_from_a_relative_workspace_root() {
3948        // `TempDir::new_in(".")` hands back an ABSOLUTE path, so rebuild the
3949        // relative form by hand — cargo runs tests with cwd at the package root,
3950        // so `./<name>` resolves to the same directory.
3951        let tmp = tempfile::TempDir::new_in(".").expect("tempdir beside cwd");
3952        let workspace_root = Path::new(".").join(tmp.path().file_name().unwrap());
3953        assert!(
3954            workspace_root.is_relative(),
3955            "fixture must reproduce the relative-root case, got {}",
3956            workspace_root.display(),
3957        );
3958
3959        write_recipe(&workspace_root, "noop-recipe");
3960        let cache_dir = workspace_root.join(".yah/cache/derive/transform");
3961        std::fs::create_dir_all(&cache_dir).unwrap();
3962        let fetch_path = workspace_root.join(".yah/cache/derive/fetch/in.bin");
3963        std::fs::create_dir_all(fetch_path.parent().unwrap()).unwrap();
3964        std::fs::write(&fetch_path, b"fetch bytes").unwrap();
3965
3966        let (capture, argv) = ArgvCapture::new(true);
3967        let executor: Arc<dyn ForgeExecutor> = capture;
3968        let transform = TransformSpec {
3969            recipe: "noop-recipe".to_string(),
3970            params: BTreeMap::new(),
3971        };
3972
3973        materialize_transform(
3974            &transform,
3975            &fetch_path,
3976            &blake3_hex(b"fetch bytes"),
3977            ZERO_SENTINEL_HEX, // bootstrap: accept whatever hash comes out
3978            &cache_dir,
3979            &workspace_root,
3980            executor.as_ref(),
3981        )
3982        .await
3983        .expect("transform must succeed");
3984
3985        // Recipe argv is ["./tool", "{{YAH_TRANSFORM_IN_0}}", "{{YAH_TRANSFORM_OUT}}"].
3986        let argv = argv.lock().await.clone();
3987        assert_eq!(argv.len(), 3, "unexpected argv shape: {argv:?}");
3988        assert!(
3989            Path::new(&argv[1]).is_absolute(),
3990            "{ENV_TRANSFORM_IN_0} binding must be absolute, got {:?}",
3991            argv[1],
3992        );
3993        assert!(
3994            Path::new(&argv[2]).is_absolute(),
3995            "{ENV_TRANSFORM_OUT} binding must be absolute, got {:?}",
3996            argv[2],
3997        );
3998    }
3999
4000    /// R546-B8 hardening: a recipe that exits 0 without writing its artifact
4001    /// must be named as such. Before this, the missing file surfaced as
4002    /// `reading transform output …: No such file or directory` on a path under
4003    /// `.yah/cache/`, which reads like cache corruption and sends the reader
4004    /// into the wrong subsystem.
4005    #[tokio::test]
4006    async fn transform_that_writes_no_output_names_the_recipe_contract() {
4007        let fx = Fixture::new(minio_slot());
4008        write_recipe(&fx.workspace_root, "writes-nothing");
4009        let cache_dir = fx.workspace_root.join(".yah/cache/derive/transform");
4010        let fetch_path = fx.workspace_root.join(".yah/cache/derive/fetch/in.bin");
4011        std::fs::create_dir_all(fetch_path.parent().unwrap()).unwrap();
4012        std::fs::write(&fetch_path, b"fetch bytes").unwrap();
4013
4014        // Exit 0, write nothing — exactly what the relative-OUT bug looked like
4015        // from the reconciler's side.
4016        let (capture, _) = ArgvCapture::new(false);
4017        let executor: Arc<dyn ForgeExecutor> = capture;
4018        let transform = TransformSpec {
4019            recipe: "writes-nothing".to_string(),
4020            params: BTreeMap::new(),
4021        };
4022
4023        let err = materialize_transform(
4024            &transform,
4025            &fetch_path,
4026            &blake3_hex(b"fetch bytes"),
4027            ZERO_SENTINEL_HEX,
4028            &cache_dir,
4029            &fx.workspace_root,
4030            executor.as_ref(),
4031        )
4032        .await
4033        .expect_err("a recipe that writes nothing must fail");
4034
4035        let msg = format!("{err:#}");
4036        assert!(msg.contains("produced no output"), "{msg}");
4037        assert!(msg.contains("writes-nothing"), "{msg}");
4038        assert!(
4039            !msg.contains("No such file or directory"),
4040            "must not leak the raw io error that reads like cache corruption: {msg}",
4041        );
4042    }
4043
4044    #[tokio::test]
4045    async fn materialize_transform_blake3_mismatch_on_output_is_hard_error() {
4046        let fx = Fixture::new(minio_slot());
4047        write_recipe(&fx.workspace_root, "wrong-output");
4048        let out_bytes = b"actual output".to_vec();
4049        let actual_hash = blake3_hex(&out_bytes);
4050        let pinned_hash = HASH_64; // deliberately wrong
4051        assert_ne!(actual_hash, pinned_hash);
4052        let (mock, _) = MockExecutor::new(out_bytes);
4053        let executor: Arc<dyn ForgeExecutor> = mock;
4054
4055        let fetch_path = fx.workspace_root.join(".yah/cache/derive/fetch/in.bin");
4056        std::fs::create_dir_all(fetch_path.parent().unwrap()).unwrap();
4057        std::fs::write(&fetch_path, b"fetch bytes").unwrap();
4058
4059        let transform = TransformSpec {
4060            recipe: "wrong-output".to_string(),
4061            params: BTreeMap::new(),
4062        };
4063        let err = materialize_transform(
4064            &transform,
4065            &fetch_path,
4066            &blake3_hex(b"fetch bytes"),
4067            pinned_hash,
4068            &fx.workspace_root.join(".yah/cache/derive/transform"),
4069            &fx.workspace_root,
4070            executor.as_ref(),
4071        )
4072        .await
4073        .expect_err("transform output mismatch must surface");
4074        let msg = format!("{err:#}");
4075        assert!(msg.contains("BLAKE3"), "{msg}");
4076        assert!(msg.contains(pinned_hash), "{msg}");
4077
4078        // R925: this is the path the `StagingOutput` guard exists for — the mock
4079        // DID write the staging output, and the mismatch bailed past the
4080        // publishing rename, so the guard must have unlinked it. The staging
4081        // name is `<derive_key>.out.<pid>.<seq>.tmp`, which no later run reuses,
4082        // so a missed cleanup accumulates without bound.
4083        //
4084        // Deliberately a directory scan, not `assert!(!cache_dir.join(format!(
4085        // "{derive_key}.tmp")).exists())`: an assertion naming one hard-coded
4086        // staging filename passes VACUOUSLY once the name carries a pid, and the
4087        // leak coverage disappears with the test still green.
4088        let cache_dir = fx.workspace_root.join(".yah/cache/derive/transform");
4089        let strays: Vec<_> = std::fs::read_dir(&cache_dir)
4090            .unwrap()
4091            .flatten()
4092            .map(|e| e.file_name().to_string_lossy().into_owned())
4093            .filter(|n| n.ends_with(".tmp"))
4094            .collect();
4095        assert!(strays.is_empty(), "staging output leaked: {strays:?}");
4096    }
4097
4098    #[tokio::test]
4099    async fn materialize_transform_recipe_failure_surfaces_stderr() {
4100        let fx = Fixture::new(minio_slot());
4101        write_recipe(&fx.workspace_root, "failing-recipe");
4102        let executor: Arc<dyn ForgeExecutor> = MockExecutor::failing("tool exploded".to_string());
4103
4104        let fetch_path = fx.workspace_root.join(".yah/cache/derive/fetch/in.bin");
4105        std::fs::create_dir_all(fetch_path.parent().unwrap()).unwrap();
4106        std::fs::write(&fetch_path, b"fetch bytes").unwrap();
4107
4108        let transform = TransformSpec {
4109            recipe: "failing-recipe".to_string(),
4110            params: BTreeMap::new(),
4111        };
4112        let err = materialize_transform(
4113            &transform,
4114            &fetch_path,
4115            &blake3_hex(b"fetch bytes"),
4116            HASH_64,
4117            &fx.workspace_root.join(".yah/cache/derive/transform"),
4118            &fx.workspace_root,
4119            executor.as_ref(),
4120        )
4121        .await
4122        .expect_err("failing recipe must surface");
4123        let msg = format!("{err:#}");
4124        assert!(msg.contains("failing-recipe"), "{msg}");
4125        assert!(msg.contains("tool exploded"), "{msg}");
4126    }
4127
4128    #[tokio::test]
4129    async fn materialize_derive_fetch_only_uses_fetch_path_for_upload() {
4130        // No transform — materialize_asset returns the fetch cache path
4131        // directly. The downstream BLAKE3 verify happens in sync_assets.
4132        let fx = Fixture::new(minio_slot());
4133        let body = b"weights v1";
4134        let hash = blake3_hex(body);
4135
4136        use axum::routing::get;
4137        let body_c = body.to_vec();
4138        let app = axum::Router::new().route(
4139            "/w.bin",
4140            get(move || {
4141                let body = body_c.clone();
4142                async move { body }
4143            }),
4144        );
4145        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
4146        let addr = listener.local_addr().unwrap();
4147        tokio::spawn(async move {
4148            axum::serve(listener, app).await.unwrap();
4149        });
4150
4151        let entry = AssetEntry {
4152            filename: "w.bin".to_string(),
4153            source: None,
4154            derive: Some(AssetDerive {
4155                fetch: FetchSource {
4156                    url: format!("http://{addr}/w.bin"),
4157                    blake3: BlakeHash(hash.clone()),
4158                    license: License::Mit,
4159                },
4160                transform: None,
4161                lock: None,
4162            }),
4163            blake3: BlakeHash(hash.clone()),
4164        };
4165        let materialized = materialize_asset(
4166            &entry,
4167            &fx.workspace_root,
4168            &fx.workload_dir(),
4169            test_executor().as_ref(),
4170        )
4171        .await
4172        .expect("fetch-only materialize");
4173        assert!(materialized
4174            .path
4175            .to_string_lossy()
4176            .contains(".yah/cache/derive/fetch/"));
4177        assert_eq!(std::fs::read(&materialized.path).unwrap(), body);
4178        assert!(
4179            materialized.discovered_fetch_hash.is_none(),
4180            "strict mode: pinned fetch.blake3 → no discovery"
4181        );
4182    }
4183
4184    #[tokio::test]
4185    async fn materialize_fetch_bootstrap_sentinel_accepts_actual_hash() {
4186        // fetch.blake3 ships as the zero sentinel; reconciler accepts the
4187        // computed hash and names the cache by the discovered value.
4188        let fx = Fixture::new(minio_slot());
4189        let body = b"discovered-upstream";
4190        let actual_hash = blake3_hex(body);
4191
4192        use axum::routing::get;
4193        let body_c = body.to_vec();
4194        let app = axum::Router::new().route(
4195            "/blob.bin",
4196            get(move || {
4197                let body = body_c.clone();
4198                async move { body }
4199            }),
4200        );
4201        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
4202        let addr = listener.local_addr().unwrap();
4203        tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
4204
4205        let fetch = FetchSource {
4206            url: format!("http://{addr}/blob.bin"),
4207            blake3: BlakeHash(ZERO_SENTINEL_HEX.to_string()),
4208            license: License::Mit,
4209        };
4210        let cache_dir = fx.workspace_root.join(".yah/cache/derive/fetch");
4211        let (path, discovered) = materialize_fetch(&fetch, &cache_dir).await.unwrap();
4212        assert!(path.exists(), "cache file must exist after bootstrap fetch");
4213        assert_eq!(
4214            discovered, actual_hash,
4215            "discovered hash = actual content hash"
4216        );
4217        assert!(
4218            path.file_name()
4219                .and_then(|s| s.to_str())
4220                .map(|s| s.starts_with(&actual_hash))
4221                .unwrap_or(false),
4222            "bootstrap mode names cache by discovered hash, got {path:?}",
4223        );
4224    }
4225
4226    #[tokio::test]
4227    async fn materialize_derive_fetch_bootstrap_surfaces_discovered_hash() {
4228        // End-to-end: AssetEntry with sentinel fetch.blake3 → materialize_asset
4229        // populates discovered_fetch_hash so sync_assets can surface it.
4230        let fx = Fixture::new(minio_slot());
4231        let body = b"upstream-bytes";
4232        let actual_hash = blake3_hex(body);
4233
4234        use axum::routing::get;
4235        let body_c = body.to_vec();
4236        let app = axum::Router::new().route(
4237            "/w.bin",
4238            get(move || {
4239                let body = body_c.clone();
4240                async move { body }
4241            }),
4242        );
4243        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
4244        let addr = listener.local_addr().unwrap();
4245        tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
4246
4247        let entry = AssetEntry {
4248            filename: "w.bin".to_string(),
4249            source: None,
4250            derive: Some(AssetDerive {
4251                fetch: FetchSource {
4252                    url: format!("http://{addr}/w.bin"),
4253                    blake3: BlakeHash(ZERO_SENTINEL_HEX.to_string()),
4254                    license: License::Mit,
4255                },
4256                transform: None,
4257                lock: None,
4258            }),
4259            // entry.blake3 still sentinel — upload-side verify in sync_assets
4260            // handles the output-hash discovery; here we only assert the
4261            // fetch-side discovery threads up correctly.
4262            blake3: BlakeHash(ZERO_SENTINEL_HEX.to_string()),
4263        };
4264        let materialized = materialize_asset(
4265            &entry,
4266            &fx.workspace_root,
4267            &fx.workload_dir(),
4268            test_executor().as_ref(),
4269        )
4270        .await
4271        .expect("bootstrap materialize");
4272        assert_eq!(
4273            materialized.discovered_fetch_hash.as_deref(),
4274            Some(actual_hash.as_str()),
4275            "fetch.blake3 sentinel → discovered hash returned",
4276        );
4277    }
4278
4279    /// W209/F4: discovered hashes write to `$YAH_OUTPUTS` with one line per
4280    /// (kind, filename) pair so the QED runner can route them through
4281    /// OutputMap into pipeline binds. Key shape encodes the filename so
4282    /// multi-asset workloads round-trip; the value is the raw blake3 hex.
4283    #[test]
4284    fn bootstrap_outputs_format_round_trip() {
4285        let dir = tempdir().unwrap();
4286        let outputs = dir.path().join("outputs");
4287        let bootstrapped = vec![
4288            BootstrappedHash {
4289                filename: "yah-desktop/whisper/distil-large-v3-coreml.tar.gz".to_string(),
4290                kind: BootstrapHashKind::Output,
4291                hash: "fb0afc9f3d966f5347c6dfd335adab12f1dc8ee6df18cf9e9ff90fe86f0416c0"
4292                    .to_string(),
4293            },
4294            BootstrappedHash {
4295                filename: "yah-desktop/whisper/distil-large-v3-coreml.tar.gz".to_string(),
4296                kind: BootstrapHashKind::Fetch,
4297                hash: "050ffe562134208781dc316181b146a725821fff005fb4ffb6de2a6ada334a9b"
4298                    .to_string(),
4299            },
4300        ];
4301        append_bootstrap_outputs(&outputs, &bootstrapped).expect("write outputs");
4302        let content = std::fs::read_to_string(&outputs).expect("read outputs");
4303        let lines: Vec<&str> = content.lines().collect();
4304        assert_eq!(lines.len(), 2, "one line per discovery");
4305        assert!(
4306            lines.contains(
4307                &"discovered_asset_blake3:yah-desktop/whisper/distil-large-v3-coreml.tar.gz=\
4308                  fb0afc9f3d966f5347c6dfd335adab12f1dc8ee6df18cf9e9ff90fe86f0416c0"
4309            ),
4310            "asset output key encodes filename: {content:?}",
4311        );
4312        assert!(
4313            lines.contains(
4314                &"discovered_fetch_blake3:yah-desktop/whisper/distil-large-v3-coreml.tar.gz=\
4315                  050ffe562134208781dc316181b146a725821fff005fb4ffb6de2a6ada334a9b"
4316            ),
4317            "fetch output key encodes filename: {content:?}",
4318        );
4319    }
4320
4321    /// Append, not truncate: a step that materializes N assets and emits
4322    /// other outputs in the same step must not clobber what's already in
4323    /// the file.
4324    #[test]
4325    fn bootstrap_outputs_append_preserves_prior_lines() {
4326        let dir = tempdir().unwrap();
4327        let outputs = dir.path().join("outputs");
4328        std::fs::write(&outputs, "prior_key=prior_value\n").unwrap();
4329        let bootstrapped = vec![BootstrappedHash {
4330            filename: "a.bin".to_string(),
4331            kind: BootstrapHashKind::Output,
4332            hash: HASH_64.to_string(),
4333        }];
4334        append_bootstrap_outputs(&outputs, &bootstrapped).expect("append");
4335        let content = std::fs::read_to_string(&outputs).expect("read");
4336        assert!(content.starts_with("prior_key=prior_value\n"), "prior kept");
4337        assert!(
4338            content.contains(&format!("discovered_asset_blake3:a.bin={HASH_64}")),
4339            "new line appended: {content:?}",
4340        );
4341    }
4342}