Skip to main content

toolkit/bootstrap/
run.rs

1use super::config::{get_gear_runtime_config, render_gear_config_for_oop};
2use super::host::{init_logging_unified, init_panic_tracing, normalize_path};
3use super::{AppConfig, RuntimeKind};
4use crate::backends::LocalProcessBackend;
5use crate::runtime::{
6    DbOptions, OopGearSpawnConfig, OopSpawnOptions, RunOptions, ShutdownOptions, run, shutdown,
7};
8use anyhow::Result;
9use figment::Figment;
10use figment::providers::Serialized;
11use std::path::{Path, PathBuf};
12use std::sync::Arc;
13use tokio_util::sync::CancellationToken;
14
15/// Spawn a signal handler task that cancels the provided token on SIGTERM/SIGINT.
16///
17/// This helper consolidates signal handling logic used by both `run_server` and `run_migrate`.
18/// The `context` parameter customizes log messages for better diagnostics.
19fn spawn_signal_handler(cancel: CancellationToken, context: &str) {
20    let context_owned = context.to_owned();
21    tokio::spawn(async move {
22        match shutdown::wait_for_shutdown().await {
23            Ok(()) => {
24                tracing::info!(target: "", "------------------");
25                tracing::info!("{}: shutdown signal received", context_owned);
26            }
27            Err(e) => {
28                tracing::warn!(
29                    error = %e,
30                    "{}: signal handler failed, falling back to ctrl_c()",
31                    context_owned
32                );
33                _ = tokio::signal::ctrl_c().await;
34            }
35        }
36        cancel.cancel();
37    });
38}
39
40/// # Errors
41///
42/// Returns an error if:
43/// - There was a critical error during initialization of the gears
44/// - Problems with the database or third-party services
45/// - An issue during runtime or shutdown
46///
47/// The TLS crypto provider is installed automatically as the first step of
48/// [`init_procedure`] (idempotent), so callers do not need to invoke
49/// [`super::init_crypto_provider`] explicitly.
50pub async fn run_server(config: AppConfig) -> Result<()> {
51    init_procedure(&config).map_err(|e| {
52        tracing::error!(error = %e, "Initialization failed");
53        e
54    })?;
55    tracing::info!("Initializing gears...");
56
57    // Generate process-level instance ID once at startup.
58    // This is shared by all gears in this process.
59    let instance_id = uuid::Uuid::new_v4();
60    tracing::info!(instance_id = %instance_id, "Generated process instance ID");
61
62    // Create root cancellation token for the entire process.
63    // This token drives shutdown for the gear runtime and all lifecycle/stateful gears.
64    let cancel = CancellationToken::new();
65
66    // Hook OS signals to the root token at the host level.
67    // This replaces the use of ShutdownOptions::Signals inside the runtime.
68    spawn_signal_handler(cancel.clone(), "server");
69
70    // Build config provider and resolve database options
71    let db_options = resolve_db_options(&config)?;
72
73    // Create OoP backend with cancellation token - it will auto-shutdown all processes on cancel
74    let oop_backend = LocalProcessBackend::new(cancel.clone());
75
76    // Build OoP spawn configuration
77    let oop_options = build_oop_spawn_options(&config, oop_backend)?;
78
79    // Run the ToolKit runtime with the root cancellation token.
80    // Shutdown is driven by the signal handler spawned above, not by ShutdownOptions::Signals.
81    // OoP gears are spawned after the start phase (once grpc-hub has bound its port).
82    //
83    // No `internal_token_provider` is set here on purpose: this is the
84    // in-process (Profile 1) host. Gear-to-gear calls resolve to LOCAL trait
85    // objects through the `ClientHub` — there is no transport, so no
86    // `X-ToolKit-Internal-Token` header/metadata is ever emitted and no outbound
87    // platform credential is needed. The provider is only consulted inside the
88    // generated REST/gRPC transport clients, which the monolith does not use for
89    // co-hosted dependencies. A genuinely remote (directory-resolved) dependency
90    // is the resolving-client path, which threads its own provider
91    // (`cpt-cf-adr-two-plane-auth`).
92    let run_options = RunOptions::new(
93        Arc::new(config),
94        db_options,
95        ShutdownOptions::Token(cancel.clone()),
96        instance_id,
97    )
98    .with_oop(oop_options);
99
100    let result = run(run_options).await;
101
102    // Graceful shutdown - flush remaining telemetry
103    #[cfg(feature = "otel")]
104    tracing_shutdown().await;
105
106    result
107}
108
109/// Run database migrations and exit.
110///
111/// This mode is designed for cloud deployment workflows where database
112/// migrations need to run as a separate step before starting the application.
113///
114/// Phases executed:
115/// - Pre-init (wire runtime internals)
116/// - DB migration (run all pending migrations)
117///
118/// The process exits after migrations complete. Any errors are reported
119/// and propagated as non-zero exit codes.
120///
121/// # Errors
122///
123/// Returns an error if:
124/// - No database configuration is found
125/// - Gear discovery fails
126/// - Pre-init phase fails
127/// - Migration phase fails
128///
129/// The TLS crypto provider is installed automatically as the first step of
130/// [`init_procedure`] (idempotent), so callers do not need to invoke
131/// [`super::init_crypto_provider`] explicitly.
132pub async fn run_migrate(config: AppConfig) -> Result<()> {
133    init_procedure(&config).map_err(|e| {
134        tracing::error!(error = %e, "Initialization failed");
135        e
136    })?;
137    tracing::info!("Starting migration mode...");
138
139    // Generate process-level instance ID for this migration run
140    let instance_id = uuid::Uuid::new_v4();
141    tracing::info!(instance_id = %instance_id, "Generated migration instance ID");
142
143    // Create cancellation token and wire it to OS signals
144    let cancel = CancellationToken::new();
145
146    // Hook OS signals to enable graceful cancellation of migrations
147    spawn_signal_handler(cancel.clone(), "migration");
148
149    // Build database options from configuration
150    let db_options = resolve_db_options(&config)?;
151
152    // Verify we have database configuration
153    if matches!(db_options, DbOptions::None) {
154        anyhow::bail!("Cannot run migrations: no database configuration found");
155    }
156
157    // Discover and build the gear registry
158    let registry = crate::registry::GearRegistry::discover_and_build()?;
159    tracing::info!(
160        gear_count = registry.gears().len(),
161        "Discovered gears for migration"
162    );
163
164    // Create the host runtime
165    let host = crate::runtime::HostRuntime::new(
166        registry,
167        Arc::new(config),
168        db_options,
169        Arc::new(crate::client_hub::ClientHub::new()),
170        cancel,
171        instance_id,
172        None, // No OoP spawning during migration
173    );
174
175    // Run only the migration phases (pre-init + DB migration)
176    let result = host.run_migration_phases().await;
177
178    // Graceful shutdown - flush remaining telemetry
179    #[cfg(feature = "otel")]
180    tracing_shutdown().await;
181
182    result?;
183
184    tracing::info!("All migrations completed successfully");
185    Ok(())
186}
187
188fn resolve_db_options(config: &AppConfig) -> Result<DbOptions> {
189    if config.database.is_none() {
190        tracing::warn!("No global database section found; running without databases");
191        return Ok(DbOptions::None);
192    }
193
194    tracing::info!("Using DbManager with Figment-based configuration");
195    let figment = Figment::new().merge(Serialized::defaults(config));
196    let db_manager = Arc::new(toolkit_db::DbManager::from_figment(
197        figment,
198        config.server.home_dir.clone(),
199    )?);
200    Ok(DbOptions::Manager(db_manager))
201}
202
203/// Build `OoP` spawn configuration from `AppConfig`.
204///
205/// This collects all gears with `type=oop` and prepares their spawn configuration.
206/// The actual spawning happens in the `HostRuntime` after the start phase.
207fn build_oop_spawn_options(
208    config: &AppConfig,
209    backend: LocalProcessBackend,
210) -> Result<Option<OopSpawnOptions>> {
211    let home_dir = PathBuf::from(&config.server.home_dir);
212    let mut gears = Vec::new();
213
214    for gear_name in config.gears.keys() {
215        if let Some(spawn_config) = try_build_oop_gear_config(config, gear_name, &home_dir)? {
216            gears.push(spawn_config);
217        }
218    }
219
220    if gears.is_empty() {
221        Ok(None)
222    } else {
223        tracing::info!(count = gears.len(), "Prepared OoP gears for spawning");
224        Ok(Some(OopSpawnOptions {
225            gears,
226            backend: Box::new(backend),
227        }))
228    }
229}
230
231/// Try to build `OoP` gear spawn config if gear is of type `OoP`
232fn try_build_oop_gear_config(
233    config: &AppConfig,
234    gear_name: &str,
235    home_dir: &Path,
236) -> Result<Option<OopGearSpawnConfig>> {
237    let Some(runtime_cfg) = get_gear_runtime_config(config, gear_name)? else {
238        return Ok(None);
239    };
240
241    if !matches!(runtime_cfg.mod_type, RuntimeKind::Oop) {
242        return Ok(None);
243    }
244
245    let exec_cfg = runtime_cfg.execution.as_ref().ok_or_else(|| {
246        anyhow::anyhow!("gear '{gear_name}' is type=oop but execution config is missing")
247    })?;
248
249    let binary = normalize_path(&exec_cfg.executable_path)?;
250    let spawn_args = exec_cfg.args.clone();
251    let env = exec_cfg.environment.clone();
252
253    // Render the complete gear config (with resolved DB)
254    let rendered_config = render_gear_config_for_oop(config, gear_name, home_dir)?;
255    let rendered_json = rendered_config.to_json()?;
256
257    tracing::debug!(
258        gear =  %gear_name,
259        "Prepared OoP gear config: db={}",
260        rendered_config.database.is_some()
261    );
262
263    Ok(Some(OopGearSpawnConfig {
264        gear_name: gear_name.to_owned(),
265        binary,
266        args: spawn_args,
267        env,
268        working_directory: exec_cfg.working_directory.clone(),
269        rendered_config_json: rendered_json,
270    }))
271}
272
273/// Initialize process-wide bootstrap state from a provided `&AppConfig`.
274///
275/// This helper performs the common startup sequence shared by server and migration modes.
276/// It does **not** load configuration; the caller is responsible for building and passing
277/// a valid `AppConfig`.
278///
279/// Steps performed:
280///
281/// - initializes tracing/logging (once, guarded by a process-wide `Once`) and metrics
282///   when OpenTelemetry is enabled
283/// - registers the panic hook used to route panics through tracing
284/// - emits a small startup span and version metadata for diagnostics
285///
286/// The rustls crypto provider is installed inside this function (idempotent
287/// `OnceLock`), so callers do not need to invoke
288/// [`super::init_crypto_provider`] separately. Direct callers of
289/// [`super::init_crypto_provider`] outside the bootstrap path (e.g. ad-hoc
290/// probe binaries that do not call `run_server` / `run_migrate`) are still
291/// supported.
292///
293/// # Errors
294///
295/// Returns an error if OpenTelemetry tracing initialization fails while
296/// tracing is enabled, or if the crypto provider installation fails.
297#[cfg_attr(not(feature = "otel"), allow(clippy::unnecessary_wraps))]
298pub fn init_procedure(config: &AppConfig) -> Result<()> {
299    // Install the rustls crypto provider FIRST — before anything that may
300    // touch TLS (OTLP exporter over HTTPS, DB connections, etc.). The call
301    // is process-wide and idempotent; calling it twice is safe.
302    super::init_crypto_provider()?;
303
304    // Build OpenTelemetry layer before logging
305    #[cfg(feature = "otel")]
306    let otel_layer = if config.opentelemetry.tracing.enabled {
307        Some(crate::telemetry::init::init_tracing(&config.opentelemetry)?)
308    } else {
309        None
310    };
311    #[cfg(not(feature = "otel"))]
312    let otel_layer = None;
313
314    // Initialize logging + otel in one Registry
315    init_logging_unified(
316        &config.logging,
317        &config.server.home_dir,
318        otel_layer,
319        config.opentelemetry.inject_trace_ids_into_logs(),
320    );
321
322    // Register custom panic hook to reroute panic backtrace into tracing.
323    init_panic_tracing();
324
325    // Initialize OpenTelemetry metrics (or confirm noop when disabled)
326    #[cfg(feature = "otel")]
327    if let Err(e) = crate::telemetry::init::init_metrics_provider(&config.opentelemetry) {
328        tracing::error!(error = %e, "OpenTelemetry metrics not initialized");
329    }
330
331    // One-time connectivity probe
332    #[cfg(feature = "otel")]
333    if config.opentelemetry.tracing.enabled
334        && let Err(e) = crate::telemetry::init::otel_connectivity_probe(&config.opentelemetry)
335    {
336        tracing::error!(error = %e, "OTLP connectivity probe failed");
337    }
338
339    // Smoke test span to confirm traces flow to Jaeger
340    tracing::info_span!("startup_check", app = config.server.name).in_scope(|| {
341        tracing::info!("startup span alive - traces should be visible in Jaeger");
342    });
343
344    tracing::info!(
345        version = env!("CARGO_PKG_VERSION"),
346        rust_version = env!("CARGO_PKG_RUST_VERSION"),
347        "{} Server starting",
348        config.server.name,
349    );
350
351    Ok(())
352}
353
354#[cfg(feature = "otel")]
355/// Flush compatibility shutdown hooks for OpenTelemetry tracing and metrics.
356///
357/// This delegates to the current telemetry shutdown helpers so callers can use a
358/// single bootstrap-level function during graceful shutdown.
359///
360/// `shutdown_metrics`/`shutdown_tracing` drain the SDK's batch processor, which
361/// blocks on exporting the final batch (real network I/O to the OTLP endpoint).
362/// Run that on a blocking-pool thread so it never stalls a Tokio worker thread.
363pub async fn tracing_shutdown() {
364    if let Err(e) = tokio::task::spawn_blocking(|| {
365        crate::telemetry::init::shutdown_metrics();
366        crate::telemetry::init::shutdown_tracing();
367    })
368    .await
369    {
370        tracing::warn!(error = %e, "Telemetry shutdown task panicked");
371    }
372}