Skip to main content

cloud/reconciler/
mod.rs

1//! Reconciler abstraction — bring a workload up against a mirror's
2//! provider slots.
3//!
4//! A reconciler is the kind-specific code that knows how to deploy one
5//! workload kind (`mesofact-static`, `container`, future `almanac`, …) to a
6//! mirror. Selection: a [`ServiceComponent`](crate::ServiceComponent)'s
7//! `kind` field picks which reconciler runs; the reconciler then dispatches
8//! on the mirror's provider slot (e.g. `mesofact-static` →
9//! `providers.static` slot → `miniflare-native` inline or `cloudflare` ref).
10//!
11//! [`MesofactStaticReconciler`] serves all four tiers through one door — the
12//! compiled Worker bundle — and varies only the object store beneath it:
13//! the `yah-s3-fs` driver at dev ([`dev_door`]), MinIO at pond ([`pond`]),
14//! R2 at cloud and ha. Which one a mirror binds is declared in
15//! `[drivers.s3]`, never branched on (W265, R584-F4).
16//!
17//!
18//! @yah:ticket(R419-F2, "Implement CloudflareWorkerReconciler (kind=cloudflare-worker)")
19//! @yah:assignee(agent:claude)
20//! @yah:at(2026-06-03T08:02:49Z)
21//! @yah:status(review)
22//! @yah:parent(R419)
23//! @yah:handoff("Landed CloudflareWorkerReconciler in crates/yah/cloud/src/reconciler/cloudflare_worker.rs and re-exported it through reconciler/mod.rs + cloud/src/lib.rs. up() validates the registry slot (must be `use = \"cloudflare\"` with non-empty `zone` + `domain`), reads workload.toml's `[build]` + `[[bindings]]`, matches each binding against a sibling mirror slot by `binding` field, runs build, reads the bundled entrypoint (default dist/index.js), idempotently lists+creates R2 buckets, deploys via deploy_worker_script with WorkerBinding::R2Bucket entries (R419-F1 surface), then attaches the custom domain via the new upsert_worker_custom_domain method on CloudflareClient. Returns RunningWorkload::adopted with public_url=https://<domain>. All config validation runs BEFORE any CF API call (R330-B5 fail-fast discipline).")
24//! @yah:handoff("Added CloudflareClient::upsert_worker_custom_domain (cloudflare.rs) — separate from upsert_worker_route because Worker Routes are zone-scoped pattern matches and Custom Domains are account-scoped hostname attachments. List-first idempotency: skips PUT when the (hostname, service, zone_id) tuple is already bound.")
25//! @yah:verify("cargo check -p cloud --lib — clean")
26//! @yah:depends_on(R419-F1)
27//!
28//! @yah:ticket(R419-F3, "Register cloudflare-worker reconciler in CLI + desktop dispatch")
29//! @yah:assignee(agent:claude)
30//! @yah:at(2026-06-03T08:02:58Z)
31//! @yah:status(review)
32//! @yah:parent(R419)
33//! @yah:handoff("Added match arm `\"cloudflare-worker\" => CloudflareWorkerReconciler::new().up(ctx)` in both dispatchers: app/yah/cli/src/cloud.rs:reconcile_component and app/yah/desktop/src/mirror_run.rs's component.kind.as_str() match. Imported CloudflareWorkerReconciler at the top of each file. Updated the desktop file-level docstring to list the new kind. Pre-existing fallback arm still produces a clean error for unknown kinds.")
34//! @yah:verify("cargo check -p cloud --lib — clean")
35//! @yah:verify("cargo check -p yah --lib --bins — clean (warnings unchanged from baseline)")
36//! @yah:verify("cargo check -p desktop --lib — clean (warnings unchanged from baseline)")
37//! @yah:depends_on(R419-F2)
38//!
39//! @yah:ticket(R419-F4, "Regression tests: misconfig fail-fast for cloudflare-worker reconciler")
40//! @yah:assignee(agent:claude)
41//! @yah:at(2026-06-03T08:03:09Z)
42//! @yah:status(review)
43//! @yah:parent(R419)
44//! @yah:handoff("Three fail-fast tests in reconciler::cloudflare_worker::tests — Fixture builds an in-tempdir yah-cr-shaped workspace (writes workload.toml + .yah/infra/providers/cloudflare.toml). up_bails_on_registry_missing_domain (case 3) drops domain off the registry slot. up_bails_on_binding_name_drift (case 2) puts binding=\"STORAGE\" in the cache slot while workload.toml binds CACHE. up_bails_on_cache_slot_missing_bucket (case 1) keeps binding=\"CACHE\" but omits bucket. Each test asserts the error message names the offending field, the slot role, the service, and the env — no CF HTTP call is made because validation runs before any client construction.")
45//! @yah:verify("cargo test -p cloud --lib reconciler::cloudflare_worker — 3 passed")
46//! @yah:verify("cargo test -p cloud --lib — 246 passed (1 pre-existing failure cloud_init::tests::embedded_template_matches_workspace_canonical is unrelated R092-F2 template drift, not caused by R419)")
47//! @yah:depends_on(R419-F2)
48//!
49//! @yah:relay(R458, "Cloud reconciler for .yah/domains/*.toml — R2 custom-domain shape")
50//! @yah:at(2026-06-05T08:40:57Z)
51//! @yah:status(open)
52//! @yah:next("F1: implement ensure_r2_custom_domain (mirror of ensure_r2_bucket) + wire into yah cloud apply as a post-services pass. Scope: domains with cdn_bucket set and no [[routes]] (today: cdn-yah-dev.toml). Worker-routed shape (yah-dev, app-yah-dev) is a separate surface.")
53//! @arch:see(.yah/domains/cdn-yah-dev.toml)
54//!
55//! @yah:ticket(R458-F1, "ensure_r2_custom_domain (CF API) + apply-time orchestration")
56//! @yah:at(2026-06-05T08:41:07Z)
57//! @yah:status(review)
58//! @yah:assignee(agent:claude)
59//! @yah:parent(R458)
60//! @arch:see(.yah/domains/cdn-yah-dev.toml)
61//! @yah:next("R458 can be archived once F1 is signed off, unless we want to keep it open for the routed-domain (Worker) reconciler shape — that's a much larger surface (DNS + Worker route management + bundle deploy) than this F1's bucket-binding.")
62//! @yah:handoff("Live-verified end-to-end. cdn.yah.dev now resolves to CF anycast IPs (104.21.43.100, 172.67.178.4) and curl HTTP/2 200s against https://cdn.yah.dev/yah-desktop/whisper/distil-large-v3-q5_1.bin (content-length 584567555 = the q5_1 bytes R422-F11 published). Second apply is idempotent — the list-first path skips the POST when the binding is already present. Files touched: (1) crates/yah/cloud/src/provider/cloudflare.rs — new R2CustomDomain output type + CloudflareClient::list_r2_custom_domains and CloudflareClient::add_r2_custom_domain methods (GET / POST /accounts/{id}/r2/buckets/{bucket}/domains/custom). The POST body needs zone_id even though the endpoint is bucket-scoped (CF rejects with 'JSON not well formed' otherwise — caught live and added on the second iteration). Re-exported through provider/mod.rs + cloud/src/lib.rs. (2) crates/yah/cloud/src/reconciler/domain.rs — new module with ensure_r2_custom_domain(account_id, bucket, domain) mirroring static_asset::ensure_r2_bucket's list-first idempotency. Resolves the parent zone id via the existing CloudflareClient::zone_id_for_name and a parent_zone_name(domain) heuristic that takes the last two labels (correct for every yah-owned zone today; the doc comment names the longest-suffix-match upgrade path for future three-label-apex zones). 3 unit tests cover apex / subdomain / deeper-subdomain. (3) crates/yah/cloud/src/reconciler/mod.rs — declared `pub mod domain` + re-exported ensure_r2_custom_domain. (4) app/yah/cli/src/cloud.rs — new DomainOutcome enum + a post-services domain pass in handle_apply that walks cfg.domains, dispatches the R2-custom-domain shape (cdn_bucket set, no [[routes]]), and routes routed-shape domains (yah-dev, app-yah-dev) into a Skipped row labelled 'has [[routes]] — Worker-routed shape'. Gated on a cloudflare provider being declared (pond-only setups print 'skip domain pass: no cloudflare provider declared'). Originally also gated on empty --service filter; that gate dropped on review — domains are workspace-scoped and the operator wants them reconciled even when narrowing services. New print_domain_summary mirrors print_apply_summary's table + JSON output. Required CF token scopes (verified live): Workers R2 Storage: Edit + Zone: Read. The existing cloudflare-api-token slot carries both.\n\nLive verification command + transcript:\n\n  $ ./target/debug/yah cloud apply --env cloud --service yah-desktop\n  ==> yah-desktop/cloud: reconciling 2 component(s)\n      component desktop (kind=binary)\n      component whisper-models (kind=static-asset)\n  ==> domain cdn-yah-dev (cdn.yah.dev): ensuring R2 custom-domain binding on bucket yah-dev\n  apply summary (cloud):\n    yah-desktop  ok       2 component(s) reconciled\n  domain summary (cloud):\n    cdn-yah-dev  ok       R2 custom domain bound\n  $ dig +short cdn.yah.dev\n  104.21.43.100\n  172.67.178.4\n  $ curl -sI https://cdn.yah.dev/yah-desktop/whisper/distil-large-v3-q5_1.bin | head -4\n  HTTP/2 200\n  content-length: 584567555\n\nUnblocks R422-T13's client-side `cdn_fallback = \"https://cdn.yah.dev/yah-desktop/whisper/{blake3}\"` — the URL now actually resolves and serves the bytes.")
63//! @yah:verify("cargo test -p cloud --lib reconciler::domain --locked  # 3 pass")
64//! @yah:verify("cargo check --workspace --locked  # clean")
65//! @yah:verify("./target/debug/yah cloud apply --env cloud --service yah-desktop  # domain summary shows cdn-yah-dev=ok, yah-dev/app-yah-dev=skipped (routed)")
66//! @yah:verify("curl -sI https://cdn.yah.dev/yah-desktop/whisper/distil-large-v3-q5_1.bin  # HTTP/2 200, content-length 584567555")
67//!
68//! @yah:ticket(R870-B25, "The cloud-init template drift guard is vacuously green, and the canonical mirror.yml it should guard is 128 lines stale")
69//! @yah:status(review)
70//! @yah:at(2026-09-11T00:24:58Z)
71//! @yah:assignee(agent:bundle-anthropic-ashguard)
72//! @yah:parent(R870)
73//! @yah:severity(high)
74//! @yah:next("OPERATIONAL QUESTION THIS RAISES, worth answering separately from the code fix: which is the provisioning path actually used, the embedded template or the repo-root canonical one? If any node was provisioned from the canonical copy since R858-F17 landed, it has no turso-backup helpers and its durability-declaring workloads will refuse to deploy. `rendered_runcmd_entries_are_all_strings` is the gate that genuinely catches the colon-space footgun and IS green — this ticket is about the twin-file guard beside it, not that one.")
75//! @yah:verify("The drift guard must FAIL on today's tree (proving it now compares something real), then pass once the canonical `.yah/infra/cloud-init/mirror.yml` is brought level with the embedded template. Asserting only the post-fix green is the weak form — it passes for the same vacuous reason it does today.")
76//! @yah:gotcha("Found by independent verification of R870-F23 phase 2 (@Ashguard:dove, session:60d4f41f), confirming a claim from @Ashguard:blade's implementation pass. NOT caused by R870-F23 — pre-existing, and F23's own change is green and unaffected. Do not read this as a phase-2 regression.")
77//! @yah:next("THE DEFECT, read not inferred. `embedded_template_matches_workspace_canonical` exists to prove the mirror.yml compiled into the binary matches the canonical copy on disk. It resolves the workspace root by walking `CARGO_MANIFEST_DIR.ancestors()`, which lands on `oss/yubaba` — an independent Cargo workspace whose `.yah/` holds only a `.gitignore`. The canonical path therefore does not exist, the test takes its `canonical_path.exists()` bootstrap branch, and asserts NOTHING. It has been green for that reason, not because the files agree.")
78//! @yah:next("WHAT THE GUARD IS MISSING, measured: the repo-root `.yah/infra/cloud-init/mirror.yml` is 128 diff-lines behind `oss/yubaba/crates/cloud/templates/mirror.yml`. It is missing the ENTIRE R858-F17 turso-backup block, and the YAML-quoting fix R870-F23 landed in the embedded copy (the two runcmd entries whose bare `: ` made cloud-init parse them as a Mapping and skip them). THE FIX IS TWO PARTS AND THE ORDER MATTERS: first make the guard non-vacuous — resolve the canonical path against the REPO root rather than the enclosing cargo workspace, or fail loudly when it cannot be found, so the bootstrap branch can no longer swallow a real absence. Then reconcile the two files. Doing only the second leaves the guard still asleep for the next drift.")
79//! @yah:handoff("GUARD MADE NON-VACUOUS FIRST, then the files reconciled, and the intermediate RED was observed. New CanonicalHome enum + locate_canonical_home() in oss/yubaba/crates/cloud/src/cloud_init.rs (non-test code, since it is a real property of the layout): the ancestor walk is keyed on `.yah/infra/` instead of `.yah/`. The monorepo root has it; oss/yubaba, whose own .yah/ holds nothing but a .gitignore, does not, and neither does the standalone export where the repo root genuinely IS oss/yubaba. embedded_template_matches_workspace_canonical now matches on that: Monorepo(root) means the canonical file is REQUIRED (missing/unreadable panics with a message naming the path and saying provisioning reads it); StandaloneExport is the only branch allowed to skip. The old `if canonical_path.exists()` bootstrap branch that swallowed a real absence is gone.")
80//! @yah:handoff("THE RED, run after step 1 and before step 2 exactly as the ticket asked: `cargo test -p yah-cloud --lib cloud_init::tests` from oss/yubaba gave 29 passed / 1 FAILED, the single failure being embedded_template_matches_workspace_canonical asserting on /Users/leif/ss/yah/.yah/infra/cloud-init/mirror.yml with the full 128-line diff in its message. That is the proof it now compares something real; the post-fix green on its own would have been indistinguishable from today's vacuous green.")
81//! @yah:handoff("THE RECONCILE WAS NOT A ONE-WAY COPY, and this is the part worth reading. The canonical copy carried a DELIBERATE CORRECTION the embedded template never received: commit 45e39f0a (2026-08-03) changed `systemctl enable --now yubaba.slice` to `systemctl start yubaba.slice`, plus the comment explaining why - app/yah/cli/resources/yubaba.slice has no [Install] section, and systemctl enable on such a unit fails. Blindly overwriting canonical with embedded (which the ticket text reads like) would have regressed that fix into the file provisioning actually ships. So the fix went BOTH ways: the slice correction was applied to oss/yubaba/crates/cloud/templates/mirror.yml first, then that file was copied over .yah/infra/cloud-init/mirror.yml. The two are now byte-identical at 237 lines; canonical gained the whole R858-F17 turso-backup block, the R870-F23 YAML double-quoting of the two colon-space runcmd entries, the litestream 0.3.13 install, and the headscale 0.23.0 + litestream-headscale.service pre-stage.")
82//! @yah:handoff("WHICH TEMPLATE PROVISIONING ACTUALLY USES - THE CANONICAL ON-DISK ONE, read not guessed. cloud_init::load_template() returns the on-disk <workspace_root>/.yah/infra/cloud-init/mirror.yml whenever it exists and falls back to the include_str! embedded copy only when it does not. provision::build_request() (oss/yubaba/crates/cloud/src/provision.rs:88) is its sole non-test caller, and its sole caller in turn is handle_provision at app/yah/cli/src/cloud.rs:14524, whose workspace_root is the camp repo root - which has the file. So the embedded template is the FALLBACK, not the live path, and the stale copy is the one every `yah cloud machine provision` has been shipping. load_template's doc comment now says so.")
83//! @yah:handoff("FLEET CONSEQUENCE, dated from git so the window is bounded. NO node was touched - code, templates and tests only, per scope. Canonical was last updated 2026-08-03 (45e39f0a). The headscale 0.23.0 pre-stage, the litestream install and litestream-headscale.service staging landed in the EMBEDDED copy only on 2026-09-05 (4bed91fe), so any node provisioned between 2026-09-05 and today silently received none of them - that is the real damage this drift did. The turso-backup block landed in the embedded copy only TODAY (a8f0d501, 2026-09-10), so no node has ever received the durability helpers from either file; that gap is fleet-wide and predates this ticket rather than being caused by the drift. Remediating existing nodes is deliberately NOT done here.")
84//! @yah:handoff("DISCOVERED WORK, fixed in this pass: the SECOND stale twin is .yah/infra/cloud-init/stand-up-yubaba.sh, mirror.yml's documented SSH transcription for LAN nodes (W257 step 6). It installed yubaba/kamaji/units from the release tarball but never the turso-backup helpers or the KAMAJI_HYDRATE_HELPER/KAMAJI_TAIL_HELPER drop-in, so every LAN node stood up by it would refuse any durability-declaring workload. Added a present-checked, non-fatal block in the script's own idiom ($SUDO install from $D, tee for the drop-in, a WARNING to stderr when the tarball predates R858-F17), using Environment= in a drop-in rather than ExecStart= flags for the reason kamaji.service's own comment gives. Verified against scripts/publish-yubaba-release.sh:326-327, which hard-asserts both binaries at exactly $STAGE_NAME/turso-backup-{hydrate,tail} - the path the script looks in. bash -n clean.")
85//! @yah:verify("cargo test -p yah-cloud --lib, run from oss/yubaba (NOT the repo root - yah-cloud is not a root workspace member and needs dev-dependencies): 1166 passed / 0 failed / 4 ignored, against the 1163/0/4 baseline. The +3 are exactly the three tests added here; nothing regressed.")
86//! @yah:verify("INTERMEDIATE RED, the verification this ticket actually asked for: after step 1 and before step 2, cargo test -p yah-cloud --lib cloud_init::tests = 29 passed / 1 failed, the single failure being embedded_template_matches_workspace_canonical naming .yah/infra/cloud-init/mirror.yml. Green afterwards.")
87//! @yah:verify("Three new tests pin the guard against going vacuous again. locate_canonical_home_finds_monorepo_root_not_the_inner_workspace and locate_canonical_home_reports_standalone_export build both tree shapes in tempdirs (the monorepo one reproduces the inner oss/yubaba/.yah/.gitignore that caused the original miss). drift_guard_cannot_go_vacuous_in_the_monorepo keys off a signal INDEPENDENT of the .yah/infra/ marker - `git subtree split --prefix=oss/yubaba` strips that prefix, so a CARGO_MANIFEST_DIR still ending in oss/yubaba/crates/cloud proves we are in the monorepo - and panics if the guard resolved StandaloneExport there.")
88//! @yah:verify("diff -u between the two mirror.yml copies is empty; both are 237 lines.")
89//! @yah:verify("bash -n .yah/infra/cloud-init/stand-up-yubaba.sh clean.")
90//! @yah:verify("rendered_runcmd_entries_are_all_strings left untouched and still green, as instructed - a different test and a different gate.")
91//! @yah:gotcha("NOT FIXED, and deliberately left as an operator call rather than decided silently: stand-up-yubaba.sh still lacks the litestream 0.3.13 install and the headscale 0.23.0 / litestream-headscale.service pre-stage that mirror.yml now carries. Those exist because R858-T4 made EVERY node a coordinator candidate, and whether the LAN/appliance class (us-west-01x, mostly no-voter) belongs in that candidate set is a fleet-topology decision, not a transcription gap. The durability helpers WERE added because their consequence is unconditional - kamaji refuses the workload on any node - while these two only matter if the node can ever own the ingress role.")
92//! @yah:gotcha("Uncommitted: this camp's git policy is `defer`, so the four changed files (oss/yubaba/crates/cloud/src/cloud_init.rs, both mirror.yml copies, .yah/infra/cloud-init/stand-up-yubaba.sh) are in the working tree for the camp's git sweep. No commit SHA to cite.")
93//! @yah:verify("INDEPENDENTLY VERIFIED BY A SECOND COURIER (session:67d56cd3) who did not implement it. GUARD IS GENUINELY HONEST, confirmed by reading the branches: `locate_canonical_home` (cloud_init.rs:208) walks ancestors for `.yah/infra/`, and the Monorepo branch (cloud_init.rs:1013-1029) `read_to_string`-panics on a missing or unreadable canonical and `assert_eq`s full trimmed content against `DEFAULT_TEMPLATE` — so BOTH a deleted canonical AND a revert to its pre-fix state now fail hard. No skip branch survives. THE STANDALONE-EXPORT DISCRIMINATOR HOLDS: `.yah/infra/` exists only at the monorepo root and nowhere under `oss/` (oss/yubaba/.yah holds one tracked .gitignore), so the real subtree-split export takes the StandaloneExport branch legitimately; the only route to that branch from inside the monorepo is deleting `.yah/infra/`, which `drift_guard_cannot_go_vacuous_in_the_monorepo` catches off an INDEPENDENT path-suffix signal. Two residual holes, both benign in DIRECTION: a standalone clone placed under some `.yah/infra/` ancestor false-FAILS rather than skipping, and renaming the oss/yubaba path would disarm only the meta-guard, not the drift guard.")
94//! @yah:verify("THE RECONCILIATION WAS NOT A ONE-WAY COPY, and that was checked rather than taken on trust — a blind embedded-to-canonical copy would have silently destroyed someone else's fix. Both copies are byte-identical at 237 lines, sha256 6bc46731e030bae8a1970bcb06cf3132323eb454ab72496d6ef9b2c6177585af. At HEAD the CANONICAL carried the 45e39f0a fix (`systemctl start yubaba.slice`, since the slice has no [Install] section) while the EMBEDDED still had `enable --now`; the working diff applies that exact hunk embedded-ward (+5/-3) against the canonical's +109/-1, proving the flow went both directions. No `enable --now yubaba.slice` remains in either file. The canonical's single deleted line was only the `chmod 0644` superseded by the version adding litestream-headscale.service. R858-F17's turso-backup block and R870-F23's two double-quoted runcmd entries are present in both, necessarily so given byte-identity. THE INTERMEDIATE RED WAS STRUCTURALLY NECESSARY, not stage-managed: the Monorepo branch compares entire trimmed contents and the two HEAD copies differed by 112 insertions / 6 deletions, so the assert could not have passed. `cloud_init::tests` is 30 tests, making the reported 29-pass/1-fail arithmetically consistent. (The \"128-line diff\" figure is 118 changed lines by numstat — cosmetic, not a defect.) COUNTS: `cargo test -p yah-cloud --lib` from oss/yubaba = 1166 passed / 0 failed / 4 ignored, the +3 over baseline being exactly the three new tests.")
95//! @yah:verify("THE FLEET FACT, CONFIRMED WITH ONE DATE CORRECTED — this is the operationally consequential part and the correction matters. PROVISIONING READS THE CANONICAL ON-DISK COPY, not the embedded one: provision.rs:88 is `cloud_init::load_template(workspace_root)`, and `load_template` (cloud_init.rs:225-232) PREFERS `.yah/infra/cloud-init/mirror.yml`, falling back to the embedded copy only when absent; sole non-test caller chain is cloud.rs:14523 with the camp repo root. So every node was provisioned from the file that was 128 lines stale. CONSEQUENCE 1, date corrected: the headscale/litestream pre-stage entered the EMBEDDED copy on 2026-09-04 in b20a1e09 — NOT 2026-09-05, which was 4bed91fe merely refining it — and the canonical never carried it, so every node provisioned since 2026-09-04 missed that pre-stage. CONSEQUENCE 2 HOLDS AS STATED: no node has ever had the turso-backup helpers. `turso-backup-hydrate` first entered the embedded template today in a8f0d501 and the canonical only in this change; grep finds no other install path (only publish-yubaba-release.sh, kamaji's consumers, both mirror.yml copies, stand-up-yubaba.sh), and `.yah/infra/machines/us-west-001.toml:42` independently records \"NO prod node has them today\". NOT VERIFIED, stated rather than smoothed: no node was touched, so the actual SET of nodes provisioned since 2026-09-04 is unconfirmed.")
96//! @yah:handoff("THIRD TWIN CLOSED — .yah/infra/cloud-init/stand-up-yubaba.sh is now both LEVEL and GUARDED. GUARD SHAPE CHOSEN: the assertion test, not generate-from-one-source, and the reason is that the twin is not a pure transcription. The script deliberately diverges from mirror.yml in ways that are CORRECT and load-bearing (enable+restart instead of `enable --now`, which is the only reason it can call itself idempotent — see its own :162-176 comment measured on us-west-013/014; write-if-absent journald ceiling so a Pi's tighter 200M bound wins; cluster-KEK install; loopback bind for no-mesh nodes; status block reading /health rather than the on-disk --version). A generator would therefore have to model the divergences, which is the whole difficulty, and mirror.yml is itself a static `{{ }}`-substituted file rather than a Rust-rendered one — so single-sourcing would mean either making mirror.yml generated (large) or parsing YAML to emit bash (fragile). Not contained; see @yah:next for what it would actually take.")
97//! @yah:handoff("THE GUARD DERIVES ITS PINS FROM THE TEMPLATE rather than restating them, which is what stops it becoming the next thing that rots. `stand_up_script_carries_the_templates_install_steps` (oss/yubaba/crates/cloud/src/cloud_init.rs, next to the mirror.yml drift guard) reads the script via the new `paths::stand_up_script`, reusing `locate_canonical_home` so StandaloneExport is the same single permitted skip. Three assertion families: (1) STAND_UP_TWIN_ANCHORS, an explicit list of artifacts a node must end up carrying, each asserted against DEFAULT_TEMPLATE AS WELL so a stale anchor goes red on the template side instead of over-constraining the script; (2) every 64-char lowercase-hex sha256 pin extracted from DEFAULT_TEMPLATE must appear in the script (the four litestream/headscale amd64+arm64 checksums; `{{YAH_YUBABA_SHA256}}` is not hex so it is not picked up, and the extractor asserts it found >=4 so a broken extractor cannot make the guard vacuous); (3) every `https://github.com/OWNER/REPO/releases/download/TAG/` prefix in the template must appear in the script — owner/repo/tag pinned, arch-templated asset filename left free. Net effect: bumping litestream or headscale in mirror.yml alone now turns this red.")
98//! @yah:handoff("THREE REDS DEMONSTRATED, not one, because the anchor half fires first and would have masked the other two. (a) Test written BEFORE the script fix: `cargo test -p yah-cloud --lib cloud_init::tests::stand_up_script` = 0 passed / 1 FAILED, naming `no litestream-headscale.service`. (b) With the script level, one hex char mutated in the script's headscale arm64 HS_SHA: FAILED, naming the missing pin 99fa9b29...e9fe. (c) The litestream URL tag moved to v0.3.14 while the template stays v0.3.13: FAILED, naming the unfetched .../download/v0.3.13/. Both mutations were reverted by hand and re-verified. SCRIPT CONTENT ADDED, in the script's own idiom: present-checked non-fatal `$SUDO install -m0644 \"$D/litestream-headscale.service\"` alongside the other three units (verified it really ships in the tarball — publish-yubaba-release.sh:287 stages it from $RESOURCES); an arch-cased litestream 0.3.13 fetch+sha256+`tar -C /usr/local/bin`; an arch-cased headscale 0.23.0 fetch+sha256+`install -m0755` into /var/lib/yah-cloud/headscale/headscale. Every failure path is a stderr WARNING naming the operational consequence, never an exit — a node without these cannot hold the ingress role, which is not a failed stand-up.")
99//! @yah:handoff("DELIBERATE CHOICE WORTH REVIEWING: the two downloads are UNCONDITIONAL, not present-checked-skip. A re-run of this script is the documented upgrade path (W257 §8), and `if [ -x /usr/local/bin/litestream ]; then skip` would recreate exactly the stale-binary trap the script's own `restart`-not-`--now` comment was written for after us-west-013/014. Cost is re-downloading ~10MB litestream + ~51MB headscale on every re-run; both are sha-pinned so the repeat is idempotent, just not free. A version-checked skip was rejected because it would depend on `headscale version` / `litestream version` output formats I did not verify. Rationale is in the script's comments. VERIFIED: `bash -n` clean; both mirror.yml copies UNTOUCHED and still byte-identical at sha256 6bc46731e030bae8a1970bcb06cf3132323eb454ab72496d6ef9b2c6177585af; `cargo test -p yah-cloud --lib` from oss/yubaba = 1167 passed / 0 failed / 4 ignored, exactly +1 over the 1166 baseline and that +1 is the new test. rustfmt --check clean on cloud_init.rs (paths.rs has ONE pre-existing unformatted hunk in `infra_source_cache_dir`, not mine, left alone). NO FLEET CONTACT of any kind — script, test and one paths.rs helper only. Uncommitted: camp git policy is `defer`, so no SHA to cite.")
100//! @yah:next("WHAT GENERATE-FROM-ONE-SOURCE WOULD ACTUALLY TAKE, having rejected it as out of scope here. The contained 20% is the third-party PINS: litestream 0.3.13 + 2 sha256s and headscale 0.23.0 + 2 sha256s now live in five places (both mirror.yml copies, stand-up-yubaba.sh, and for headscale also `cloud::mesh::HEADSCALE_VERSION` and `yubaba::DEFAULT_HEADSCALE_VERSION`, which mirror.yml's own comment admits are kept in lockstep by convention with no dependency edge). Making those Rust constants the one source and rendering them into mirror.yml as `{{ LITESTREAM_BLOCK }}` / `{{ HEADSCALE_BLOCK }}` would collapse five to one, but it requires touching BOTH mirror.yml copies in lockstep and a matching emitter for the script. The remaining 80% — the install STEPS — is not contained: the script's correct divergences (idempotent enable+restart, write-if-absent ceilings, KEK, loopback bind) mean a generator must model divergence, and the script's real destination is the `yah cloud machine bootstrap` command W242 Phase 1 already plans, which would emit both from one Rust model. That is the right home for this, and it is a relay, not a hunk.")
101//! @yah:gotcha("SUPERSEDES the earlier gotcha beginning \"NOT FIXED, and deliberately left as an operator call\" — that entry is now STALE and should be read as history, not state. The litestream 0.3.13 install and the headscale 0.23.0 / litestream-headscale.service pre-stage ARE now in stand-up-yubaba.sh, added on the relay leader's explicit instruction in this follow-on pass. The fleet-topology question that entry deferred has therefore been ANSWERED IN THE AFFIRMATIVE BY DEFAULT: every LAN/appliance node stood up by this script from here on is provisioned as a coordinator candidate (headscale binary staged at /var/lib/yah-cloud/headscale/headscale, replication unit laid down, neither started — leader.rs still decides who runs it). If that is NOT wanted for the us-west-01x class, the lever is an opt-out env guard in the script, not reverting it, because reverting now goes red against `stand_up_script_carries_the_templates_install_steps`. Cost of the affirmative answer is bounded and staging-only: ~61MB fetched per run and two staged-but-inert artifacts.")
102//! @yah:verify("THIRD-TWIN GUARD: `cargo test -p yah-cloud --lib cloud_init::` from oss/yubaba = 31 passed / 0 failed (30 before, +1 = stand_up_script_carries_the_templates_install_steps). Full lib suite = 1167 passed / 0 failed / 4 ignored vs the 1166/0/4 baseline. bash -n .yah/infra/cloud-init/stand-up-yubaba.sh clean. Both mirror.yml copies still sha256 6bc46731e030bae8a1970bcb06cf3132323eb454ab72496d6ef9b2c6177585af — the guard on THEM was not disturbed and is still green.")
103//! @yah:handoff("ACCEPTED BY THE RELAY LEADER (@Ashguard:hydra, session:39386823). Both halves landed in the required order — guard made honest FIRST, then the two mirror.yml copies reconciled — plus a third twin (stand-up-yubaba.sh) closed that the fix itself exposed. Implemented by @Ashguard:polaris (session:9d2e59a3), independently verified by session:67d56cd3, third-twin follow-on by session:039a0f82. The operationally consequential finding is in the verify entries: provisioning reads the CANONICAL on-disk template, and that was the copy which had drifted.")
104//! @yah:verify("COLUMN MOVED VIA `yah board move R870-B25 review` AFTER BOTH `board.review` VERBS REFUSED — recorded because the next agent will hit it too and the error message actively misleads. The MCP verb and `yah board review` both fail with \"ticket 'R870-B25' not found — it may have been archived or may not exist in this camp\", which is false: the CLI's fresh scan saw it fine, anchored at oss/yubaba/crates/cloud/src/reconciler/mod.rs:67. ROOT CAUSE, grounded not guessed: `arch.review_ticket` is a daemon-only gated write with NO in-process fallback — already filed as R606-T3 and stated verbatim at app/yah/cli/src/camp.rs:975, \"Board reads/updates fall back in-process; review/move/etc do not — that asymmetry is the bug\" — and the daemon resolves transition targets from its IN-MEMORY store rather than from disk (crates/yah/camp-service/src/service.rs:4199-4216, which emits that exact string). So \"exists on disk\" and \"the daemon can transition it\" are independent facts. `yah board move` turns out to be on the falls-back side despite that note, which is why it works. Contributing: the daemon is version-skewed and degraded (0.8.36+74874f3e-dirty vs CLI 0.8.37+e896d28a-dirty) and refused read probes with EAGAIN, the signature R606-S2 pinned to the 500ms fast-path RPC floor under load. ALSO CONFIRMED, since it was the other candidate explanation: `oss/yubaba` is NOT a subcamp — only cheers, mesofact, turso-backup and xlb carry `.yah/camp.toml` — so no `--path` is needed and the filing location was never the problem.")
105
106use std::collections::{BTreeMap, VecDeque};
107use std::path::{Path, PathBuf};
108use std::sync::Arc;
109
110use anyhow::{Context, Result};
111use async_trait::async_trait;
112use serde::{Deserialize, Serialize};
113use tokio::sync::{oneshot, Mutex as AsyncMutex};
114
115use workload_spec::{NamespaceId, TenantId};
116
117use crate::{GitSource, MirrorConfig, MirrorProviderSlot, ServiceComponent, ServiceConfig};
118
119pub mod bundle_store;
120pub(crate) mod cf_creds;
121pub mod cloudflare_worker;
122pub mod container;
123pub mod derive_cache_prune;
124pub mod dev_door;
125pub mod domain;
126pub mod headscale;
127pub mod ingress;
128pub mod ingress_verify;
129// R918-F5 — the local-process reconciler is Unix-only *by design*, not by
130// accident of which syscalls it happened to reach for. Its contract is
131// "replace the predecessor": `kill(pid, 0)` liveness plus a SIGTERM-grace-
132// SIGKILL ladder. A build with that ladder removed would still spawn, and
133// would silently double-spawn instead of reaping — exactly the failure that
134// looks inert in a diff and strands processes on a node. So the module is
135// gated whole rather than half-ported. A non-unix `yah` is a client (it talks
136// to a camp and submits builds); it supervises no workloads, so it needs no
137// local-process reconciler. Reversing this means implementing
138// OpenProcess/TerminateProcess here, at the point someone actually wants a
139// Windows fleet node — see .yah/docs/working/W352-windows-and-macos-build-targets.md.
140#[cfg(unix)]
141pub mod device;
142#[cfg(unix)]
143pub mod local_process;
144pub mod mesofact_bundle;
145pub mod mesofact_static;
146// `pub(crate)` rather than private: `sanitize_ident` is the crate's ONE mesh-ident
147// normalizer, and R870-F23's `inner_door::component_workload_ident` has to fold
148// its derived ident the same way `local_process` folds its own. A second copy
149// would be a second normalizer that can drift.
150pub(crate) mod native_support;
151pub mod pg_driver;
152pub mod s3_driver;
153pub mod smtp_driver;
154pub mod pond;
155pub mod pond_door;
156pub mod pond_publish;
157pub mod publish_beacon;
158pub mod r2_publish;
159pub mod service_discovery;
160pub mod static_asset;
161pub mod static_asset_prune;
162pub mod sync_status;
163
164#[cfg(test)]
165mod lowering_golden;
166
167pub use bundle_store::{publish_bundle_to_r2, PublishReport as BundlePublishReport};
168pub use cloudflare_worker::CloudflareWorkerReconciler;
169pub use container::{ContainerOptions, ContainerReconciler};
170pub use derive_cache_prune::{
171    collect_live_derive_hashes, compute_derive_cache_candidates, execute_derive_cache_prune,
172    DeriveCacheLiveHashes, DerivePruneCandidate,
173};
174pub use domain::{
175    deploy_domain_passway, diff_apex_records, ensure_passway_apex, ensure_r2_custom_domain,
176    list_live_apex_records, plan_domain_passway, plan_passway_apex, public_origins, ApexRecordDiff,
177    DomainPasswayPlan, LiveApexRecord, PasswayApexOutcome, PasswayOrigin,
178};
179pub use headscale::{
180    DeclaredHeadscale, DeclaredPolicy, DeclaredPreauthKey, HeadscaleReconciler,
181    WORKLOAD_KIND as HEADSCALE_WORKLOAD_KIND,
182};
183pub use ingress::{
184    collate_front_doors, declared as ingress_declared, ensure_tunnel_ingress, machine_mesh_addrs,
185    plan_ingress, publish_tunnel_ingress, resolve_ingress_candidates, resolve_ingress_placements,
186    Collation, IngressPlan, IngressRule, NodeFrontDoor, PlannedEdge, TunnelIngressOutcome,
187};
188pub use ingress_verify::{
189    apply_public_path, resolve_upstreams_reporting, verify_collation, BeaconFetch, DialOutcome,
190    EndpointCheck, PublicReadings, RuleResolution, RuleResolutions, RuleVerdict,
191    UndeclaredDoorProbe, VerifyFinding, VerifyReport, undeclared_door_probes,
192};
193#[cfg(unix)]
194pub use local_process::LocalProcessReconciler;
195#[cfg(unix)]
196pub use device::DeviceReconciler;
197pub use mesofact_bundle::{
198    resolve_bundle_machines, BundleSlot, MesofactBundleReconciler, RevalidateSlot,
199    SLOT_ROLE as BUNDLE_SLOT_ROLE,
200};
201pub use dev_door::{sync_dev_door, up_dev_door, DEFAULT_DEV_DOOR_PORT};
202pub use mesofact_static::MesofactStaticReconciler;
203pub use pond::{PondOptions, PondState};
204pub use pond_door::{
205    door_env, door_state_dir, ensure_pond_cert_as, plan_pond_door, pond_hostname,
206    resolve_passway_binary, spawn_pond_door, CertPair, PondDoorPlan, DEFAULT_DOOR_PORT, POND_TLD,
207};
208// R918-F5 — `is_root` reads an effective uid; `ensure_pond_cert` is the wrapper
209// that folds it in. Both unix-only. `ensure_pond_cert_as` above is the portable
210// half, with the privilege decision injected.
211#[cfg(unix)]
212pub use pond_door::{ensure_pond_cert, is_root};
213pub use pond_publish::{derive_minio_key, publish_to_pond, PondPublishReport};
214pub use r2_publish::{
215    publish_to_r2, R2PublishReport, R2PurgeOpts, R2_ACCESS_KEY_ENV, R2_ACCESS_KEY_SLOT,
216    R2_SECRET_KEY_ENV, R2_SECRET_KEY_SLOT,
217};
218pub use service_discovery::{
219    DiscoveredRecord, RecordVisibility, ServiceRecordFanout, UnknownReason,
220};
221pub use static_asset::StaticAssetReconciler;
222pub use static_asset_prune::{
223    compute_live_set, compute_prune_candidates, execute_prune, load_service_and_mirror,
224    PruneCandidate, PruneOutcome, PruneReport,
225};
226pub use sync_status::{
227    compute_cell, compute_service, new_sync_id, summarize, CellStatus, DriftEntry, HealthState,
228    MirrorObservation, Runtime, ServiceStatus, StatusSummary, SyncHistoryEntry, SyncOutcome,
229    SyncState, WireContainerStatus,
230};
231
232// ─── Log buffer ─────────────────────────────────────────────────────────────
233
234const LOG_CAP: usize = 500;
235
236#[derive(Debug, Default)]
237struct LogRing {
238    lines: VecDeque<String>,
239    /// Monotonically increasing total lines ever pushed (never decrements).
240    total: usize,
241}
242
243/// Bounded ring buffer for child-process stdout/stderr (R263-F3).
244/// Shared between the reader tasks and the Tauri `mirror_run_logs` command.
245#[derive(Debug, Clone, Default)]
246pub struct LogBuffer(Arc<AsyncMutex<LogRing>>);
247
248impl LogBuffer {
249    pub fn new() -> Self {
250        Self::default()
251    }
252
253    /// Append a line; drops the oldest entry when over capacity.
254    pub async fn push(&self, line: String) {
255        let mut ring = self.0.lock().await;
256        ring.total += 1;
257        ring.lines.push_back(line);
258        if ring.lines.len() > LOG_CAP {
259            ring.lines.pop_front();
260        }
261    }
262
263    /// Return lines not yet seen by the caller.
264    ///
265    /// `since` is the `total` cursor from the previous call (0 = nothing
266    /// seen yet). Returns `(new_lines, new_cursor)`. Pass `new_cursor` back
267    /// on the next call to receive only incremental output.
268    pub async fn since(&self, since: usize) -> (Vec<String>, usize) {
269        let ring = self.0.lock().await;
270        let oldest = ring.total.saturating_sub(ring.lines.len());
271        let skip = since.saturating_sub(oldest);
272        let new_lines: Vec<String> = ring.lines.iter().skip(skip).cloned().collect();
273        (new_lines, ring.total)
274    }
275
276    /// Current write cursor — the `total` [`Self::since`] would hand back
277    /// right now if nothing more were pushed. Lets a producer mark a boundary
278    /// (e.g. "everything before this point was the build phase") for a later
279    /// reader to seek past without re-reading lines it doesn't want.
280    pub async fn cursor(&self) -> usize {
281        self.0.lock().await.total
282    }
283}
284
285/// Shared cell for one phase-boundary cursor a reconciler can publish
286/// mid-`up()`, so a poller sees a multi-phase bring-up's internal transition
287/// (e.g. build → run) before the whole call returns — the same problem
288/// [`LogBuffer`] solves for output, for a single position instead of a ring.
289/// `None` until the reconciler reaches that phase; a caller registers one
290/// before calling `up()` to observe it live, same pattern as
291/// [`LogBuffer::clone`]-and-hand-in.
292#[derive(Debug, Clone, Default)]
293pub struct PhaseCursor(Arc<AsyncMutex<Option<usize>>>);
294
295impl PhaseCursor {
296    pub fn new() -> Self {
297        Self::default()
298    }
299
300    pub async fn set(&self, cursor: usize) {
301        *self.0.lock().await = Some(cursor);
302    }
303
304    pub async fn get(&self) -> Option<usize> {
305        *self.0.lock().await
306    }
307}
308
309/// The `(tenant, namespace)` a bring-up is scoped to (W206). Reconcilers that
310/// touch a credentialed provider resolve it at this scope
311/// ([`CfProvider::resolve_scoped`](super::reconciler::cf_creds)) so a namespace's
312/// Cloudflare zone/account/keystore slots come from its own scope rather than the
313/// workspace-global defaults. Defaults to the singleton `(default, default)`,
314/// which collapses every scoped lookup back to the historical global slots — so
315/// single-namespace deployments are unaffected.
316#[derive(Debug, Clone)]
317pub struct ProviderScope {
318    pub tenant: TenantId,
319    pub namespace: NamespaceId,
320}
321
322impl ProviderScope {
323    /// The degenerate single-tenant / single-namespace scope. Scoped provider
324    /// lookups made against it resolve to the pre-W206 global keystore slots.
325    pub fn singleton() -> Self {
326        Self {
327            tenant: TenantId::singleton(),
328            namespace: NamespaceId::singleton(),
329        }
330    }
331}
332
333impl Default for ProviderScope {
334    fn default() -> Self {
335        Self::singleton()
336    }
337}
338
339/// Inputs a reconciler sees for one bring-up.
340pub struct ReconcileCtx<'a> {
341    /// Workspace root (parent of `.yah/`). Used to resolve relative paths
342    /// on the component.
343    pub workspace_root: &'a Path,
344    /// Service that owns the component.
345    pub service: &'a ServiceConfig,
346    /// Component being brought up.
347    pub component: &'a ServiceComponent,
348    /// Mirror manifest the bring-up targets.
349    pub mirror: &'a MirrorConfig,
350    /// Environment name (file stem of `mirrors/<env>.toml`).
351    pub env: &'a str,
352    /// `(tenant, namespace)` this bring-up is scoped to (W206). Credentialed
353    /// providers resolve at this scope; defaults to [`ProviderScope::singleton`].
354    pub scope: ProviderScope,
355}
356
357impl<'a> ReconcileCtx<'a> {
358    /// Absolute path to the component's workload directory (the parent of
359    /// `workload.toml`).
360    ///
361    /// In-tree components resolve to `<workspace_root>/<path>`. For
362    /// `git`-sourced components (R561-F1, "BYO git") this points into the
363    /// local clone — `<source_cache>/<subdir>/<path>` — which is empty until
364    /// [`materialize`](Self::materialize) runs (approach A: clone-at-reconcile,
365    /// so config load + validation stay offline).
366    pub fn workload_dir(&self) -> PathBuf {
367        match &self.component.git {
368            None => self.workspace_root.join(&self.component.path),
369            Some(git) => {
370                let mut dir = self.source_cache_dir();
371                if let Some(subdir) = &git.subdir {
372                    dir = dir.join(subdir);
373                }
374                dir.join(&self.component.path)
375            }
376        }
377    }
378
379    /// Root of the local clone for a `git`-sourced component:
380    /// `<workspace_root>/.yah/infra/state/sources/<service>/<component_id>`.
381    fn source_cache_dir(&self) -> PathBuf {
382        self.workspace_root
383            .join(".yah/infra/state/sources")
384            .join(&self.service.name)
385            .join(&self.component.id)
386    }
387
388    /// Ensure a `git`-sourced component's code is present locally before build
389    /// (R561-F1, approach A). No-op for in-tree components. Idempotent: clones
390    /// on the first call, fetches + re-checks-out the pinned ref thereafter.
391    ///
392    /// Reconcilers MUST call this at the top of [`up`](Reconciler::up) before
393    /// reading [`workload_dir`](Self::workload_dir) for a remote component.
394    pub async fn materialize(&self) -> Result<()> {
395        let Some(git) = &self.component.git else {
396            return Ok(());
397        };
398        materialize_git_source(git, &self.source_cache_dir())
399            .await
400            .with_context(|| {
401                format!(
402                    "materializing git source {}@{} for {}/{}",
403                    git.repo, git.r#ref, self.service.name, self.component.id
404                )
405            })
406    }
407
408    /// Read `<workload_dir>/workload.toml` and extract just the `kind`
409    /// discriminator.
410    ///
411    /// Why this and not the strongly-typed [`workload_spec::Workload`]
412    /// parse: reconcilers only need the kind to
413    /// dispatch; per-kind tooling (e.g. `mesofact-dev`'s
414    /// `WatchOptions::from_workload`) does its own parsing for the
415    /// build/out_dir fields it cares about.
416    pub fn workload_kind(&self) -> Result<String> {
417        workload_kind(&self.workload_dir())
418    }
419
420    /// Look up a provider slot by role (e.g. `"static"`, `"compute"`).
421    ///
422    /// Tries the component-qualified key first (`"<role>:<component id>"`)
423    /// before falling back to the bare role. A mirror role is normally
424    /// service-wide — one `providers.static` slot serves every static-kind
425    /// component — but a service can declare more than one component with
426    /// the same role (e.g. two `mesofact-static`/`mesofact-spa` components
427    /// sharing one mirror), and those need distinct ports to ever both come
428    /// up. The qualified key is how a mirror opts a specific component out
429    /// of sharing the bare-role slot:
430    ///
431    /// ```toml
432    /// [providers."static:site"]
433    /// kind = "miniflare-native"
434    /// port = 4331
435    ///
436    /// [providers."static:app"]
437    /// kind = "miniflare-native"
438    /// port = 4332
439    /// ```
440    ///
441    /// A mirror with only the bare role (the common, single-component case)
442    /// is unaffected — the qualified lookup misses and falls through.
443    pub fn slot(&self, role: &str) -> Option<&'a MirrorProviderSlot> {
444        let qualified = format!("{role}:{}", self.component.id);
445        self.mirror
446            .providers
447            .get(qualified.as_str())
448            .or_else(|| self.mirror.providers.get(role))
449    }
450
451    /// This environment's `[build.<component id>]` override, if the mirror
452    /// declares one that changes anything (R905).
453    ///
454    /// Keyed by component id rather than by role: a build belongs to the
455    /// project that declares it, and two components sharing a role still build
456    /// with two different commands.
457    pub fn build_override(&self) -> Option<&'a crate::config::MirrorBuildOverride> {
458        self.mirror.build_override(&self.component.id)
459    }
460}
461
462/// An explicit teardown hook for a workload this process did not spawn as a
463/// child — see [`RunningWorkload::with_teardown`] (R714-B1).
464///
465/// Boxed rather than a generic parameter because [`RunningWorkload`] is stored
466/// in heterogeneous collections (the desktop's mirror registry) and cannot
467/// carry a type parameter.
468/// `Sync` on the boxed closure is load-bearing, not belt-and-braces: without
469/// it `RunningWorkload` stops being `Sync`, so `&RunningWorkload` stops being
470/// `Send`, and every desktop `#[tauri::command]` that holds one across an
471/// `.await` (`mirror_run_logs` iterates the registry's handles) fails to
472/// compile with "future cannot be sent between threads safely". The captures a
473/// teardown needs — a container name and a workspace root — are `Sync` anyway.
474type TeardownFuture = std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send>>;
475type TeardownFn = Box<dyn FnOnce() -> TeardownFuture + Send + Sync + 'static>;
476
477/// Newtype so [`RunningWorkload`] can keep its `#[derive(Debug)]` — a boxed
478/// closure is not `Debug`.
479struct Teardown(TeardownFn);
480
481impl std::fmt::Debug for Teardown {
482    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
483        f.write_str("Teardown(<fn>)")
484    }
485}
486
487/// Handle to a workload that's been brought up. Owns the lifecycle: drop
488/// or call [`RunningWorkload::shutdown`] to take it back down.
489#[derive(Debug)]
490pub struct RunningWorkload {
491    /// Workload kind that was reconciled (e.g. `"mesofact-static"`).
492    pub kind: String,
493    /// Slot role this workload occupies on the mirror (e.g. `"static"`).
494    pub slot: String,
495    /// Local URL the workload exposes, when applicable. `None` for
496    /// workloads that publish to a non-local artifact store (e.g. R2).
497    pub dev_url: Option<String>,
498    /// Notional URL of the deployed artifact in the production case
499    /// (e.g. `https://yah.dev`). `None` until Cloudflare/R2 wiring lands.
500    pub public_url: Option<String>,
501    /// Secondary local UI console, if the workload exposes one (e.g. MinIO
502    /// console on a pond tier). Surfaced as a separate chip in the Services
503    /// matrix next to `dev_url`.
504    pub console_url: Option<String>,
505    /// Ring buffer for stdout/stderr from the workload's child process.
506    /// `None` for workloads that don't capture stdio (e.g. container-backed).
507    pub log_buffer: Option<LogBuffer>,
508
509    /// R546-B12: human-readable lines the reconciler wants the operator to see
510    /// on THIS run — what it actually did, not what is configured. A clean
511    /// static-asset reconcile used to print nothing at all, so a successful
512    /// publish and a successful no-op were indistinguishable at the apply
513    /// surface, and the one view that would have disambiguated them (`yah cloud
514    /// status`) was itself blind.
515    ///
516    /// Deliberately NOT on [`RunningWorkloadSummary`]: that type is serialized
517    /// across the process boundary to the desktop UI, and this is per-run
518    /// console output, not state the UI should cache.
519    pub notes: Vec<String>,
520
521    /// Sender that signals the supervisor task to tear down. Closing the
522    /// channel (drop) is equivalent to sending — supervisor exits on
523    /// channel close.
524    shutdown: Option<oneshot::Sender<()>>,
525    /// Joinable task that owns any child process and reaps it on signal.
526    supervisor: Option<tokio::task::JoinHandle<Result<()>>>,
527    /// R714-B1: teardown for a workload that runs OUTSIDE this process, so
528    /// there is no child to reap and no supervisor to signal.
529    ///
530    /// Deliberately run from [`RunningWorkload::shutdown`] only, never from
531    /// `Drop`. The two are different intents and conflating them breaks both
532    /// directions: a container the desktop started is meant to outlive the
533    /// desktop (the next launch re-adopts it with `adopt_only`), so quitting
534    /// the app must not `docker rm -f` it; but the ■ button IS an explicit
535    /// stop, and must.
536    teardown: Option<Teardown>,
537
538    /// R715-F4: where to ask this workload for a live status document, when
539    /// it declared a process-control channel (W315). `None` for a workload
540    /// with no channel — polling is then simply skipped, not an error.
541    ///
542    /// The desktop holds `RunningWorkload` in-process (it links this crate
543    /// directly, unlike the ad-hoc `run.spawn` path which lives behind the
544    /// camp daemon's socket), so a live poll is [`crate::proc_control::fetch_status`]
545    /// called straight against this endpoint — no RPC hop, no handle registry.
546    control: Option<crate::proc_control::ControlEndpoint>,
547
548    /// [`LogBuffer`] cursor marking the end of the build phase, for a
549    /// component that was compiled before being spawned (`local-process`
550    /// with `cargo_package` set). Lines before this index in `log_buffer` are
551    /// `cargo build` output; lines at or after are the spawned process's own
552    /// stdout/stderr. `None` when nothing was built (no `cargo_package`, or a
553    /// reconciler that predates this field).
554    pub build_log_end: Option<usize>,
555}
556
557impl RunningWorkload {
558    /// Create a handle for a workload that's already running externally
559    /// (e.g. embedded in yah-camp). No subprocess is owned; shutdown is a
560    /// no-op so the caller can call `shutdown()` uniformly.
561    pub fn adopted(
562        kind: impl Into<String>,
563        slot: impl Into<String>,
564        dev_url: Option<String>,
565    ) -> Self {
566        Self {
567            kind: kind.into(),
568            slot: slot.into(),
569            dev_url,
570            public_url: None,
571            console_url: None,
572            log_buffer: None,
573            notes: Vec::new(),
574            shutdown: None,
575            supervisor: None,
576            teardown: None,
577            control: None,
578            build_log_end: None,
579        }
580    }
581
582    /// Attach an explicit teardown to a handle for an externally-running
583    /// workload (R714-B1).
584    ///
585    /// `adopted()` alone gives a handle whose `shutdown()` is a documented
586    /// no-op. That is right for workloads another process owns and will keep
587    /// re-asserting (pond containers under camp's yubaba), and wrong for ones
588    /// this process started and nobody else will ever stop — for those the ■
589    /// button reported success while the container kept running.
590    ///
591    /// `teardown` runs on `shutdown()` and NOT on `Drop`; see [`Self::teardown`].
592    pub fn with_teardown<F, Fut>(mut self, teardown: F) -> Self
593    where
594        F: FnOnce() -> Fut + Send + Sync + 'static,
595        Fut: std::future::Future<Output = Result<()>> + Send + 'static,
596    {
597        self.teardown = Some(Teardown(Box::new(move || Box::pin(teardown()))));
598        self
599    }
600
601    /// Whether this handle carries a real teardown — i.e. whether a failed
602    /// [`Self::shutdown`] means the workload is **still running**.
603    ///
604    /// The stop path needs this to decide between an operator-facing error and
605    /// a log line, and R875-B1 is why it is a property of the handle rather
606    /// than a list of kinds at the call site. That list started as
607    /// `kind == "container"` when R714-B1 gave containers a teardown, and was
608    /// silently wrong the moment a second kind grew one: an adopted
609    /// mesofact-dev whose teardown failed reported a successful stop with the
610    /// server still serving. A handle knows whether it owns the workload; a
611    /// string comparison at the call site only knows what it was last taught.
612    pub fn owns_teardown(&self) -> bool {
613        self.teardown.is_some()
614    }
615
616    /// Attach per-run operator-facing lines (R546-B12). See [`Self::notes`].
617    pub fn with_notes(mut self, notes: Vec<String>) -> Self {
618        self.notes = notes;
619        self
620    }
621
622    /// Record where to poll this workload's process-control channel, when it
623    /// declared one (R715-F4). A no-op (leaves `control: None`) when the
624    /// argument is `None` — callers can pass the reconciler's resolved
625    /// endpoint straight through without an `if let`.
626    pub fn with_control(mut self, control: Option<crate::proc_control::ControlEndpoint>) -> Self {
627        self.control = control;
628        self
629    }
630
631    /// Ask this workload's process-control channel for a live status
632    /// document (R715-F4). `None` when no channel was declared, or when the
633    /// endpoint didn't answer — both are "nothing new to show", not errors;
634    /// see [`crate::proc_control::fetch_status`] for why a poll failure isn't
635    /// itself meaningful.
636    pub async fn poll_control(&self) -> Option<crate::proc_control::ProcStatus> {
637        let endpoint = self.control.as_ref()?;
638        crate::proc_control::fetch_status(endpoint).await.ok()
639    }
640
641    /// Whether the supervisor task that owns this workload's child process is
642    /// still running. `spawn_native_log_supervisor` (native_support.rs) exits
643    /// as soon as `NativeRuntime::get_workload` reports a terminal state —
644    /// which itself comes from a real `child.wait()` in kamaji's native
645    /// backend, so this catches a crash (segfault, panic, `SIGKILL`) the same
646    /// way it catches a clean exit, not just an unresponsive process.
647    ///
648    /// A workload with no supervisor (`RunningWorkload::adopted` — runs
649    /// outside this process, e.g. a pond container camp re-asserts) has
650    /// nothing to check here and reports alive unconditionally; its liveness
651    /// is whatever tracks it, not this handle.
652    pub fn is_alive(&self) -> bool {
653        self.supervisor.as_ref().is_none_or(|h| !h.is_finished())
654    }
655
656    /// Record where the build phase ends in `log_buffer`, for a component
657    /// that was compiled before being spawned. See [`Self::build_log_end`].
658    pub fn with_build_log_end(mut self, cursor: Option<usize>) -> Self {
659        self.build_log_end = cursor;
660        self
661    }
662
663    /// Set the public URL for a published workload (e.g. `"https://yah.dev"`).
664    pub fn with_public_url(mut self, url: impl Into<String>) -> Self {
665        self.public_url = Some(url.into());
666        self
667    }
668
669    /// Set the console URL for a workload that exposes a secondary local UI
670    /// (e.g. MinIO console on a pond tier).
671    pub fn with_console_url(mut self, url: impl Into<String>) -> Self {
672        self.console_url = Some(url.into());
673        self
674    }
675
676    /// Gracefully tear down: signal the supervisor, await its exit, then run
677    /// any explicit teardown hook.
678    ///
679    /// The hook runs LAST and its error propagates. A stop that could not tear
680    /// the workload down must surface as an error, never as a silent success —
681    /// that silence is the whole of R714-B1.
682    pub async fn shutdown(mut self) -> Result<()> {
683        if let Some(tx) = self.shutdown.take() {
684            let _ = tx.send(());
685        }
686        if let Some(handle) = self.supervisor.take() {
687            handle
688                .await
689                .context("joining workload supervisor")?
690                .context("workload supervisor")?;
691        }
692        if let Some(Teardown(hook)) = self.teardown.take() {
693            hook().await.context("workload teardown")?;
694        }
695        Ok(())
696    }
697}
698
699impl Drop for RunningWorkload {
700    fn drop(&mut self) {
701        // Best-effort signal. The supervisor task is detached and will
702        // reap its child when it observes the closed channel.
703        if let Some(tx) = self.shutdown.take() {
704            let _ = tx.send(());
705        }
706    }
707}
708
709/// Bring one workload up. Each impl handles one [`ServiceComponent::kind`].
710#[async_trait]
711pub trait Reconciler: Send + Sync {
712    /// Workload kind this reconciler handles (matches `ServiceComponent.kind`).
713    fn kind(&self) -> &'static str;
714
715    /// Bring the workload up. Returns a handle whose lifecycle is tied to
716    /// the mirror being up.
717    async fn up(&self, ctx: ReconcileCtx<'_>) -> Result<RunningWorkload>;
718}
719
720/// Read `<workload_dir>/workload.toml` and return just the `kind` field.
721/// See [`ReconcileCtx::workload_kind`] for why we don't deserialize through
722/// the strong types yet.
723pub fn workload_kind(workload_dir: &Path) -> Result<String> {
724    let path = workload_dir.join("workload.toml");
725    let src =
726        std::fs::read_to_string(&path).with_context(|| format!("reading {}", path.display()))?;
727    let value: toml::Value =
728        toml::from_str(&src).with_context(|| format!("parsing {}", path.display()))?;
729    let kind = value
730        .get("kind")
731        .and_then(|v| v.as_str())
732        .with_context(|| format!("{}: missing `kind` field", path.display()))?;
733    Ok(kind.to_string())
734}
735
736/// Shallow-clone (or update) a [`GitSource`] into `dir` (R561-F1). Idempotent:
737/// clones when `dir/.git` is absent, otherwise fetches the pinned ref and
738/// force-checks-it-out. Uses the system `git` so it inherits the operator's
739/// credential helpers / SSH agent — no in-process git library.
740///
741/// `pub` (R615-T3 / W274): `yah infra sync` reuses this verbatim for
742/// `InfraSourceKind::Git` sources rather than a second shallow-clone-or-pull
743/// implementation — same "one git-source shape, reused" discipline R615-F1
744/// already applied to the type.
745pub async fn materialize_git_source(git: &GitSource, dir: &Path) -> Result<()> {
746    use tokio::process::Command;
747
748    async fn run_git(args: &[&std::ffi::OsStr]) -> Result<()> {
749        let out = Command::new("git")
750            .args(args)
751            .output()
752            .await
753            .context("spawning git")?;
754        if !out.status.success() {
755            anyhow::bail!(
756                "git {} failed: {}",
757                args.iter()
758                    .map(|a| a.to_string_lossy())
759                    .collect::<Vec<_>>()
760                    .join(" "),
761                String::from_utf8_lossy(&out.stderr).trim()
762            );
763        }
764        Ok(())
765    }
766
767    use std::ffi::OsStr;
768    let dir_os = dir.as_os_str();
769    let r#ref = git.r#ref.as_str();
770
771    if dir.join(".git").is_dir() {
772        // Existing checkout — update to the pinned ref.
773        run_git(&[
774            OsStr::new("-C"),
775            dir_os,
776            OsStr::new("fetch"),
777            OsStr::new("--depth"),
778            OsStr::new("1"),
779            OsStr::new("origin"),
780            OsStr::new(r#ref),
781        ])
782        .await?;
783        run_git(&[
784            OsStr::new("-C"),
785            dir_os,
786            OsStr::new("checkout"),
787            OsStr::new("--force"),
788            OsStr::new("FETCH_HEAD"),
789        ])
790        .await?;
791    } else {
792        if let Some(parent) = dir.parent() {
793            tokio::fs::create_dir_all(parent)
794                .await
795                .with_context(|| format!("creating {}", parent.display()))?;
796        }
797        // `--branch` accepts a branch or tag. Pinning to a bare commit SHA is a
798        // follow-up (needs clone-then-fetch); the common case is a branch/tag.
799        run_git(&[
800            OsStr::new("clone"),
801            OsStr::new("--depth"),
802            OsStr::new("1"),
803            OsStr::new("--branch"),
804            OsStr::new(r#ref),
805            OsStr::new(git.repo.as_str()),
806            dir_os,
807        ])
808        .await?;
809    }
810    Ok(())
811}
812
813/// Build a `RunningWorkload` from the pieces a reconciler produces.
814pub(crate) fn into_running(
815    kind: impl Into<String>,
816    slot: impl Into<String>,
817    dev_url: Option<String>,
818    public_url: Option<String>,
819    log_buffer: Option<LogBuffer>,
820    shutdown: oneshot::Sender<()>,
821    supervisor: tokio::task::JoinHandle<Result<()>>,
822) -> RunningWorkload {
823    RunningWorkload {
824        kind: kind.into(),
825        slot: slot.into(),
826        dev_url,
827        public_url,
828        console_url: None,
829        log_buffer,
830        notes: Vec::new(),
831        shutdown: Some(shutdown),
832        supervisor: Some(supervisor),
833        teardown: None,
834        control: None,
835        build_log_end: None,
836    }
837}
838
839/// Wait for a TCP port to start accepting connections. Returns `true` if
840/// the port came up within `timeout`, `false` otherwise. Useful for
841/// reconcilers that spawn a server and need to know when it's reachable
842/// before reporting success.
843///
844/// R918-F5 — unix-only: its only caller is [`local_process`], gated above.
845#[cfg(unix)]
846pub(crate) async fn wait_for_port(
847    addr: std::net::SocketAddr,
848    timeout: std::time::Duration,
849) -> bool {
850    let deadline = tokio::time::Instant::now() + timeout;
851    loop {
852        if tokio::net::TcpStream::connect(addr).await.is_ok() {
853            return true;
854        }
855        if tokio::time::Instant::now() >= deadline {
856            return false;
857        }
858        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
859    }
860}
861
862// `wait_for_http_ready` lived here pre-R374-F3 to back the MinIO health
863// probe in pond's bring-up path. That logic moved to
864// `local_driver::pond_minio::wait_for_http_ready` so yubaba + cloud share
865// it. The mesofact-static reconciler arm uses [`wait_for_port`] for
866// dev-tier port readiness; nothing else needs an HTTP-level probe today.
867
868/// Pluck a `u16` out of a [`MirrorProviderSlot`]'s inline `fields` map.
869/// Returns `None` if the key is absent or out of range.
870pub(crate) fn slot_field_u16(fields: &BTreeMap<String, toml::Value>, key: &str) -> Option<u16> {
871    fields
872        .get(key)
873        .and_then(|v| v.as_integer())
874        .and_then(|n| u16::try_from(n).ok())
875}
876
877/// Serializable summary of a running workload — what the desktop / CLI
878/// hands to the UI. Subset of [`RunningWorkload`] that's safe to cross
879/// process boundaries.
880#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
881pub struct RunningWorkloadSummary {
882    pub kind: String,
883    pub slot: String,
884    pub dev_url: Option<String>,
885    pub public_url: Option<String>,
886    pub console_url: Option<String>,
887    /// Operator-facing lines from [`RunningWorkload::notes`] (R546-B12).
888    ///
889    /// These used to stop here — the summary dropped them, so every note a
890    /// reconciler attached was written into a struct nobody read. They matter
891    /// most for a workload with no `dev_url` to click (R715-T2): the notes are
892    /// then the only structured thing the Run tab has to show about it.
893    #[serde(default, skip_serializing_if = "Vec::is_empty")]
894    pub notes: Vec<String>,
895    /// See [`RunningWorkload::build_log_end`]. Lets a log-tail consumer skip
896    /// straight to run-phase output, or show everything from 0 when the
897    /// operator wants the build log too (e.g. after a compile failure).
898    #[serde(default, skip_serializing_if = "Option::is_none")]
899    pub build_log_end: Option<usize>,
900}
901
902impl From<&RunningWorkload> for RunningWorkloadSummary {
903    fn from(r: &RunningWorkload) -> Self {
904        Self {
905            kind: r.kind.clone(),
906            slot: r.slot.clone(),
907            dev_url: r.dev_url.clone(),
908            public_url: r.public_url.clone(),
909            console_url: r.console_url.clone(),
910            notes: r.notes.clone(),
911            build_log_end: r.build_log_end,
912        }
913    }
914}
915
916#[cfg(test)]
917mod source_seam_tests {
918    //! R561-F1 — the BYO-git source seam: path resolution + materialization.
919    use super::*;
920    use std::collections::BTreeMap;
921
922    fn component(git: Option<GitSource>) -> ServiceComponent {
923        ServiceComponent {
924            mount: None,
925            id: "site".into(),
926            kind: "mesofact-static".into(),
927            path: "site".into(),
928            role: "static".into(),
929            publishes: None,
930            wave: 0,
931            git,
932            deploy: Default::default(),
933        }
934    }
935
936    fn service(comp: ServiceComponent) -> ServiceConfig {
937        ServiceConfig {
938            schema_version: 1,
939            name: "scrabcake".into(),
940            address: crate::config::ServiceAddress::front_door("scrabcake.example"),
941            description: None,
942            components: vec![comp],
943            db: crate::DbCatalog::default(),
944        }
945    }
946
947    fn mirror() -> MirrorConfig {
948        MirrorConfig {
949            schema_version: 1,
950            shape: crate::MirrorShape::Local,
951            providers: BTreeMap::new(),
952            ingress: Default::default(),
953            ingress_machines: Vec::new(),
954            drivers: Default::default(),
955            asset_aliases: BTreeMap::new(),
956            build: Default::default(),
957        }
958    }
959
960    fn ctx<'a>(ws: &'a Path, svc: &'a ServiceConfig, mir: &'a MirrorConfig) -> ReconcileCtx<'a> {
961        ReconcileCtx {
962            workspace_root: ws,
963            service: svc,
964            component: &svc.components[0],
965            mirror: mir,
966            env: "dev",
967            scope: ProviderScope::singleton(),
968        }
969    }
970
971    #[test]
972    fn workload_dir_in_tree_joins_workspace_root() {
973        let svc = service(component(None));
974        let mir = mirror();
975        assert_eq!(
976            ctx(Path::new("/ws"), &svc, &mir).workload_dir(),
977            Path::new("/ws/site")
978        );
979    }
980
981    #[test]
982    fn workload_dir_git_resolves_into_source_cache_with_subdir() {
983        let git = GitSource {
984            repo: "https://example.com/r.git".into(),
985            r#ref: "main".into(),
986            subdir: Some("apps".into()),
987        };
988        let svc = service(component(Some(git)));
989        let mir = mirror();
990        assert_eq!(
991            ctx(Path::new("/ws"), &svc, &mir).workload_dir(),
992            Path::new("/ws/.yah/infra/state/sources/scrabcake/site/apps/site")
993        );
994    }
995
996    #[tokio::test]
997    async fn materialize_is_noop_for_in_tree_component() {
998        let svc = service(component(None));
999        let mir = mirror();
1000        // No git source → Ok, and nothing is written under the workspace.
1001        ctx(Path::new("/nonexistent-ws"), &svc, &mir)
1002            .materialize()
1003            .await
1004            .unwrap();
1005    }
1006
1007    #[tokio::test]
1008    async fn materialize_clones_git_source_offline() {
1009        fn git(args: &[&str], cwd: &Path) {
1010            let out = std::process::Command::new("git")
1011                .args(args)
1012                .current_dir(cwd)
1013                .env("GIT_AUTHOR_NAME", "t")
1014                .env("GIT_AUTHOR_EMAIL", "t@t")
1015                .env("GIT_COMMITTER_NAME", "t")
1016                .env("GIT_COMMITTER_EMAIL", "t@t")
1017                .output()
1018                .unwrap();
1019            assert!(
1020                out.status.success(),
1021                "git {args:?}: {}",
1022                String::from_utf8_lossy(&out.stderr)
1023            );
1024        }
1025
1026        let tmp = tempfile::tempdir().unwrap();
1027        let src = tmp.path().join("src-repo");
1028        std::fs::create_dir_all(&src).unwrap();
1029        git(&["init", "-b", "main"], &src);
1030        std::fs::write(src.join("hello.txt"), "hi").unwrap();
1031        git(&["add", "."], &src);
1032        git(&["commit", "-m", "init"], &src);
1033
1034        let source = GitSource {
1035            repo: format!("file://{}", src.display()),
1036            r#ref: "main".into(),
1037            subdir: None,
1038        };
1039        let dest = tmp.path().join("cache");
1040
1041        // First call clones.
1042        materialize_git_source(&source, &dest).await.unwrap();
1043        assert!(dest.join("hello.txt").is_file());
1044
1045        // Second call takes the update path and stays green (idempotent).
1046        materialize_git_source(&source, &dest).await.unwrap();
1047        assert!(dest.join("hello.txt").is_file());
1048    }
1049}
1050
1051#[cfg(test)]
1052mod teardown_tests {
1053    //! R714-B1 — the explicit-teardown contract on [`RunningWorkload`].
1054    //!
1055    //! These are about WHEN the hook runs, not what it does. The bug being
1056    //! fixed was a `shutdown()` that reported success having done nothing, and
1057    //! the trap in fixing it is a `Drop` that tears down a container which is
1058    //! supposed to survive the process.
1059    use super::*;
1060    use std::sync::atomic::{AtomicUsize, Ordering};
1061
1062    fn counting() -> (RunningWorkload, Arc<AtomicUsize>) {
1063        let hits = Arc::new(AtomicUsize::new(0));
1064        let seen = hits.clone();
1065        let w = RunningWorkload::adopted("container", "compute", None)
1066            .with_teardown(move || {
1067                let seen = seen.clone();
1068                async move {
1069                    seen.fetch_add(1, Ordering::SeqCst);
1070                    Ok(())
1071                }
1072            });
1073        (w, hits)
1074    }
1075
1076    /// Attaching a teardown must not cost `RunningWorkload` its auto traits.
1077    /// The desktop stores these handles in a shared registry and its Tauri
1078    /// commands iterate them across `.await` points, which needs `Sync` — a
1079    /// hook that is `Send` but not `Sync` takes it away here and surfaces two
1080    /// crates over as "future cannot be sent between threads safely", with a
1081    /// span pointing at `mirror_run_logs` rather than at this file.
1082    #[test]
1083    fn a_teardown_does_not_cost_the_handle_send_or_sync() {
1084        fn assert_send_sync<T: Send + Sync>() {}
1085        assert_send_sync::<RunningWorkload>();
1086    }
1087
1088    #[tokio::test]
1089    async fn shutdown_runs_the_teardown() {
1090        let (w, hits) = counting();
1091        w.shutdown().await.unwrap();
1092        assert_eq!(hits.load(Ordering::SeqCst), 1);
1093    }
1094
1095    #[tokio::test]
1096    async fn dropping_the_handle_does_NOT_run_the_teardown() {
1097        // A container this process started is meant to outlive it — the next
1098        // launch re-adopts it. Quitting the app must not `docker rm -f` it.
1099        let (w, hits) = counting();
1100        drop(w);
1101        // Yield so a stray spawned task would have had a chance to run.
1102        tokio::task::yield_now().await;
1103        assert_eq!(hits.load(Ordering::SeqCst), 0);
1104    }
1105
1106    #[tokio::test]
1107    async fn a_failing_teardown_makes_shutdown_fail() {
1108        // The whole of R714-B1: a stop that could not tear the workload down
1109        // must not report success.
1110        let w = RunningWorkload::adopted("container", "compute", None)
1111            .with_teardown(|| async { anyhow::bail!("docker stop refused") });
1112        let err = w.shutdown().await.unwrap_err();
1113        assert!(format!("{err:#}").contains("docker stop refused"), "{err:#}");
1114    }
1115
1116    #[tokio::test]
1117    async fn an_adopted_handle_without_a_teardown_still_shuts_down_cleanly() {
1118        // Pond containers are owned by camp's yubaba and stop through
1119        // `workload.stop`; their no-op shutdown is correct and must stay.
1120        let w = RunningWorkload::adopted("mesofact-static", "static", None);
1121        w.shutdown().await.unwrap();
1122    }
1123
1124    /// R875-B1. The stop path decides "still running, tell the operator" vs
1125    /// "untidy supervisor, log it" from this, so it has to answer for the
1126    /// handle in front of it rather than for the kind string it carries — the
1127    /// two disagreed for every adopted mesofact-dev.
1128    #[test]
1129    fn owns_teardown_tracks_the_hook_not_the_kind() {
1130        let (owned, _) = counting();
1131        assert!(owned.owns_teardown());
1132        assert!(
1133            !RunningWorkload::adopted("container", "compute", None).owns_teardown(),
1134            "a bare adopted handle owns nothing, whatever its kind says"
1135        );
1136        assert!(
1137            RunningWorkload::adopted("mesofact-static", "static", None)
1138                .with_teardown(|| async { Ok(()) })
1139                .owns_teardown(),
1140            "a non-container kind with a teardown owns its workload"
1141        );
1142    }
1143
1144    #[tokio::test]
1145    async fn the_teardown_runs_after_the_supervisor_is_joined() {
1146        // Ordering matters for a workload that has both: reap the child first,
1147        // then remove the container it was talking to.
1148        let order = Arc::new(AsyncMutex::new(Vec::<&'static str>::new()));
1149        let (tx, rx) = oneshot::channel::<()>();
1150        let sup_order = order.clone();
1151        let supervisor = tokio::spawn(async move {
1152            let _ = rx.await;
1153            sup_order.lock().await.push("supervisor");
1154            Ok(())
1155        });
1156        let hook_order = order.clone();
1157        let w = into_running("container", "compute", None, None, None, tx, supervisor)
1158            .with_teardown(move || {
1159                let hook_order = hook_order.clone();
1160                async move {
1161                    hook_order.lock().await.push("teardown");
1162                    Ok(())
1163                }
1164            });
1165
1166        w.shutdown().await.unwrap();
1167        assert_eq!(*order.lock().await, vec!["supervisor", "teardown"]);
1168    }
1169}