Skip to main content

toolkit/runtime/
host_runtime.rs

1//! Host Runtime - orchestrates the full `ToolKit` lifecycle
2//!
3//! This gear contains the `HostRuntime` type that owns and coordinates
4//! the execution of all lifecycle phases.
5//!
6//! High-level phase order:
7//! - `pre_init` (system gears only)
8//! - DB migrations (gears with DB capability)
9//! - `init` (all gears)
10//! - proxy-wiring (`#[toolkit::consumes]` clients; feature-gated)
11//! - `post_init` (system gears only; runs after *all* `init` complete)
12//! - REST wiring (gears with REST capability; requires a single REST host)
13//! - gRPC registration (gears with gRPC capability; requires a single gRPC hub)
14//! - start/stop (stateful gears)
15//! - `OoP` spawn / wait / stop (host-only orchestration)
16//!
17//! Both lifecycle paths — in-process (`run_phases_internal`) and `OoP`
18//! (`run_oop_serving`) — run these in the same relative order. Proxy-wiring in
19//! particular must stay before `start`: a gear resolving a consumed contract
20//! during its own `start` has to find the client in the `ClientHub` regardless
21//! of profile. The `OoP` path additionally has no REST-wiring, directory-register
22//! or spawn phases (`oop_serve` owns those concerns).
23
24use axum::Router;
25use std::collections::HashSet;
26use std::sync::Arc;
27
28use tokio_util::sync::CancellationToken;
29use uuid::Uuid;
30
31use crate::backends::OopSpawnConfig;
32use crate::client_hub::ClientHub;
33use crate::config::ConfigProvider;
34use crate::context::GearContextBuilder;
35use crate::registry::{
36    ApiGatewayCap, GearEntry, GearRegistry, GrpcHubCap, RegistryError, RestApiCap, RunnableCap,
37    SystemCap,
38};
39use crate::runtime::{GearManager, GrpcInstallerStore, OopSpawnOptions, SystemContext};
40
41#[cfg(feature = "db")]
42use crate::registry::DatabaseCap;
43
44/// How the runtime should provide DBs to gears.
45#[derive(Clone)]
46pub enum DbOptions {
47    /// No database integration. `GearCtx::db()` will be `None`, `db_required()` will error.
48    None,
49    /// Use a `DbManager` to handle database connections with Figment-based configuration.
50    #[cfg(feature = "db")]
51    Manager(Arc<toolkit_db::DbManager>),
52}
53
54/// Runtime execution mode that determines which phases to run.
55#[derive(Debug, Clone, Copy, PartialEq, Eq)]
56pub enum RunMode {
57    /// Run all phases and wait for shutdown signal (normal application mode).
58    Full,
59    /// Run only pre-init and DB migration phases, then exit (for cloud deployments).
60    MigrateOnly,
61}
62
63/// Environment variable name for passing directory endpoint to `OoP` gears.
64pub const TOOLKIT_DIRECTORY_ENDPOINT_ENV: &str = "TOOLKIT_DIRECTORY_ENDPOINT";
65
66/// Environment variable name for passing rendered gear config to `OoP` gears.
67pub const TOOLKIT_MODULE_CONFIG_ENV: &str = "TOOLKIT_MODULE_CONFIG";
68
69/// Default shutdown deadline for graceful gear stop (35 seconds).
70///
71/// This is intentionally 5 seconds longer than `WithLifecycle::stop_timeout` (30s default)
72/// to ensure deterministic behavior: the lifecycle's internal timeout fires first,
73/// and the runtime deadline acts as a hard backstop.
74pub const DEFAULT_SHUTDOWN_DEADLINE: std::time::Duration = std::time::Duration::from_secs(35);
75
76/// `HostRuntime` owns the lifecycle orchestration for `ToolKit`.
77///
78/// It encapsulates all runtime state and drives gears through the full lifecycle (see gear docs).
79/// Read a consumer's ADR-0004 static-endpoint override for `dep_gear` from
80/// `gears.<owner_gear>.config.consumer_wiring.<dep_gear>` (a base endpoint URI
81/// string). This is the dev/test escape hatch that bypasses service discovery;
82/// returns `None` when unset.
83fn static_endpoint_override(
84    cfg: &dyn ConfigProvider,
85    owner_gear: &str,
86    dep_gear: &str,
87) -> Option<String> {
88    cfg.get_gear_config(owner_gear)?
89        .get("config")?
90        .get("consumer_wiring")?
91        .get(dep_gear)?
92        .as_str()
93        .map(str::to_owned)
94}
95
96pub struct HostRuntime {
97    registry: GearRegistry,
98    ctx_builder: GearContextBuilder,
99    instance_id: Uuid,
100    gear_manager: Arc<GearManager>,
101    grpc_installers: Arc<GrpcInstallerStore>,
102    client_hub: Arc<ClientHub>,
103    /// Per-gear config, retained for the proxy-wiring phase to read a consumer's
104    /// static-endpoint override (dev/test escape hatch, ADR-0004).
105    gears_cfg: Arc<dyn ConfigProvider>,
106    /// Process-level dependency-resolution + draining signal, published in
107    /// `client_hub` for the `/readyz` probe and updated by the proxy-wiring
108    /// readiness loop + draining watcher. In-process, an
109    /// [`ReadinessHealthcheck`](super::readiness::ReadinessHealthcheck) leaf
110    /// bridges it into the gateway's healthcheck registry.
111    dep_checker: Arc<super::readiness::DependencyChecker>,
112    /// Set once the in-process directory-register phase has advertised this
113    /// process's REST providers, so shutdown deregisters them exactly once — and
114    /// only in the in-process host path. In `OoP` serving, presence + deregister
115    /// is owned by `oop_serve`, so this stays `false` and avoids a double
116    /// deregister.
117    rest_providers_registered: std::sync::atomic::AtomicBool,
118    cancel: CancellationToken,
119    #[allow(dead_code)]
120    db_options: DbOptions,
121    /// `OoP` gear spawn configuration and backend
122    oop_options: Option<OopSpawnOptions>,
123    /// Maximum time allowed for graceful shutdown before hard-stop signal is sent.
124    shutdown_deadline: std::time::Duration,
125}
126
127impl HostRuntime {
128    /// Create a new `HostRuntime` instance.
129    ///
130    /// This prepares all runtime components but does not start any lifecycle phases.
131    pub fn new(
132        registry: GearRegistry,
133        gears_cfg: Arc<dyn ConfigProvider>,
134        db_options: DbOptions,
135        client_hub: Arc<ClientHub>,
136        cancel: CancellationToken,
137        instance_id: Uuid,
138        oop_options: Option<OopSpawnOptions>,
139    ) -> Self {
140        // Create runtime-owned components for system gears
141        let gear_manager = Arc::new(GearManager::new());
142        let grpc_installers = Arc::new(GrpcInstallerStore::new());
143
144        // Process-level dependency/draining signal, published so the gateway's
145        // /readyz handler can fetch it (concrete-type key). Created before any
146        // phase runs.
147        let dep_checker = Arc::new(super::readiness::DependencyChecker::new());
148        client_hub.register::<super::readiness::DependencyChecker>(dep_checker.clone());
149
150        // Build the context builder that will resolve per-gear DbHandles
151        let ctx_builder = GearContextBuilder::new(
152            instance_id,
153            gears_cfg.clone(),
154            client_hub.clone(),
155            cancel.clone(),
156        );
157        #[cfg(feature = "db")]
158        let ctx_builder = match &db_options {
159            DbOptions::Manager(mgr) => ctx_builder.with_db_manager(mgr.clone()),
160            DbOptions::None => ctx_builder,
161        };
162
163        Self {
164            registry,
165            ctx_builder,
166            instance_id,
167            gear_manager,
168            grpc_installers,
169            client_hub,
170            gears_cfg,
171            dep_checker,
172            rest_providers_registered: std::sync::atomic::AtomicBool::new(false),
173            cancel,
174            db_options,
175            oop_options,
176            shutdown_deadline: DEFAULT_SHUTDOWN_DEADLINE,
177        }
178    }
179
180    /// Set a custom shutdown deadline for graceful gear stop.
181    ///
182    /// This is the maximum time the runtime will wait for each gear to stop gracefully
183    /// before sending the hard-stop signal (cancelling the deadline token).
184    ///
185    /// # Relationship with `WithLifecycle::stop_timeout`
186    ///
187    /// When using `WithLifecycle`, its `stop_timeout` (default 30s) races against this
188    /// `shutdown_deadline` (also default 30s). To ensure deterministic behavior:
189    ///
190    /// - `WithLifecycle::stop_timeout` should be **less than** `shutdown_deadline`
191    /// - This allows the lifecycle's internal timeout to trigger first for graceful cleanup
192    /// - The runtime's `deadline_token` then acts as a hard backstop
193    ///
194    /// Example: `stop_timeout = 25s`, `shutdown_deadline = 30s`
195    #[must_use]
196    pub fn with_shutdown_deadline(mut self, deadline: std::time::Duration) -> Self {
197        self.shutdown_deadline = deadline;
198        self
199    }
200
201    /// `PRE_INIT` phase: wire runtime internals into system gears.
202    ///
203    /// This phase runs before init and only for gears with the "system" capability.
204    ///
205    /// # Errors
206    /// Returns `RegistryError` if system wiring fails.
207    pub fn run_pre_init_phase(&self) -> Result<(), RegistryError> {
208        tracing::info!("Phase: pre_init");
209
210        let sys_ctx = SystemContext::new(
211            self.instance_id,
212            Arc::clone(&self.gear_manager),
213            Arc::clone(&self.grpc_installers),
214        );
215
216        for entry in self.registry.gears() {
217            // Check for cancellation before processing each gear
218            if self.cancel.is_cancelled() {
219                tracing::warn!("Pre-init phase cancelled by signal");
220                return Err(RegistryError::Cancelled);
221            }
222
223            if let Some(sys_mod) = entry.caps.query::<SystemCap>() {
224                tracing::debug!(gear = entry.name, "Running system pre_init");
225                sys_mod
226                    .pre_init(&sys_ctx)
227                    .map_err(|e| RegistryError::PreInit {
228                        gear: entry.name,
229                        source: e,
230                    })?;
231            }
232        }
233
234        Ok(())
235    }
236
237    /// Helper: resolve context for a gear with error mapping.
238    #[cfg(feature = "db")]
239    async fn gear_context(
240        &self,
241        gear_name: &'static str,
242    ) -> Result<crate::context::GearCtx, RegistryError> {
243        self.ctx_builder
244            .for_gear(gear_name)
245            .await
246            .map_err(|e| RegistryError::DbMigrate {
247                gear: gear_name,
248                source: e,
249            })
250    }
251
252    /// Helper: extract DB handle and gear if both exist.
253    #[cfg(feature = "db")]
254    async fn db_migration_target(
255        &self,
256        gear_name: &'static str,
257        ctx: &crate::context::GearCtx,
258        db_gear: Option<Arc<dyn crate::contracts::DatabaseCapability>>,
259    ) -> Result<
260        Option<(
261            toolkit_db::Db,
262            Arc<dyn crate::contracts::DatabaseCapability>,
263        )>,
264        RegistryError,
265    > {
266        let Some(dbm) = db_gear else {
267            return Ok(None);
268        };
269
270        // Important: DB migrations require access to the underlying `Db`, not just `DBProvider`.
271        // `GearCtx` intentionally exposes only `DBProvider` for better DX and to reduce mistakes.
272        // So the runtime resolves the `Db` directly from its `DbManager`.
273        let db = match &self.db_options {
274            DbOptions::None => None,
275            #[cfg(feature = "db")]
276            DbOptions::Manager(mgr) => {
277                mgr.get(gear_name)
278                    .await
279                    .map_err(|e| RegistryError::DbMigrate {
280                        gear: gear_name,
281                        source: e.into(),
282                    })?
283            }
284        };
285
286        _ = ctx; // ctx is kept for parity/error context; DB is resolved from manager above.
287        Ok(db.map(|db| (db, dbm)))
288    }
289
290    /// Helper: run migrations for a single gear using the new migration runner.
291    ///
292    /// This collects migrations from the gear and executes them via the
293    /// runtime's privileged connection. Gears never see the raw connection.
294    #[cfg(feature = "db")]
295    async fn migrate_gear(
296        gear_name: &'static str,
297        db: &toolkit_db::Db,
298        db_gear: Arc<dyn crate::contracts::DatabaseCapability>,
299    ) -> Result<(), RegistryError> {
300        // Collect migrations from the gear
301        let migrations = db_gear.migrations();
302
303        if migrations.is_empty() {
304            tracing::debug!(gear = gear_name, "No migrations to run");
305            return Ok(());
306        }
307
308        tracing::debug!(
309            gear = gear_name,
310            count = migrations.len(),
311            "Running DB migrations"
312        );
313
314        // Execute migrations using the migration runner
315        let result =
316            toolkit_db::migration_runner::run_migrations_for_gear(db, gear_name, migrations)
317                .await
318                .map_err(|e| RegistryError::DbMigrate {
319                    gear: gear_name,
320                    source: anyhow::Error::new(e),
321                })?;
322
323        tracing::info!(
324            gear = gear_name,
325            applied = result.applied,
326            skipped = result.skipped,
327            "DB migrations completed"
328        );
329
330        Ok(())
331    }
332
333    /// DB MIGRATION phase: run migrations for all gears with DB capability.
334    ///
335    /// Runs before init, with system gears processed first.
336    ///
337    /// Gears provide migrations via `DatabaseCapability::migrations()`.
338    /// The runtime executes them with a privileged connection that gears
339    /// never receive directly. Each gear gets a separate migration history
340    /// table, preventing cross-gear interference.
341    #[cfg(feature = "db")]
342    async fn run_db_phase(&self) -> Result<(), RegistryError> {
343        tracing::info!("Phase: db (before init)");
344
345        for entry in self.registry.gears_by_system_priority() {
346            // Check for cancellation before processing each gear
347            if self.cancel.is_cancelled() {
348                tracing::warn!("DB migration phase cancelled by signal");
349                return Err(RegistryError::Cancelled);
350            }
351
352            let ctx = self.gear_context(entry.name).await?;
353            let db_gear = entry.caps.query::<DatabaseCap>();
354
355            match self
356                .db_migration_target(entry.name, &ctx, db_gear.clone())
357                .await?
358            {
359                Some((db, dbm)) => {
360                    Self::migrate_gear(entry.name, &db, dbm).await?;
361                }
362                None if db_gear.is_some() => {
363                    tracing::debug!(
364                        gear = entry.name,
365                        "Gear has DbGear trait but no DB handle (no config)"
366                    );
367                }
368                None => {}
369            }
370        }
371
372        Ok(())
373    }
374
375    /// INIT phase: initialize all gears in topological order.
376    ///
377    /// System gears initialize first, followed by user gears.
378    async fn run_init_phase(&self) -> Result<(), RegistryError> {
379        tracing::info!("Phase: init");
380
381        for entry in self.registry.gears_by_system_priority() {
382            let ctx =
383                self.ctx_builder
384                    .for_gear(entry.name)
385                    .await
386                    .map_err(|e| RegistryError::Init {
387                        gear: entry.name,
388                        source: e,
389                    })?;
390            tracing::info!(gear = entry.name, "Initializing a gear...");
391            entry
392                .core
393                .init(&ctx)
394                .await
395                .map_err(|e| RegistryError::Init {
396                    gear: entry.name,
397                    source: e,
398                })?;
399            tracing::info!(gear = entry.name, "Initialized a gear.");
400        }
401
402        Ok(())
403    }
404
405    /// `POST_INIT` phase: optional hook after ALL gears completed `init()`.
406    ///
407    /// Consumer proxy-wiring phase (eventual readiness).
408    ///
409    /// Runs after init (compile-time / local registrations) and before
410    /// post-init. Replays each `ConsumerRegistration` emitted by
411    /// `#[toolkit::consumes]`: if a compile-time impl is already in the
412    /// `ClientHub` it wins (the wiring closure short-circuits); otherwise a
413    /// directory-resolving REST client is registered under the contract trait.
414    ///
415    /// Non-blocking: endpoint discovery is lazy/per-call inside the resolving
416    /// client, so this phase never waits on provider availability (ADR-0007).
417    /// A no-op when no consumer is registered, preserving the phase-order
418    /// invariants relied on by existing tests.
419    #[allow(
420        clippy::unused_async,
421        reason = "kept async for symmetry with the other `run_*_phase` steps awaited in sequence by `run_gear_phases`; the awaited work runs in a spawned readiness-probe task"
422    )]
423    async fn run_proxy_wiring_phase(&self) -> Result<(), RegistryError> {
424        use crate::discovery::{
425            ConsumerRegistration, DirectoryEndpointResolver, NullEndpointResolver,
426        };
427        use toolkit_contract::runtime::resolving::EndpointResolver;
428
429        let regs: Vec<&ConsumerRegistration> = inventory::iter::<ConsumerRegistration>
430            .into_iter()
431            .collect();
432        if regs.is_empty() {
433            return Ok(());
434        }
435        tracing::info!(
436            count = regs.len(),
437            "Phase: proxy-wiring (consumer discovery)"
438        );
439
440        // Without a DirectoryClient we cannot resolve *remote* providers, but we
441        // must NOT silently skip wiring: co-located consumers still short-circuit
442        // to their local impl, and remote consumers must register as unresolved
443        // readiness gates so `/readyz` stays 503 (a misconfigured consumer must
444        // not report Ready). A null resolver makes every remote lookup `Ok(None)`
445        // so the loop below classifies deps correctly without a directory.
446        let (resolver, have_directory): (Arc<dyn EndpointResolver>, bool) =
447            if let Ok(dir) = self.client_hub.get::<dyn crate::DirectoryClient>() {
448                (Arc::new(DirectoryEndpointResolver::new(dir)), true)
449            } else {
450                tracing::error!(
451                    consumers = regs.len(),
452                    "proxy-wiring: no DirectoryClient in ClientHub; remote consumer \
453                     dependencies cannot be resolved and will gate /readyz (503). \
454                     Co-located (local) dependencies are unaffected."
455                );
456                (Arc::new(NullEndpointResolver), false)
457            };
458
459        // Wire each consumer contract. The outcome distinguishes a co-located
460        // local impl (hub short-circuit) from a directory-resolving REST client.
461        // Every dep is registered as a readiness gate; local ones are marked
462        // resolved immediately, and only remote ones gate readiness + get the
463        // background directory-resolve loop (ADR-0007: startup-gating + sticky).
464        // `owner_gear` is derived by `#[toolkit::consumes]` from the struct
465        // ident, while `#[toolkit::gear(name = ...)]` sets the registry name
466        // independently. When they diverge, wiring still works (the loop below
467        // does not filter on owner) but the static-override config key silently
468        // resolves to nothing. Say so rather than leaving the operator to wonder
469        // why `consumer_wiring` is ignored.
470        let known_gears: std::collections::HashSet<&str> =
471            self.registry.gears().iter().map(GearEntry::name).collect();
472        for reg in &regs {
473            if !known_gears.contains(reg.owner_gear) {
474                tracing::warn!(
475                    owner = reg.owner_gear,
476                    dep = reg.dep_gear,
477                    "proxy-wiring: consumer's owner gear name does not match any registered gear; \
478                     the `gears.{}.config.consumer_wiring.{}` static override will never resolve. \
479                     Rename the gear to the kebab-case of its struct ident.",
480                    reg.owner_gear,
481                    reg.dep_gear,
482                );
483            }
484        }
485
486        let mut remote_deps: Vec<String> = Vec::new();
487        for reg in &regs {
488            // ADR-0004 static-endpoint override (dev/test escape hatch): if the
489            // consumer's config declares a fixed endpoint for this dep, wire it
490            // directly (bypassing discovery) via a `StaticEndpointResolver`. A
491            // fixed endpoint needs no probe loop, so it is readiness-resolved
492            // immediately. Emitted at `warn!` — it must not be used in production.
493            let static_override =
494                static_endpoint_override(self.gears_cfg.as_ref(), reg.owner_gear, reg.dep_gear);
495            let (reg_resolver, is_static): (Arc<dyn EndpointResolver>, bool) =
496                if let Some(endpoint) = &static_override {
497                    tracing::warn!(
498                        owner = reg.owner_gear,
499                        dep = reg.dep_gear,
500                        endpoint = %endpoint,
501                        "proxy-wiring: STATIC endpoint override in use (ADR-0004 dev/test \
502                         escape hatch) - bypasses service discovery; MUST NOT be used in \
503                         production"
504                    );
505                    (
506                        Arc::new(crate::discovery::StaticEndpointResolver::new(
507                            endpoint.clone(),
508                        )),
509                        true,
510                    )
511                } else {
512                    (Arc::clone(&resolver), false)
513                };
514
515            let outcome = (reg.wire)(&self.client_hub, reg_resolver).map_err(|source| {
516                RegistryError::ProxyWiring {
517                    gear: reg.owner_gear,
518                    source,
519                }
520            })?;
521            self.dep_checker.register_dep(reg.dep_gear.to_owned());
522            match outcome {
523                // Local impl won the hub short-circuit — resolved.
524                crate::discovery::WireOutcome::Local => {
525                    self.dep_checker.mark_resolved(reg.dep_gear);
526                }
527                // Static override → fixed endpoint, always resolvable — no probe.
528                crate::discovery::WireOutcome::Remote if is_static => {
529                    self.dep_checker.mark_resolved(reg.dep_gear);
530                }
531                // Directory-resolved remote → gate readiness + background probe.
532                crate::discovery::WireOutcome::Remote => remote_deps.push(reg.dep_gear.to_owned()),
533            }
534            tracing::debug!(
535                owner = reg.owner_gear,
536                dep = reg.dep_gear,
537                outcome = ?outcome,
538                static_override = is_static,
539                "wired consumer contract"
540            );
541        }
542
543        // Only remote deps need directory resolution; local wins are already
544        // resolved above. Without a directory the null resolver can never resolve
545        // them, so skip the probe entirely and leave those deps gating /readyz.
546        if !have_directory || remote_deps.is_empty() {
547            return Ok(());
548        }
549
550        let readiness = Arc::clone(&self.dep_checker);
551        let cancel = self.cancel.clone();
552        tokio::spawn(async move {
553            const BASE: std::time::Duration = std::time::Duration::from_millis(100);
554            const MAX: std::time::Duration = std::time::Duration::from_secs(30);
555            let mut pending = remote_deps;
556            let mut backoff = BASE;
557            while !pending.is_empty() {
558                let mut still_pending = Vec::new();
559                for dep in pending {
560                    match resolver.resolve_endpoint(&dep).await {
561                        Ok(Some(_)) => {
562                            readiness.mark_resolved(&dep);
563                            tracing::info!(dep = %dep, "readiness: dependency resolved");
564                        }
565                        // No live instance yet — expected during startup churn.
566                        Ok(None) => still_pending.push(dep),
567                        // Genuine directory-backend failure: surface it (a stuck
568                        // Starting / 503 pod otherwise has no diagnostic trail).
569                        Err(e) => {
570                            tracing::warn!(dep = %dep, error = %e, "readiness: directory lookup failed");
571                            still_pending.push(dep);
572                        }
573                    }
574                }
575                pending = still_pending;
576                if pending.is_empty() {
577                    break;
578                }
579                tokio::select! {
580                    () = cancel.cancelled() => break,
581                    () = tokio::time::sleep(backoff) => {}
582                }
583                backoff = (backoff * 2).min(MAX);
584            }
585        });
586
587        Ok(())
588    }
589
590    /// This provides a global barrier between initialization-time registration
591    /// and subsequent phases that may rely on a fully-populated runtime registry.
592    /// `init` -> proxy-wiring -> `post_init`, the segment both lifecycle paths
593    /// share.
594    ///
595    /// Extracted so the two paths cannot disagree on it. They did once: the
596    /// `OoP` path ran proxy-wiring *after* `start`, because the call was dropped
597    /// where a retired untyped `resolve_deps` stopgap used to sit rather than
598    /// chosen deliberately. A gear resolving a consumed contract during its own
599    /// `start` then found nothing in the `ClientHub` under Profile 2/3 while
600    /// working fine in Profile 1.
601    ///
602    /// Everything after this segment legitimately differs — the in-process path
603    /// composes a REST router, the `OoP` path leaves that to `oop_serve`.
604    ///
605    /// Keeping this a single private method is what enforces the invariant: if
606    /// a future change inlines the three calls back into both paths, this method
607    /// becomes dead code and the workspace's `-D warnings` build fails.
608    async fn run_init_wiring_post_init(&self) -> Result<(), RegistryError> {
609        self.run_init_phase().await?;
610
611        self.run_proxy_wiring_phase().await?;
612
613        self.run_post_init_phase().await
614    }
615
616    ///
617    /// System gears run first, followed by user gears, preserving topo order.
618    async fn run_post_init_phase(&self) -> Result<(), RegistryError> {
619        tracing::info!("Phase: post_init");
620
621        let sys_ctx = SystemContext::new(
622            self.instance_id,
623            Arc::clone(&self.gear_manager),
624            Arc::clone(&self.grpc_installers),
625        );
626
627        for entry in self.registry.gears_by_system_priority() {
628            if let Some(sys_mod) = entry.caps.query::<SystemCap>() {
629                sys_mod
630                    .post_init(&sys_ctx)
631                    .await
632                    .map_err(|e| RegistryError::PostInit {
633                        gear: entry.name,
634                        source: e,
635                    })?;
636            }
637        }
638
639        Ok(())
640    }
641
642    /// REST phase: compose the router against the REST host.
643    ///
644    /// This is a synchronous phase that builds the final Router by:
645    /// 1. Preparing the host gear
646    /// 2. Registering all REST providers
647    /// 3. Finalizing with `OpenAPI` endpoints
648    async fn run_rest_phase(&self) -> Result<Router, RegistryError> {
649        tracing::info!("Phase: rest (sync)");
650
651        let mut router = Router::new();
652
653        // Find host(s) and whether any rest gears exist
654        let host_count = self
655            .registry
656            .gears()
657            .iter()
658            .filter(|e| e.caps.has::<ApiGatewayCap>())
659            .count();
660
661        match host_count {
662            0 => {
663                return if self
664                    .registry
665                    .gears()
666                    .iter()
667                    .any(|e| e.caps.has::<RestApiCap>())
668                {
669                    Err(RegistryError::RestRequiresHost)
670                } else {
671                    Ok(router)
672                };
673            }
674            1 => { /* proceed */ }
675            _ => return Err(RegistryError::MultipleRestHosts),
676        }
677
678        // Resolve the single host entry and its gear context
679        let host_idx = self
680            .registry
681            .gears()
682            .iter()
683            .position(|e| e.caps.has::<ApiGatewayCap>())
684            .ok_or(RegistryError::RestHostNotFoundAfterValidation)?;
685        let host_entry = &self.registry.gears()[host_idx];
686        let Some(host) = host_entry.caps.query::<ApiGatewayCap>() else {
687            return Err(RegistryError::RestHostMissingFromEntry);
688        };
689        let host_ctx = self
690            .ctx_builder
691            .for_gear(host_entry.name)
692            .await
693            .map_err(|e| RegistryError::RestPrepare {
694                gear: host_entry.name,
695                source: e,
696            })?;
697
698        // use host as the registry
699        let registry: &dyn crate::contracts::OpenApiRegistry = host.as_registry();
700
701        // Healthcheck registry, passed explicitly to the REST host and providers below
702        // (not via ClientHub). Seeded with the host's shutdown token so in-flight checks
703        // are aborted on shutdown.
704        let hc_registry = Arc::new(
705            crate::healthcheck::RestHealthcheckRegistry::with_cancellation(
706                host_ctx.cancellation_token().clone(),
707            ),
708        );
709
710        // Bridge the process-level eventual-readiness state into the served
711        // probe: while any consumed dependency is unresolved (or the process is
712        // draining) this synthetic check reports Unhealthy, so the gateway's
713        // `/readyz` returns 503 until all `#[toolkit::consumes]` deps are wired.
714        // (`/healthz` is a static liveness handler and is unaffected.)
715        hc_registry.register(
716            "readiness",
717            Arc::new(super::readiness::ReadinessHealthcheck::new(
718                self.dep_checker.clone(),
719            )),
720        );
721
722        // 1) Host prepare: base Router / global middlewares / basic OAS meta
723        router = host
724            .rest_prepare(&host_ctx, router, hc_registry.clone())
725            .map_err(|source| RegistryError::RestPrepare {
726                gear: host_entry.name,
727                source,
728            })?;
729
730        // 2) Register all REST providers (in the current discovery order)
731        for e in self.registry.gears() {
732            if let Some(rest) = e.caps.query::<RestApiCap>() {
733                let ctx = self.ctx_builder.for_gear(e.name).await.map_err(|err| {
734                    RegistryError::RestRegister {
735                        gear: e.name,
736                        source: err,
737                    }
738                })?;
739
740                router = rest
741                    .register_rest(&ctx, router, registry)
742                    .map_err(|source| RegistryError::RestRegister {
743                        gear: e.name,
744                        source,
745                    })?;
746
747                // Register the gear's readiness healthcheck after successful route registration.
748                if let Some(hc) = rest.healthcheck(&ctx) {
749                    hc_registry.register(e.name, hc);
750                }
751            }
752        }
753
754        // 3) Host finalize: attach /openapi.json and /docs, persist Router if needed (no server start)
755        router = host
756            .rest_finalize(&host_ctx, router, hc_registry)
757            .map_err(|source| RegistryError::RestFinalize {
758                gear: host_entry.name,
759                source,
760            })?;
761
762        Ok(router)
763    }
764
765    /// gRPC registration phase: collect services from all grpc gears.
766    ///
767    /// Services are stored in the installer store for the `grpc-hub` to consume during start.
768    async fn run_grpc_phase(&self) -> Result<(), RegistryError> {
769        tracing::info!("Phase: grpc (registration)");
770
771        // If no grpc_hub and no grpc_services, skip the phase
772        if self.registry.grpc_hub.is_none() && self.registry.grpc_services.is_empty() {
773            return Ok(());
774        }
775
776        // If there are grpc_services but no hub, that's an error
777        if self.registry.grpc_hub.is_none() && !self.registry.grpc_services.is_empty() {
778            return Err(RegistryError::GrpcRequiresHub);
779        }
780
781        // If there's a hub, collect all services grouped by gear and hand them off to the installer store
782        if let Some(hub_name) = &self.registry.grpc_hub {
783            let mut gears_data = Vec::new();
784            let mut seen = HashSet::new();
785
786            // Collect services from all grpc gears
787            for (gear_name, service_gear) in &self.registry.grpc_services {
788                let ctx = self.ctx_builder.for_gear(gear_name).await.map_err(|err| {
789                    RegistryError::GrpcRegister {
790                        gear: gear_name.clone(),
791                        source: err,
792                    }
793                })?;
794
795                let installers = service_gear
796                    .get_grpc_services(&ctx)
797                    .await
798                    .map_err(|source| RegistryError::GrpcRegister {
799                        gear: gear_name.clone(),
800                        source,
801                    })?;
802
803                for reg in &installers {
804                    if !seen.insert(reg.service_name) {
805                        return Err(RegistryError::GrpcRegister {
806                            gear: gear_name.clone(),
807                            source: anyhow::anyhow!(
808                                "Duplicate gRPC service name: {}",
809                                reg.service_name
810                            ),
811                        });
812                    }
813                }
814
815                gears_data.push(crate::runtime::GearInstallers {
816                    gear_name: gear_name.clone(),
817                    installers,
818                });
819            }
820
821            self.grpc_installers
822                .set(crate::runtime::GrpcInstallerData { gears: gears_data })
823                .map_err(|source| RegistryError::GrpcRegister {
824                    gear: hub_name.clone(),
825                    source,
826                })?;
827        }
828
829        Ok(())
830    }
831
832    /// START phase: start all stateful gears.
833    ///
834    /// System gears start first, followed by user gears.
835    async fn run_start_phase(&self) -> Result<(), RegistryError> {
836        tracing::info!("Phase: start");
837
838        for e in self.registry.gears_by_system_priority() {
839            if let Some(s) = e.caps.query::<RunnableCap>() {
840                tracing::debug!(
841                    gear = e.name,
842                    is_system = e.caps.has::<SystemCap>(),
843                    "Starting stateful gear"
844                );
845                s.start(self.cancel.clone())
846                    .await
847                    .map_err(|source| RegistryError::Start {
848                        gear: e.name,
849                        source,
850                    })?;
851                tracing::info!(gear = e.name, "Started gear");
852            }
853        }
854
855        Ok(())
856    }
857
858    /// Stop a single gear, logging errors but continuing execution.
859    async fn stop_one_gear(entry: &GearEntry, cancel: CancellationToken) {
860        if let Some(s) = entry.caps.query::<RunnableCap>() {
861            match s.stop(cancel).await {
862                Err(err) => {
863                    tracing::warn!(gear =  entry.name, error = %err, "Failed to stop gear");
864                }
865                _ => {
866                    tracing::info!(gear = entry.name, "Stopped gear");
867                }
868            }
869        }
870    }
871
872    /// STOP phase: stop all stateful gears in reverse order.
873    ///
874    /// # Two-Phase Shutdown Contract
875    ///
876    /// This phase implements a proper two-phase shutdown for **each gear**:
877    ///
878    /// 1. **Graceful stop request**: Each gear's `stop(deadline_token)` is called with a
879    ///    *fresh* cancellation token (not the already-cancelled root token). Gears should
880    ///    interpret this as "please stop gracefully".
881    ///
882    /// 2. **Hard-stop deadline**: After `shutdown_deadline` expires **for that gear**,
883    ///    its `deadline_token` is cancelled. Gears should interpret this as "abort immediately".
884    ///
885    /// Each gear gets its own independent deadline — if gear A takes 25s to stop,
886    /// gear B still gets the full `shutdown_deadline` for its graceful shutdown.
887    ///
888    /// This allows gears to implement real graceful shutdown:
889    /// - Request cooperative shutdown of child tasks
890    /// - Wait for them to finish gracefully
891    /// - If `deadline_token` fires, switch to hard-abort mode
892    ///
893    /// Errors are logged but do not fail the shutdown process.
894    /// Note: `OoP` gears are stopped automatically by the backend when the
895    /// cancellation token is triggered.
896    async fn run_stop_phase(&self) -> Result<(), RegistryError> {
897        tracing::info!("Phase: stop");
898
899        // Drop our REST providers from the directory first, so consumers stop
900        // resolving an endpoint that is about to disappear.
901        self.deregister_rest_providers().await;
902
903        let deadline = self.shutdown_deadline;
904
905        // Stop all gears in reverse order, each with its own independent deadline
906        for e in self.registry.gears().iter().rev() {
907            let gear_name = e.name;
908
909            // Create a fresh deadline token for THIS gear
910            // Each gear gets the full shutdown_deadline independently
911            let deadline_token = CancellationToken::new();
912            let deadline_token_for_timeout = deadline_token.clone();
913
914            // Spawn a task to cancel this gear's deadline token after shutdown_deadline
915            let deadline_task = tokio::spawn(async move {
916                tokio::time::sleep(deadline).await;
917                tracing::warn!(
918                    gear = gear_name,
919                    deadline_secs = deadline.as_secs(),
920                    "Gear shutdown deadline reached, sending hard-stop signal"
921                );
922                deadline_token_for_timeout.cancel();
923            });
924
925            // Stop this gear with its own deadline token
926            // The gear can observe the token transition from uncancelled→cancelled
927            Self::stop_one_gear(e, deadline_token).await;
928
929            // Cancel the deadline task and await it to ensure full cleanup
930            deadline_task.abort();
931            #[allow(clippy::let_underscore_must_use)]
932            let _ = deadline_task.await;
933        }
934
935        Ok(())
936    }
937
938    /// Run the stop phase with a watchdog that force-exits the process if the
939    /// stop phase hangs on a blocking syscall. The watchdog is disarmed whether
940    /// the stop phase succeeds or fails, so a failing stop phase does not leave
941    /// the watchdog active and eventually force-exit the process.
942    async fn run_stop_phase_guarded(&self) -> Result<(), RegistryError> {
943        let gear_count = u32::try_from(self.registry.gears().len().max(1)).unwrap_or(1);
944        let stop_timeout = self
945            .shutdown_deadline
946            .checked_mul(gear_count)
947            .and_then(|d| d.checked_add(std::time::Duration::from_secs(5)))
948            .unwrap_or(self.shutdown_deadline);
949
950        // Use a channel to arm/disarm the watchdog. If the lifecycle future is
951        // dropped before the stop phase finishes (e.g. an outer timeout), the
952        // sender is dropped and the watchdog exits without killing the process.
953        let (disarm_tx, disarm_rx) = std::sync::mpsc::channel::<()>();
954        std::thread::spawn(move || {
955            match disarm_rx.recv_timeout(stop_timeout) {
956                Ok(()) | Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
957                    // Stop phase completed, or the lifecycle future was cancelled.
958                }
959                Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
960                    tracing::warn!(
961                        timeout_secs = stop_timeout.as_secs(),
962                        "shutdown: stop phase timed out, force exiting"
963                    );
964                    std::process::exit(1);
965                }
966            }
967        });
968
969        let stop_result = self.run_stop_phase().await;
970        // Disarm the watchdog before propagating the stop-phase result. This runs
971        // for both success and failure so a failing stop phase does not leave the
972        // watchdog armed and eventually force-exit the process.
973        let _ = disarm_tx.send(()).ok();
974
975        stop_result
976    }
977
978    /// `OoP` SPAWN phase: spawn out-of-process gears after start phase.
979    ///
980    /// This phase runs after `grpc-hub` is already listening, so we can pass
981    /// the real directory endpoint to `OoP` gears.
982    async fn run_oop_spawn_phase(&self) -> Result<(), RegistryError> {
983        let oop_opts = match &self.oop_options {
984            Some(opts) if !opts.gears.is_empty() => opts,
985            _ => return Ok(()),
986        };
987
988        tracing::info!("Phase: oop_spawn");
989
990        // Wait for grpc_hub to publish its endpoint (it runs async in start phase)
991        let directory_endpoint = self.wait_for_grpc_hub_endpoint().await;
992
993        for gear_cfg in &oop_opts.gears {
994            // Build environment with directory endpoint and rendered config
995            // Note: User controls --config via execution.args in master config
996            let mut env = gear_cfg.env.clone();
997            env.insert(
998                TOOLKIT_MODULE_CONFIG_ENV.to_owned(),
999                gear_cfg.rendered_config_json.clone(),
1000            );
1001            if let Some(ref endpoint) = directory_endpoint {
1002                env.insert(TOOLKIT_DIRECTORY_ENDPOINT_ENV.to_owned(), endpoint.clone());
1003            }
1004
1005            // Use args from execution config as-is (user controls --config via args)
1006            let args = gear_cfg.args.clone();
1007
1008            let spawn_config = OopSpawnConfig {
1009                gear_name: gear_cfg.gear_name.clone(),
1010                binary: gear_cfg.binary.clone(),
1011                args,
1012                env,
1013                working_directory: gear_cfg.working_directory.clone(),
1014            };
1015
1016            oop_opts
1017                .backend
1018                .spawn(spawn_config)
1019                .await
1020                .map_err(|e| RegistryError::OopSpawn {
1021                    gear: gear_cfg.gear_name.clone(),
1022                    source: e,
1023                })?;
1024
1025            tracing::info!(
1026                gear =  %gear_cfg.gear_name,
1027                directory_endpoint = ?directory_endpoint,
1028                "Spawned OoP gear via backend"
1029            );
1030        }
1031
1032        Ok(())
1033    }
1034
1035    /// Wait for `grpc-hub` to publish its bound endpoint.
1036    ///
1037    /// Polls the `GrpcHubGear::bound_endpoint()` with a short interval until available or timeout.
1038    /// Returns None if no `grpc-hub` is running or if it times out.
1039    async fn wait_for_grpc_hub_endpoint(&self) -> Option<String> {
1040        const POLL_INTERVAL: std::time::Duration = std::time::Duration::from_millis(10);
1041        const MAX_WAIT: std::time::Duration = std::time::Duration::from_secs(5);
1042
1043        // Find grpc_hub in registry
1044        let grpc_hub = self
1045            .registry
1046            .gears()
1047            .iter()
1048            .find_map(|e| e.caps.query::<GrpcHubCap>());
1049
1050        let Some(hub) = grpc_hub else {
1051            return None; // No grpc_hub registered
1052        };
1053
1054        let start = std::time::Instant::now();
1055
1056        loop {
1057            if let Some(endpoint) = hub.bound_endpoint() {
1058                tracing::debug!(
1059                    endpoint = %endpoint,
1060                    elapsed_ms = start.elapsed().as_millis(),
1061                    "gRPC hub endpoint available"
1062                );
1063                return Some(endpoint);
1064            }
1065
1066            if start.elapsed() > MAX_WAIT {
1067                tracing::warn!("Timed out waiting for gRPC hub to bind");
1068                return None;
1069            }
1070
1071            tokio::time::sleep(POLL_INTERVAL).await;
1072        }
1073    }
1074
1075    /// Wait for the REST host gateway to publish its bound endpoint.
1076    ///
1077    /// The gateway binds its listener asynchronously in the start phase, so the
1078    /// bound endpoint may not be set the instant the start phase returns; poll
1079    /// with a short interval until available or timeout.
1080    async fn wait_for_rest_endpoint(
1081        &self,
1082        host: &Arc<dyn crate::contracts::ApiGatewayCapability>,
1083    ) -> Option<String> {
1084        const POLL_INTERVAL: std::time::Duration = std::time::Duration::from_millis(10);
1085        const MAX_WAIT: std::time::Duration = std::time::Duration::from_secs(5);
1086
1087        let start = std::time::Instant::now();
1088        loop {
1089            if let Some(endpoint) = host.bound_endpoint() {
1090                return Some(endpoint);
1091            }
1092            if start.elapsed() > MAX_WAIT {
1093                tracing::warn!("Timed out waiting for REST host to bind");
1094                return None;
1095            }
1096            tokio::time::sleep(POLL_INTERVAL).await;
1097        }
1098    }
1099
1100    /// Directory-register phase (eventual readiness, provider side).
1101    ///
1102    /// After the REST server has bound, advertise every in-process REST provider
1103    /// gear in the service directory under its own gear name, pointing at the
1104    /// shared gateway endpoint. Consumers resolving a provider gear name then
1105    /// receive this endpoint and the gateway routes to the provider's handlers.
1106    ///
1107    /// No-op when there is no REST host or no REST provider gears, so non-REST
1108    /// deployments and existing tests are unaffected. Registers through the
1109    /// `DirectoryClient` in the `ClientHub`, which uniformly targets the
1110    /// in-process directory (`LocalDirectoryClient`) or the central directory
1111    /// (`DirectoryGrpcClient` for `OoP`) depending on what the host wired.
1112    async fn run_directory_register_phase(&self) -> Result<(), RegistryError> {
1113        let rest_gears = self.rest_provider_gears();
1114        if rest_gears.is_empty() {
1115            return Ok(());
1116        }
1117
1118        let Some(host) = self
1119            .registry
1120            .gears()
1121            .iter()
1122            .find_map(|e| e.caps.query::<ApiGatewayCap>())
1123        else {
1124            return Ok(()); // no REST host serving the routes
1125        };
1126
1127        let Some(endpoint) = self.wait_for_rest_endpoint(&host).await else {
1128            tracing::warn!(
1129                "directory-register: REST host endpoint unavailable; skipping REST provider registration"
1130            );
1131            return Ok(());
1132        };
1133
1134        let Ok(dir) = self.client_hub.get::<dyn crate::DirectoryClient>() else {
1135            tracing::debug!(
1136                "directory-register: no DirectoryClient in ClientHub; skipping REST provider registration"
1137            );
1138            return Ok(());
1139        };
1140
1141        let instance_id = self.instance_id.to_string();
1142        for gear in rest_gears {
1143            // The directory keys instances by (gear, instance_id) and replaces
1144            // wholesale. grpc-hub may have already registered this same
1145            // (gear, instance_id) with gRPC services during the start phase, so
1146            // carry those forward instead of clobbering them to empty — adding
1147            // the REST endpoint must augment, not replace, the entry.
1148            let (grpc_services, version) = match dir.list_instances(gear).await {
1149                Ok(insts) => insts
1150                    .into_iter()
1151                    .find(|i| i.instance_id == instance_id)
1152                    .map(|i| (i.grpc_services, i.version))
1153                    .unwrap_or_default(),
1154                Err(_) => (Vec::new(), None),
1155            };
1156            let info = crate::RegisterInstanceInfo {
1157                gear: gear.to_owned(),
1158                instance_id: instance_id.clone(),
1159                grpc_services,
1160                version,
1161                rest_endpoint: Some(crate::ServiceEndpoint::new(endpoint.clone())),
1162                // OpenAPI spec is published separately (grpc-hub start phase);
1163                // the REST-augmentation registration does not carry it.
1164                openapi_spec: None,
1165            };
1166            match dir.register_instance(info).await {
1167                Ok(()) => {
1168                    tracing::info!(gear, endpoint = %endpoint, "registered REST provider in directory");
1169                }
1170                Err(e) => {
1171                    tracing::warn!(gear, error = %e, "directory-register: failed to register REST provider");
1172                }
1173            }
1174        }
1175        // Mark that this (in-process host) process advertised its REST providers,
1176        // so the stop phase deregisters them exactly once. `OoP` serving never
1177        // runs this phase, so its deregister is owned solely by `oop_serve`.
1178        self.rest_providers_registered
1179            .store(true, std::sync::atomic::Ordering::SeqCst);
1180        Ok(())
1181    }
1182
1183    /// Names of all gears that provide a REST API (have `RestApiCap`), excluding
1184    /// the REST host gateway itself (`ApiGatewayCap`) — the gateway is the
1185    /// transport, not a contract provider, so it must not be advertised in the
1186    /// directory under its own gear name.
1187    fn rest_provider_gears(&self) -> Vec<&'static str> {
1188        self.registry
1189            .gears()
1190            .iter()
1191            .filter(|e| e.caps.has::<RestApiCap>() && !e.caps.has::<ApiGatewayCap>())
1192            .map(|e| e.name)
1193            .collect()
1194    }
1195
1196    /// Deregister this process's REST providers from the directory on shutdown,
1197    /// so consumers stop resolving an endpoint that is going away. Best-effort.
1198    async fn deregister_rest_providers(&self) {
1199        // Only the in-process host path registers REST providers in the directory
1200        // (via `run_directory_register_phase`). In `OoP` serving, `oop_serve`
1201        // owns presence + deregister, so skip here to avoid a double deregister.
1202        if !self
1203            .rest_providers_registered
1204            .load(std::sync::atomic::Ordering::SeqCst)
1205        {
1206            return;
1207        }
1208        let rest_gears = self.rest_provider_gears();
1209        if rest_gears.is_empty() {
1210            return;
1211        }
1212        let Ok(dir) = self.client_hub.get::<dyn crate::DirectoryClient>() else {
1213            return;
1214        };
1215        let instance_id = self.instance_id.to_string();
1216        for gear in rest_gears {
1217            if let Err(e) = dir.deregister_instance(gear, &instance_id).await {
1218                tracing::warn!(gear, error = %e, "directory-deregister: failed to deregister REST provider");
1219            }
1220        }
1221    }
1222
1223    /// Run the full gear lifecycle (all phases).
1224    ///
1225    /// This is the standard entry point for normal application execution.
1226    /// It runs all phases from pre-init through shutdown.
1227    ///
1228    /// # Errors
1229    ///
1230    /// Returns an error if any gear phase fails during execution.
1231    pub async fn run_gear_phases(self) -> anyhow::Result<()> {
1232        self.run_phases_internal(RunMode::Full).await
1233    }
1234
1235    /// Run only the migration phases (pre-init + DB migration).
1236    ///
1237    /// This is designed for cloud deployment workflows where database migrations
1238    /// need to run as a separate step before starting the application.
1239    /// The process exits after migrations complete.
1240    ///
1241    /// # Errors
1242    ///
1243    /// Returns an error if pre-init or migration phases fail.
1244    pub async fn run_migration_phases(self) -> anyhow::Result<()> {
1245        self.run_phases_internal(RunMode::MigrateOnly).await
1246    }
1247
1248    /// Internal implementation that runs gear phases based on the mode.
1249    ///
1250    /// This private method contains the actual phase execution logic and is called
1251    /// by both `run_gear_phases()` and `run_migration_phases()`.
1252    ///
1253    /// # Modes
1254    ///
1255    /// - `RunMode::Full`: Executes all phases and waits for shutdown signal
1256    /// - `RunMode::MigrateOnly`: Executes only pre-init and DB migration phases, then exits
1257    ///
1258    /// # Phases (Full Mode)
1259    ///
1260    /// 1. Pre-init (system gears only)
1261    /// 2. DB migration (all gears with database capability)
1262    /// 3. Init (all gears)
1263    /// 4. Post-init (system gears only)
1264    /// 5. REST (gears with REST capability)
1265    /// 6. gRPC (gears with gRPC capability)
1266    /// 7. Start (runnable gears)
1267    /// 8. `OoP` spawn (out-of-process gears)
1268    /// 9. Wait for cancellation
1269    /// 10. Stop (runnable gears in reverse order)
1270    async fn run_phases_internal(self, mode: RunMode) -> anyhow::Result<()> {
1271        // Log execution mode
1272        match mode {
1273            RunMode::Full => {
1274                tracing::info!("Running full lifecycle (all phases)");
1275            }
1276            RunMode::MigrateOnly => {
1277                tracing::info!("Running in migration mode (pre-init + db phases only)");
1278            }
1279        }
1280
1281        // 1. Pre-init phase (before init, only for system gears)
1282        self.run_pre_init_phase()?;
1283
1284        // 2. DB migration phase (system gears first)
1285        #[cfg(feature = "db")]
1286        {
1287            self.run_db_phase().await?;
1288        }
1289        #[cfg(not(feature = "db"))]
1290        {
1291            // No DB integration in this build.
1292        }
1293
1294        // Exit early if running in migration-only mode
1295        if mode == RunMode::MigrateOnly {
1296            tracing::info!("Migration phases completed successfully");
1297            return Ok(());
1298        }
1299
1300        // 3. Init -> proxy-wiring -> post-init (shared with the OoP path)
1301        self.run_init_wiring_post_init().await?;
1302
1303        // 5. REST phase (synchronous router composition)
1304        let _router = self.run_rest_phase().await?;
1305
1306        // 6. gRPC registration phase
1307        self.run_grpc_phase().await?;
1308
1309        // 7. Start phase
1310        self.run_start_phase().await?;
1311
1312        // Draining watcher: flip readiness to Draining the moment shutdown
1313        // begins so /readyz reports 503 and the orchestrator drains the pod out
1314        // of the load balancer before the stop phase tears gears down.
1315        {
1316            let readiness = Arc::clone(&self.dep_checker);
1317            let cancel = self.cancel.clone();
1318            tokio::spawn(async move {
1319                cancel.cancelled().await;
1320                readiness.set_draining(true);
1321            });
1322        }
1323
1324        // 7b. Directory-register phase: advertise in-process REST providers in
1325        //     the directory once the gateway has bound its listener.
1326        self.run_directory_register_phase().await?;
1327
1328        // 8. OoP spawn phase (after grpc_hub is running)
1329        self.run_oop_spawn_phase().await?;
1330
1331        // 9. Wait for cancellation
1332        self.cancel.cancelled().await;
1333
1334        // 10. Stop phase with hard timeout.
1335        //     Blocking stop implementations are guarded by a watchdog thread so
1336        //     a hang cannot block shutdown, whether in the in-process or OoP path.
1337        self.run_stop_phase_guarded().await?;
1338        Ok(())
1339    }
1340}
1341
1342/// Out-of-process HTTP serving lifecycle (`cpt-cf-component-oop-bootstrap`).
1343#[cfg(feature = "bootstrap")]
1344impl HostRuntime {
1345    /// Compose a **host-less** REST router from all `RestApiCap` gears, plus the
1346    /// gear's generated `OpenAPI` document (serialized JSON).
1347    ///
1348    /// Unlike [`run_rest_phase`](Self::run_rest_phase), this does not require an
1349    /// `ApiGatewayCap` host: `OoP` gears serve their own routes directly.
1350    async fn compose_oop_router(
1351        &self,
1352        options: &crate::runtime::OopServeOptions,
1353        hc_registry: &Arc<crate::healthcheck::RestHealthcheckRegistry>,
1354    ) -> anyhow::Result<(Router, String)> {
1355        use crate::api::{OpenApiInfo, OpenApiRegistryImpl};
1356        use anyhow::Context as _;
1357
1358        let registry = OpenApiRegistryImpl::new();
1359        let mut router = Router::new();
1360
1361        for entry in self.registry.gears() {
1362            if let Some(rest) = entry.caps.query::<RestApiCap>() {
1363                let ctx = self
1364                    .ctx_builder
1365                    .for_gear(entry.name)
1366                    .await
1367                    .with_context(|| format!("OoP router: build context for '{}'", entry.name))?;
1368                router = rest
1369                    .register_rest(&ctx, router, &registry)
1370                    .with_context(|| format!("OoP router: register_rest for '{}'", entry.name))?;
1371
1372                // Register the gear's readiness healthcheck (the same mechanism
1373                // the api-gateway host uses), so /readyz reflects it identically
1374                // whether the gear runs in-process or OoP.
1375                if let Some(hc) = rest.healthcheck(&ctx) {
1376                    hc_registry.register(entry.name, hc);
1377                }
1378            }
1379        }
1380
1381        let info = OpenApiInfo {
1382            title: options.gear_name.clone(),
1383            version: options
1384                .version
1385                .clone()
1386                .unwrap_or_else(|| "0.0.0".to_owned()),
1387            description: None,
1388            servers: vec![],
1389        };
1390        let openapi = registry
1391            .build_openapi(&info)
1392            .context("OoP router: build OpenAPI document")?;
1393        let json = serde_json::to_string(&openapi).context("OoP router: serialize OpenAPI")?;
1394
1395        Ok((router, json))
1396    }
1397
1398    /// Run the full `OoP` gear lifecycle: phases (`pre_init` … `start`), then
1399    /// serve the composed router with framework probes, background
1400    /// self-registration, dependency resolution, and graceful drain, then the
1401    /// `stop` phase.
1402    ///
1403    /// # Errors
1404    /// Returns an error if any lifecycle phase or the HTTP server fails.
1405    pub async fn run_oop_serving(
1406        self,
1407        options: crate::runtime::OopServeOptions,
1408    ) -> anyhow::Result<()> {
1409        use crate::runtime::ReadinessState;
1410
1411        tracing::info!("Running OoP serving lifecycle");
1412
1413        // Make the directory client reachable to the proxy-wiring phase, which
1414        // resolves remote `#[toolkit::consumes]` providers through the
1415        // `ClientHub`. The `OoP` presence loop uses `options.directory`
1416        // directly; consumer wiring reads it here.
1417        if self.client_hub.get::<dyn crate::DirectoryClient>().is_err() {
1418            self.client_hub
1419                .register::<dyn crate::DirectoryClient>(Arc::clone(&options.directory));
1420        }
1421
1422        // Shared gear healthcheck registry (same mechanism as the api-gateway
1423        // host path). Seeded with the root cancellation token so in-flight
1424        // checks are aborted on shutdown. Populated during router composition.
1425        let hc_registry = Arc::new(
1426            crate::healthcheck::RestHealthcheckRegistry::with_cancellation(self.cancel.clone()),
1427        );
1428
1429        // Readiness gates on `#[toolkit::consumes]` dependencies only — the same
1430        // policy as the in-process host path — via the shared `DependencyChecker`
1431        // fed by the proxy-wiring phase below. `deps = [...]`-only declarations
1432        // remain for topo-sort ordering but do NOT gate `/readyz`. The
1433        // healthcheck registry supplies the per-gear readiness dimension; both
1434        // feed the `/readyz` aggregate.
1435        let readiness = ReadinessState::from_checker(
1436            Arc::clone(&self.dep_checker),
1437            Arc::clone(&hc_registry),
1438            options.healthcheck_timeout,
1439        );
1440
1441        // Bind the HTTP server and serve probes BEFORE the (possibly slow)
1442        // lifecycle phases, so the kubelet's liveness probe (`/healthz`) passes
1443        // immediately instead of getting connection-refused during `start()`.
1444        // Gear routes reply `503 starting` until attached below.
1445        let mut server = super::oop_serve::OopHttpServer::start(
1446            Arc::clone(&readiness),
1447            options,
1448            self.cancel.clone(),
1449        )
1450        .await?;
1451
1452        // Lifecycle phases up to start, then wire consumers (typed
1453        // directory-resolving clients feed the shared `DependencyChecker`),
1454        // then compose the host-less REST router + OpenAPI spec (collecting each
1455        // gear's healthcheck into the shared registry). Grouped so a failure
1456        // tears the probe server down cleanly.
1457        let mut started = false;
1458        let composed: anyhow::Result<(Router, String)> = async {
1459            self.run_pre_init_phase()?;
1460            #[cfg(feature = "db")]
1461            self.run_db_phase().await?;
1462            // Init -> proxy-wiring -> post-init, shared with the in-process path
1463            // so a gear resolving a consumed contract during `start` finds the
1464            // client in the hub under both profiles.
1465            self.run_init_wiring_post_init().await?;
1466            self.run_grpc_phase().await?;
1467            self.run_start_phase().await?;
1468            started = true;
1469            // The gear lifecycle has populated the ClientHub. If an in-process
1470            // authn stack (e.g. a linked authn-resolver gear) registered a
1471            // DynBearerAuthenticator bridge, install the tenant plane now —
1472            // before `attach` layers security_context_middleware. (The platform
1473            // plane is built eagerly at bootstrap from `oop_http.internal_auth`.)
1474            server.resolve_bearer_authenticator(&self.client_hub);
1475            self.compose_oop_router(server.options(), &hc_registry)
1476                .await
1477        }
1478        .await;
1479
1480        let serve_result = match composed {
1481            Ok((gear_router, openapi_json)) => {
1482                // Publish gear routes (they go live) + start directory presence.
1483                // Dependency resolution already ran in the proxy-wiring phase.
1484                server.attach(gear_router, openapi_json);
1485                // Serve until cancelled, drain, then deregister.
1486                server.join().await
1487            }
1488            Err(e) => {
1489                tracing::error!(error = %e, "OoP startup failed before serving gear routes");
1490                // Tear down the probe server that is already bound.
1491                self.cancel.cancel();
1492                if let Err(join_err) = server.join().await {
1493                    tracing::warn!(error = %join_err, "OoP probe server teardown after startup failure errored");
1494                }
1495                Err(e)
1496            }
1497        };
1498
1499        // Stop phase runs only if start completed successfully. Errors are logged
1500        // but do not fail the shutdown process (same contract as run_stop_phase).
1501        if started && let Err(e) = self.run_stop_phase_guarded().await {
1502            tracing::warn!(error = %e, "OoP stop phase reported an error");
1503        }
1504
1505        serve_result
1506    }
1507}
1508
1509#[cfg(test)]
1510#[cfg(feature = "bootstrap")]
1511#[cfg_attr(coverage_nightly, coverage(off))]
1512#[path = "host_runtime_oop_tests.rs"]
1513mod host_runtime_oop_tests;
1514
1515#[cfg(test)]
1516#[cfg_attr(coverage_nightly, coverage(off))]
1517mod tests {
1518    use super::*;
1519    use crate::context::GearCtx;
1520    use crate::contracts::{Gear, RunnableCapability, SystemCapability};
1521    use crate::registry::RegistryBuilder;
1522    use std::sync::Arc;
1523    use std::sync::atomic::{AtomicUsize, Ordering};
1524    use tokio::sync::Mutex;
1525
1526    #[derive(Default)]
1527    #[allow(dead_code)]
1528    struct DummyCore;
1529    #[async_trait::async_trait]
1530    impl Gear for DummyCore {
1531        async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
1532            Ok(())
1533        }
1534    }
1535
1536    struct StopOrderTracker {
1537        my_order: usize,
1538        stop_order: Arc<AtomicUsize>,
1539    }
1540
1541    impl StopOrderTracker {
1542        fn new(counter: &Arc<AtomicUsize>, stop_order: Arc<AtomicUsize>) -> Self {
1543            let my_order = counter.fetch_add(1, Ordering::SeqCst);
1544            Self {
1545                my_order,
1546                stop_order,
1547            }
1548        }
1549    }
1550
1551    #[async_trait::async_trait]
1552    impl Gear for StopOrderTracker {
1553        async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
1554            Ok(())
1555        }
1556    }
1557
1558    #[async_trait::async_trait]
1559    impl RunnableCapability for StopOrderTracker {
1560        async fn start(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
1561            Ok(())
1562        }
1563        async fn stop(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
1564            let order = self.stop_order.fetch_add(1, Ordering::SeqCst);
1565            tracing::info!(my_order = self.my_order, stop_order = order, "Gear stopped");
1566            Ok(())
1567        }
1568    }
1569
1570    #[tokio::test]
1571    async fn test_stop_phase_reverse_order() {
1572        let counter = Arc::new(AtomicUsize::new(0));
1573        let stop_order = Arc::new(AtomicUsize::new(0));
1574
1575        let gear_a = Arc::new(StopOrderTracker::new(&counter, stop_order.clone()));
1576        let gear_b = Arc::new(StopOrderTracker::new(&counter, stop_order.clone()));
1577        let gear_c = Arc::new(StopOrderTracker::new(&counter, stop_order.clone()));
1578
1579        let mut builder = RegistryBuilder::default();
1580        builder.register_core_with_meta("a", &[], gear_a.clone() as Arc<dyn Gear>);
1581        builder.register_core_with_meta("b", &["a"], gear_b.clone() as Arc<dyn Gear>);
1582        builder.register_core_with_meta("c", &["b"], gear_c.clone() as Arc<dyn Gear>);
1583
1584        builder.register_stateful_with_meta("a", gear_a.clone() as Arc<dyn RunnableCapability>);
1585        builder.register_stateful_with_meta("b", gear_b.clone() as Arc<dyn RunnableCapability>);
1586        builder.register_stateful_with_meta("c", gear_c.clone() as Arc<dyn RunnableCapability>);
1587
1588        let registry = builder.build_topo_sorted().unwrap();
1589
1590        // Verify gear order is a -> b -> c
1591        let gear_names: Vec<_> = registry.gears().iter().map(|m| m.name).collect();
1592        assert_eq!(gear_names, vec!["a", "b", "c"]);
1593
1594        let client_hub = Arc::new(ClientHub::new());
1595        let cancel = CancellationToken::new();
1596        let config_provider: Arc<dyn ConfigProvider> = Arc::new(EmptyConfigProvider);
1597
1598        let runtime = HostRuntime::new(
1599            registry,
1600            config_provider,
1601            DbOptions::None,
1602            client_hub,
1603            cancel.clone(),
1604            Uuid::new_v4(),
1605            None,
1606        );
1607
1608        // Run stop phase
1609        runtime.run_stop_phase().await.unwrap();
1610
1611        // Verify gears stopped in reverse order: c (stop_order=0), b (stop_order=1), a (stop_order=2)
1612        // Gear order is: a=0, b=1, c=2
1613        // Stop order should be: c=0, b=1, a=2
1614        assert_eq!(stop_order.load(Ordering::SeqCst), 3);
1615    }
1616
1617    #[tokio::test]
1618    async fn test_stop_phase_continues_on_error() {
1619        struct FailingGear {
1620            should_fail: bool,
1621            stopped: Arc<AtomicUsize>,
1622        }
1623
1624        #[async_trait::async_trait]
1625        impl Gear for FailingGear {
1626            async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
1627                Ok(())
1628            }
1629        }
1630
1631        #[async_trait::async_trait]
1632        impl RunnableCapability for FailingGear {
1633            async fn start(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
1634                Ok(())
1635            }
1636            async fn stop(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
1637                self.stopped.fetch_add(1, Ordering::SeqCst);
1638                if self.should_fail {
1639                    anyhow::bail!("Intentional failure")
1640                }
1641                Ok(())
1642            }
1643        }
1644
1645        let stopped = Arc::new(AtomicUsize::new(0));
1646        let gear_a = Arc::new(FailingGear {
1647            should_fail: false,
1648            stopped: stopped.clone(),
1649        });
1650        let gear_b = Arc::new(FailingGear {
1651            should_fail: true,
1652            stopped: stopped.clone(),
1653        });
1654        let gear_c = Arc::new(FailingGear {
1655            should_fail: false,
1656            stopped: stopped.clone(),
1657        });
1658
1659        let mut builder = RegistryBuilder::default();
1660        builder.register_core_with_meta("a", &[], gear_a.clone() as Arc<dyn Gear>);
1661        builder.register_core_with_meta("b", &["a"], gear_b.clone() as Arc<dyn Gear>);
1662        builder.register_core_with_meta("c", &["b"], gear_c.clone() as Arc<dyn Gear>);
1663
1664        builder.register_stateful_with_meta("a", gear_a.clone() as Arc<dyn RunnableCapability>);
1665        builder.register_stateful_with_meta("b", gear_b.clone() as Arc<dyn RunnableCapability>);
1666        builder.register_stateful_with_meta("c", gear_c.clone() as Arc<dyn RunnableCapability>);
1667
1668        let registry = builder.build_topo_sorted().unwrap();
1669
1670        let client_hub = Arc::new(ClientHub::new());
1671        let cancel = CancellationToken::new();
1672        let config_provider: Arc<dyn ConfigProvider> = Arc::new(EmptyConfigProvider);
1673
1674        let runtime = HostRuntime::new(
1675            registry,
1676            config_provider,
1677            DbOptions::None,
1678            client_hub,
1679            cancel.clone(),
1680            Uuid::new_v4(),
1681            None,
1682        );
1683
1684        // Run stop phase - should not fail even though gear_b fails
1685        runtime.run_stop_phase().await.unwrap();
1686
1687        // All gears should have attempted to stop
1688        assert_eq!(stopped.load(Ordering::SeqCst), 3);
1689    }
1690
1691    struct EmptyConfigProvider;
1692    impl ConfigProvider for EmptyConfigProvider {
1693        fn get_gear_config(&self, _gear_name: &str) -> Option<&serde_json::Value> {
1694            None
1695        }
1696    }
1697
1698    #[test]
1699    fn static_endpoint_override_reads_nested_consumer_wiring_key() {
1700        struct MapCfg(std::collections::HashMap<String, serde_json::Value>);
1701        impl ConfigProvider for MapCfg {
1702            fn get_gear_config(&self, gear: &str) -> Option<&serde_json::Value> {
1703                self.0.get(gear)
1704            }
1705        }
1706        let mut map = std::collections::HashMap::new();
1707        map.insert(
1708            "orders".to_owned(),
1709            serde_json::json!({
1710                "config": { "consumer_wiring": { "billing": "http://localhost:8081" } }
1711            }),
1712        );
1713        let cfg = MapCfg(map);
1714
1715        // Present override is read from `config.consumer_wiring.<dep>`.
1716        assert_eq!(
1717            super::static_endpoint_override(&cfg, "orders", "billing").as_deref(),
1718            Some("http://localhost:8081")
1719        );
1720        // Absent dep / owner → None (falls through to directory resolution).
1721        assert_eq!(
1722            super::static_endpoint_override(&cfg, "orders", "inventory"),
1723            None
1724        );
1725        assert_eq!(
1726            super::static_endpoint_override(&cfg, "warehouse", "billing"),
1727            None
1728        );
1729        assert_eq!(
1730            super::static_endpoint_override(&EmptyConfigProvider, "orders", "billing"),
1731            None
1732        );
1733    }
1734
1735    /// The override is keyed by the *gear name* (kebab), which is what
1736    /// `#[toolkit::consumes]` puts in `ConsumerRegistration::owner_gear`.
1737    /// Regression guard: the macro used to emit `stringify!(StructIdent)`, so
1738    /// the lookup asked for `gears.ApiContractsConsumer` — a key no config
1739    /// declares — and the escape hatch could never fire.
1740    #[test]
1741    fn static_endpoint_override_is_keyed_by_kebab_gear_name() {
1742        struct MapCfg(std::collections::HashMap<String, serde_json::Value>);
1743        impl ConfigProvider for MapCfg {
1744            fn get_gear_config(&self, gear: &str) -> Option<&serde_json::Value> {
1745                self.0.get(gear)
1746            }
1747        }
1748        let mut map = std::collections::HashMap::new();
1749        map.insert(
1750            "api-contracts-consumer".to_owned(),
1751            serde_json::json!({
1752                "config": { "consumer_wiring": { "api-contracts": "http://localhost:9099" } }
1753            }),
1754        );
1755        let cfg = MapCfg(map);
1756
1757        assert_eq!(
1758            super::static_endpoint_override(&cfg, "api-contracts-consumer", "api-contracts")
1759                .as_deref(),
1760            Some("http://localhost:9099"),
1761        );
1762        // The pre-fix value — the Rust struct ident — must NOT resolve.
1763        assert_eq!(
1764            super::static_endpoint_override(&cfg, "ApiContractsConsumer", "api-contracts"),
1765            None,
1766        );
1767    }
1768
1769    #[tokio::test]
1770    async fn test_post_init_runs_after_all_init_and_system_first() {
1771        #[derive(Clone)]
1772        struct TrackHooks {
1773            name: &'static str,
1774            events: Arc<Mutex<Vec<String>>>,
1775        }
1776
1777        #[async_trait::async_trait]
1778        impl Gear for TrackHooks {
1779            async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
1780                self.events.lock().await.push(format!("init:{}", self.name));
1781                Ok(())
1782            }
1783        }
1784
1785        #[async_trait::async_trait]
1786        impl SystemCapability for TrackHooks {
1787            fn pre_init(&self, _sys: &crate::runtime::SystemContext) -> anyhow::Result<()> {
1788                Ok(())
1789            }
1790
1791            async fn post_init(&self, _sys: &crate::runtime::SystemContext) -> anyhow::Result<()> {
1792                self.events
1793                    .lock()
1794                    .await
1795                    .push(format!("post_init:{}", self.name));
1796                Ok(())
1797            }
1798        }
1799
1800        let events = Arc::new(Mutex::new(Vec::<String>::new()));
1801        let sys_a = Arc::new(TrackHooks {
1802            name: "sys_a",
1803            events: events.clone(),
1804        });
1805        let user_b = Arc::new(TrackHooks {
1806            name: "user_b",
1807            events: events.clone(),
1808        });
1809        let user_c = Arc::new(TrackHooks {
1810            name: "user_c",
1811            events: events.clone(),
1812        });
1813
1814        let mut builder = RegistryBuilder::default();
1815        builder.register_core_with_meta("sys_a", &[], sys_a.clone() as Arc<dyn Gear>);
1816        builder.register_core_with_meta("user_b", &["sys_a"], user_b.clone() as Arc<dyn Gear>);
1817        builder.register_core_with_meta("user_c", &["user_b"], user_c.clone() as Arc<dyn Gear>);
1818        builder.register_system_with_meta("sys_a", sys_a.clone() as Arc<dyn SystemCapability>);
1819
1820        let registry = builder.build_topo_sorted().unwrap();
1821
1822        let client_hub = Arc::new(ClientHub::new());
1823        let cancel = CancellationToken::new();
1824        let config_provider: Arc<dyn ConfigProvider> = Arc::new(EmptyConfigProvider);
1825
1826        let runtime = HostRuntime::new(
1827            registry,
1828            config_provider,
1829            DbOptions::None,
1830            client_hub,
1831            cancel,
1832            Uuid::new_v4(),
1833            None,
1834        );
1835
1836        // Run init phase for all gears, then post_init as a separate barrier phase.
1837        runtime.run_init_phase().await.unwrap();
1838        runtime.run_post_init_phase().await.unwrap();
1839
1840        let events = events.lock().await.clone();
1841        let first_post_init = events
1842            .iter()
1843            .position(|e| e.starts_with("post_init:"))
1844            .expect("expected post_init events");
1845        assert!(
1846            events[..first_post_init]
1847                .iter()
1848                .all(|e| e.starts_with("init:")),
1849            "expected all init events before post_init, got: {events:?}"
1850        );
1851
1852        // system-first order within each phase
1853        assert_eq!(
1854            events,
1855            vec![
1856                "init:sys_a",
1857                "init:user_b",
1858                "init:user_c",
1859                "post_init:sys_a",
1860            ]
1861        );
1862    }
1863
1864    /// The in-process and `OoP` paths must agree on where proxy-wiring sits.
1865    ///
1866    /// They did not: the `OoP` path ran it *after* `start`, so a gear resolving
1867    /// a consumed contract during its own `start` found nothing in the
1868    /// `ClientHub` under Profile 2/3 while working under Profile 1. Both paths
1869    /// now go through `run_init_wiring_post_init`, which pins the order.
1870    ///
1871    /// This test covers the two observable endpoints of that segment. The
1872    /// wiring step between them is a no-op here (no `#[toolkit::consumes]`
1873    /// registration is linked into this test binary), so what actually stops
1874    /// the paths diverging again is the shared method itself — inlining it back
1875    /// into both call sites makes it dead code and fails the `-D warnings`
1876    /// build.
1877    #[tokio::test]
1878    async fn init_wiring_post_init_runs_as_one_ordered_segment() {
1879        #[derive(Clone)]
1880        struct TrackHooks {
1881            events: Arc<Mutex<Vec<String>>>,
1882        }
1883
1884        #[async_trait::async_trait]
1885        impl Gear for TrackHooks {
1886            async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
1887                self.events.lock().await.push("init".to_owned());
1888                Ok(())
1889            }
1890        }
1891
1892        #[async_trait::async_trait]
1893        impl SystemCapability for TrackHooks {
1894            fn pre_init(&self, _sys: &crate::runtime::SystemContext) -> anyhow::Result<()> {
1895                Ok(())
1896            }
1897
1898            async fn post_init(&self, _sys: &crate::runtime::SystemContext) -> anyhow::Result<()> {
1899                self.events.lock().await.push("post_init".to_owned());
1900                Ok(())
1901            }
1902        }
1903
1904        let events = Arc::new(Mutex::new(Vec::<String>::new()));
1905        let gear = Arc::new(TrackHooks {
1906            events: events.clone(),
1907        });
1908
1909        let mut builder = RegistryBuilder::default();
1910        builder.register_core_with_meta("sys", &[], gear.clone() as Arc<dyn Gear>);
1911        builder.register_system_with_meta("sys", gear.clone() as Arc<dyn SystemCapability>);
1912        let registry = builder.build_topo_sorted().unwrap();
1913
1914        let runtime = HostRuntime::new(
1915            registry,
1916            Arc::new(EmptyConfigProvider) as Arc<dyn ConfigProvider>,
1917            DbOptions::None,
1918            Arc::new(ClientHub::new()),
1919            CancellationToken::new(),
1920            Uuid::new_v4(),
1921            None,
1922        );
1923
1924        runtime.run_init_wiring_post_init().await.unwrap();
1925
1926        assert_eq!(events.lock().await.clone(), vec!["init", "post_init"]);
1927    }
1928
1929    #[tokio::test]
1930    async fn test_stop_phase_provides_fresh_deadline_token() {
1931        use std::sync::atomic::AtomicBool;
1932
1933        struct TokenCheckGear {
1934            stop_was_called: AtomicBool,
1935            token_was_cancelled_on_entry: AtomicBool,
1936        }
1937
1938        #[async_trait::async_trait]
1939        impl Gear for TokenCheckGear {
1940            async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
1941                Ok(())
1942            }
1943        }
1944
1945        #[async_trait::async_trait]
1946        impl RunnableCapability for TokenCheckGear {
1947            async fn start(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
1948                Ok(())
1949            }
1950            async fn stop(&self, deadline_token: CancellationToken) -> anyhow::Result<()> {
1951                // Record that stop() was called
1952                self.stop_was_called.store(true, Ordering::SeqCst);
1953                // Record whether the token was already cancelled when stop() was called
1954                self.token_was_cancelled_on_entry
1955                    .store(deadline_token.is_cancelled(), Ordering::SeqCst);
1956                Ok(())
1957            }
1958        }
1959
1960        let gear = Arc::new(TokenCheckGear {
1961            stop_was_called: AtomicBool::new(false),
1962            // Default to true to detect if stop() was never called
1963            token_was_cancelled_on_entry: AtomicBool::new(true),
1964        });
1965
1966        let mut builder = RegistryBuilder::default();
1967        builder.register_core_with_meta("test", &[], gear.clone() as Arc<dyn Gear>);
1968        builder.register_stateful_with_meta("test", gear.clone() as Arc<dyn RunnableCapability>);
1969
1970        let registry = builder.build_topo_sorted().unwrap();
1971        let client_hub = Arc::new(ClientHub::new());
1972        let cancel = CancellationToken::new();
1973        let config_provider: Arc<dyn ConfigProvider> = Arc::new(EmptyConfigProvider);
1974
1975        let runtime = HostRuntime::new(
1976            registry,
1977            config_provider,
1978            DbOptions::None,
1979            client_hub,
1980            cancel.clone(),
1981            Uuid::new_v4(),
1982            None,
1983        );
1984
1985        // Run stop phase - the deadline token should NOT be cancelled
1986        runtime.run_stop_phase().await.unwrap();
1987
1988        // First, verify stop() was actually called (guards against silent registration failures)
1989        assert!(
1990            gear.stop_was_called.load(Ordering::SeqCst),
1991            "stop() was never called - gear may not have been registered correctly"
1992        );
1993
1994        // The token should NOT have been cancelled when stop() was called
1995        // This is the key fix: gears get a fresh token, not the already-cancelled root token
1996        assert!(
1997            !gear.token_was_cancelled_on_entry.load(Ordering::SeqCst),
1998            "deadline_token should NOT be cancelled when stop() is called - this enables graceful shutdown"
1999        );
2000    }
2001
2002    #[tokio::test]
2003    async fn test_stop_phase_graceful_shutdown_completes_before_deadline() {
2004        use std::sync::atomic::AtomicBool;
2005        use std::time::Duration;
2006
2007        struct GracefulGear {
2008            graceful_completed: AtomicBool,
2009            deadline_fired: AtomicBool,
2010        }
2011
2012        #[async_trait::async_trait]
2013        impl Gear for GracefulGear {
2014            async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
2015                Ok(())
2016            }
2017        }
2018
2019        #[async_trait::async_trait]
2020        impl RunnableCapability for GracefulGear {
2021            async fn start(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
2022                Ok(())
2023            }
2024            async fn stop(&self, deadline_token: CancellationToken) -> anyhow::Result<()> {
2025                // Simulate graceful shutdown that completes quickly (10ms)
2026                tokio::select! {
2027                    () = tokio::time::sleep(Duration::from_millis(10)) => {
2028                        self.graceful_completed.store(true, Ordering::SeqCst);
2029                    }
2030                    () = deadline_token.cancelled() => {
2031                        self.deadline_fired.store(true, Ordering::SeqCst);
2032                    }
2033                }
2034                Ok(())
2035            }
2036        }
2037
2038        let gear = Arc::new(GracefulGear {
2039            graceful_completed: AtomicBool::new(false),
2040            deadline_fired: AtomicBool::new(false),
2041        });
2042
2043        let mut builder = RegistryBuilder::default();
2044        builder.register_core_with_meta("test", &[], gear.clone() as Arc<dyn Gear>);
2045        builder.register_stateful_with_meta("test", gear.clone() as Arc<dyn RunnableCapability>);
2046
2047        let registry = builder.build_topo_sorted().unwrap();
2048        let client_hub = Arc::new(ClientHub::new());
2049        let cancel = CancellationToken::new();
2050        let config_provider: Arc<dyn ConfigProvider> = Arc::new(EmptyConfigProvider);
2051
2052        // Use a long deadline (5s) - gear should complete gracefully before this
2053        let runtime = HostRuntime::new(
2054            registry,
2055            config_provider,
2056            DbOptions::None,
2057            client_hub,
2058            cancel.clone(),
2059            Uuid::new_v4(),
2060            None,
2061        )
2062        .with_shutdown_deadline(Duration::from_secs(5));
2063
2064        runtime.run_stop_phase().await.unwrap();
2065
2066        // Graceful shutdown should have completed
2067        assert!(
2068            gear.graceful_completed.load(Ordering::SeqCst),
2069            "graceful shutdown should complete"
2070        );
2071        // Deadline should NOT have fired (gear finished before deadline)
2072        assert!(
2073            !gear.deadline_fired.load(Ordering::SeqCst),
2074            "deadline should not fire when graceful shutdown completes quickly"
2075        );
2076    }
2077
2078    #[tokio::test]
2079    async fn test_stop_phase_deadline_fires_for_slow_gear() {
2080        use std::sync::atomic::AtomicBool;
2081        use std::time::Duration;
2082
2083        struct SlowGear {
2084            graceful_completed: AtomicBool,
2085            deadline_fired: AtomicBool,
2086        }
2087
2088        #[async_trait::async_trait]
2089        impl Gear for SlowGear {
2090            async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
2091                Ok(())
2092            }
2093        }
2094
2095        #[async_trait::async_trait]
2096        impl RunnableCapability for SlowGear {
2097            async fn start(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
2098                Ok(())
2099            }
2100            async fn stop(&self, deadline_token: CancellationToken) -> anyhow::Result<()> {
2101                // Simulate slow graceful shutdown (would take 10s, but deadline is 100ms)
2102                tokio::select! {
2103                    () = tokio::time::sleep(Duration::from_secs(10)) => {
2104                        self.graceful_completed.store(true, Ordering::SeqCst);
2105                    }
2106                    () = deadline_token.cancelled() => {
2107                        self.deadline_fired.store(true, Ordering::SeqCst);
2108                    }
2109                }
2110                Ok(())
2111            }
2112        }
2113
2114        let gear = Arc::new(SlowGear {
2115            graceful_completed: AtomicBool::new(false),
2116            deadline_fired: AtomicBool::new(false),
2117        });
2118
2119        let mut builder = RegistryBuilder::default();
2120        builder.register_core_with_meta("test", &[], gear.clone() as Arc<dyn Gear>);
2121        builder.register_stateful_with_meta("test", gear.clone() as Arc<dyn RunnableCapability>);
2122
2123        let registry = builder.build_topo_sorted().unwrap();
2124        let client_hub = Arc::new(ClientHub::new());
2125        let cancel = CancellationToken::new();
2126        let config_provider: Arc<dyn ConfigProvider> = Arc::new(EmptyConfigProvider);
2127
2128        // Use a short deadline (100ms) - gear should be interrupted by deadline
2129        let runtime = HostRuntime::new(
2130            registry,
2131            config_provider,
2132            DbOptions::None,
2133            client_hub,
2134            cancel.clone(),
2135            Uuid::new_v4(),
2136            None,
2137        )
2138        .with_shutdown_deadline(Duration::from_millis(100));
2139
2140        runtime.run_stop_phase().await.unwrap();
2141
2142        // Graceful shutdown should NOT have completed (deadline fired first)
2143        assert!(
2144            !gear.graceful_completed.load(Ordering::SeqCst),
2145            "graceful shutdown should not complete when deadline fires first"
2146        );
2147        // Deadline should have fired
2148        assert!(
2149            gear.deadline_fired.load(Ordering::SeqCst),
2150            "deadline should fire for slow gears"
2151        );
2152    }
2153}