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