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