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 → `local-static` inline or `cloudflare` ref).
10//!
11//! T3 ships [`MesofactStaticReconciler`] with the `local-static` path
12//! wired (spawn `mesofact-dev` as a child process). The Cloudflare path is
13//! a stub — the production reconciler lands once `mesofact-publisher`
14//! integration is on the roadmap.
15//!
16//!
17//! @yah:ticket(R419-F2, "Implement CloudflareWorkerReconciler (kind=cloudflare-worker)")
18//! @yah:assignee(agent:claude)
19//! @yah:at(2026-06-03T08:02:49Z)
20//! @yah:status(review)
21//! @yah:parent(R419)
22//! @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).")
23//! @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.")
24//! @yah:verify("cargo check -p cloud --lib — clean")
25//! @yah:depends_on(R419-F1)
26//!
27//! @yah:ticket(R419-F3, "Register cloudflare-worker reconciler in CLI + desktop dispatch")
28//! @yah:assignee(agent:claude)
29//! @yah:at(2026-06-03T08:02:58Z)
30//! @yah:status(review)
31//! @yah:parent(R419)
32//! @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.")
33//! @yah:verify("cargo check -p cloud --lib — clean")
34//! @yah:verify("cargo check -p yah --lib --bins — clean (warnings unchanged from baseline)")
35//! @yah:verify("cargo check -p desktop --lib — clean (warnings unchanged from baseline)")
36//! @yah:depends_on(R419-F2)
37//!
38//! @yah:ticket(R419-F4, "Regression tests: misconfig fail-fast for cloudflare-worker reconciler")
39//! @yah:assignee(agent:claude)
40//! @yah:at(2026-06-03T08:03:09Z)
41//! @yah:status(review)
42//! @yah:parent(R419)
43//! @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.")
44//! @yah:verify("cargo test -p cloud --lib reconciler::cloudflare_worker — 3 passed")
45//! @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)")
46//! @yah:depends_on(R419-F2)
47//!
48//! @yah:relay(R458, "Cloud reconciler for .yah/domains/*.toml — R2 custom-domain shape")
49//! @yah:at(2026-06-05T08:40:57Z)
50//! @yah:status(open)
51//! @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.")
52//! @arch:see(.yah/domains/cdn-yah-dev.toml)
53//!
54//! @yah:ticket(R458-F1, "ensure_r2_custom_domain (CF API) + apply-time orchestration")
55//! @yah:at(2026-06-05T08:41:07Z)
56//! @yah:status(review)
57//! @yah:assignee(agent:claude)
58//! @yah:parent(R458)
59//! @arch:see(.yah/domains/cdn-yah-dev.toml)
60//! @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.")
61//! @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.")
62//! @yah:verify("cargo test -p cloud --lib reconciler::domain --locked  # 3 pass")
63//! @yah:verify("cargo check --workspace --locked  # clean")
64//! @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)")
65//! @yah:verify("curl -sI https://cdn.yah.dev/yah-desktop/whisper/distil-large-v3-q5_1.bin  # HTTP/2 200, content-length 584567555")
66//!
67//! @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")
68//! @yah:at(2026-09-11T00:24:58Z)
69//! @yah:status(open)
70//! @yah:assignee(agent:bundle-anthropic-ashguard)
71//! @yah:parent(R870)
72//! @yah:severity(high)
73//! @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.")
74//! @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.")
75//! @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.")
76//! @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.")
77//! @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.")
78
79use std::collections::{BTreeMap, VecDeque};
80use std::path::{Path, PathBuf};
81use std::sync::Arc;
82
83use anyhow::{Context, Result};
84use async_trait::async_trait;
85use serde::{Deserialize, Serialize};
86use tokio::sync::{oneshot, Mutex as AsyncMutex};
87
88use workload_spec::{NamespaceId, TenantId};
89
90use crate::{GitSource, MirrorConfig, MirrorProviderSlot, ServiceComponent, ServiceConfig};
91
92pub mod bundle_store;
93pub(crate) mod cf_creds;
94pub mod cloudflare_worker;
95pub mod container;
96pub mod derive_cache_prune;
97pub mod domain;
98pub mod headscale;
99pub mod ingress;
100pub mod ingress_verify;
101pub mod local_process;
102pub mod mesofact_bundle;
103pub mod mesofact_static;
104// `pub(crate)` rather than private: `sanitize_ident` is the crate's ONE mesh-ident
105// normalizer, and R870-F23's `inner_door::component_workload_ident` has to fold
106// its derived ident the same way `local_process` folds its own. A second copy
107// would be a second normalizer that can drift.
108pub(crate) mod native_support;
109pub mod pg_driver;
110pub mod pond;
111pub mod pond_door;
112pub mod pond_publish;
113pub mod publish_beacon;
114pub mod r2_publish;
115pub mod service_discovery;
116pub mod static_asset;
117pub mod static_asset_prune;
118pub mod sync_status;
119
120#[cfg(test)]
121mod lowering_golden;
122
123pub use bundle_store::{publish_bundle_to_r2, PublishReport as BundlePublishReport};
124pub use cloudflare_worker::CloudflareWorkerReconciler;
125pub use container::{ContainerOptions, ContainerReconciler};
126pub use derive_cache_prune::{
127    collect_live_derive_hashes, compute_derive_cache_candidates, execute_derive_cache_prune,
128    DeriveCacheLiveHashes, DerivePruneCandidate,
129};
130pub use domain::{
131    deploy_domain_passway, diff_apex_records, ensure_passway_apex, ensure_r2_custom_domain,
132    list_live_apex_records, plan_domain_passway, plan_passway_apex, public_origins, ApexRecordDiff,
133    DomainPasswayPlan, LiveApexRecord, PasswayApexOutcome, PasswayOrigin,
134};
135pub use headscale::{
136    DeclaredHeadscale, DeclaredPolicy, DeclaredPreauthKey, HeadscaleReconciler,
137    WORKLOAD_KIND as HEADSCALE_WORKLOAD_KIND,
138};
139pub use ingress::{
140    collate_front_doors, declared as ingress_declared, ensure_tunnel_ingress, machine_mesh_addrs,
141    plan_ingress, publish_tunnel_ingress, resolve_ingress_candidates, resolve_ingress_placements,
142    Collation, IngressPlan, IngressRule, NodeFrontDoor, PlannedEdge, TunnelIngressOutcome,
143};
144pub use ingress_verify::{
145    apply_public_path, resolve_upstreams_reporting, verify_collation, BeaconFetch, DialOutcome,
146    EndpointCheck, PublicReadings, RuleResolution, RuleResolutions, RuleVerdict, VerifyFinding,
147    VerifyReport,
148};
149pub use local_process::LocalProcessReconciler;
150pub use mesofact_bundle::{
151    resolve_bundle_machines, BundleSlot, MesofactBundleReconciler, RevalidateSlot,
152    SLOT_ROLE as BUNDLE_SLOT_ROLE,
153};
154pub use mesofact_static::{LocalStaticOptions, MesofactStaticReconciler};
155pub use pond::{PondOptions, PondState};
156pub use pond_door::{
157    door_env, door_state_dir, ensure_pond_cert, ensure_pond_cert_as, is_root, plan_pond_door,
158    pond_hostname, resolve_passway_binary, spawn_pond_door, CertPair, PondDoorPlan,
159    DEFAULT_DOOR_PORT, POND_TLD,
160};
161pub use pond_publish::{derive_minio_key, publish_to_pond, PondPublishReport};
162pub use r2_publish::{
163    publish_to_r2, R2PublishReport, R2PurgeOpts, R2_ACCESS_KEY_ENV, R2_ACCESS_KEY_SLOT,
164    R2_SECRET_KEY_ENV, R2_SECRET_KEY_SLOT,
165};
166pub use service_discovery::{
167    DiscoveredRecord, RecordVisibility, ServiceRecordFanout, UnknownReason,
168};
169pub use static_asset::StaticAssetReconciler;
170pub use static_asset_prune::{
171    compute_live_set, compute_prune_candidates, execute_prune, load_service_and_mirror,
172    PruneCandidate, PruneOutcome, PruneReport,
173};
174pub use sync_status::{
175    compute_cell, compute_service, new_sync_id, summarize, CellStatus, DriftEntry, HealthState,
176    MirrorObservation, Runtime, ServiceStatus, StatusSummary, SyncHistoryEntry, SyncOutcome,
177    SyncState, WireContainerStatus,
178};
179
180// ─── Log buffer ─────────────────────────────────────────────────────────────
181
182const LOG_CAP: usize = 500;
183
184#[derive(Debug, Default)]
185struct LogRing {
186    lines: VecDeque<String>,
187    /// Monotonically increasing total lines ever pushed (never decrements).
188    total: usize,
189}
190
191/// Bounded ring buffer for child-process stdout/stderr (R263-F3).
192/// Shared between the reader tasks and the Tauri `mirror_run_logs` command.
193#[derive(Debug, Clone, Default)]
194pub struct LogBuffer(Arc<AsyncMutex<LogRing>>);
195
196impl LogBuffer {
197    pub fn new() -> Self {
198        Self::default()
199    }
200
201    /// Append a line; drops the oldest entry when over capacity.
202    pub async fn push(&self, line: String) {
203        let mut ring = self.0.lock().await;
204        ring.total += 1;
205        ring.lines.push_back(line);
206        if ring.lines.len() > LOG_CAP {
207            ring.lines.pop_front();
208        }
209    }
210
211    /// Return lines not yet seen by the caller.
212    ///
213    /// `since` is the `total` cursor from the previous call (0 = nothing
214    /// seen yet). Returns `(new_lines, new_cursor)`. Pass `new_cursor` back
215    /// on the next call to receive only incremental output.
216    pub async fn since(&self, since: usize) -> (Vec<String>, usize) {
217        let ring = self.0.lock().await;
218        let oldest = ring.total.saturating_sub(ring.lines.len());
219        let skip = since.saturating_sub(oldest);
220        let new_lines: Vec<String> = ring.lines.iter().skip(skip).cloned().collect();
221        (new_lines, ring.total)
222    }
223
224    /// Current write cursor — the `total` [`Self::since`] would hand back
225    /// right now if nothing more were pushed. Lets a producer mark a boundary
226    /// (e.g. "everything before this point was the build phase") for a later
227    /// reader to seek past without re-reading lines it doesn't want.
228    pub async fn cursor(&self) -> usize {
229        self.0.lock().await.total
230    }
231}
232
233/// Shared cell for one phase-boundary cursor a reconciler can publish
234/// mid-`up()`, so a poller sees a multi-phase bring-up's internal transition
235/// (e.g. build → run) before the whole call returns — the same problem
236/// [`LogBuffer`] solves for output, for a single position instead of a ring.
237/// `None` until the reconciler reaches that phase; a caller registers one
238/// before calling `up()` to observe it live, same pattern as
239/// [`LogBuffer::clone`]-and-hand-in.
240#[derive(Debug, Clone, Default)]
241pub struct PhaseCursor(Arc<AsyncMutex<Option<usize>>>);
242
243impl PhaseCursor {
244    pub fn new() -> Self {
245        Self::default()
246    }
247
248    pub async fn set(&self, cursor: usize) {
249        *self.0.lock().await = Some(cursor);
250    }
251
252    pub async fn get(&self) -> Option<usize> {
253        *self.0.lock().await
254    }
255}
256
257/// The `(tenant, namespace)` a bring-up is scoped to (W206). Reconcilers that
258/// touch a credentialed provider resolve it at this scope
259/// ([`CfProvider::resolve_scoped`](super::reconciler::cf_creds)) so a namespace's
260/// Cloudflare zone/account/keystore slots come from its own scope rather than the
261/// workspace-global defaults. Defaults to the singleton `(default, default)`,
262/// which collapses every scoped lookup back to the historical global slots — so
263/// single-namespace deployments are unaffected.
264#[derive(Debug, Clone)]
265pub struct ProviderScope {
266    pub tenant: TenantId,
267    pub namespace: NamespaceId,
268}
269
270impl ProviderScope {
271    /// The degenerate single-tenant / single-namespace scope. Scoped provider
272    /// lookups made against it resolve to the pre-W206 global keystore slots.
273    pub fn singleton() -> Self {
274        Self {
275            tenant: TenantId::singleton(),
276            namespace: NamespaceId::singleton(),
277        }
278    }
279}
280
281impl Default for ProviderScope {
282    fn default() -> Self {
283        Self::singleton()
284    }
285}
286
287/// Inputs a reconciler sees for one bring-up.
288pub struct ReconcileCtx<'a> {
289    /// Workspace root (parent of `.yah/`). Used to resolve relative paths
290    /// on the component.
291    pub workspace_root: &'a Path,
292    /// Service that owns the component.
293    pub service: &'a ServiceConfig,
294    /// Component being brought up.
295    pub component: &'a ServiceComponent,
296    /// Mirror manifest the bring-up targets.
297    pub mirror: &'a MirrorConfig,
298    /// Environment name (file stem of `mirrors/<env>.toml`).
299    pub env: &'a str,
300    /// `(tenant, namespace)` this bring-up is scoped to (W206). Credentialed
301    /// providers resolve at this scope; defaults to [`ProviderScope::singleton`].
302    pub scope: ProviderScope,
303}
304
305impl<'a> ReconcileCtx<'a> {
306    /// Absolute path to the component's workload directory (the parent of
307    /// `workload.toml`).
308    ///
309    /// In-tree components resolve to `<workspace_root>/<path>`. For
310    /// `git`-sourced components (R561-F1, "BYO git") this points into the
311    /// local clone — `<source_cache>/<subdir>/<path>` — which is empty until
312    /// [`materialize`](Self::materialize) runs (approach A: clone-at-reconcile,
313    /// so config load + validation stay offline).
314    pub fn workload_dir(&self) -> PathBuf {
315        match &self.component.git {
316            None => self.workspace_root.join(&self.component.path),
317            Some(git) => {
318                let mut dir = self.source_cache_dir();
319                if let Some(subdir) = &git.subdir {
320                    dir = dir.join(subdir);
321                }
322                dir.join(&self.component.path)
323            }
324        }
325    }
326
327    /// Root of the local clone for a `git`-sourced component:
328    /// `<workspace_root>/.yah/infra/state/sources/<service>/<component_id>`.
329    fn source_cache_dir(&self) -> PathBuf {
330        self.workspace_root
331            .join(".yah/infra/state/sources")
332            .join(&self.service.name)
333            .join(&self.component.id)
334    }
335
336    /// Ensure a `git`-sourced component's code is present locally before build
337    /// (R561-F1, approach A). No-op for in-tree components. Idempotent: clones
338    /// on the first call, fetches + re-checks-out the pinned ref thereafter.
339    ///
340    /// Reconcilers MUST call this at the top of [`up`](Reconciler::up) before
341    /// reading [`workload_dir`](Self::workload_dir) for a remote component.
342    pub async fn materialize(&self) -> Result<()> {
343        let Some(git) = &self.component.git else {
344            return Ok(());
345        };
346        materialize_git_source(git, &self.source_cache_dir())
347            .await
348            .with_context(|| {
349                format!(
350                    "materializing git source {}@{} for {}/{}",
351                    git.repo, git.r#ref, self.service.name, self.component.id
352                )
353            })
354    }
355
356    /// Read `<workload_dir>/workload.toml` and extract just the `kind`
357    /// discriminator.
358    ///
359    /// Why this and not the strongly-typed [`workload_spec::Workload`]
360    /// parse: the on-disk `schema_version = 1` form predates the
361    /// `SchemaVersion::V1` enum and won't round-trip through the strong
362    /// types until B3 lands (see `crates/yah/cloud/src/config.rs` test
363    /// `web_workload_round_trips`). Reconcilers only need the kind to
364    /// dispatch; per-kind tooling (e.g. `mesofact-dev`'s
365    /// `WatchOptions::from_workload`) does its own parsing for the
366    /// build/out_dir fields it cares about.
367    pub fn workload_kind(&self) -> Result<String> {
368        workload_kind(&self.workload_dir())
369    }
370
371    /// Look up a provider slot by role (e.g. `"static"`, `"compute"`).
372    ///
373    /// Tries the component-qualified key first (`"<role>:<component id>"`)
374    /// before falling back to the bare role. A mirror role is normally
375    /// service-wide — one `providers.static` slot serves every static-kind
376    /// component — but a service can declare more than one component with
377    /// the same role (e.g. two `mesofact-static`/`mesofact-spa` components
378    /// sharing one mirror), and those need distinct ports to ever both come
379    /// up. The qualified key is how a mirror opts a specific component out
380    /// of sharing the bare-role slot:
381    ///
382    /// ```toml
383    /// [providers."static:site"]
384    /// kind = "local-static"
385    /// port = 4331
386    ///
387    /// [providers."static:app"]
388    /// kind = "local-static"
389    /// port = 4332
390    /// ```
391    ///
392    /// A mirror with only the bare role (the common, single-component case)
393    /// is unaffected — the qualified lookup misses and falls through.
394    pub fn slot(&self, role: &str) -> Option<&'a MirrorProviderSlot> {
395        let qualified = format!("{role}:{}", self.component.id);
396        self.mirror
397            .providers
398            .get(qualified.as_str())
399            .or_else(|| self.mirror.providers.get(role))
400    }
401}
402
403/// An explicit teardown hook for a workload this process did not spawn as a
404/// child — see [`RunningWorkload::with_teardown`] (R714-B1).
405///
406/// Boxed rather than a generic parameter because [`RunningWorkload`] is stored
407/// in heterogeneous collections (the desktop's mirror registry) and cannot
408/// carry a type parameter.
409/// `Sync` on the boxed closure is load-bearing, not belt-and-braces: without
410/// it `RunningWorkload` stops being `Sync`, so `&RunningWorkload` stops being
411/// `Send`, and every desktop `#[tauri::command]` that holds one across an
412/// `.await` (`mirror_run_logs` iterates the registry's handles) fails to
413/// compile with "future cannot be sent between threads safely". The captures a
414/// teardown needs — a container name and a workspace root — are `Sync` anyway.
415type TeardownFuture = std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send>>;
416type TeardownFn = Box<dyn FnOnce() -> TeardownFuture + Send + Sync + 'static>;
417
418/// Newtype so [`RunningWorkload`] can keep its `#[derive(Debug)]` — a boxed
419/// closure is not `Debug`.
420struct Teardown(TeardownFn);
421
422impl std::fmt::Debug for Teardown {
423    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
424        f.write_str("Teardown(<fn>)")
425    }
426}
427
428/// Handle to a workload that's been brought up. Owns the lifecycle: drop
429/// or call [`RunningWorkload::shutdown`] to take it back down.
430#[derive(Debug)]
431pub struct RunningWorkload {
432    /// Workload kind that was reconciled (e.g. `"mesofact-static"`).
433    pub kind: String,
434    /// Slot role this workload occupies on the mirror (e.g. `"static"`).
435    pub slot: String,
436    /// Local URL the workload exposes, when applicable. `None` for
437    /// workloads that publish to a non-local artifact store (e.g. R2).
438    pub dev_url: Option<String>,
439    /// Notional URL of the deployed artifact in the production case
440    /// (e.g. `https://yah.dev`). `None` until Cloudflare/R2 wiring lands.
441    pub public_url: Option<String>,
442    /// Secondary local UI console, if the workload exposes one (e.g. MinIO
443    /// console on a pond tier). Surfaced as a separate chip in the Services
444    /// matrix next to `dev_url`.
445    pub console_url: Option<String>,
446    /// Ring buffer for stdout/stderr from the workload's child process.
447    /// `None` for workloads that don't capture stdio (e.g. container-backed).
448    pub log_buffer: Option<LogBuffer>,
449
450    /// R546-B12: human-readable lines the reconciler wants the operator to see
451    /// on THIS run — what it actually did, not what is configured. A clean
452    /// static-asset reconcile used to print nothing at all, so a successful
453    /// publish and a successful no-op were indistinguishable at the apply
454    /// surface, and the one view that would have disambiguated them (`yah cloud
455    /// status`) was itself blind.
456    ///
457    /// Deliberately NOT on [`RunningWorkloadSummary`]: that type is serialized
458    /// across the process boundary to the desktop UI, and this is per-run
459    /// console output, not state the UI should cache.
460    pub notes: Vec<String>,
461
462    /// Sender that signals the supervisor task to tear down. Closing the
463    /// channel (drop) is equivalent to sending — supervisor exits on
464    /// channel close.
465    shutdown: Option<oneshot::Sender<()>>,
466    /// Joinable task that owns any child process and reaps it on signal.
467    supervisor: Option<tokio::task::JoinHandle<Result<()>>>,
468    /// R714-B1: teardown for a workload that runs OUTSIDE this process, so
469    /// there is no child to reap and no supervisor to signal.
470    ///
471    /// Deliberately run from [`RunningWorkload::shutdown`] only, never from
472    /// `Drop`. The two are different intents and conflating them breaks both
473    /// directions: a container the desktop started is meant to outlive the
474    /// desktop (the next launch re-adopts it with `adopt_only`), so quitting
475    /// the app must not `docker rm -f` it; but the ■ button IS an explicit
476    /// stop, and must.
477    teardown: Option<Teardown>,
478
479    /// R715-F4: where to ask this workload for a live status document, when
480    /// it declared a process-control channel (W315). `None` for a workload
481    /// with no channel — polling is then simply skipped, not an error.
482    ///
483    /// The desktop holds `RunningWorkload` in-process (it links this crate
484    /// directly, unlike the ad-hoc `run.spawn` path which lives behind the
485    /// camp daemon's socket), so a live poll is [`crate::proc_control::fetch_status`]
486    /// called straight against this endpoint — no RPC hop, no handle registry.
487    control: Option<crate::proc_control::ControlEndpoint>,
488
489    /// [`LogBuffer`] cursor marking the end of the build phase, for a
490    /// component that was compiled before being spawned (`local-process`
491    /// with `cargo_package` set). Lines before this index in `log_buffer` are
492    /// `cargo build` output; lines at or after are the spawned process's own
493    /// stdout/stderr. `None` when nothing was built (no `cargo_package`, or a
494    /// reconciler that predates this field).
495    pub build_log_end: Option<usize>,
496}
497
498impl RunningWorkload {
499    /// Create a handle for a workload that's already running externally
500    /// (e.g. embedded in yah-camp). No subprocess is owned; shutdown is a
501    /// no-op so the caller can call `shutdown()` uniformly.
502    pub fn adopted(
503        kind: impl Into<String>,
504        slot: impl Into<String>,
505        dev_url: Option<String>,
506    ) -> Self {
507        Self {
508            kind: kind.into(),
509            slot: slot.into(),
510            dev_url,
511            public_url: None,
512            console_url: None,
513            log_buffer: None,
514            notes: Vec::new(),
515            shutdown: None,
516            supervisor: None,
517            teardown: None,
518            control: None,
519            build_log_end: None,
520        }
521    }
522
523    /// Attach an explicit teardown to a handle for an externally-running
524    /// workload (R714-B1).
525    ///
526    /// `adopted()` alone gives a handle whose `shutdown()` is a documented
527    /// no-op. That is right for workloads another process owns and will keep
528    /// re-asserting (pond containers under camp's yubaba), and wrong for ones
529    /// this process started and nobody else will ever stop — for those the ■
530    /// button reported success while the container kept running.
531    ///
532    /// `teardown` runs on `shutdown()` and NOT on `Drop`; see [`Self::teardown`].
533    pub fn with_teardown<F, Fut>(mut self, teardown: F) -> Self
534    where
535        F: FnOnce() -> Fut + Send + Sync + 'static,
536        Fut: std::future::Future<Output = Result<()>> + Send + 'static,
537    {
538        self.teardown = Some(Teardown(Box::new(move || Box::pin(teardown()))));
539        self
540    }
541
542    /// Whether this handle carries a real teardown — i.e. whether a failed
543    /// [`Self::shutdown`] means the workload is **still running**.
544    ///
545    /// The stop path needs this to decide between an operator-facing error and
546    /// a log line, and R875-B1 is why it is a property of the handle rather
547    /// than a list of kinds at the call site. That list started as
548    /// `kind == "container"` when R714-B1 gave containers a teardown, and was
549    /// silently wrong the moment a second kind grew one: an adopted
550    /// mesofact-dev whose teardown failed reported a successful stop with the
551    /// server still serving. A handle knows whether it owns the workload; a
552    /// string comparison at the call site only knows what it was last taught.
553    pub fn owns_teardown(&self) -> bool {
554        self.teardown.is_some()
555    }
556
557    /// Attach per-run operator-facing lines (R546-B12). See [`Self::notes`].
558    pub fn with_notes(mut self, notes: Vec<String>) -> Self {
559        self.notes = notes;
560        self
561    }
562
563    /// Record where to poll this workload's process-control channel, when it
564    /// declared one (R715-F4). A no-op (leaves `control: None`) when the
565    /// argument is `None` — callers can pass the reconciler's resolved
566    /// endpoint straight through without an `if let`.
567    pub fn with_control(mut self, control: Option<crate::proc_control::ControlEndpoint>) -> Self {
568        self.control = control;
569        self
570    }
571
572    /// Ask this workload's process-control channel for a live status
573    /// document (R715-F4). `None` when no channel was declared, or when the
574    /// endpoint didn't answer — both are "nothing new to show", not errors;
575    /// see [`crate::proc_control::fetch_status`] for why a poll failure isn't
576    /// itself meaningful.
577    pub async fn poll_control(&self) -> Option<crate::proc_control::ProcStatus> {
578        let endpoint = self.control.as_ref()?;
579        crate::proc_control::fetch_status(endpoint).await.ok()
580    }
581
582    /// Whether the supervisor task that owns this workload's child process is
583    /// still running. `spawn_native_log_supervisor` (native_support.rs) exits
584    /// as soon as `NativeRuntime::get_workload` reports a terminal state —
585    /// which itself comes from a real `child.wait()` in kamaji's native
586    /// backend, so this catches a crash (segfault, panic, `SIGKILL`) the same
587    /// way it catches a clean exit, not just an unresponsive process.
588    ///
589    /// A workload with no supervisor (`RunningWorkload::adopted` — runs
590    /// outside this process, e.g. a pond container camp re-asserts) has
591    /// nothing to check here and reports alive unconditionally; its liveness
592    /// is whatever tracks it, not this handle.
593    pub fn is_alive(&self) -> bool {
594        self.supervisor.as_ref().is_none_or(|h| !h.is_finished())
595    }
596
597    /// Record where the build phase ends in `log_buffer`, for a component
598    /// that was compiled before being spawned. See [`Self::build_log_end`].
599    pub fn with_build_log_end(mut self, cursor: Option<usize>) -> Self {
600        self.build_log_end = cursor;
601        self
602    }
603
604    /// Set the public URL for a published workload (e.g. `"https://yah.dev"`).
605    pub fn with_public_url(mut self, url: impl Into<String>) -> Self {
606        self.public_url = Some(url.into());
607        self
608    }
609
610    /// Set the console URL for a workload that exposes a secondary local UI
611    /// (e.g. MinIO console on a pond tier).
612    pub fn with_console_url(mut self, url: impl Into<String>) -> Self {
613        self.console_url = Some(url.into());
614        self
615    }
616
617    /// Gracefully tear down: signal the supervisor, await its exit, then run
618    /// any explicit teardown hook.
619    ///
620    /// The hook runs LAST and its error propagates. A stop that could not tear
621    /// the workload down must surface as an error, never as a silent success —
622    /// that silence is the whole of R714-B1.
623    pub async fn shutdown(mut self) -> Result<()> {
624        if let Some(tx) = self.shutdown.take() {
625            let _ = tx.send(());
626        }
627        if let Some(handle) = self.supervisor.take() {
628            handle
629                .await
630                .context("joining workload supervisor")?
631                .context("workload supervisor")?;
632        }
633        if let Some(Teardown(hook)) = self.teardown.take() {
634            hook().await.context("workload teardown")?;
635        }
636        Ok(())
637    }
638}
639
640impl Drop for RunningWorkload {
641    fn drop(&mut self) {
642        // Best-effort signal. The supervisor task is detached and will
643        // reap its child when it observes the closed channel.
644        if let Some(tx) = self.shutdown.take() {
645            let _ = tx.send(());
646        }
647    }
648}
649
650/// Bring one workload up. Each impl handles one [`ServiceComponent::kind`].
651#[async_trait]
652pub trait Reconciler: Send + Sync {
653    /// Workload kind this reconciler handles (matches `ServiceComponent.kind`).
654    fn kind(&self) -> &'static str;
655
656    /// Bring the workload up. Returns a handle whose lifecycle is tied to
657    /// the mirror being up.
658    async fn up(&self, ctx: ReconcileCtx<'_>) -> Result<RunningWorkload>;
659}
660
661/// Read `<workload_dir>/workload.toml` and return just the `kind` field.
662/// See [`ReconcileCtx::workload_kind`] for why we don't deserialize through
663/// the strong types yet.
664pub fn workload_kind(workload_dir: &Path) -> Result<String> {
665    let path = workload_dir.join("workload.toml");
666    let src =
667        std::fs::read_to_string(&path).with_context(|| format!("reading {}", path.display()))?;
668    let value: toml::Value =
669        toml::from_str(&src).with_context(|| format!("parsing {}", path.display()))?;
670    let kind = value
671        .get("kind")
672        .and_then(|v| v.as_str())
673        .with_context(|| format!("{}: missing `kind` field", path.display()))?;
674    Ok(kind.to_string())
675}
676
677/// Shallow-clone (or update) a [`GitSource`] into `dir` (R561-F1). Idempotent:
678/// clones when `dir/.git` is absent, otherwise fetches the pinned ref and
679/// force-checks-it-out. Uses the system `git` so it inherits the operator's
680/// credential helpers / SSH agent — no in-process git library.
681///
682/// `pub` (R615-T3 / W274): `yah infra sync` reuses this verbatim for
683/// `InfraSourceKind::Git` sources rather than a second shallow-clone-or-pull
684/// implementation — same "one git-source shape, reused" discipline R615-F1
685/// already applied to the type.
686pub async fn materialize_git_source(git: &GitSource, dir: &Path) -> Result<()> {
687    use tokio::process::Command;
688
689    async fn run_git(args: &[&std::ffi::OsStr]) -> Result<()> {
690        let out = Command::new("git")
691            .args(args)
692            .output()
693            .await
694            .context("spawning git")?;
695        if !out.status.success() {
696            anyhow::bail!(
697                "git {} failed: {}",
698                args.iter()
699                    .map(|a| a.to_string_lossy())
700                    .collect::<Vec<_>>()
701                    .join(" "),
702                String::from_utf8_lossy(&out.stderr).trim()
703            );
704        }
705        Ok(())
706    }
707
708    use std::ffi::OsStr;
709    let dir_os = dir.as_os_str();
710    let r#ref = git.r#ref.as_str();
711
712    if dir.join(".git").is_dir() {
713        // Existing checkout — update to the pinned ref.
714        run_git(&[
715            OsStr::new("-C"),
716            dir_os,
717            OsStr::new("fetch"),
718            OsStr::new("--depth"),
719            OsStr::new("1"),
720            OsStr::new("origin"),
721            OsStr::new(r#ref),
722        ])
723        .await?;
724        run_git(&[
725            OsStr::new("-C"),
726            dir_os,
727            OsStr::new("checkout"),
728            OsStr::new("--force"),
729            OsStr::new("FETCH_HEAD"),
730        ])
731        .await?;
732    } else {
733        if let Some(parent) = dir.parent() {
734            tokio::fs::create_dir_all(parent)
735                .await
736                .with_context(|| format!("creating {}", parent.display()))?;
737        }
738        // `--branch` accepts a branch or tag. Pinning to a bare commit SHA is a
739        // follow-up (needs clone-then-fetch); the common case is a branch/tag.
740        run_git(&[
741            OsStr::new("clone"),
742            OsStr::new("--depth"),
743            OsStr::new("1"),
744            OsStr::new("--branch"),
745            OsStr::new(r#ref),
746            OsStr::new(git.repo.as_str()),
747            dir_os,
748        ])
749        .await?;
750    }
751    Ok(())
752}
753
754/// Build a `RunningWorkload` from the pieces a reconciler produces.
755pub(crate) fn into_running(
756    kind: impl Into<String>,
757    slot: impl Into<String>,
758    dev_url: Option<String>,
759    public_url: Option<String>,
760    log_buffer: Option<LogBuffer>,
761    shutdown: oneshot::Sender<()>,
762    supervisor: tokio::task::JoinHandle<Result<()>>,
763) -> RunningWorkload {
764    RunningWorkload {
765        kind: kind.into(),
766        slot: slot.into(),
767        dev_url,
768        public_url,
769        console_url: None,
770        log_buffer,
771        notes: Vec::new(),
772        shutdown: Some(shutdown),
773        supervisor: Some(supervisor),
774        teardown: None,
775        control: None,
776        build_log_end: None,
777    }
778}
779
780/// Wait for a TCP port to start accepting connections. Returns `true` if
781/// the port came up within `timeout`, `false` otherwise. Useful for
782/// reconcilers that spawn a server and need to know when it's reachable
783/// before reporting success.
784pub(crate) async fn wait_for_port(
785    addr: std::net::SocketAddr,
786    timeout: std::time::Duration,
787) -> bool {
788    let deadline = tokio::time::Instant::now() + timeout;
789    loop {
790        if tokio::net::TcpStream::connect(addr).await.is_ok() {
791            return true;
792        }
793        if tokio::time::Instant::now() >= deadline {
794            return false;
795        }
796        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
797    }
798}
799
800// `wait_for_http_ready` lived here pre-R374-F3 to back the MinIO health
801// probe in pond's bring-up path. That logic moved to
802// `local_driver::pond_minio::wait_for_http_ready` so yubaba + cloud share
803// it. The mesofact-static reconciler arm uses [`wait_for_port`] for
804// dev-tier port readiness; nothing else needs an HTTP-level probe today.
805
806/// Pluck a `u16` out of a [`MirrorProviderSlot`]'s inline `fields` map.
807/// Returns `None` if the key is absent or out of range.
808pub(crate) fn slot_field_u16(fields: &BTreeMap<String, toml::Value>, key: &str) -> Option<u16> {
809    fields
810        .get(key)
811        .and_then(|v| v.as_integer())
812        .and_then(|n| u16::try_from(n).ok())
813}
814
815/// Serializable summary of a running workload — what the desktop / CLI
816/// hands to the UI. Subset of [`RunningWorkload`] that's safe to cross
817/// process boundaries.
818#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
819pub struct RunningWorkloadSummary {
820    pub kind: String,
821    pub slot: String,
822    pub dev_url: Option<String>,
823    pub public_url: Option<String>,
824    pub console_url: Option<String>,
825    /// Operator-facing lines from [`RunningWorkload::notes`] (R546-B12).
826    ///
827    /// These used to stop here — the summary dropped them, so every note a
828    /// reconciler attached was written into a struct nobody read. They matter
829    /// most for a workload with no `dev_url` to click (R715-T2): the notes are
830    /// then the only structured thing the Run tab has to show about it.
831    #[serde(default, skip_serializing_if = "Vec::is_empty")]
832    pub notes: Vec<String>,
833    /// See [`RunningWorkload::build_log_end`]. Lets a log-tail consumer skip
834    /// straight to run-phase output, or show everything from 0 when the
835    /// operator wants the build log too (e.g. after a compile failure).
836    #[serde(default, skip_serializing_if = "Option::is_none")]
837    pub build_log_end: Option<usize>,
838}
839
840impl From<&RunningWorkload> for RunningWorkloadSummary {
841    fn from(r: &RunningWorkload) -> Self {
842        Self {
843            kind: r.kind.clone(),
844            slot: r.slot.clone(),
845            dev_url: r.dev_url.clone(),
846            public_url: r.public_url.clone(),
847            console_url: r.console_url.clone(),
848            notes: r.notes.clone(),
849            build_log_end: r.build_log_end,
850        }
851    }
852}
853
854#[cfg(test)]
855mod source_seam_tests {
856    //! R561-F1 — the BYO-git source seam: path resolution + materialization.
857    use super::*;
858    use std::collections::BTreeMap;
859
860    fn component(git: Option<GitSource>) -> ServiceComponent {
861        ServiceComponent {
862            mount: None,
863            id: "site".into(),
864            kind: "mesofact-static".into(),
865            path: "site".into(),
866            role: "static".into(),
867            publishes: None,
868            wave: 0,
869            git,
870            deploy: Default::default(),
871        }
872    }
873
874    fn service(comp: ServiceComponent) -> ServiceConfig {
875        ServiceConfig {
876            schema_version: 1,
877            name: "scrabcake".into(),
878            domain: "scrabcake.example".into(),
879            components: vec![comp],
880            db: crate::DbCatalog::default(),
881        }
882    }
883
884    fn mirror() -> MirrorConfig {
885        MirrorConfig {
886            schema_version: 1,
887            shape: crate::MirrorShape::Local,
888            providers: BTreeMap::new(),
889            ingress: Default::default(),
890            ingress_machines: Vec::new(),
891            drivers: Default::default(),
892            asset_aliases: BTreeMap::new(),
893        }
894    }
895
896    fn ctx<'a>(ws: &'a Path, svc: &'a ServiceConfig, mir: &'a MirrorConfig) -> ReconcileCtx<'a> {
897        ReconcileCtx {
898            workspace_root: ws,
899            service: svc,
900            component: &svc.components[0],
901            mirror: mir,
902            env: "dev",
903            scope: ProviderScope::singleton(),
904        }
905    }
906
907    #[test]
908    fn workload_dir_in_tree_joins_workspace_root() {
909        let svc = service(component(None));
910        let mir = mirror();
911        assert_eq!(
912            ctx(Path::new("/ws"), &svc, &mir).workload_dir(),
913            Path::new("/ws/site")
914        );
915    }
916
917    #[test]
918    fn workload_dir_git_resolves_into_source_cache_with_subdir() {
919        let git = GitSource {
920            repo: "https://example.com/r.git".into(),
921            r#ref: "main".into(),
922            subdir: Some("apps".into()),
923        };
924        let svc = service(component(Some(git)));
925        let mir = mirror();
926        assert_eq!(
927            ctx(Path::new("/ws"), &svc, &mir).workload_dir(),
928            Path::new("/ws/.yah/infra/state/sources/scrabcake/site/apps/site")
929        );
930    }
931
932    #[tokio::test]
933    async fn materialize_is_noop_for_in_tree_component() {
934        let svc = service(component(None));
935        let mir = mirror();
936        // No git source → Ok, and nothing is written under the workspace.
937        ctx(Path::new("/nonexistent-ws"), &svc, &mir)
938            .materialize()
939            .await
940            .unwrap();
941    }
942
943    #[tokio::test]
944    async fn materialize_clones_git_source_offline() {
945        fn git(args: &[&str], cwd: &Path) {
946            let out = std::process::Command::new("git")
947                .args(args)
948                .current_dir(cwd)
949                .env("GIT_AUTHOR_NAME", "t")
950                .env("GIT_AUTHOR_EMAIL", "t@t")
951                .env("GIT_COMMITTER_NAME", "t")
952                .env("GIT_COMMITTER_EMAIL", "t@t")
953                .output()
954                .unwrap();
955            assert!(
956                out.status.success(),
957                "git {args:?}: {}",
958                String::from_utf8_lossy(&out.stderr)
959            );
960        }
961
962        let tmp = tempfile::tempdir().unwrap();
963        let src = tmp.path().join("src-repo");
964        std::fs::create_dir_all(&src).unwrap();
965        git(&["init", "-b", "main"], &src);
966        std::fs::write(src.join("hello.txt"), "hi").unwrap();
967        git(&["add", "."], &src);
968        git(&["commit", "-m", "init"], &src);
969
970        let source = GitSource {
971            repo: format!("file://{}", src.display()),
972            r#ref: "main".into(),
973            subdir: None,
974        };
975        let dest = tmp.path().join("cache");
976
977        // First call clones.
978        materialize_git_source(&source, &dest).await.unwrap();
979        assert!(dest.join("hello.txt").is_file());
980
981        // Second call takes the update path and stays green (idempotent).
982        materialize_git_source(&source, &dest).await.unwrap();
983        assert!(dest.join("hello.txt").is_file());
984    }
985}
986
987#[cfg(test)]
988mod teardown_tests {
989    //! R714-B1 — the explicit-teardown contract on [`RunningWorkload`].
990    //!
991    //! These are about WHEN the hook runs, not what it does. The bug being
992    //! fixed was a `shutdown()` that reported success having done nothing, and
993    //! the trap in fixing it is a `Drop` that tears down a container which is
994    //! supposed to survive the process.
995    use super::*;
996    use std::sync::atomic::{AtomicUsize, Ordering};
997
998    fn counting() -> (RunningWorkload, Arc<AtomicUsize>) {
999        let hits = Arc::new(AtomicUsize::new(0));
1000        let seen = hits.clone();
1001        let w = RunningWorkload::adopted("container", "compute", None)
1002            .with_teardown(move || {
1003                let seen = seen.clone();
1004                async move {
1005                    seen.fetch_add(1, Ordering::SeqCst);
1006                    Ok(())
1007                }
1008            });
1009        (w, hits)
1010    }
1011
1012    /// Attaching a teardown must not cost `RunningWorkload` its auto traits.
1013    /// The desktop stores these handles in a shared registry and its Tauri
1014    /// commands iterate them across `.await` points, which needs `Sync` — a
1015    /// hook that is `Send` but not `Sync` takes it away here and surfaces two
1016    /// crates over as "future cannot be sent between threads safely", with a
1017    /// span pointing at `mirror_run_logs` rather than at this file.
1018    #[test]
1019    fn a_teardown_does_not_cost_the_handle_send_or_sync() {
1020        fn assert_send_sync<T: Send + Sync>() {}
1021        assert_send_sync::<RunningWorkload>();
1022    }
1023
1024    #[tokio::test]
1025    async fn shutdown_runs_the_teardown() {
1026        let (w, hits) = counting();
1027        w.shutdown().await.unwrap();
1028        assert_eq!(hits.load(Ordering::SeqCst), 1);
1029    }
1030
1031    #[tokio::test]
1032    async fn dropping_the_handle_does_NOT_run_the_teardown() {
1033        // A container this process started is meant to outlive it — the next
1034        // launch re-adopts it. Quitting the app must not `docker rm -f` it.
1035        let (w, hits) = counting();
1036        drop(w);
1037        // Yield so a stray spawned task would have had a chance to run.
1038        tokio::task::yield_now().await;
1039        assert_eq!(hits.load(Ordering::SeqCst), 0);
1040    }
1041
1042    #[tokio::test]
1043    async fn a_failing_teardown_makes_shutdown_fail() {
1044        // The whole of R714-B1: a stop that could not tear the workload down
1045        // must not report success.
1046        let w = RunningWorkload::adopted("container", "compute", None)
1047            .with_teardown(|| async { anyhow::bail!("docker stop refused") });
1048        let err = w.shutdown().await.unwrap_err();
1049        assert!(format!("{err:#}").contains("docker stop refused"), "{err:#}");
1050    }
1051
1052    #[tokio::test]
1053    async fn an_adopted_handle_without_a_teardown_still_shuts_down_cleanly() {
1054        // Pond containers are owned by camp's yubaba and stop through
1055        // `workload.stop`; their no-op shutdown is correct and must stay.
1056        let w = RunningWorkload::adopted("mesofact-static", "static", None);
1057        w.shutdown().await.unwrap();
1058    }
1059
1060    /// R875-B1. The stop path decides "still running, tell the operator" vs
1061    /// "untidy supervisor, log it" from this, so it has to answer for the
1062    /// handle in front of it rather than for the kind string it carries — the
1063    /// two disagreed for every adopted mesofact-dev.
1064    #[test]
1065    fn owns_teardown_tracks_the_hook_not_the_kind() {
1066        let (owned, _) = counting();
1067        assert!(owned.owns_teardown());
1068        assert!(
1069            !RunningWorkload::adopted("container", "compute", None).owns_teardown(),
1070            "a bare adopted handle owns nothing, whatever its kind says"
1071        );
1072        assert!(
1073            RunningWorkload::adopted("mesofact-static", "static", None)
1074                .with_teardown(|| async { Ok(()) })
1075                .owns_teardown(),
1076            "a non-container kind with a teardown owns its workload"
1077        );
1078    }
1079
1080    #[tokio::test]
1081    async fn the_teardown_runs_after_the_supervisor_is_joined() {
1082        // Ordering matters for a workload that has both: reap the child first,
1083        // then remove the container it was talking to.
1084        let order = Arc::new(AsyncMutex::new(Vec::<&'static str>::new()));
1085        let (tx, rx) = oneshot::channel::<()>();
1086        let sup_order = order.clone();
1087        let supervisor = tokio::spawn(async move {
1088            let _ = rx.await;
1089            sup_order.lock().await.push("supervisor");
1090            Ok(())
1091        });
1092        let hook_order = order.clone();
1093        let w = into_running("container", "compute", None, None, None, tx, supervisor)
1094            .with_teardown(move || {
1095                let hook_order = hook_order.clone();
1096                async move {
1097                    hook_order.lock().await.push("teardown");
1098                    Ok(())
1099                }
1100            });
1101
1102        w.shutdown().await.unwrap();
1103        assert_eq!(*order.lock().await, vec!["supervisor", "teardown"]);
1104    }
1105}