1use 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#[derive(Clone)]
46pub enum DbOptions {
47 None,
49 #[cfg(feature = "db")]
51 Manager(Arc<toolkit_db::DbManager>),
52}
53
54#[derive(Debug, Clone, Copy, PartialEq, Eq)]
56pub enum RunMode {
57 Full,
59 MigrateOnly,
61}
62
63pub const TOOLKIT_DIRECTORY_ENDPOINT_ENV: &str = "TOOLKIT_DIRECTORY_ENDPOINT";
65
66pub const TOOLKIT_MODULE_CONFIG_ENV: &str = "TOOLKIT_MODULE_CONFIG";
68
69pub const DEFAULT_SHUTDOWN_DEADLINE: std::time::Duration = std::time::Duration::from_secs(35);
75
76fn 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 gears_cfg: Arc<dyn ConfigProvider>,
106 dep_checker: Arc<super::readiness::DependencyChecker>,
112 rest_providers_registered: std::sync::atomic::AtomicBool,
118 cancel: CancellationToken,
119 #[allow(dead_code)]
120 db_options: DbOptions,
121 oop_options: Option<OopSpawnOptions>,
123 shutdown_deadline: std::time::Duration,
125}
126
127impl HostRuntime {
128 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 let gear_manager = Arc::new(GearManager::new());
142 let grpc_installers = Arc::new(GrpcInstallerStore::new());
143
144 let dep_checker = Arc::new(super::readiness::DependencyChecker::new());
148 client_hub.register::<super::readiness::DependencyChecker>(dep_checker.clone());
149
150 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 #[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 pub fn run_pre_init_phase(&self) -> Result<(), RegistryError> {
208 tracing::info!("Phase: pre_init");
209
210 let sys_ctx = SystemContext::new(
211 self.instance_id,
212 Arc::clone(&self.gear_manager),
213 Arc::clone(&self.grpc_installers),
214 );
215
216 for entry in self.registry.gears() {
217 if self.cancel.is_cancelled() {
219 tracing::warn!("Pre-init phase cancelled by signal");
220 return Err(RegistryError::Cancelled);
221 }
222
223 if let Some(sys_mod) = entry.caps.query::<SystemCap>() {
224 tracing::debug!(gear = entry.name, "Running system pre_init");
225 sys_mod
226 .pre_init(&sys_ctx)
227 .map_err(|e| RegistryError::PreInit {
228 gear: entry.name,
229 source: e,
230 })?;
231 }
232 }
233
234 Ok(())
235 }
236
237 #[cfg(feature = "db")]
239 async fn gear_context(
240 &self,
241 gear_name: &'static str,
242 ) -> Result<crate::context::GearCtx, RegistryError> {
243 self.ctx_builder
244 .for_gear(gear_name)
245 .await
246 .map_err(|e| RegistryError::DbMigrate {
247 gear: gear_name,
248 source: e,
249 })
250 }
251
252 #[cfg(feature = "db")]
254 async fn db_migration_target(
255 &self,
256 gear_name: &'static str,
257 ctx: &crate::context::GearCtx,
258 db_gear: Option<Arc<dyn crate::contracts::DatabaseCapability>>,
259 ) -> Result<
260 Option<(
261 toolkit_db::Db,
262 Arc<dyn crate::contracts::DatabaseCapability>,
263 )>,
264 RegistryError,
265 > {
266 let Some(dbm) = db_gear else {
267 return Ok(None);
268 };
269
270 let db = match &self.db_options {
274 DbOptions::None => None,
275 #[cfg(feature = "db")]
276 DbOptions::Manager(mgr) => {
277 mgr.get(gear_name)
278 .await
279 .map_err(|e| RegistryError::DbMigrate {
280 gear: gear_name,
281 source: e.into(),
282 })?
283 }
284 };
285
286 _ = ctx; Ok(db.map(|db| (db, dbm)))
288 }
289
290 #[cfg(feature = "db")]
295 async fn migrate_gear(
296 gear_name: &'static str,
297 db: &toolkit_db::Db,
298 db_gear: Arc<dyn crate::contracts::DatabaseCapability>,
299 ) -> Result<(), RegistryError> {
300 let migrations = db_gear.migrations();
302
303 if migrations.is_empty() {
304 tracing::debug!(gear = gear_name, "No migrations to run");
305 return Ok(());
306 }
307
308 tracing::debug!(
309 gear = gear_name,
310 count = migrations.len(),
311 "Running DB migrations"
312 );
313
314 let result =
316 toolkit_db::migration_runner::run_migrations_for_gear(db, gear_name, migrations)
317 .await
318 .map_err(|e| RegistryError::DbMigrate {
319 gear: gear_name,
320 source: anyhow::Error::new(e),
321 })?;
322
323 tracing::info!(
324 gear = gear_name,
325 applied = result.applied,
326 skipped = result.skipped,
327 "DB migrations completed"
328 );
329
330 Ok(())
331 }
332
333 #[cfg(feature = "db")]
342 async fn run_db_phase(&self) -> Result<(), RegistryError> {
343 tracing::info!("Phase: db (before init)");
344
345 for entry in self.registry.gears_by_system_priority() {
346 if self.cancel.is_cancelled() {
348 tracing::warn!("DB migration phase cancelled by signal");
349 return Err(RegistryError::Cancelled);
350 }
351
352 let ctx = self.gear_context(entry.name).await?;
353 let db_gear = entry.caps.query::<DatabaseCap>();
354
355 match self
356 .db_migration_target(entry.name, &ctx, db_gear.clone())
357 .await?
358 {
359 Some((db, dbm)) => {
360 Self::migrate_gear(entry.name, &db, dbm).await?;
361 }
362 None if db_gear.is_some() => {
363 tracing::debug!(
364 gear = entry.name,
365 "Gear has DbGear trait but no DB handle (no config)"
366 );
367 }
368 None => {}
369 }
370 }
371
372 Ok(())
373 }
374
375 async fn run_init_phase(&self) -> Result<(), RegistryError> {
379 tracing::info!("Phase: init");
380
381 for entry in self.registry.gears_by_system_priority() {
382 let ctx =
383 self.ctx_builder
384 .for_gear(entry.name)
385 .await
386 .map_err(|e| RegistryError::Init {
387 gear: entry.name,
388 source: e,
389 })?;
390 tracing::info!(gear = entry.name, "Initializing a gear...");
391 entry
392 .core
393 .init(&ctx)
394 .await
395 .map_err(|e| RegistryError::Init {
396 gear: entry.name,
397 source: e,
398 })?;
399 tracing::info!(gear = entry.name, "Initialized a gear.");
400 }
401
402 Ok(())
403 }
404
405 #[allow(
420 clippy::unused_async,
421 reason = "kept async for symmetry with the other `run_*_phase` steps awaited in sequence by `run_gear_phases`; the awaited work runs in a spawned readiness-probe task"
422 )]
423 async fn run_proxy_wiring_phase(&self) -> Result<(), RegistryError> {
424 use crate::discovery::{
425 ConsumerRegistration, DirectoryEndpointResolver, NullEndpointResolver,
426 };
427 use toolkit_contract::runtime::resolving::EndpointResolver;
428
429 let regs: Vec<&ConsumerRegistration> = inventory::iter::<ConsumerRegistration>
430 .into_iter()
431 .collect();
432 if regs.is_empty() {
433 return Ok(());
434 }
435 tracing::info!(
436 count = regs.len(),
437 "Phase: proxy-wiring (consumer discovery)"
438 );
439
440 let (resolver, have_directory): (Arc<dyn EndpointResolver>, bool) =
447 if let Ok(dir) = self.client_hub.get::<dyn crate::DirectoryClient>() {
448 (Arc::new(DirectoryEndpointResolver::new(dir)), true)
449 } else {
450 tracing::error!(
451 consumers = regs.len(),
452 "proxy-wiring: no DirectoryClient in ClientHub; remote consumer \
453 dependencies cannot be resolved and will gate /readyz (503). \
454 Co-located (local) dependencies are unaffected."
455 );
456 (Arc::new(NullEndpointResolver), false)
457 };
458
459 let known_gears: std::collections::HashSet<&str> =
471 self.registry.gears().iter().map(GearEntry::name).collect();
472 for reg in ®s {
473 if !known_gears.contains(reg.owner_gear) {
474 tracing::warn!(
475 owner = reg.owner_gear,
476 dep = reg.dep_gear,
477 "proxy-wiring: consumer's owner gear name does not match any registered gear; \
478 the `gears.{}.config.consumer_wiring.{}` static override will never resolve. \
479 Rename the gear to the kebab-case of its struct ident.",
480 reg.owner_gear,
481 reg.dep_gear,
482 );
483 }
484 }
485
486 let mut remote_deps: Vec<String> = Vec::new();
487 for reg in ®s {
488 let static_override =
494 static_endpoint_override(self.gears_cfg.as_ref(), reg.owner_gear, reg.dep_gear);
495 let (reg_resolver, is_static): (Arc<dyn EndpointResolver>, bool) =
496 if let Some(endpoint) = &static_override {
497 tracing::warn!(
498 owner = reg.owner_gear,
499 dep = reg.dep_gear,
500 endpoint = %endpoint,
501 "proxy-wiring: STATIC endpoint override in use (ADR-0004 dev/test \
502 escape hatch) - bypasses service discovery; MUST NOT be used in \
503 production"
504 );
505 (
506 Arc::new(crate::discovery::StaticEndpointResolver::new(
507 endpoint.clone(),
508 )),
509 true,
510 )
511 } else {
512 (Arc::clone(&resolver), false)
513 };
514
515 let outcome = (reg.wire)(&self.client_hub, reg_resolver).map_err(|source| {
516 RegistryError::ProxyWiring {
517 gear: reg.owner_gear,
518 source,
519 }
520 })?;
521 self.dep_checker.register_dep(reg.dep_gear.to_owned());
522 match outcome {
523 crate::discovery::WireOutcome::Local => {
525 self.dep_checker.mark_resolved(reg.dep_gear);
526 }
527 crate::discovery::WireOutcome::Remote if is_static => {
529 self.dep_checker.mark_resolved(reg.dep_gear);
530 }
531 crate::discovery::WireOutcome::Remote => remote_deps.push(reg.dep_gear.to_owned()),
533 }
534 tracing::debug!(
535 owner = reg.owner_gear,
536 dep = reg.dep_gear,
537 outcome = ?outcome,
538 static_override = is_static,
539 "wired consumer contract"
540 );
541 }
542
543 if !have_directory || remote_deps.is_empty() {
547 return Ok(());
548 }
549
550 let readiness = Arc::clone(&self.dep_checker);
551 let cancel = self.cancel.clone();
552 tokio::spawn(async move {
553 const BASE: std::time::Duration = std::time::Duration::from_millis(100);
554 const MAX: std::time::Duration = std::time::Duration::from_secs(30);
555 let mut pending = remote_deps;
556 let mut backoff = BASE;
557 while !pending.is_empty() {
558 let mut still_pending = Vec::new();
559 for dep in pending {
560 match resolver.resolve_endpoint(&dep).await {
561 Ok(Some(_)) => {
562 readiness.mark_resolved(&dep);
563 tracing::info!(dep = %dep, "readiness: dependency resolved");
564 }
565 Ok(None) => still_pending.push(dep),
567 Err(e) => {
570 tracing::warn!(dep = %dep, error = %e, "readiness: directory lookup failed");
571 still_pending.push(dep);
572 }
573 }
574 }
575 pending = still_pending;
576 if pending.is_empty() {
577 break;
578 }
579 tokio::select! {
580 () = cancel.cancelled() => break,
581 () = tokio::time::sleep(backoff) => {}
582 }
583 backoff = (backoff * 2).min(MAX);
584 }
585 });
586
587 Ok(())
588 }
589
590 async fn run_init_wiring_post_init(&self) -> Result<(), RegistryError> {
609 self.run_init_phase().await?;
610
611 self.run_proxy_wiring_phase().await?;
612
613 self.run_post_init_phase().await
614 }
615
616 async fn run_post_init_phase(&self) -> Result<(), RegistryError> {
619 tracing::info!("Phase: post_init");
620
621 let sys_ctx = SystemContext::new(
622 self.instance_id,
623 Arc::clone(&self.gear_manager),
624 Arc::clone(&self.grpc_installers),
625 );
626
627 for entry in self.registry.gears_by_system_priority() {
628 if let Some(sys_mod) = entry.caps.query::<SystemCap>() {
629 sys_mod
630 .post_init(&sys_ctx)
631 .await
632 .map_err(|e| RegistryError::PostInit {
633 gear: entry.name,
634 source: e,
635 })?;
636 }
637 }
638
639 Ok(())
640 }
641
642 async fn run_rest_phase(&self) -> Result<Router, RegistryError> {
649 tracing::info!("Phase: rest (sync)");
650
651 let mut router = Router::new();
652
653 let host_count = self
655 .registry
656 .gears()
657 .iter()
658 .filter(|e| e.caps.has::<ApiGatewayCap>())
659 .count();
660
661 match host_count {
662 0 => {
663 return if self
664 .registry
665 .gears()
666 .iter()
667 .any(|e| e.caps.has::<RestApiCap>())
668 {
669 Err(RegistryError::RestRequiresHost)
670 } else {
671 Ok(router)
672 };
673 }
674 1 => { }
675 _ => return Err(RegistryError::MultipleRestHosts),
676 }
677
678 let host_idx = self
680 .registry
681 .gears()
682 .iter()
683 .position(|e| e.caps.has::<ApiGatewayCap>())
684 .ok_or(RegistryError::RestHostNotFoundAfterValidation)?;
685 let host_entry = &self.registry.gears()[host_idx];
686 let Some(host) = host_entry.caps.query::<ApiGatewayCap>() else {
687 return Err(RegistryError::RestHostMissingFromEntry);
688 };
689 let host_ctx = self
690 .ctx_builder
691 .for_gear(host_entry.name)
692 .await
693 .map_err(|e| RegistryError::RestPrepare {
694 gear: host_entry.name,
695 source: e,
696 })?;
697
698 let registry: &dyn crate::contracts::OpenApiRegistry = host.as_registry();
700
701 let hc_registry = Arc::new(
705 crate::healthcheck::RestHealthcheckRegistry::with_cancellation(
706 host_ctx.cancellation_token().clone(),
707 ),
708 );
709
710 hc_registry.register(
716 "readiness",
717 Arc::new(super::readiness::ReadinessHealthcheck::new(
718 self.dep_checker.clone(),
719 )),
720 );
721
722 router = host
724 .rest_prepare(&host_ctx, router, hc_registry.clone())
725 .map_err(|source| RegistryError::RestPrepare {
726 gear: host_entry.name,
727 source,
728 })?;
729
730 for e in self.registry.gears() {
732 if let Some(rest) = e.caps.query::<RestApiCap>() {
733 let ctx = self.ctx_builder.for_gear(e.name).await.map_err(|err| {
734 RegistryError::RestRegister {
735 gear: e.name,
736 source: err,
737 }
738 })?;
739
740 router = rest
741 .register_rest(&ctx, router, registry)
742 .map_err(|source| RegistryError::RestRegister {
743 gear: e.name,
744 source,
745 })?;
746
747 if let Some(hc) = rest.healthcheck(&ctx) {
749 hc_registry.register(e.name, hc);
750 }
751 }
752 }
753
754 router = host
756 .rest_finalize(&host_ctx, router, hc_registry)
757 .map_err(|source| RegistryError::RestFinalize {
758 gear: host_entry.name,
759 source,
760 })?;
761
762 Ok(router)
763 }
764
765 async fn run_grpc_phase(&self) -> Result<(), RegistryError> {
769 tracing::info!("Phase: grpc (registration)");
770
771 if self.registry.grpc_hub.is_none() && self.registry.grpc_services.is_empty() {
773 return Ok(());
774 }
775
776 if self.registry.grpc_hub.is_none() && !self.registry.grpc_services.is_empty() {
778 return Err(RegistryError::GrpcRequiresHub);
779 }
780
781 if let Some(hub_name) = &self.registry.grpc_hub {
783 let mut gears_data = Vec::new();
784 let mut seen = HashSet::new();
785
786 for (gear_name, service_gear) in &self.registry.grpc_services {
788 let ctx = self.ctx_builder.for_gear(gear_name).await.map_err(|err| {
789 RegistryError::GrpcRegister {
790 gear: gear_name.clone(),
791 source: err,
792 }
793 })?;
794
795 let installers = service_gear
796 .get_grpc_services(&ctx)
797 .await
798 .map_err(|source| RegistryError::GrpcRegister {
799 gear: gear_name.clone(),
800 source,
801 })?;
802
803 for reg in &installers {
804 if !seen.insert(reg.service_name) {
805 return Err(RegistryError::GrpcRegister {
806 gear: gear_name.clone(),
807 source: anyhow::anyhow!(
808 "Duplicate gRPC service name: {}",
809 reg.service_name
810 ),
811 });
812 }
813 }
814
815 gears_data.push(crate::runtime::GearInstallers {
816 gear_name: gear_name.clone(),
817 installers,
818 });
819 }
820
821 self.grpc_installers
822 .set(crate::runtime::GrpcInstallerData { gears: gears_data })
823 .map_err(|source| RegistryError::GrpcRegister {
824 gear: hub_name.clone(),
825 source,
826 })?;
827 }
828
829 Ok(())
830 }
831
832 async fn run_start_phase(&self) -> Result<(), RegistryError> {
836 tracing::info!("Phase: start");
837
838 for e in self.registry.gears_by_system_priority() {
839 if let Some(s) = e.caps.query::<RunnableCap>() {
840 tracing::debug!(
841 gear = e.name,
842 is_system = e.caps.has::<SystemCap>(),
843 "Starting stateful gear"
844 );
845 s.start(self.cancel.clone())
846 .await
847 .map_err(|source| RegistryError::Start {
848 gear: e.name,
849 source,
850 })?;
851 tracing::info!(gear = e.name, "Started gear");
852 }
853 }
854
855 Ok(())
856 }
857
858 async fn stop_one_gear(entry: &GearEntry, cancel: CancellationToken) {
860 if let Some(s) = entry.caps.query::<RunnableCap>() {
861 match s.stop(cancel).await {
862 Err(err) => {
863 tracing::warn!(gear = entry.name, error = %err, "Failed to stop gear");
864 }
865 _ => {
866 tracing::info!(gear = entry.name, "Stopped gear");
867 }
868 }
869 }
870 }
871
872 async fn run_stop_phase(&self) -> Result<(), RegistryError> {
897 tracing::info!("Phase: stop");
898
899 self.deregister_rest_providers().await;
902
903 let deadline = self.shutdown_deadline;
904
905 for e in self.registry.gears().iter().rev() {
907 let gear_name = e.name;
908
909 let deadline_token = CancellationToken::new();
912 let deadline_token_for_timeout = deadline_token.clone();
913
914 let deadline_task = tokio::spawn(async move {
916 tokio::time::sleep(deadline).await;
917 tracing::warn!(
918 gear = gear_name,
919 deadline_secs = deadline.as_secs(),
920 "Gear shutdown deadline reached, sending hard-stop signal"
921 );
922 deadline_token_for_timeout.cancel();
923 });
924
925 Self::stop_one_gear(e, deadline_token).await;
928
929 deadline_task.abort();
931 #[allow(clippy::let_underscore_must_use)]
932 let _ = deadline_task.await;
933 }
934
935 Ok(())
936 }
937
938 async fn run_stop_phase_guarded(&self) -> Result<(), RegistryError> {
943 let gear_count = u32::try_from(self.registry.gears().len().max(1)).unwrap_or(1);
944 let stop_timeout = self
945 .shutdown_deadline
946 .checked_mul(gear_count)
947 .and_then(|d| d.checked_add(std::time::Duration::from_secs(5)))
948 .unwrap_or(self.shutdown_deadline);
949
950 let (disarm_tx, disarm_rx) = std::sync::mpsc::channel::<()>();
954 std::thread::spawn(move || {
955 match disarm_rx.recv_timeout(stop_timeout) {
956 Ok(()) | Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
957 }
959 Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
960 tracing::warn!(
961 timeout_secs = stop_timeout.as_secs(),
962 "shutdown: stop phase timed out, force exiting"
963 );
964 std::process::exit(1);
965 }
966 }
967 });
968
969 let stop_result = self.run_stop_phase().await;
970 let _ = disarm_tx.send(()).ok();
974
975 stop_result
976 }
977
978 async fn run_oop_spawn_phase(&self) -> Result<(), RegistryError> {
983 let oop_opts = match &self.oop_options {
984 Some(opts) if !opts.gears.is_empty() => opts,
985 _ => return Ok(()),
986 };
987
988 tracing::info!("Phase: oop_spawn");
989
990 let directory_endpoint = self.wait_for_grpc_hub_endpoint().await;
992
993 for gear_cfg in &oop_opts.gears {
994 let mut env = gear_cfg.env.clone();
997 env.insert(
998 TOOLKIT_MODULE_CONFIG_ENV.to_owned(),
999 gear_cfg.rendered_config_json.clone(),
1000 );
1001 if let Some(ref endpoint) = directory_endpoint {
1002 env.insert(TOOLKIT_DIRECTORY_ENDPOINT_ENV.to_owned(), endpoint.clone());
1003 }
1004
1005 let args = gear_cfg.args.clone();
1007
1008 let spawn_config = OopSpawnConfig {
1009 gear_name: gear_cfg.gear_name.clone(),
1010 binary: gear_cfg.binary.clone(),
1011 args,
1012 env,
1013 working_directory: gear_cfg.working_directory.clone(),
1014 };
1015
1016 oop_opts
1017 .backend
1018 .spawn(spawn_config)
1019 .await
1020 .map_err(|e| RegistryError::OopSpawn {
1021 gear: gear_cfg.gear_name.clone(),
1022 source: e,
1023 })?;
1024
1025 tracing::info!(
1026 gear = %gear_cfg.gear_name,
1027 directory_endpoint = ?directory_endpoint,
1028 "Spawned OoP gear via backend"
1029 );
1030 }
1031
1032 Ok(())
1033 }
1034
1035 async fn wait_for_grpc_hub_endpoint(&self) -> Option<String> {
1040 const POLL_INTERVAL: std::time::Duration = std::time::Duration::from_millis(10);
1041 const MAX_WAIT: std::time::Duration = std::time::Duration::from_secs(5);
1042
1043 let grpc_hub = self
1045 .registry
1046 .gears()
1047 .iter()
1048 .find_map(|e| e.caps.query::<GrpcHubCap>());
1049
1050 let Some(hub) = grpc_hub else {
1051 return None; };
1053
1054 let start = std::time::Instant::now();
1055
1056 loop {
1057 if let Some(endpoint) = hub.bound_endpoint() {
1058 tracing::debug!(
1059 endpoint = %endpoint,
1060 elapsed_ms = start.elapsed().as_millis(),
1061 "gRPC hub endpoint available"
1062 );
1063 return Some(endpoint);
1064 }
1065
1066 if start.elapsed() > MAX_WAIT {
1067 tracing::warn!("Timed out waiting for gRPC hub to bind");
1068 return None;
1069 }
1070
1071 tokio::time::sleep(POLL_INTERVAL).await;
1072 }
1073 }
1074
1075 async fn wait_for_rest_endpoint(
1081 &self,
1082 host: &Arc<dyn crate::contracts::ApiGatewayCapability>,
1083 ) -> Option<String> {
1084 const POLL_INTERVAL: std::time::Duration = std::time::Duration::from_millis(10);
1085 const MAX_WAIT: std::time::Duration = std::time::Duration::from_secs(5);
1086
1087 let start = std::time::Instant::now();
1088 loop {
1089 if let Some(endpoint) = host.bound_endpoint() {
1090 return Some(endpoint);
1091 }
1092 if start.elapsed() > MAX_WAIT {
1093 tracing::warn!("Timed out waiting for REST host to bind");
1094 return None;
1095 }
1096 tokio::time::sleep(POLL_INTERVAL).await;
1097 }
1098 }
1099
1100 async fn run_directory_register_phase(&self) -> Result<(), RegistryError> {
1113 let rest_gears = self.rest_provider_gears();
1114 if rest_gears.is_empty() {
1115 return Ok(());
1116 }
1117
1118 let Some(host) = self
1119 .registry
1120 .gears()
1121 .iter()
1122 .find_map(|e| e.caps.query::<ApiGatewayCap>())
1123 else {
1124 return Ok(()); };
1126
1127 let Some(endpoint) = self.wait_for_rest_endpoint(&host).await else {
1128 tracing::warn!(
1129 "directory-register: REST host endpoint unavailable; skipping REST provider registration"
1130 );
1131 return Ok(());
1132 };
1133
1134 let Ok(dir) = self.client_hub.get::<dyn crate::DirectoryClient>() else {
1135 tracing::debug!(
1136 "directory-register: no DirectoryClient in ClientHub; skipping REST provider registration"
1137 );
1138 return Ok(());
1139 };
1140
1141 let instance_id = self.instance_id.to_string();
1142 for gear in rest_gears {
1143 let (grpc_services, version) = match dir.list_instances(gear).await {
1149 Ok(insts) => insts
1150 .into_iter()
1151 .find(|i| i.instance_id == instance_id)
1152 .map(|i| (i.grpc_services, i.version))
1153 .unwrap_or_default(),
1154 Err(_) => (Vec::new(), None),
1155 };
1156 let info = crate::RegisterInstanceInfo {
1157 gear: gear.to_owned(),
1158 instance_id: instance_id.clone(),
1159 grpc_services,
1160 version,
1161 rest_endpoint: Some(crate::ServiceEndpoint::new(endpoint.clone())),
1162 openapi_spec: None,
1165 };
1166 match dir.register_instance(info).await {
1167 Ok(()) => {
1168 tracing::info!(gear, endpoint = %endpoint, "registered REST provider in directory");
1169 }
1170 Err(e) => {
1171 tracing::warn!(gear, error = %e, "directory-register: failed to register REST provider");
1172 }
1173 }
1174 }
1175 self.rest_providers_registered
1179 .store(true, std::sync::atomic::Ordering::SeqCst);
1180 Ok(())
1181 }
1182
1183 fn rest_provider_gears(&self) -> Vec<&'static str> {
1188 self.registry
1189 .gears()
1190 .iter()
1191 .filter(|e| e.caps.has::<RestApiCap>() && !e.caps.has::<ApiGatewayCap>())
1192 .map(|e| e.name)
1193 .collect()
1194 }
1195
1196 async fn deregister_rest_providers(&self) {
1199 if !self
1203 .rest_providers_registered
1204 .load(std::sync::atomic::Ordering::SeqCst)
1205 {
1206 return;
1207 }
1208 let rest_gears = self.rest_provider_gears();
1209 if rest_gears.is_empty() {
1210 return;
1211 }
1212 let Ok(dir) = self.client_hub.get::<dyn crate::DirectoryClient>() else {
1213 return;
1214 };
1215 let instance_id = self.instance_id.to_string();
1216 for gear in rest_gears {
1217 if let Err(e) = dir.deregister_instance(gear, &instance_id).await {
1218 tracing::warn!(gear, error = %e, "directory-deregister: failed to deregister REST provider");
1219 }
1220 }
1221 }
1222
1223 pub async fn run_gear_phases(self) -> anyhow::Result<()> {
1232 self.run_phases_internal(RunMode::Full).await
1233 }
1234
1235 pub async fn run_migration_phases(self) -> anyhow::Result<()> {
1245 self.run_phases_internal(RunMode::MigrateOnly).await
1246 }
1247
1248 async fn run_phases_internal(self, mode: RunMode) -> anyhow::Result<()> {
1271 match mode {
1273 RunMode::Full => {
1274 tracing::info!("Running full lifecycle (all phases)");
1275 }
1276 RunMode::MigrateOnly => {
1277 tracing::info!("Running in migration mode (pre-init + db phases only)");
1278 }
1279 }
1280
1281 self.run_pre_init_phase()?;
1283
1284 #[cfg(feature = "db")]
1286 {
1287 self.run_db_phase().await?;
1288 }
1289 #[cfg(not(feature = "db"))]
1290 {
1291 }
1293
1294 if mode == RunMode::MigrateOnly {
1296 tracing::info!("Migration phases completed successfully");
1297 return Ok(());
1298 }
1299
1300 self.run_init_wiring_post_init().await?;
1302
1303 let _router = self.run_rest_phase().await?;
1305
1306 self.run_grpc_phase().await?;
1308
1309 self.run_start_phase().await?;
1311
1312 {
1316 let readiness = Arc::clone(&self.dep_checker);
1317 let cancel = self.cancel.clone();
1318 tokio::spawn(async move {
1319 cancel.cancelled().await;
1320 readiness.set_draining(true);
1321 });
1322 }
1323
1324 self.run_directory_register_phase().await?;
1327
1328 self.run_oop_spawn_phase().await?;
1330
1331 self.cancel.cancelled().await;
1333
1334 self.run_stop_phase_guarded().await?;
1338 Ok(())
1339 }
1340}
1341
1342#[cfg(feature = "bootstrap")]
1344impl HostRuntime {
1345 async fn compose_oop_router(
1351 &self,
1352 options: &crate::runtime::OopServeOptions,
1353 hc_registry: &Arc<crate::healthcheck::RestHealthcheckRegistry>,
1354 ) -> anyhow::Result<(Router, String)> {
1355 use crate::api::{OpenApiInfo, OpenApiRegistryImpl};
1356 use anyhow::Context as _;
1357
1358 let registry = OpenApiRegistryImpl::new();
1359 let mut router = Router::new();
1360
1361 for entry in self.registry.gears() {
1362 if let Some(rest) = entry.caps.query::<RestApiCap>() {
1363 let ctx = self
1364 .ctx_builder
1365 .for_gear(entry.name)
1366 .await
1367 .with_context(|| format!("OoP router: build context for '{}'", entry.name))?;
1368 router = rest
1369 .register_rest(&ctx, router, ®istry)
1370 .with_context(|| format!("OoP router: register_rest for '{}'", entry.name))?;
1371
1372 if let Some(hc) = rest.healthcheck(&ctx) {
1376 hc_registry.register(entry.name, hc);
1377 }
1378 }
1379 }
1380
1381 let info = OpenApiInfo {
1382 title: options.gear_name.clone(),
1383 version: options
1384 .version
1385 .clone()
1386 .unwrap_or_else(|| "0.0.0".to_owned()),
1387 description: None,
1388 servers: vec![],
1389 };
1390 let openapi = registry
1391 .build_openapi(&info)
1392 .context("OoP router: build OpenAPI document")?;
1393 let json = serde_json::to_string(&openapi).context("OoP router: serialize OpenAPI")?;
1394
1395 Ok((router, json))
1396 }
1397
1398 pub async fn run_oop_serving(
1406 self,
1407 options: crate::runtime::OopServeOptions,
1408 ) -> anyhow::Result<()> {
1409 use crate::runtime::ReadinessState;
1410
1411 tracing::info!("Running OoP serving lifecycle");
1412
1413 if self.client_hub.get::<dyn crate::DirectoryClient>().is_err() {
1418 self.client_hub
1419 .register::<dyn crate::DirectoryClient>(Arc::clone(&options.directory));
1420 }
1421
1422 let hc_registry = Arc::new(
1426 crate::healthcheck::RestHealthcheckRegistry::with_cancellation(self.cancel.clone()),
1427 );
1428
1429 let readiness = ReadinessState::from_checker(
1436 Arc::clone(&self.dep_checker),
1437 Arc::clone(&hc_registry),
1438 options.healthcheck_timeout,
1439 );
1440
1441 let mut server = super::oop_serve::OopHttpServer::start(
1446 Arc::clone(&readiness),
1447 options,
1448 self.cancel.clone(),
1449 )
1450 .await?;
1451
1452 let mut started = false;
1458 let composed: anyhow::Result<(Router, String)> = async {
1459 self.run_pre_init_phase()?;
1460 #[cfg(feature = "db")]
1461 self.run_db_phase().await?;
1462 self.run_init_wiring_post_init().await?;
1466 self.run_grpc_phase().await?;
1467 self.run_start_phase().await?;
1468 started = true;
1469 server.resolve_bearer_authenticator(&self.client_hub);
1475 self.compose_oop_router(server.options(), &hc_registry)
1476 .await
1477 }
1478 .await;
1479
1480 let serve_result = match composed {
1481 Ok((gear_router, openapi_json)) => {
1482 server.attach(gear_router, openapi_json);
1485 server.join().await
1487 }
1488 Err(e) => {
1489 tracing::error!(error = %e, "OoP startup failed before serving gear routes");
1490 self.cancel.cancel();
1492 if let Err(join_err) = server.join().await {
1493 tracing::warn!(error = %join_err, "OoP probe server teardown after startup failure errored");
1494 }
1495 Err(e)
1496 }
1497 };
1498
1499 if started && let Err(e) = self.run_stop_phase_guarded().await {
1502 tracing::warn!(error = %e, "OoP stop phase reported an error");
1503 }
1504
1505 serve_result
1506 }
1507}
1508
1509#[cfg(test)]
1510#[cfg(feature = "bootstrap")]
1511#[cfg_attr(coverage_nightly, coverage(off))]
1512#[path = "host_runtime_oop_tests.rs"]
1513mod host_runtime_oop_tests;
1514
1515#[cfg(test)]
1516#[cfg_attr(coverage_nightly, coverage(off))]
1517mod tests {
1518 use super::*;
1519 use crate::context::GearCtx;
1520 use crate::contracts::{Gear, RunnableCapability, SystemCapability};
1521 use crate::registry::RegistryBuilder;
1522 use std::sync::Arc;
1523 use std::sync::atomic::{AtomicUsize, Ordering};
1524 use tokio::sync::Mutex;
1525
1526 #[derive(Default)]
1527 #[allow(dead_code)]
1528 struct DummyCore;
1529 #[async_trait::async_trait]
1530 impl Gear for DummyCore {
1531 async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
1532 Ok(())
1533 }
1534 }
1535
1536 struct StopOrderTracker {
1537 my_order: usize,
1538 stop_order: Arc<AtomicUsize>,
1539 }
1540
1541 impl StopOrderTracker {
1542 fn new(counter: &Arc<AtomicUsize>, stop_order: Arc<AtomicUsize>) -> Self {
1543 let my_order = counter.fetch_add(1, Ordering::SeqCst);
1544 Self {
1545 my_order,
1546 stop_order,
1547 }
1548 }
1549 }
1550
1551 #[async_trait::async_trait]
1552 impl Gear for StopOrderTracker {
1553 async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
1554 Ok(())
1555 }
1556 }
1557
1558 #[async_trait::async_trait]
1559 impl RunnableCapability for StopOrderTracker {
1560 async fn start(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
1561 Ok(())
1562 }
1563 async fn stop(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
1564 let order = self.stop_order.fetch_add(1, Ordering::SeqCst);
1565 tracing::info!(my_order = self.my_order, stop_order = order, "Gear stopped");
1566 Ok(())
1567 }
1568 }
1569
1570 #[tokio::test]
1571 async fn test_stop_phase_reverse_order() {
1572 let counter = Arc::new(AtomicUsize::new(0));
1573 let stop_order = Arc::new(AtomicUsize::new(0));
1574
1575 let gear_a = Arc::new(StopOrderTracker::new(&counter, stop_order.clone()));
1576 let gear_b = Arc::new(StopOrderTracker::new(&counter, stop_order.clone()));
1577 let gear_c = Arc::new(StopOrderTracker::new(&counter, stop_order.clone()));
1578
1579 let mut builder = RegistryBuilder::default();
1580 builder.register_core_with_meta("a", &[], gear_a.clone() as Arc<dyn Gear>);
1581 builder.register_core_with_meta("b", &["a"], gear_b.clone() as Arc<dyn Gear>);
1582 builder.register_core_with_meta("c", &["b"], gear_c.clone() as Arc<dyn Gear>);
1583
1584 builder.register_stateful_with_meta("a", gear_a.clone() as Arc<dyn RunnableCapability>);
1585 builder.register_stateful_with_meta("b", gear_b.clone() as Arc<dyn RunnableCapability>);
1586 builder.register_stateful_with_meta("c", gear_c.clone() as Arc<dyn RunnableCapability>);
1587
1588 let registry = builder.build_topo_sorted().unwrap();
1589
1590 let gear_names: Vec<_> = registry.gears().iter().map(|m| m.name).collect();
1592 assert_eq!(gear_names, vec!["a", "b", "c"]);
1593
1594 let client_hub = Arc::new(ClientHub::new());
1595 let cancel = CancellationToken::new();
1596 let config_provider: Arc<dyn ConfigProvider> = Arc::new(EmptyConfigProvider);
1597
1598 let runtime = HostRuntime::new(
1599 registry,
1600 config_provider,
1601 DbOptions::None,
1602 client_hub,
1603 cancel.clone(),
1604 Uuid::new_v4(),
1605 None,
1606 );
1607
1608 runtime.run_stop_phase().await.unwrap();
1610
1611 assert_eq!(stop_order.load(Ordering::SeqCst), 3);
1615 }
1616
1617 #[tokio::test]
1618 async fn test_stop_phase_continues_on_error() {
1619 struct FailingGear {
1620 should_fail: bool,
1621 stopped: Arc<AtomicUsize>,
1622 }
1623
1624 #[async_trait::async_trait]
1625 impl Gear for FailingGear {
1626 async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
1627 Ok(())
1628 }
1629 }
1630
1631 #[async_trait::async_trait]
1632 impl RunnableCapability for FailingGear {
1633 async fn start(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
1634 Ok(())
1635 }
1636 async fn stop(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
1637 self.stopped.fetch_add(1, Ordering::SeqCst);
1638 if self.should_fail {
1639 anyhow::bail!("Intentional failure")
1640 }
1641 Ok(())
1642 }
1643 }
1644
1645 let stopped = Arc::new(AtomicUsize::new(0));
1646 let gear_a = Arc::new(FailingGear {
1647 should_fail: false,
1648 stopped: stopped.clone(),
1649 });
1650 let gear_b = Arc::new(FailingGear {
1651 should_fail: true,
1652 stopped: stopped.clone(),
1653 });
1654 let gear_c = Arc::new(FailingGear {
1655 should_fail: false,
1656 stopped: stopped.clone(),
1657 });
1658
1659 let mut builder = RegistryBuilder::default();
1660 builder.register_core_with_meta("a", &[], gear_a.clone() as Arc<dyn Gear>);
1661 builder.register_core_with_meta("b", &["a"], gear_b.clone() as Arc<dyn Gear>);
1662 builder.register_core_with_meta("c", &["b"], gear_c.clone() as Arc<dyn Gear>);
1663
1664 builder.register_stateful_with_meta("a", gear_a.clone() as Arc<dyn RunnableCapability>);
1665 builder.register_stateful_with_meta("b", gear_b.clone() as Arc<dyn RunnableCapability>);
1666 builder.register_stateful_with_meta("c", gear_c.clone() as Arc<dyn RunnableCapability>);
1667
1668 let registry = builder.build_topo_sorted().unwrap();
1669
1670 let client_hub = Arc::new(ClientHub::new());
1671 let cancel = CancellationToken::new();
1672 let config_provider: Arc<dyn ConfigProvider> = Arc::new(EmptyConfigProvider);
1673
1674 let runtime = HostRuntime::new(
1675 registry,
1676 config_provider,
1677 DbOptions::None,
1678 client_hub,
1679 cancel.clone(),
1680 Uuid::new_v4(),
1681 None,
1682 );
1683
1684 runtime.run_stop_phase().await.unwrap();
1686
1687 assert_eq!(stopped.load(Ordering::SeqCst), 3);
1689 }
1690
1691 struct EmptyConfigProvider;
1692 impl ConfigProvider for EmptyConfigProvider {
1693 fn get_gear_config(&self, _gear_name: &str) -> Option<&serde_json::Value> {
1694 None
1695 }
1696 }
1697
1698 #[test]
1699 fn static_endpoint_override_reads_nested_consumer_wiring_key() {
1700 struct MapCfg(std::collections::HashMap<String, serde_json::Value>);
1701 impl ConfigProvider for MapCfg {
1702 fn get_gear_config(&self, gear: &str) -> Option<&serde_json::Value> {
1703 self.0.get(gear)
1704 }
1705 }
1706 let mut map = std::collections::HashMap::new();
1707 map.insert(
1708 "orders".to_owned(),
1709 serde_json::json!({
1710 "config": { "consumer_wiring": { "billing": "http://localhost:8081" } }
1711 }),
1712 );
1713 let cfg = MapCfg(map);
1714
1715 assert_eq!(
1717 super::static_endpoint_override(&cfg, "orders", "billing").as_deref(),
1718 Some("http://localhost:8081")
1719 );
1720 assert_eq!(
1722 super::static_endpoint_override(&cfg, "orders", "inventory"),
1723 None
1724 );
1725 assert_eq!(
1726 super::static_endpoint_override(&cfg, "warehouse", "billing"),
1727 None
1728 );
1729 assert_eq!(
1730 super::static_endpoint_override(&EmptyConfigProvider, "orders", "billing"),
1731 None
1732 );
1733 }
1734
1735 #[test]
1741 fn static_endpoint_override_is_keyed_by_kebab_gear_name() {
1742 struct MapCfg(std::collections::HashMap<String, serde_json::Value>);
1743 impl ConfigProvider for MapCfg {
1744 fn get_gear_config(&self, gear: &str) -> Option<&serde_json::Value> {
1745 self.0.get(gear)
1746 }
1747 }
1748 let mut map = std::collections::HashMap::new();
1749 map.insert(
1750 "api-contracts-consumer".to_owned(),
1751 serde_json::json!({
1752 "config": { "consumer_wiring": { "api-contracts": "http://localhost:9099" } }
1753 }),
1754 );
1755 let cfg = MapCfg(map);
1756
1757 assert_eq!(
1758 super::static_endpoint_override(&cfg, "api-contracts-consumer", "api-contracts")
1759 .as_deref(),
1760 Some("http://localhost:9099"),
1761 );
1762 assert_eq!(
1764 super::static_endpoint_override(&cfg, "ApiContractsConsumer", "api-contracts"),
1765 None,
1766 );
1767 }
1768
1769 #[tokio::test]
1770 async fn test_post_init_runs_after_all_init_and_system_first() {
1771 #[derive(Clone)]
1772 struct TrackHooks {
1773 name: &'static str,
1774 events: Arc<Mutex<Vec<String>>>,
1775 }
1776
1777 #[async_trait::async_trait]
1778 impl Gear for TrackHooks {
1779 async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
1780 self.events.lock().await.push(format!("init:{}", self.name));
1781 Ok(())
1782 }
1783 }
1784
1785 #[async_trait::async_trait]
1786 impl SystemCapability for TrackHooks {
1787 fn pre_init(&self, _sys: &crate::runtime::SystemContext) -> anyhow::Result<()> {
1788 Ok(())
1789 }
1790
1791 async fn post_init(&self, _sys: &crate::runtime::SystemContext) -> anyhow::Result<()> {
1792 self.events
1793 .lock()
1794 .await
1795 .push(format!("post_init:{}", self.name));
1796 Ok(())
1797 }
1798 }
1799
1800 let events = Arc::new(Mutex::new(Vec::<String>::new()));
1801 let sys_a = Arc::new(TrackHooks {
1802 name: "sys_a",
1803 events: events.clone(),
1804 });
1805 let user_b = Arc::new(TrackHooks {
1806 name: "user_b",
1807 events: events.clone(),
1808 });
1809 let user_c = Arc::new(TrackHooks {
1810 name: "user_c",
1811 events: events.clone(),
1812 });
1813
1814 let mut builder = RegistryBuilder::default();
1815 builder.register_core_with_meta("sys_a", &[], sys_a.clone() as Arc<dyn Gear>);
1816 builder.register_core_with_meta("user_b", &["sys_a"], user_b.clone() as Arc<dyn Gear>);
1817 builder.register_core_with_meta("user_c", &["user_b"], user_c.clone() as Arc<dyn Gear>);
1818 builder.register_system_with_meta("sys_a", sys_a.clone() as Arc<dyn SystemCapability>);
1819
1820 let registry = builder.build_topo_sorted().unwrap();
1821
1822 let client_hub = Arc::new(ClientHub::new());
1823 let cancel = CancellationToken::new();
1824 let config_provider: Arc<dyn ConfigProvider> = Arc::new(EmptyConfigProvider);
1825
1826 let runtime = HostRuntime::new(
1827 registry,
1828 config_provider,
1829 DbOptions::None,
1830 client_hub,
1831 cancel,
1832 Uuid::new_v4(),
1833 None,
1834 );
1835
1836 runtime.run_init_phase().await.unwrap();
1838 runtime.run_post_init_phase().await.unwrap();
1839
1840 let events = events.lock().await.clone();
1841 let first_post_init = events
1842 .iter()
1843 .position(|e| e.starts_with("post_init:"))
1844 .expect("expected post_init events");
1845 assert!(
1846 events[..first_post_init]
1847 .iter()
1848 .all(|e| e.starts_with("init:")),
1849 "expected all init events before post_init, got: {events:?}"
1850 );
1851
1852 assert_eq!(
1854 events,
1855 vec![
1856 "init:sys_a",
1857 "init:user_b",
1858 "init:user_c",
1859 "post_init:sys_a",
1860 ]
1861 );
1862 }
1863
1864 #[tokio::test]
1878 async fn init_wiring_post_init_runs_as_one_ordered_segment() {
1879 #[derive(Clone)]
1880 struct TrackHooks {
1881 events: Arc<Mutex<Vec<String>>>,
1882 }
1883
1884 #[async_trait::async_trait]
1885 impl Gear for TrackHooks {
1886 async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
1887 self.events.lock().await.push("init".to_owned());
1888 Ok(())
1889 }
1890 }
1891
1892 #[async_trait::async_trait]
1893 impl SystemCapability for TrackHooks {
1894 fn pre_init(&self, _sys: &crate::runtime::SystemContext) -> anyhow::Result<()> {
1895 Ok(())
1896 }
1897
1898 async fn post_init(&self, _sys: &crate::runtime::SystemContext) -> anyhow::Result<()> {
1899 self.events.lock().await.push("post_init".to_owned());
1900 Ok(())
1901 }
1902 }
1903
1904 let events = Arc::new(Mutex::new(Vec::<String>::new()));
1905 let gear = Arc::new(TrackHooks {
1906 events: events.clone(),
1907 });
1908
1909 let mut builder = RegistryBuilder::default();
1910 builder.register_core_with_meta("sys", &[], gear.clone() as Arc<dyn Gear>);
1911 builder.register_system_with_meta("sys", gear.clone() as Arc<dyn SystemCapability>);
1912 let registry = builder.build_topo_sorted().unwrap();
1913
1914 let runtime = HostRuntime::new(
1915 registry,
1916 Arc::new(EmptyConfigProvider) as Arc<dyn ConfigProvider>,
1917 DbOptions::None,
1918 Arc::new(ClientHub::new()),
1919 CancellationToken::new(),
1920 Uuid::new_v4(),
1921 None,
1922 );
1923
1924 runtime.run_init_wiring_post_init().await.unwrap();
1925
1926 assert_eq!(events.lock().await.clone(), vec!["init", "post_init"]);
1927 }
1928
1929 #[tokio::test]
1930 async fn test_stop_phase_provides_fresh_deadline_token() {
1931 use std::sync::atomic::AtomicBool;
1932
1933 struct TokenCheckGear {
1934 stop_was_called: AtomicBool,
1935 token_was_cancelled_on_entry: AtomicBool,
1936 }
1937
1938 #[async_trait::async_trait]
1939 impl Gear for TokenCheckGear {
1940 async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
1941 Ok(())
1942 }
1943 }
1944
1945 #[async_trait::async_trait]
1946 impl RunnableCapability for TokenCheckGear {
1947 async fn start(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
1948 Ok(())
1949 }
1950 async fn stop(&self, deadline_token: CancellationToken) -> anyhow::Result<()> {
1951 self.stop_was_called.store(true, Ordering::SeqCst);
1953 self.token_was_cancelled_on_entry
1955 .store(deadline_token.is_cancelled(), Ordering::SeqCst);
1956 Ok(())
1957 }
1958 }
1959
1960 let gear = Arc::new(TokenCheckGear {
1961 stop_was_called: AtomicBool::new(false),
1962 token_was_cancelled_on_entry: AtomicBool::new(true),
1964 });
1965
1966 let mut builder = RegistryBuilder::default();
1967 builder.register_core_with_meta("test", &[], gear.clone() as Arc<dyn Gear>);
1968 builder.register_stateful_with_meta("test", gear.clone() as Arc<dyn RunnableCapability>);
1969
1970 let registry = builder.build_topo_sorted().unwrap();
1971 let client_hub = Arc::new(ClientHub::new());
1972 let cancel = CancellationToken::new();
1973 let config_provider: Arc<dyn ConfigProvider> = Arc::new(EmptyConfigProvider);
1974
1975 let runtime = HostRuntime::new(
1976 registry,
1977 config_provider,
1978 DbOptions::None,
1979 client_hub,
1980 cancel.clone(),
1981 Uuid::new_v4(),
1982 None,
1983 );
1984
1985 runtime.run_stop_phase().await.unwrap();
1987
1988 assert!(
1990 gear.stop_was_called.load(Ordering::SeqCst),
1991 "stop() was never called - gear may not have been registered correctly"
1992 );
1993
1994 assert!(
1997 !gear.token_was_cancelled_on_entry.load(Ordering::SeqCst),
1998 "deadline_token should NOT be cancelled when stop() is called - this enables graceful shutdown"
1999 );
2000 }
2001
2002 #[tokio::test]
2003 async fn test_stop_phase_graceful_shutdown_completes_before_deadline() {
2004 use std::sync::atomic::AtomicBool;
2005 use std::time::Duration;
2006
2007 struct GracefulGear {
2008 graceful_completed: AtomicBool,
2009 deadline_fired: AtomicBool,
2010 }
2011
2012 #[async_trait::async_trait]
2013 impl Gear for GracefulGear {
2014 async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
2015 Ok(())
2016 }
2017 }
2018
2019 #[async_trait::async_trait]
2020 impl RunnableCapability for GracefulGear {
2021 async fn start(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
2022 Ok(())
2023 }
2024 async fn stop(&self, deadline_token: CancellationToken) -> anyhow::Result<()> {
2025 tokio::select! {
2027 () = tokio::time::sleep(Duration::from_millis(10)) => {
2028 self.graceful_completed.store(true, Ordering::SeqCst);
2029 }
2030 () = deadline_token.cancelled() => {
2031 self.deadline_fired.store(true, Ordering::SeqCst);
2032 }
2033 }
2034 Ok(())
2035 }
2036 }
2037
2038 let gear = Arc::new(GracefulGear {
2039 graceful_completed: AtomicBool::new(false),
2040 deadline_fired: AtomicBool::new(false),
2041 });
2042
2043 let mut builder = RegistryBuilder::default();
2044 builder.register_core_with_meta("test", &[], gear.clone() as Arc<dyn Gear>);
2045 builder.register_stateful_with_meta("test", gear.clone() as Arc<dyn RunnableCapability>);
2046
2047 let registry = builder.build_topo_sorted().unwrap();
2048 let client_hub = Arc::new(ClientHub::new());
2049 let cancel = CancellationToken::new();
2050 let config_provider: Arc<dyn ConfigProvider> = Arc::new(EmptyConfigProvider);
2051
2052 let runtime = HostRuntime::new(
2054 registry,
2055 config_provider,
2056 DbOptions::None,
2057 client_hub,
2058 cancel.clone(),
2059 Uuid::new_v4(),
2060 None,
2061 )
2062 .with_shutdown_deadline(Duration::from_secs(5));
2063
2064 runtime.run_stop_phase().await.unwrap();
2065
2066 assert!(
2068 gear.graceful_completed.load(Ordering::SeqCst),
2069 "graceful shutdown should complete"
2070 );
2071 assert!(
2073 !gear.deadline_fired.load(Ordering::SeqCst),
2074 "deadline should not fire when graceful shutdown completes quickly"
2075 );
2076 }
2077
2078 #[tokio::test]
2079 async fn test_stop_phase_deadline_fires_for_slow_gear() {
2080 use std::sync::atomic::AtomicBool;
2081 use std::time::Duration;
2082
2083 struct SlowGear {
2084 graceful_completed: AtomicBool,
2085 deadline_fired: AtomicBool,
2086 }
2087
2088 #[async_trait::async_trait]
2089 impl Gear for SlowGear {
2090 async fn init(&self, _ctx: &GearCtx) -> anyhow::Result<()> {
2091 Ok(())
2092 }
2093 }
2094
2095 #[async_trait::async_trait]
2096 impl RunnableCapability for SlowGear {
2097 async fn start(&self, _cancel: CancellationToken) -> anyhow::Result<()> {
2098 Ok(())
2099 }
2100 async fn stop(&self, deadline_token: CancellationToken) -> anyhow::Result<()> {
2101 tokio::select! {
2103 () = tokio::time::sleep(Duration::from_secs(10)) => {
2104 self.graceful_completed.store(true, Ordering::SeqCst);
2105 }
2106 () = deadline_token.cancelled() => {
2107 self.deadline_fired.store(true, Ordering::SeqCst);
2108 }
2109 }
2110 Ok(())
2111 }
2112 }
2113
2114 let gear = Arc::new(SlowGear {
2115 graceful_completed: AtomicBool::new(false),
2116 deadline_fired: AtomicBool::new(false),
2117 });
2118
2119 let mut builder = RegistryBuilder::default();
2120 builder.register_core_with_meta("test", &[], gear.clone() as Arc<dyn Gear>);
2121 builder.register_stateful_with_meta("test", gear.clone() as Arc<dyn RunnableCapability>);
2122
2123 let registry = builder.build_topo_sorted().unwrap();
2124 let client_hub = Arc::new(ClientHub::new());
2125 let cancel = CancellationToken::new();
2126 let config_provider: Arc<dyn ConfigProvider> = Arc::new(EmptyConfigProvider);
2127
2128 let runtime = HostRuntime::new(
2130 registry,
2131 config_provider,
2132 DbOptions::None,
2133 client_hub,
2134 cancel.clone(),
2135 Uuid::new_v4(),
2136 None,
2137 )
2138 .with_shutdown_deadline(Duration::from_millis(100));
2139
2140 runtime.run_stop_phase().await.unwrap();
2141
2142 assert!(
2144 !gear.graceful_completed.load(Ordering::SeqCst),
2145 "graceful shutdown should not complete when deadline fires first"
2146 );
2147 assert!(
2149 gear.deadline_fired.load(Ordering::SeqCst),
2150 "deadline should fire for slow gears"
2151 );
2152 }
2153}