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.
101///
102/// `#[non_exhaustive]`: construct via [`RunOptions::new`] + the `with_*`
103/// builders rather than a struct literal, so adding a new optional field does
104/// not break every call site (and out-of-crate consumers get a deprecation-free
105/// upgrade path). The four arguments to `new` are the always-required fields;
106/// everything else defaults to "off".
107#[non_exhaustive]
108pub struct RunOptions {
109 /// Provider of gear config sections (raw JSON by gear name).
110 pub gears_cfg: Arc<dyn ConfigProvider>,
111 /// DB strategy: none, or `DbManager`.
112 pub db: DbOptions,
113 /// Shutdown strategy.
114 pub shutdown: ShutdownOptions,
115 /// Pre-registered clients to inject into the `ClientHub` before gear initialization.
116 ///
117 /// This is useful for `OoP` bootstrap where clients (like `DirectoryGrpcClient`)
118 /// are created before calling `run()` and need to be available in the `ClientHub`.
119 pub clients: Vec<ClientRegistration>,
120 /// Process-level instance ID.
121 ///
122 /// This is a unique identifier for this process instance, generated once at bootstrap
123 /// (either in `run_oop_with_options` for `OoP` gears or in the main host).
124 /// It is propagated to all gears via `GearCtx::instance_id()` and `SystemContext::instance_id()`.
125 pub instance_id: Uuid,
126 /// `OoP` gear spawn configuration.
127 ///
128 /// These gears are spawned after the start phase, once `grpc-hub` is running
129 /// and the real directory endpoint is known.
130 pub oop: Option<OopSpawnOptions>,
131 /// Maximum time allowed for each gear's graceful shutdown before hard-stop.
132 ///
133 /// If `None`, uses [`DEFAULT_SHUTDOWN_DEADLINE`](crate::runtime::DEFAULT_SHUTDOWN_DEADLINE) (35 seconds).
134 ///
135 /// See `HostRuntime::with_shutdown_deadline` for details on the relationship
136 /// with `WithLifecycle::stop_timeout`.
137 pub shutdown_deadline: Option<std::time::Duration>,
138 /// Process-wide platform-plane credential source (the selected
139 /// `InternalCredential`), threaded into every [`GearCtx`] so
140 /// `#[toolkit::provides]`-generated clients attach `X-ToolKit-Internal-Token`
141 /// on platform-plane (`PlatformSecurityContext`) methods. `None` (Profile 1
142 /// / in-process, or no credential configured) attaches nothing.
143 pub internal_token_provider: Option<toolkit_contract::runtime::config::InternalTokenProvider>,
144}
145
146impl RunOptions {
147 /// Create `RunOptions` with the four always-required fields; all optional
148 /// fields default to "off" (`clients` empty, no `oop`, default shutdown
149 /// deadline, no platform credential). Layer options on with the `with_*`
150 /// builders.
151 #[must_use]
152 pub fn new(
153 gears_cfg: Arc<dyn ConfigProvider>,
154 db: DbOptions,
155 shutdown: ShutdownOptions,
156 instance_id: Uuid,
157 ) -> Self {
158 Self {
159 gears_cfg,
160 db,
161 shutdown,
162 clients: Vec::new(),
163 instance_id,
164 oop: None,
165 shutdown_deadline: None,
166 internal_token_provider: None,
167 }
168 }
169
170 /// Pre-register clients to inject into the `ClientHub` before gear init.
171 #[must_use]
172 pub fn with_clients(mut self, clients: Vec<ClientRegistration>) -> Self {
173 self.clients = clients;
174 self
175 }
176
177 /// Set the `OoP` gear spawn configuration.
178 #[must_use]
179 pub fn with_oop(mut self, oop: Option<OopSpawnOptions>) -> Self {
180 self.oop = oop;
181 self
182 }
183
184 /// Override the per-gear graceful-shutdown deadline.
185 #[must_use]
186 pub fn with_shutdown_deadline(mut self, deadline: Option<std::time::Duration>) -> Self {
187 self.shutdown_deadline = deadline;
188 self
189 }
190
191 /// Set the process-wide platform-plane credential source (see the field).
192 #[must_use]
193 pub fn with_internal_token_provider(
194 mut self,
195 provider: Option<toolkit_contract::runtime::config::InternalTokenProvider>,
196 ) -> Self {
197 self.internal_token_provider = provider;
198 self
199 }
200}
201
202/// Construct a `HostRuntime` by wiring shutdown, discovering gears, building
203/// the `ClientHub`, and injecting pre-registered clients.
204///
205/// Shared by [`run`] and [`run_oop_serving`]. `opts` is consumed; the returned
206/// `HostRuntime` owns the root lifecycle `CancellationToken`.
207fn build_host_runtime(opts: RunOptions) -> anyhow::Result<HostRuntime> {
208 // 1. Prepare cancellation token based on shutdown options.
209 let cancel = match &opts.shutdown {
210 ShutdownOptions::Token(t) => t.clone(),
211 _ => CancellationToken::new(),
212 };
213
214 // 2. Spawn shutdown waiter (Signals / Future).
215 match opts.shutdown {
216 ShutdownOptions::Signals => {
217 let c = cancel.clone();
218 tokio::spawn(async move {
219 match shutdown::wait_for_shutdown().await {
220 Ok(()) => {
221 tracing::info!(target: "", "------------------");
222 tracing::info!("shutdown: signal received");
223 }
224 Err(e) => {
225 tracing::warn!(
226 error = %e,
227 "shutdown: primary waiter failed; falling back to ctrl_c()"
228 );
229 _ = tokio::signal::ctrl_c().await;
230 }
231 }
232 c.cancel();
233 });
234 }
235 ShutdownOptions::Future(waiter) => {
236 let c = cancel.clone();
237 tokio::spawn(async move {
238 waiter.await;
239 tracing::info!("shutdown: external future completed");
240 c.cancel();
241 });
242 }
243 ShutdownOptions::Token(_) => {
244 tracing::info!("shutdown: external token will control lifecycle");
245 }
246 }
247
248 // 3. Discover gears.
249 let registry = GearRegistry::discover_and_build()?;
250
251 // 4. Build shared `ClientHub` and inject pre-registered clients.
252 let hub = Arc::new(ClientHub::default());
253 for registration in opts.clients {
254 registration.apply(&hub);
255 }
256
257 // 5. Instantiate `HostRuntime`.
258 let mut host = HostRuntime::new(
259 registry,
260 opts.gears_cfg.clone(),
261 opts.db,
262 hub,
263 cancel,
264 opts.instance_id,
265 opts.oop,
266 );
267 if let Some(deadline) = opts.shutdown_deadline {
268 host = host.with_shutdown_deadline(deadline);
269 }
270 host = host.with_internal_token_provider(opts.internal_token_provider);
271
272 Ok(host)
273}
274
275/// Run the full gear lifecycle.
276///
277/// Discovers gears, wires shutdown, and delegates phase execution to [`HostRuntime`].
278///
279/// # Errors
280/// Returns an error if any lifecycle phase fails.
281pub async fn run(opts: RunOptions) -> anyhow::Result<()> {
282 let host = build_host_runtime(opts)?;
283 host.run_gear_phases().await
284}
285
286/// Run an out-of-process gear and serve its HTTP routes.
287///
288/// Uses the same setup as [`run`], then drives the `OoP` serving lifecycle:
289/// lifecycle phases, an Axum HTTP server with framework probes, background
290/// self-registration, dependency resolution, and graceful drain.
291///
292/// # Errors
293/// Returns an error if discovery, any lifecycle phase, or the HTTP server fails.
294#[cfg(feature = "bootstrap")]
295pub async fn run_oop_serving(
296 opts: RunOptions,
297 serve: crate::runtime::OopServeOptions,
298) -> anyhow::Result<()> {
299 let host = build_host_runtime(opts)?;
300 host.run_oop_serving(serve).await
301}