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
67use std::collections::{BTreeMap, VecDeque};
68use std::path::{Path, PathBuf};
69use std::sync::Arc;
70
71use anyhow::{Context, Result};
72use async_trait::async_trait;
73use serde::{Deserialize, Serialize};
74use tokio::sync::{oneshot, Mutex as AsyncMutex};
75
76use workload_spec::{NamespaceId, TenantId};
77
78use crate::{GitSource, MirrorConfig, MirrorProviderSlot, ServiceComponent, ServiceConfig};
79
80pub mod bundle_store;
81pub(crate) mod cf_creds;
82pub mod cloudflare_worker;
83pub mod container;
84pub mod derive_cache_prune;
85pub mod domain;
86pub mod ingress;
87pub mod local_process;
88pub mod mesofact_bundle;
89pub mod mesofact_static;
90mod native_support;
91pub mod pg_driver;
92pub mod pond;
93pub mod pond_publish;
94pub mod publish_beacon;
95pub mod r2_publish;
96pub mod static_asset;
97pub mod static_asset_prune;
98pub mod sync_status;
99
100#[cfg(test)]
101mod lowering_golden;
102
103pub use bundle_store::{publish_bundle_to_r2, PublishReport as BundlePublishReport};
104pub use cloudflare_worker::CloudflareWorkerReconciler;
105pub use container::{ContainerOptions, ContainerReconciler};
106pub use derive_cache_prune::{
107    collect_live_derive_hashes, compute_derive_cache_candidates, execute_derive_cache_prune,
108    DeriveCacheLiveHashes, DerivePruneCandidate,
109};
110pub use domain::ensure_r2_custom_domain;
111pub use ingress::{
112    declared as ingress_declared, ensure_tunnel_ingress, plan_ingress, publish_tunnel_ingress,
113    IngressPlan, IngressRule, TunnelIngressOutcome,
114};
115pub use local_process::LocalProcessReconciler;
116pub use mesofact_bundle::{
117    resolve_bundle_machines, BundleSlot, MesofactBundleReconciler, RevalidateSlot,
118    SLOT_ROLE as BUNDLE_SLOT_ROLE,
119};
120pub use mesofact_static::{LocalStaticOptions, MesofactStaticReconciler};
121pub use pond::{PondOptions, PondState};
122pub use pond_publish::{derive_minio_key, publish_to_pond, PondPublishReport};
123pub use r2_publish::{
124    publish_to_r2, R2PublishReport, R2PurgeOpts, R2_ACCESS_KEY_ENV, R2_ACCESS_KEY_SLOT,
125    R2_SECRET_KEY_ENV, R2_SECRET_KEY_SLOT,
126};
127pub use static_asset::StaticAssetReconciler;
128pub use static_asset_prune::{
129    compute_live_set, compute_prune_candidates, execute_prune, load_service_and_mirror,
130    PruneCandidate, PruneOutcome, PruneReport,
131};
132pub use sync_status::{
133    compute_cell, compute_service, new_sync_id, summarize, CellStatus, DriftEntry, HealthState,
134    MirrorObservation, Runtime, ServiceStatus, StatusSummary, SyncHistoryEntry, SyncOutcome,
135    SyncState, WireContainerStatus,
136};
137
138// ─── Log buffer ─────────────────────────────────────────────────────────────
139
140const LOG_CAP: usize = 500;
141
142#[derive(Debug, Default)]
143struct LogRing {
144    lines: VecDeque<String>,
145    /// Monotonically increasing total lines ever pushed (never decrements).
146    total: usize,
147}
148
149/// Bounded ring buffer for child-process stdout/stderr (R263-F3).
150/// Shared between the reader tasks and the Tauri `mirror_run_logs` command.
151#[derive(Debug, Clone, Default)]
152pub struct LogBuffer(Arc<AsyncMutex<LogRing>>);
153
154impl LogBuffer {
155    pub fn new() -> Self {
156        Self::default()
157    }
158
159    /// Append a line; drops the oldest entry when over capacity.
160    pub async fn push(&self, line: String) {
161        let mut ring = self.0.lock().await;
162        ring.total += 1;
163        ring.lines.push_back(line);
164        if ring.lines.len() > LOG_CAP {
165            ring.lines.pop_front();
166        }
167    }
168
169    /// Return lines not yet seen by the caller.
170    ///
171    /// `since` is the `total` cursor from the previous call (0 = nothing
172    /// seen yet). Returns `(new_lines, new_cursor)`. Pass `new_cursor` back
173    /// on the next call to receive only incremental output.
174    pub async fn since(&self, since: usize) -> (Vec<String>, usize) {
175        let ring = self.0.lock().await;
176        let oldest = ring.total.saturating_sub(ring.lines.len());
177        let skip = since.saturating_sub(oldest);
178        let new_lines: Vec<String> = ring.lines.iter().skip(skip).cloned().collect();
179        (new_lines, ring.total)
180    }
181}
182
183/// The `(tenant, namespace)` a bring-up is scoped to (W206). Reconcilers that
184/// touch a credentialed provider resolve it at this scope
185/// ([`CfProvider::resolve_scoped`](super::reconciler::cf_creds)) so a namespace's
186/// Cloudflare zone/account/keystore slots come from its own scope rather than the
187/// workspace-global defaults. Defaults to the singleton `(default, default)`,
188/// which collapses every scoped lookup back to the historical global slots — so
189/// single-namespace deployments are unaffected.
190#[derive(Debug, Clone)]
191pub struct ProviderScope {
192    pub tenant: TenantId,
193    pub namespace: NamespaceId,
194}
195
196impl ProviderScope {
197    /// The degenerate single-tenant / single-namespace scope. Scoped provider
198    /// lookups made against it resolve to the pre-W206 global keystore slots.
199    pub fn singleton() -> Self {
200        Self {
201            tenant: TenantId::singleton(),
202            namespace: NamespaceId::singleton(),
203        }
204    }
205}
206
207impl Default for ProviderScope {
208    fn default() -> Self {
209        Self::singleton()
210    }
211}
212
213/// Inputs a reconciler sees for one bring-up.
214pub struct ReconcileCtx<'a> {
215    /// Workspace root (parent of `.yah/`). Used to resolve relative paths
216    /// on the component.
217    pub workspace_root: &'a Path,
218    /// Service that owns the component.
219    pub service: &'a ServiceConfig,
220    /// Component being brought up.
221    pub component: &'a ServiceComponent,
222    /// Mirror manifest the bring-up targets.
223    pub mirror: &'a MirrorConfig,
224    /// Environment name (file stem of `mirrors/<env>.toml`).
225    pub env: &'a str,
226    /// `(tenant, namespace)` this bring-up is scoped to (W206). Credentialed
227    /// providers resolve at this scope; defaults to [`ProviderScope::singleton`].
228    pub scope: ProviderScope,
229}
230
231impl<'a> ReconcileCtx<'a> {
232    /// Absolute path to the component's workload directory (the parent of
233    /// `workload.toml`).
234    ///
235    /// In-tree components resolve to `<workspace_root>/<path>`. For
236    /// `git`-sourced components (R561-F1, "BYO git") this points into the
237    /// local clone — `<source_cache>/<subdir>/<path>` — which is empty until
238    /// [`materialize`](Self::materialize) runs (approach A: clone-at-reconcile,
239    /// so config load + validation stay offline).
240    pub fn workload_dir(&self) -> PathBuf {
241        match &self.component.git {
242            None => self.workspace_root.join(&self.component.path),
243            Some(git) => {
244                let mut dir = self.source_cache_dir();
245                if let Some(subdir) = &git.subdir {
246                    dir = dir.join(subdir);
247                }
248                dir.join(&self.component.path)
249            }
250        }
251    }
252
253    /// Root of the local clone for a `git`-sourced component:
254    /// `<workspace_root>/.yah/infra/state/sources/<service>/<component_id>`.
255    fn source_cache_dir(&self) -> PathBuf {
256        self.workspace_root
257            .join(".yah/infra/state/sources")
258            .join(&self.service.name)
259            .join(&self.component.id)
260    }
261
262    /// Ensure a `git`-sourced component's code is present locally before build
263    /// (R561-F1, approach A). No-op for in-tree components. Idempotent: clones
264    /// on the first call, fetches + re-checks-out the pinned ref thereafter.
265    ///
266    /// Reconcilers MUST call this at the top of [`up`](Reconciler::up) before
267    /// reading [`workload_dir`](Self::workload_dir) for a remote component.
268    pub async fn materialize(&self) -> Result<()> {
269        let Some(git) = &self.component.git else {
270            return Ok(());
271        };
272        materialize_git_source(git, &self.source_cache_dir())
273            .await
274            .with_context(|| {
275                format!(
276                    "materializing git source {}@{} for {}/{}",
277                    git.repo, git.r#ref, self.service.name, self.component.id
278                )
279            })
280    }
281
282    /// Read `<workload_dir>/workload.toml` and extract just the `kind`
283    /// discriminator.
284    ///
285    /// Why this and not the strongly-typed [`workload_spec::Workload`]
286    /// parse: the on-disk `schema_version = 1` form predates the
287    /// `SchemaVersion::V1` enum and won't round-trip through the strong
288    /// types until B3 lands (see `crates/yah/cloud/src/config.rs` test
289    /// `web_workload_round_trips`). Reconcilers only need the kind to
290    /// dispatch; per-kind tooling (e.g. `mesofact-dev`'s
291    /// `WatchOptions::from_workload`) does its own parsing for the
292    /// build/out_dir fields it cares about.
293    pub fn workload_kind(&self) -> Result<String> {
294        workload_kind(&self.workload_dir())
295    }
296
297    /// Look up a provider slot by role (e.g. `"static"`, `"compute"`).
298    pub fn slot(&self, role: &str) -> Option<&'a MirrorProviderSlot> {
299        self.mirror.providers.get(role)
300    }
301}
302
303/// An explicit teardown hook for a workload this process did not spawn as a
304/// child — see [`RunningWorkload::with_teardown`] (R714-B1).
305///
306/// Boxed rather than a generic parameter because [`RunningWorkload`] is stored
307/// in heterogeneous collections (the desktop's mirror registry) and cannot
308/// carry a type parameter.
309/// `Sync` on the boxed closure is load-bearing, not belt-and-braces: without
310/// it `RunningWorkload` stops being `Sync`, so `&RunningWorkload` stops being
311/// `Send`, and every desktop `#[tauri::command]` that holds one across an
312/// `.await` (`mirror_run_logs` iterates the registry's handles) fails to
313/// compile with "future cannot be sent between threads safely". The captures a
314/// teardown needs — a container name and a workspace root — are `Sync` anyway.
315type TeardownFuture = std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send>>;
316type TeardownFn = Box<dyn FnOnce() -> TeardownFuture + Send + Sync + 'static>;
317
318/// Newtype so [`RunningWorkload`] can keep its `#[derive(Debug)]` — a boxed
319/// closure is not `Debug`.
320struct Teardown(TeardownFn);
321
322impl std::fmt::Debug for Teardown {
323    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
324        f.write_str("Teardown(<fn>)")
325    }
326}
327
328/// Handle to a workload that's been brought up. Owns the lifecycle: drop
329/// or call [`RunningWorkload::shutdown`] to take it back down.
330#[derive(Debug)]
331pub struct RunningWorkload {
332    /// Workload kind that was reconciled (e.g. `"mesofact-static"`).
333    pub kind: String,
334    /// Slot role this workload occupies on the mirror (e.g. `"static"`).
335    pub slot: String,
336    /// Local URL the workload exposes, when applicable. `None` for
337    /// workloads that publish to a non-local artifact store (e.g. R2).
338    pub dev_url: Option<String>,
339    /// Notional URL of the deployed artifact in the production case
340    /// (e.g. `https://yah.dev`). `None` until Cloudflare/R2 wiring lands.
341    pub public_url: Option<String>,
342    /// Secondary local UI console, if the workload exposes one (e.g. MinIO
343    /// console on a pond tier). Surfaced as a separate chip in the Services
344    /// matrix next to `dev_url`.
345    pub console_url: Option<String>,
346    /// Ring buffer for stdout/stderr from the workload's child process.
347    /// `None` for workloads that don't capture stdio (e.g. container-backed).
348    pub log_buffer: Option<LogBuffer>,
349
350    /// R546-B12: human-readable lines the reconciler wants the operator to see
351    /// on THIS run — what it actually did, not what is configured. A clean
352    /// static-asset reconcile used to print nothing at all, so a successful
353    /// publish and a successful no-op were indistinguishable at the apply
354    /// surface, and the one view that would have disambiguated them (`yah cloud
355    /// status`) was itself blind.
356    ///
357    /// Deliberately NOT on [`RunningWorkloadSummary`]: that type is serialized
358    /// across the process boundary to the desktop UI, and this is per-run
359    /// console output, not state the UI should cache.
360    pub notes: Vec<String>,
361
362    /// Sender that signals the supervisor task to tear down. Closing the
363    /// channel (drop) is equivalent to sending — supervisor exits on
364    /// channel close.
365    shutdown: Option<oneshot::Sender<()>>,
366    /// Joinable task that owns any child process and reaps it on signal.
367    supervisor: Option<tokio::task::JoinHandle<Result<()>>>,
368    /// R714-B1: teardown for a workload that runs OUTSIDE this process, so
369    /// there is no child to reap and no supervisor to signal.
370    ///
371    /// Deliberately run from [`RunningWorkload::shutdown`] only, never from
372    /// `Drop`. The two are different intents and conflating them breaks both
373    /// directions: a container the desktop started is meant to outlive the
374    /// desktop (the next launch re-adopts it with `adopt_only`), so quitting
375    /// the app must not `docker rm -f` it; but the ■ button IS an explicit
376    /// stop, and must.
377    teardown: Option<Teardown>,
378}
379
380impl RunningWorkload {
381    /// Create a handle for a workload that's already running externally
382    /// (e.g. embedded in yah-camp). No subprocess is owned; shutdown is a
383    /// no-op so the caller can call `shutdown()` uniformly.
384    pub fn adopted(
385        kind: impl Into<String>,
386        slot: impl Into<String>,
387        dev_url: Option<String>,
388    ) -> Self {
389        Self {
390            kind: kind.into(),
391            slot: slot.into(),
392            dev_url,
393            public_url: None,
394            console_url: None,
395            log_buffer: None,
396            notes: Vec::new(),
397            shutdown: None,
398            supervisor: None,
399            teardown: None,
400        }
401    }
402
403    /// Attach an explicit teardown to a handle for an externally-running
404    /// workload (R714-B1).
405    ///
406    /// `adopted()` alone gives a handle whose `shutdown()` is a documented
407    /// no-op. That is right for workloads another process owns and will keep
408    /// re-asserting (pond containers under camp's yubaba), and wrong for ones
409    /// this process started and nobody else will ever stop — for those the ■
410    /// button reported success while the container kept running.
411    ///
412    /// `teardown` runs on `shutdown()` and NOT on `Drop`; see [`Self::teardown`].
413    pub fn with_teardown<F, Fut>(mut self, teardown: F) -> Self
414    where
415        F: FnOnce() -> Fut + Send + Sync + 'static,
416        Fut: std::future::Future<Output = Result<()>> + Send + 'static,
417    {
418        self.teardown = Some(Teardown(Box::new(move || Box::pin(teardown()))));
419        self
420    }
421
422    /// Attach per-run operator-facing lines (R546-B12). See [`Self::notes`].
423    pub fn with_notes(mut self, notes: Vec<String>) -> Self {
424        self.notes = notes;
425        self
426    }
427
428    /// Set the public URL for a published workload (e.g. `"https://yah.dev"`).
429    pub fn with_public_url(mut self, url: impl Into<String>) -> Self {
430        self.public_url = Some(url.into());
431        self
432    }
433
434    /// Set the console URL for a workload that exposes a secondary local UI
435    /// (e.g. MinIO console on a pond tier).
436    pub fn with_console_url(mut self, url: impl Into<String>) -> Self {
437        self.console_url = Some(url.into());
438        self
439    }
440
441    /// Gracefully tear down: signal the supervisor, await its exit, then run
442    /// any explicit teardown hook.
443    ///
444    /// The hook runs LAST and its error propagates. A stop that could not tear
445    /// the workload down must surface as an error, never as a silent success —
446    /// that silence is the whole of R714-B1.
447    pub async fn shutdown(mut self) -> Result<()> {
448        if let Some(tx) = self.shutdown.take() {
449            let _ = tx.send(());
450        }
451        if let Some(handle) = self.supervisor.take() {
452            handle
453                .await
454                .context("joining workload supervisor")?
455                .context("workload supervisor")?;
456        }
457        if let Some(Teardown(hook)) = self.teardown.take() {
458            hook().await.context("workload teardown")?;
459        }
460        Ok(())
461    }
462}
463
464impl Drop for RunningWorkload {
465    fn drop(&mut self) {
466        // Best-effort signal. The supervisor task is detached and will
467        // reap its child when it observes the closed channel.
468        if let Some(tx) = self.shutdown.take() {
469            let _ = tx.send(());
470        }
471    }
472}
473
474/// Bring one workload up. Each impl handles one [`ServiceComponent::kind`].
475#[async_trait]
476pub trait Reconciler: Send + Sync {
477    /// Workload kind this reconciler handles (matches `ServiceComponent.kind`).
478    fn kind(&self) -> &'static str;
479
480    /// Bring the workload up. Returns a handle whose lifecycle is tied to
481    /// the mirror being up.
482    async fn up(&self, ctx: ReconcileCtx<'_>) -> Result<RunningWorkload>;
483}
484
485/// Read `<workload_dir>/workload.toml` and return just the `kind` field.
486/// See [`ReconcileCtx::workload_kind`] for why we don't deserialize through
487/// the strong types yet.
488pub fn workload_kind(workload_dir: &Path) -> Result<String> {
489    let path = workload_dir.join("workload.toml");
490    let src =
491        std::fs::read_to_string(&path).with_context(|| format!("reading {}", path.display()))?;
492    let value: toml::Value =
493        toml::from_str(&src).with_context(|| format!("parsing {}", path.display()))?;
494    let kind = value
495        .get("kind")
496        .and_then(|v| v.as_str())
497        .with_context(|| format!("{}: missing `kind` field", path.display()))?;
498    Ok(kind.to_string())
499}
500
501/// Shallow-clone (or update) a [`GitSource`] into `dir` (R561-F1). Idempotent:
502/// clones when `dir/.git` is absent, otherwise fetches the pinned ref and
503/// force-checks-it-out. Uses the system `git` so it inherits the operator's
504/// credential helpers / SSH agent — no in-process git library.
505///
506/// `pub` (R615-T3 / W274): `yah infra sync` reuses this verbatim for
507/// `InfraSourceKind::Git` sources rather than a second shallow-clone-or-pull
508/// implementation — same "one git-source shape, reused" discipline R615-F1
509/// already applied to the type.
510pub async fn materialize_git_source(git: &GitSource, dir: &Path) -> Result<()> {
511    use tokio::process::Command;
512
513    async fn run_git(args: &[&std::ffi::OsStr]) -> Result<()> {
514        let out = Command::new("git")
515            .args(args)
516            .output()
517            .await
518            .context("spawning git")?;
519        if !out.status.success() {
520            anyhow::bail!(
521                "git {} failed: {}",
522                args.iter()
523                    .map(|a| a.to_string_lossy())
524                    .collect::<Vec<_>>()
525                    .join(" "),
526                String::from_utf8_lossy(&out.stderr).trim()
527            );
528        }
529        Ok(())
530    }
531
532    use std::ffi::OsStr;
533    let dir_os = dir.as_os_str();
534    let r#ref = git.r#ref.as_str();
535
536    if dir.join(".git").is_dir() {
537        // Existing checkout — update to the pinned ref.
538        run_git(&[
539            OsStr::new("-C"),
540            dir_os,
541            OsStr::new("fetch"),
542            OsStr::new("--depth"),
543            OsStr::new("1"),
544            OsStr::new("origin"),
545            OsStr::new(r#ref),
546        ])
547        .await?;
548        run_git(&[
549            OsStr::new("-C"),
550            dir_os,
551            OsStr::new("checkout"),
552            OsStr::new("--force"),
553            OsStr::new("FETCH_HEAD"),
554        ])
555        .await?;
556    } else {
557        if let Some(parent) = dir.parent() {
558            tokio::fs::create_dir_all(parent)
559                .await
560                .with_context(|| format!("creating {}", parent.display()))?;
561        }
562        // `--branch` accepts a branch or tag. Pinning to a bare commit SHA is a
563        // follow-up (needs clone-then-fetch); the common case is a branch/tag.
564        run_git(&[
565            OsStr::new("clone"),
566            OsStr::new("--depth"),
567            OsStr::new("1"),
568            OsStr::new("--branch"),
569            OsStr::new(r#ref),
570            OsStr::new(git.repo.as_str()),
571            dir_os,
572        ])
573        .await?;
574    }
575    Ok(())
576}
577
578/// Build a `RunningWorkload` from the pieces a reconciler produces.
579pub(crate) fn into_running(
580    kind: impl Into<String>,
581    slot: impl Into<String>,
582    dev_url: Option<String>,
583    public_url: Option<String>,
584    log_buffer: Option<LogBuffer>,
585    shutdown: oneshot::Sender<()>,
586    supervisor: tokio::task::JoinHandle<Result<()>>,
587) -> RunningWorkload {
588    RunningWorkload {
589        kind: kind.into(),
590        slot: slot.into(),
591        dev_url,
592        public_url,
593        console_url: None,
594        log_buffer,
595        notes: Vec::new(),
596        shutdown: Some(shutdown),
597        supervisor: Some(supervisor),
598        teardown: None,
599    }
600}
601
602/// Wait for a TCP port to start accepting connections. Returns `true` if
603/// the port came up within `timeout`, `false` otherwise. Useful for
604/// reconcilers that spawn a server and need to know when it's reachable
605/// before reporting success.
606pub(crate) async fn wait_for_port(
607    addr: std::net::SocketAddr,
608    timeout: std::time::Duration,
609) -> bool {
610    let deadline = tokio::time::Instant::now() + timeout;
611    loop {
612        if tokio::net::TcpStream::connect(addr).await.is_ok() {
613            return true;
614        }
615        if tokio::time::Instant::now() >= deadline {
616            return false;
617        }
618        tokio::time::sleep(std::time::Duration::from_millis(50)).await;
619    }
620}
621
622// `wait_for_http_ready` lived here pre-R374-F3 to back the MinIO health
623// probe in pond's bring-up path. That logic moved to
624// `local_driver::pond_minio::wait_for_http_ready` so yubaba + cloud share
625// it. The mesofact-static reconciler arm uses [`wait_for_port`] for
626// dev-tier port readiness; nothing else needs an HTTP-level probe today.
627
628/// Pluck a `u16` out of a [`MirrorProviderSlot`]'s inline `fields` map.
629/// Returns `None` if the key is absent or out of range.
630pub(crate) fn slot_field_u16(fields: &BTreeMap<String, toml::Value>, key: &str) -> Option<u16> {
631    fields
632        .get(key)
633        .and_then(|v| v.as_integer())
634        .and_then(|n| u16::try_from(n).ok())
635}
636
637/// Serializable summary of a running workload — what the desktop / CLI
638/// hands to the UI. Subset of [`RunningWorkload`] that's safe to cross
639/// process boundaries.
640#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
641pub struct RunningWorkloadSummary {
642    pub kind: String,
643    pub slot: String,
644    pub dev_url: Option<String>,
645    pub public_url: Option<String>,
646    pub console_url: Option<String>,
647}
648
649impl From<&RunningWorkload> for RunningWorkloadSummary {
650    fn from(r: &RunningWorkload) -> Self {
651        Self {
652            kind: r.kind.clone(),
653            slot: r.slot.clone(),
654            dev_url: r.dev_url.clone(),
655            public_url: r.public_url.clone(),
656            console_url: r.console_url.clone(),
657        }
658    }
659}
660
661#[cfg(test)]
662mod source_seam_tests {
663    //! R561-F1 — the BYO-git source seam: path resolution + materialization.
664    use super::*;
665    use std::collections::BTreeMap;
666
667    fn component(git: Option<GitSource>) -> ServiceComponent {
668        ServiceComponent {
669            id: "site".into(),
670            kind: "mesofact-static".into(),
671            path: "site".into(),
672            role: "static".into(),
673            publishes: None,
674            wave: 0,
675            git,
676        }
677    }
678
679    fn service(comp: ServiceComponent) -> ServiceConfig {
680        ServiceConfig {
681            schema_version: 1,
682            name: "scrabcake".into(),
683            domain: "scrabcake.example".into(),
684            components: vec![comp],
685            db: crate::DbCatalog::default(),
686        }
687    }
688
689    fn mirror() -> MirrorConfig {
690        MirrorConfig {
691            schema_version: 1,
692            shape: crate::MirrorShape::Local,
693            providers: BTreeMap::new(),
694            ingress: Default::default(),
695            drivers: Default::default(),
696            asset_aliases: BTreeMap::new(),
697        }
698    }
699
700    fn ctx<'a>(ws: &'a Path, svc: &'a ServiceConfig, mir: &'a MirrorConfig) -> ReconcileCtx<'a> {
701        ReconcileCtx {
702            workspace_root: ws,
703            service: svc,
704            component: &svc.components[0],
705            mirror: mir,
706            env: "dev",
707            scope: ProviderScope::singleton(),
708        }
709    }
710
711    #[test]
712    fn workload_dir_in_tree_joins_workspace_root() {
713        let svc = service(component(None));
714        let mir = mirror();
715        assert_eq!(
716            ctx(Path::new("/ws"), &svc, &mir).workload_dir(),
717            Path::new("/ws/site")
718        );
719    }
720
721    #[test]
722    fn workload_dir_git_resolves_into_source_cache_with_subdir() {
723        let git = GitSource {
724            repo: "https://example.com/r.git".into(),
725            r#ref: "main".into(),
726            subdir: Some("apps".into()),
727        };
728        let svc = service(component(Some(git)));
729        let mir = mirror();
730        assert_eq!(
731            ctx(Path::new("/ws"), &svc, &mir).workload_dir(),
732            Path::new("/ws/.yah/infra/state/sources/scrabcake/site/apps/site")
733        );
734    }
735
736    #[tokio::test]
737    async fn materialize_is_noop_for_in_tree_component() {
738        let svc = service(component(None));
739        let mir = mirror();
740        // No git source → Ok, and nothing is written under the workspace.
741        ctx(Path::new("/nonexistent-ws"), &svc, &mir)
742            .materialize()
743            .await
744            .unwrap();
745    }
746
747    #[tokio::test]
748    async fn materialize_clones_git_source_offline() {
749        fn git(args: &[&str], cwd: &Path) {
750            let out = std::process::Command::new("git")
751                .args(args)
752                .current_dir(cwd)
753                .env("GIT_AUTHOR_NAME", "t")
754                .env("GIT_AUTHOR_EMAIL", "t@t")
755                .env("GIT_COMMITTER_NAME", "t")
756                .env("GIT_COMMITTER_EMAIL", "t@t")
757                .output()
758                .unwrap();
759            assert!(
760                out.status.success(),
761                "git {args:?}: {}",
762                String::from_utf8_lossy(&out.stderr)
763            );
764        }
765
766        let tmp = tempfile::tempdir().unwrap();
767        let src = tmp.path().join("src-repo");
768        std::fs::create_dir_all(&src).unwrap();
769        git(&["init", "-b", "main"], &src);
770        std::fs::write(src.join("hello.txt"), "hi").unwrap();
771        git(&["add", "."], &src);
772        git(&["commit", "-m", "init"], &src);
773
774        let source = GitSource {
775            repo: format!("file://{}", src.display()),
776            r#ref: "main".into(),
777            subdir: None,
778        };
779        let dest = tmp.path().join("cache");
780
781        // First call clones.
782        materialize_git_source(&source, &dest).await.unwrap();
783        assert!(dest.join("hello.txt").is_file());
784
785        // Second call takes the update path and stays green (idempotent).
786        materialize_git_source(&source, &dest).await.unwrap();
787        assert!(dest.join("hello.txt").is_file());
788    }
789}
790
791#[cfg(test)]
792mod teardown_tests {
793    //! R714-B1 — the explicit-teardown contract on [`RunningWorkload`].
794    //!
795    //! These are about WHEN the hook runs, not what it does. The bug being
796    //! fixed was a `shutdown()` that reported success having done nothing, and
797    //! the trap in fixing it is a `Drop` that tears down a container which is
798    //! supposed to survive the process.
799    use super::*;
800    use std::sync::atomic::{AtomicUsize, Ordering};
801
802    fn counting() -> (RunningWorkload, Arc<AtomicUsize>) {
803        let hits = Arc::new(AtomicUsize::new(0));
804        let seen = hits.clone();
805        let w = RunningWorkload::adopted("container", "compute", None)
806            .with_teardown(move || {
807                let seen = seen.clone();
808                async move {
809                    seen.fetch_add(1, Ordering::SeqCst);
810                    Ok(())
811                }
812            });
813        (w, hits)
814    }
815
816    /// Attaching a teardown must not cost `RunningWorkload` its auto traits.
817    /// The desktop stores these handles in a shared registry and its Tauri
818    /// commands iterate them across `.await` points, which needs `Sync` — a
819    /// hook that is `Send` but not `Sync` takes it away here and surfaces two
820    /// crates over as "future cannot be sent between threads safely", with a
821    /// span pointing at `mirror_run_logs` rather than at this file.
822    #[test]
823    fn a_teardown_does_not_cost_the_handle_send_or_sync() {
824        fn assert_send_sync<T: Send + Sync>() {}
825        assert_send_sync::<RunningWorkload>();
826    }
827
828    #[tokio::test]
829    async fn shutdown_runs_the_teardown() {
830        let (w, hits) = counting();
831        w.shutdown().await.unwrap();
832        assert_eq!(hits.load(Ordering::SeqCst), 1);
833    }
834
835    #[tokio::test]
836    async fn dropping_the_handle_does_NOT_run_the_teardown() {
837        // A container this process started is meant to outlive it — the next
838        // launch re-adopts it. Quitting the app must not `docker rm -f` it.
839        let (w, hits) = counting();
840        drop(w);
841        // Yield so a stray spawned task would have had a chance to run.
842        tokio::task::yield_now().await;
843        assert_eq!(hits.load(Ordering::SeqCst), 0);
844    }
845
846    #[tokio::test]
847    async fn a_failing_teardown_makes_shutdown_fail() {
848        // The whole of R714-B1: a stop that could not tear the workload down
849        // must not report success.
850        let w = RunningWorkload::adopted("container", "compute", None)
851            .with_teardown(|| async { anyhow::bail!("docker stop refused") });
852        let err = w.shutdown().await.unwrap_err();
853        assert!(format!("{err:#}").contains("docker stop refused"), "{err:#}");
854    }
855
856    #[tokio::test]
857    async fn an_adopted_handle_without_a_teardown_still_shuts_down_cleanly() {
858        // Pond containers are owned by camp's yubaba and stop through
859        // `workload.stop`; their no-op shutdown is correct and must stay.
860        let w = RunningWorkload::adopted("mesofact-static", "static", None);
861        w.shutdown().await.unwrap();
862    }
863
864    #[tokio::test]
865    async fn the_teardown_runs_after_the_supervisor_is_joined() {
866        // Ordering matters for a workload that has both: reap the child first,
867        // then remove the container it was talking to.
868        let order = Arc::new(AsyncMutex::new(Vec::<&'static str>::new()));
869        let (tx, rx) = oneshot::channel::<()>();
870        let sup_order = order.clone();
871        let supervisor = tokio::spawn(async move {
872            let _ = rx.await;
873            sup_order.lock().await.push("supervisor");
874            Ok(())
875        });
876        let hook_order = order.clone();
877        let w = into_running("container", "compute", None, None, None, tx, supervisor)
878            .with_teardown(move || {
879                let hook_order = hook_order.clone();
880                async move {
881                    hook_order.lock().await.push("teardown");
882                    Ok(())
883                }
884            });
885
886        w.shutdown().await.unwrap();
887        assert_eq!(*order.lock().await, vec!["supervisor", "teardown"]);
888    }
889}