Skip to main content

toolkit/runtime/
runner.rs

1//! `ToolKit` runtime runner.
2//!
3//! Supported DB modes:
4//!   - `DbOptions::None` — gears get no DB in their contexts.
5//!   - `DbOptions::Manager` — gears use `GearContextBuilder` to resolve per-gear `DbHandles`.
6//!
7//! Design notes:
8//! - We use **`GearContextBuilder`** to resolve per-gear `DbHandles` at runtime.
9//! - Phase order is orchestrated by `HostRuntime` (see `runtime/host_runtime.rs` docs).
10//! - Gears receive a fully-scoped `GearCtx` with a resolved Option<DbHandle>.
11//! - Shutdown can be driven by OS signals, an external `CancellationToken`,
12//!   or an arbitrary future.
13//! - Pre-registered clients can be injected into the `ClientHub` via `RunOptions::clients`.
14//! - `OoP` gears are spawned after the start phase so that `grpc-hub` is already running
15//!   and the real directory endpoint is known.
16
17use crate::backends::OopBackend;
18use crate::client_hub::ClientHub;
19use crate::config::ConfigProvider;
20use crate::registry::GearRegistry;
21use crate::runtime::shutdown;
22use crate::runtime::{DbOptions, HostRuntime};
23use std::collections::HashMap;
24use std::path::PathBuf;
25use std::{future::Future, pin::Pin, sync::Arc};
26use tokio_util::sync::CancellationToken;
27use uuid::Uuid;
28
29/// A type-erased client registration for injecting clients into the `ClientHub`.
30///
31/// This is used to pass pre-created clients (like gRPC clients) from bootstrap code
32/// into the runtime's `ClientHub` before gears are initialized.
33pub struct ClientRegistration {
34    /// Callback that registers the client into the hub.
35    register_fn: Box<dyn FnOnce(&ClientHub) + Send>,
36}
37
38impl ClientRegistration {
39    /// Create a new client registration for a trait object type.
40    ///
41    /// # Example
42    /// ```ignore
43    /// let api: Arc<dyn DirectoryClient> = Arc::new(client);
44    /// ClientRegistration::new::<dyn DirectoryClient>(api)
45    /// ```
46    pub fn new<T>(client: Arc<T>) -> Self
47    where
48        T: ?Sized + Send + Sync + 'static,
49    {
50        Self {
51            register_fn: Box::new(move |hub| {
52                hub.register::<T>(client);
53            }),
54        }
55    }
56
57    /// Execute the registration against the given hub.
58    pub(crate) fn apply(self, hub: &ClientHub) {
59        (self.register_fn)(hub);
60    }
61}
62
63/// How the runtime should decide when to stop.
64pub enum ShutdownOptions {
65    /// Listen for OS signals (Ctrl+C / SIGTERM).
66    Signals,
67    /// An external `CancellationToken` controls the lifecycle.
68    Token(CancellationToken),
69    /// An arbitrary future; when it completes, we initiate shutdown.
70    Future(Pin<Box<dyn Future<Output = ()> + Send>>),
71}
72
73/// Configuration for a single `OoP` gear to be spawned.
74#[derive(Clone)]
75pub struct OopGearSpawnConfig {
76    /// Name of the gear (e.g., "calculator").
77    pub gear_name: String,
78    /// Path to the gear executable.
79    pub binary: PathBuf,
80    /// Command-line arguments passed to the gear binary.
81    ///
82    /// Note: the user controls `--config` via `execution.args` in the master config.
83    pub args: Vec<String>,
84    /// Environment variables to set for the spawned process.
85    pub env: HashMap<String, String>,
86    /// Working directory for the spawned process.
87    pub working_directory: Option<String>,
88    /// Rendered gear configuration JSON, injected as the `TOOLKIT_MODULE_CONFIG` environment variable.
89    pub rendered_config_json: String,
90}
91
92/// Options for spawning `OoP` gears.
93pub struct OopSpawnOptions {
94    /// Gears to spawn after the start phase, once shared infrastructure is ready.
95    pub gears: Vec<OopGearSpawnConfig>,
96    /// Backend used to spawn gear processes (e.g. `LocalProcessBackend`).
97    pub backend: Box<dyn OopBackend>,
98}
99
100/// Options for running the `ToolKit` runner.
101pub struct RunOptions {
102    /// Provider of gear config sections (raw JSON by gear name).
103    pub gears_cfg: Arc<dyn ConfigProvider>,
104    /// DB strategy: none, or `DbManager`.
105    pub db: DbOptions,
106    /// Shutdown strategy.
107    pub shutdown: ShutdownOptions,
108    /// Pre-registered clients to inject into the `ClientHub` before gear initialization.
109    ///
110    /// This is useful for `OoP` bootstrap where clients (like `DirectoryGrpcClient`)
111    /// are created before calling `run()` and need to be available in the `ClientHub`.
112    pub clients: Vec<ClientRegistration>,
113    /// Process-level instance ID.
114    ///
115    /// This is a unique identifier for this process instance, generated once at bootstrap
116    /// (either in `run_oop_with_options` for `OoP` gears or in the main host).
117    /// It is propagated to all gears via `GearCtx::instance_id()` and `SystemContext::instance_id()`.
118    pub instance_id: Uuid,
119    /// `OoP` gear spawn configuration.
120    ///
121    /// These gears are spawned after the start phase, once `grpc-hub` is running
122    /// and the real directory endpoint is known.
123    pub oop: Option<OopSpawnOptions>,
124    /// Maximum time allowed for each gear's graceful shutdown before hard-stop.
125    ///
126    /// If `None`, uses `DEFAULT_SHUTDOWN_DEADLINE` (30 seconds).
127    ///
128    /// See `HostRuntime::with_shutdown_deadline` for details on the relationship
129    /// with `WithLifecycle::stop_timeout`.
130    pub shutdown_deadline: Option<std::time::Duration>,
131}
132
133/// Construct a `HostRuntime` by wiring shutdown, discovering gears, building
134/// the `ClientHub`, and injecting pre-registered clients.
135///
136/// Shared by [`run`] and [`run_oop_serving`]. `opts` is consumed; the returned
137/// `HostRuntime` owns the root lifecycle `CancellationToken`.
138fn build_host_runtime(opts: RunOptions) -> anyhow::Result<HostRuntime> {
139    // 1. Prepare cancellation token based on shutdown options.
140    let cancel = match &opts.shutdown {
141        ShutdownOptions::Token(t) => t.clone(),
142        _ => CancellationToken::new(),
143    };
144
145    // 2. Spawn shutdown waiter (Signals / Future).
146    match opts.shutdown {
147        ShutdownOptions::Signals => {
148            let c = cancel.clone();
149            tokio::spawn(async move {
150                match shutdown::wait_for_shutdown().await {
151                    Ok(()) => {
152                        tracing::info!(target: "", "------------------");
153                        tracing::info!("shutdown: signal received");
154                    }
155                    Err(e) => {
156                        tracing::warn!(
157                            error = %e,
158                            "shutdown: primary waiter failed; falling back to ctrl_c()"
159                        );
160                        _ = tokio::signal::ctrl_c().await;
161                    }
162                }
163                c.cancel();
164            });
165        }
166        ShutdownOptions::Future(waiter) => {
167            let c = cancel.clone();
168            tokio::spawn(async move {
169                waiter.await;
170                tracing::info!("shutdown: external future completed");
171                c.cancel();
172            });
173        }
174        ShutdownOptions::Token(_) => {
175            tracing::info!("shutdown: external token will control lifecycle");
176        }
177    }
178
179    // 3. Discover gears.
180    let registry = GearRegistry::discover_and_build()?;
181
182    // 4. Build shared `ClientHub` and inject pre-registered clients.
183    let hub = Arc::new(ClientHub::default());
184    for registration in opts.clients {
185        registration.apply(&hub);
186    }
187
188    // 5. Instantiate `HostRuntime`.
189    let mut host = HostRuntime::new(
190        registry,
191        opts.gears_cfg.clone(),
192        opts.db,
193        hub,
194        cancel,
195        opts.instance_id,
196        opts.oop,
197    );
198    if let Some(deadline) = opts.shutdown_deadline {
199        host = host.with_shutdown_deadline(deadline);
200    }
201
202    Ok(host)
203}
204
205/// Run the full gear lifecycle.
206///
207/// Discovers gears, wires shutdown, and delegates phase execution to [`HostRuntime`].
208///
209/// # Errors
210/// Returns an error if any lifecycle phase fails.
211pub async fn run(opts: RunOptions) -> anyhow::Result<()> {
212    let host = build_host_runtime(opts)?;
213    host.run_gear_phases().await
214}
215
216/// Run an out-of-process gear and serve its HTTP routes.
217///
218/// Uses the same setup as [`run`], then drives the `OoP` serving lifecycle:
219/// lifecycle phases, an Axum HTTP server with framework probes, background
220/// self-registration, dependency resolution, and graceful drain.
221///
222/// # Errors
223/// Returns an error if discovery, any lifecycle phase, or the HTTP server fails.
224#[cfg(feature = "bootstrap")]
225pub async fn run_oop_serving(
226    opts: RunOptions,
227    serve: crate::runtime::OopServeOptions,
228) -> anyhow::Result<()> {
229    let host = build_host_runtime(opts)?;
230    host.run_oop_serving(serve).await
231}