Skip to main content

toolkit/bootstrap/
oop.rs

1//! Out-of-process gear bootstrap library
2//!
3//! This gear provides reusable functionality for bootstrapping `OoP` (out-of-process)
4//! `ToolKit` gears in local (non-k8s) environments.
5//!
6//! ## Features
7//!
8//! - Configuration loading using `toolkit-bootstrap`
9//! - Logging initialization with tracing
10//! - gRPC connection to `DirectoryService`
11//! - Gear instance registration
12//! - Heartbeat management
13//! - Gear lifecycle execution
14//!
15//! ## Shutdown Model
16//!
17//! Shutdown is driven by a single root `CancellationToken` per process:
18//! - OS signals (SIGTERM, SIGINT, Ctrl+C) are hooked at bootstrap level
19//! - The root token is passed to `RunOptions::Token` for gear runtime shutdown
20//! - Background tasks (like heartbeat) use child tokens derived from the root
21//!
22//! On shutdown, the gear deregisters itself from the `DirectoryService` before exiting.
23//!
24//! ## Example
25//!
26//! ```rust,no_run
27//! use toolkit::bootstrap::oop::{OopRunOptions, run_oop_with_options};
28//!
29//! #[tokio::main]
30//! async fn main() -> anyhow::Result<()> {
31//!     let opts = OopRunOptions {
32//!         gear_name: "my_gear".to_string(),
33//!         instance_id: None,
34//!         directory_endpoint: "http://127.0.0.1:50051".to_string(),
35//!         config_path: None,
36//!         verbose: 0,
37//!         print_config: false,
38//!         heartbeat_interval_secs: 5,
39//!         version: None,
40//!     };
41//!
42//!     run_oop_with_options(opts).await
43//! }
44//! ```
45
46use anyhow::{Context, Result};
47use figment::{Figment, providers::Serialized};
48use std::path::{Path, PathBuf};
49use std::sync::Arc;
50use std::time::Duration;
51use tokio::time::sleep;
52use tokio_util::sync::CancellationToken;
53use tracing::{debug, error, info, warn};
54use uuid::Uuid;
55
56use super::config::{
57    AppConfig, CliArgs, LoggingConfig, RenderedDbConfig, RenderedGearConfig,
58    TOOLKIT_MODULE_CONFIG_ENV,
59};
60use crate::bootstrap::host::{init_logging_unified, init_panic_tracing};
61use crate::runtime::{
62    ClientRegistration, DbOptions, OopServeOptions, RunOptions, ShutdownOptions,
63    TOOLKIT_DIRECTORY_ENDPOINT_ENV, run, run_oop_serving, shutdown,
64};
65use cf_system_sdks::directory::{DirectoryClient, DirectoryGrpcClient};
66
67/// Configuration options for `OoP` gear bootstrap
68#[derive(Debug, Clone)]
69pub struct OopRunOptions {
70    /// Logical gear name (e.g., "`file-parser`")
71    pub gear_name: String,
72
73    /// Instance ID (defaults to a random UUID if None)
74    pub instance_id: Option<Uuid>,
75
76    /// Directory service gRPC endpoint (e.g., "<http://127.0.0.1:50051>")
77    pub directory_endpoint: String,
78
79    /// Path to configuration file
80    pub config_path: Option<PathBuf>,
81
82    /// Log verbosity level (0=default, 1=debug, 2=trace)
83    pub verbose: u8,
84
85    /// Print effective configuration and exit
86    pub print_config: bool,
87
88    /// Heartbeat interval in seconds (default: 5)
89    pub heartbeat_interval_secs: u64,
90
91    /// Gear version (used for `DirectoryService` registration and `OpenAPI` version).
92    /// Defaults to `None`.
93    pub version: Option<String>,
94}
95
96impl Default for OopRunOptions {
97    fn default() -> Self {
98        // Check for config path in environment variable as fallback
99        let config_path = std::env::var("TOOLKIT_CONFIG_PATH").ok().map(PathBuf::from);
100
101        // Check for directory endpoint in environment variable (set by master host)
102        // This is the preferred way to get the endpoint when spawned by master host
103        let directory_endpoint = std::env::var(TOOLKIT_DIRECTORY_ENDPOINT_ENV)
104            .unwrap_or_else(|_| "http://127.0.0.1:50051".to_owned());
105
106        Self {
107            gear_name: String::new(),
108            instance_id: None,
109            directory_endpoint,
110            config_path,
111            verbose: 0,
112            print_config: false,
113            heartbeat_interval_secs: 5,
114            version: None,
115        }
116    }
117}
118
119/// Builds the final configuration and `DbOptions` for an `OoP` gear.
120///
121/// Configuration merge strategy (for each section):
122/// - **Database**: field-by-field merge using `DbManager` (master as base, local as override)
123/// - **Logging**: key-by-key merge (each subsystem key is overridden by local)
124/// - **Config**: local completely replaces master if present
125///
126/// The local config file (--config) can override any settings from master's `TOOLKIT_MODULE_CONFIG`.
127///
128/// For database, the merge happens at 3 levels:
129/// 1. Global database.servers.* from master
130/// 2. Gear's database section from master (gears.<name>.database)
131/// 3. Gear's database section from local --config (overrides master)
132#[tracing::instrument(
133    level = "debug",
134    skip(local_config, rendered_config),
135    fields(
136        has_rendered = rendered_config.is_some(),
137        has_local_db = local_config.database.is_some()
138    )
139)]
140fn build_oop_config_and_db(
141    local_config: &AppConfig,
142    gear_name: &str,
143    rendered_config: Option<&RenderedGearConfig>,
144) -> Result<(AppConfig, LoggingConfig, DbOptions)> {
145    let home_dir = PathBuf::from(&local_config.server.home_dir);
146
147    // Build final_config for gear's "config" section
148    let final_config = if let Some(rendered) = rendered_config {
149        // TOOLKIT_MODULE_CONFIG exists: use rendered config as BASE, local config as OVERRIDE
150        let mut config = local_config.clone();
151
152        // Get or create the gear entry
153        let gear_entry = config
154            .gears
155            .entry(gear_name.to_owned())
156            .or_insert_with(|| serde_json::json!({}));
157
158        // Merge rendered.config as base, local gear config as override
159        if let Some(obj) = gear_entry.as_object_mut() {
160            // If local doesn't have "config" section, use rendered entirely
161            // If local has "config" section, it takes precedence (local overrides master)
162            if !obj.contains_key("config") || obj["config"].is_null() {
163                obj.insert("config".to_owned(), rendered.config.clone());
164            }
165            // If local has "config", it already overrides - no action needed
166        }
167
168        debug!(
169            gear =  %gear_name,
170            has_rendered_db = %rendered.database.is_some(),
171            has_rendered_logging = %rendered.logging.is_some(),
172            "Using rendered config from master as base, local config as override"
173        );
174
175        config
176    } else {
177        // No TOOLKIT_MODULE_CONFIG: use local config entirely (standalone mode)
178        debug!(
179            gear =  %gear_name,
180            "No rendered config from master, using local config entirely (standalone mode)"
181        );
182        local_config.clone()
183    };
184
185    // Merge logging: master logging (base) + local logging (override by key)
186    let final_logging = merge_logging_configs(
187        rendered_config.as_ref().and_then(|r| r.logging.as_ref()),
188        &local_config.logging,
189    );
190
191    // Build DbOptions using Figment merge + DbManager
192    // This allows field-by-field merge: master db config (base) -> local db config (override)
193    let db_options = build_merged_db_options(
194        &home_dir,
195        gear_name,
196        rendered_config.as_ref().and_then(|r| r.database.as_ref()),
197        local_config,
198    )?;
199
200    Ok((final_config, final_logging, db_options))
201}
202
203/// Merges logging configurations: master as base, local as override (by key).
204///
205/// Each key in the logging `HashMap` (e.g., "default", "calculator", "sqlx")
206/// is overridden by local if present.
207fn merge_logging_configs(master: Option<&LoggingConfig>, local: &LoggingConfig) -> LoggingConfig {
208    master
209        .cloned()
210        .unwrap_or_default()
211        .into_iter()
212        .chain(local.clone())
213        .collect()
214}
215
216/// Builds `DbOptions` by merging rendered config from master with local config.
217///
218/// Uses Figment to merge configurations and `DbManager` to handle the actual
219/// database connection setup with field-by-field merge logic.
220fn build_merged_db_options(
221    home_dir: &Path,
222    gear_name: &str,
223    rendered_db: Option<&RenderedDbConfig>,
224    local_config: &AppConfig,
225) -> Result<DbOptions> {
226    // Check if we have any database configuration
227    let has_rendered_db = rendered_db.is_some_and(|db| db.gear.is_some() || db.global.is_some());
228    let has_local_db = local_config.database.is_some()
229        || local_config
230            .gears
231            .get(gear_name)
232            .and_then(|m| m.get("database"))
233            .is_some();
234
235    if !has_rendered_db && !has_local_db {
236        debug!(
237            gear =  %gear_name,
238            "No database config available"
239        );
240        return Ok(DbOptions::None);
241    }
242
243    // Build a merged configuration for DbManager:
244    // 1. Start with rendered config from master (global servers + gear db)
245    // 2. Overlay local config (local can override any field)
246
247    let mut merged_config = serde_json::Map::new();
248
249    // Step 1: Add rendered database config from master as base
250    if let Some(rendered) = rendered_db {
251        // Add global servers from master
252        if let Some(ref global) = rendered.global {
253            let global_json = serde_json::to_value(global)
254                .context("Failed to serialize rendered global db config")?;
255            merged_config.insert("database".to_owned(), global_json);
256        }
257
258        // Add gear's database config from master
259        if let Some(ref gear_db) = rendered.gear {
260            let gear_db_json = serde_json::to_value(gear_db)
261                .context("Failed to serialize rendered gear db config")?;
262
263            let mut gears = serde_json::Map::new();
264            let mut gear_entry = serde_json::Map::new();
265            gear_entry.insert("database".to_owned(), gear_db_json);
266            gears.insert(gear_name.to_owned(), serde_json::Value::Object(gear_entry));
267            merged_config.insert("gears".to_owned(), serde_json::Value::Object(gears));
268        }
269    }
270
271    // Step 2: Overlay local config (local overrides master)
272    // Local global database config
273    if let Some(ref local_db) = local_config.database {
274        let local_db_json =
275            serde_json::to_value(local_db).context("Failed to serialize local global db config")?;
276
277        // Merge with existing or replace
278        if let Some(existing) = merged_config.get_mut("database") {
279            merge_json_objects(existing, &local_db_json);
280        } else {
281            merged_config.insert("database".to_owned(), local_db_json);
282        }
283    }
284
285    // Local gear database config
286    if let Some(local_gear) = local_config.gears.get(gear_name)
287        && let Some(local_gear_db) = local_gear.get("database")
288    {
289        let gears = merged_config
290            .entry("gears".to_owned())
291            .or_insert_with(|| serde_json::Value::Object(serde_json::Map::new()));
292
293        if let Some(gears_obj) = gears.as_object_mut() {
294            let gear_entry = gears_obj
295                .entry(gear_name.to_owned())
296                .or_insert_with(|| serde_json::Value::Object(serde_json::Map::new()));
297
298            if let Some(gear_obj) = gear_entry.as_object_mut() {
299                if let Some(existing_db) = gear_obj.get_mut("database") {
300                    merge_json_objects(existing_db, local_gear_db);
301                } else {
302                    gear_obj.insert("database".to_owned(), local_gear_db.clone());
303                }
304            }
305        }
306    }
307
308    debug!(
309        gear =  %gear_name,
310        has_rendered = %rendered_db.is_some(),
311        has_local_global = %local_config.database.is_some(),
312        "Building DbManager with merged config"
313    );
314
315    // Create DbManager from merged Figment
316    let figment = Figment::new().merge(Serialized::defaults(serde_json::Value::Object(
317        merged_config,
318    )));
319    let db_manager = Arc::new(
320        toolkit_db::DbManager::from_figment(figment, home_dir.to_path_buf())
321            .context("Failed to create DbManager from merged config")?,
322    );
323
324    Ok(DbOptions::Manager(db_manager))
325}
326
327/// Recursively merges source JSON object into target.
328/// Source values override target values for matching keys.
329fn merge_json_objects(target: &mut serde_json::Value, source: &serde_json::Value) {
330    if let (Some(target_obj), Some(source_obj)) = (target.as_object_mut(), source.as_object()) {
331        for (key, value) in source_obj {
332            if let Some(target_value) = target_obj.get_mut(key) {
333                // Recursively merge objects, otherwise replace
334                if target_value.is_object() && value.is_object() {
335                    merge_json_objects(target_value, value);
336                } else {
337                    *target_value = value.clone();
338                }
339            } else {
340                target_obj.insert(key.clone(), value.clone());
341            }
342        }
343    } else {
344        // If target is not an object, replace entirely
345        *target = source.clone();
346    }
347}
348
349/// Run an out-of-process gear with the given options
350///
351/// This function:
352/// 1. Creates a root `CancellationToken` for the process
353/// 2. Hooks OS signals (SIGTERM, SIGINT, Ctrl+C) to trigger cancellation
354/// 3. Loads configuration and initializes logging
355/// 4. Connects to the `DirectoryService`
356/// 5. Registers the gear instance
357/// 6. Starts a background heartbeat loop (using a child token)
358/// 7. Runs the gear lifecycle with `ShutdownOptions::Token`
359/// 8. Deregisters from `DirectoryService` on shutdown
360///
361/// ## Shutdown Model
362///
363/// A single root cancellation token drives shutdown for the entire process.
364/// OS signals are hooked at this bootstrap level (not via `ShutdownOptions::Signals`).
365/// The heartbeat loop and gear runtime both observe this token tree.
366///
367/// # Arguments
368///
369/// * `opts` - Bootstrap configuration options
370///
371/// # Returns
372///
373/// * `Ok(())` - If the gear lifecycle completed successfully
374/// * `Err(e)` - If any step failed
375///
376/// # Example
377///
378/// ```rust,no_run
379/// use toolkit::bootstrap::oop::{OopRunOptions, run_oop_with_options};
380///
381/// #[tokio::main]
382/// async fn main() -> anyhow::Result<()> {
383///     let opts = OopRunOptions {
384///         gear_name: "file-parser".to_string(),
385///         instance_id: None,
386///         directory_endpoint: "http://127.0.0.1:50051".to_string(),
387///         config_path: None,
388///         verbose: 1,
389///         print_config: false,
390///         heartbeat_interval_secs: 5,
391///         version: None,
392///     };
393///
394///     run_oop_with_options(opts).await
395/// }
396/// ```
397///
398/// # Errors
399/// Returns an error if the `OoP` gear fails to start or run.
400#[tracing::instrument(
401    level = "info",
402    name = "oop_bootstrap",
403    skip(opts),
404    fields(
405        gear =  %opts.gear_name,
406        directory = %opts.directory_endpoint
407    )
408)]
409pub async fn run_oop_with_options(opts: OopRunOptions) -> Result<()> {
410    // Generate instance ID if not provided
411    let instance_id = opts.instance_id.unwrap_or_else(Uuid::new_v4);
412
413    // Create root cancellation token for the entire process.
414    // This token drives shutdown for the gear runtime and all background tasks.
415    let cancel = CancellationToken::new();
416
417    // Hook OS signals to the root token at bootstrap level.
418    // This replaces the use of ShutdownOptions::Signals inside the runtime.
419    let cancel_for_signals = cancel.clone();
420    tokio::spawn(async move {
421        match shutdown::wait_for_shutdown().await {
422            Ok(()) => {
423                info!(target: "", "------------------");
424                info!("shutdown: signal received in OoP bootstrap");
425            }
426            Err(e) => {
427                warn!(
428                    error = %e,
429                    "shutdown: primary waiter failed in OoP bootstrap, falling back to ctrl_c()"
430                );
431                _ = tokio::signal::ctrl_c().await;
432            }
433        }
434        cancel_for_signals.cancel();
435    });
436
437    // Prepare CLI args for AppConfig loading
438    let args = CliArgs {
439        config: opts
440            .config_path
441            .as_ref()
442            .map(|p| p.to_string_lossy().to_string()),
443        print_config: opts.print_config,
444        verbose: opts.verbose,
445        mock: false,
446    };
447
448    // Load configuration
449    let mut config = AppConfig::load_or_default(opts.config_path.as_ref())?;
450    config.apply_cli_overrides(args.verbose);
451
452    // Try to read rendered gear config from master host via env var BEFORE logging init
453    // so we can use the tracing config from master for OTEL
454    let rendered_config = match std::env::var(TOOLKIT_MODULE_CONFIG_ENV) {
455        Ok(json) => RenderedGearConfig::from_json(&json).ok(),
456        Err(_) => None,
457    };
458
459    // Build final config by merging:
460    // 1. Rendered config from master host (base)
461    // 2. Local config file (override)
462    // This also merges logging configuration for proper initialization
463    let (final_config, merged_logging, db_options) =
464        build_oop_config_and_db(&config, &opts.gear_name, rendered_config.as_ref())?;
465
466    // Use OpenTelemetry config from rendered (master) config only.
467    // OoP gears do not fall back to local config for telemetry — if the master
468    // does not provide an opentelemetry section, telemetry is skipped entirely.
469    #[cfg(feature = "otel")]
470    let otel_cfg = rendered_config
471        .as_ref()
472        .and_then(|rc| rc.opentelemetry.as_ref());
473
474    // Initialize OTEL tracing layer (if tracing is enabled)
475    #[cfg(feature = "otel")]
476    let otel_layer = otel_cfg
477        .filter(|cfg| cfg.tracing.enabled)
478        .map(crate::telemetry::init_tracing)
479        .transpose()?;
480    #[cfg(not(feature = "otel"))]
481    let otel_layer = None;
482
483    // Initialize OpenTelemetry metrics provider (if configured and enabled).
484    // Store error to log after logging is initialized.
485    #[cfg(feature = "otel")]
486    let metrics_init_error = otel_cfg
487        .filter(|cfg| cfg.metrics.enabled)
488        .and_then(|cfg| crate::telemetry::init::init_metrics_provider(cfg).err());
489
490    // Initialize logging with MERGED config (master base + local override)
491    init_logging_unified(&merged_logging, &config.server.home_dir, otel_layer);
492
493    // Now that logging is available, report deferred metrics init error
494    #[cfg(feature = "otel")]
495    if let Some(e) = metrics_init_error {
496        tracing::error!(error = %e, "OpenTelemetry metrics not initialized (OoP)");
497    }
498
499    // Register custom panic hook to reroute panic backtrace into tracing.
500    init_panic_tracing();
501
502    // Now we can log - report what we received from master
503    if let Some(ref rc) = rendered_config {
504        info!(
505            env_var = TOOLKIT_MODULE_CONFIG_ENV,
506            has_database = rc.database.is_some(),
507            has_config = !rc.config.is_null(),
508            has_logging = rc.logging.is_some(),
509            has_opentelemetry = rc.opentelemetry.is_some(),
510            "Received rendered config from master host"
511        );
512    } else if std::env::var(TOOLKIT_MODULE_CONFIG_ENV).is_ok() {
513        warn!(
514            env_var = TOOLKIT_MODULE_CONFIG_ENV,
515            "Failed to parse rendered config from master host, using local config only"
516        );
517    } else {
518        debug!(
519            env_var = TOOLKIT_MODULE_CONFIG_ENV,
520            "No rendered config from master host, using local config only"
521        );
522    }
523
524    info!(
525        gear =  %opts.gear_name,
526        instance_id = %instance_id,
527        directory_endpoint = %opts.directory_endpoint,
528        "OoP gear bootstrap starting"
529    );
530
531    // Print config and exit if requested
532    if opts.print_config {
533        print_config(&config);
534        return Ok(());
535    }
536
537    // Connect to DirectoryService. When the gear configures a platform-plane
538    // credential, attach it (as `x-toolkit-internal-token`) to every outbound
539    // system call via an InternalAuthInterceptor (`cpt-cf-adr-two-plane-auth`).
540    info!(
541        "Connecting to directory service at {}",
542        opts.directory_endpoint
543    );
544    let internal_auth_cfg = final_config
545        .oop_http
546        .as_ref()
547        .and_then(|h| h.internal_auth.as_ref());
548    // Build the outbound DirectoryService interceptor and the
549    // `InternalTokenProvider` (for `#[toolkit::provides]` clients) from one
550    // credential source so they can't diverge after a rotation — see
551    // `build_platform_credentials`. `None` => no outbound platform credential.
552    let (directory_client, internal_token_provider) = if let Some(cfg) = internal_auth_cfg {
553        let (interceptor, provider) = build_platform_credentials(cfg, &cancel).await?;
554        info!("Attaching platform-plane credential to outbound DirectoryService calls");
555        let client =
556            DirectoryGrpcClient::connect_with_interceptor(&opts.directory_endpoint, interceptor)
557                .await?;
558        (client, provider)
559    } else {
560        (
561            DirectoryGrpcClient::connect(&opts.directory_endpoint).await?,
562            None,
563        )
564    };
565    let directory_api: Arc<dyn DirectoryClient> = Arc::new(directory_client);
566
567    info!("Successfully connected to directory service");
568
569    // Capture OoP HTTP config (if any) before moving the config into the provider.
570    let oop_http = final_config.oop_http.clone();
571
572    // Build config provider for gears
573    let config_provider = Arc::new(final_config);
574
575    // The DirectoryClient (gRPC client) is injected into the ClientHub so gears can access it.
576    // `oop` is left unset: OoP gears don't spawn other OoP gears.
577    let run_options = RunOptions::new(
578        config_provider,
579        db_options,
580        ShutdownOptions::Token(cancel.clone()),
581        instance_id,
582    )
583    .with_clients(vec![ClientRegistration::new::<dyn DirectoryClient>(
584        Arc::clone(&directory_api),
585    )])
586    .with_internal_token_provider(internal_token_provider);
587
588    // When `oop_http` is configured, run the HTTP-serving lifecycle:
589    // Axum server + probes + directory presence (registration + heartbeat +
590    // self-heal, owned by `presence_loop`) + dependency resolution + drain.
591    // Otherwise fall back to the legacy gRPC-only lifecycle.
592    let result = if let Some(http_cfg) = oop_http {
593        info!("Starting OoP HTTP-serving lifecycle");
594        let serve = build_oop_serve_options(
595            &http_cfg,
596            &opts.gear_name,
597            instance_id,
598            opts.version.clone(),
599            Duration::from_secs(opts.heartbeat_interval_secs),
600            Arc::clone(&directory_api),
601        )
602        .await?;
603        run_oop_serving(run_options, serve).await
604    } else {
605        info!("Starting gear lifecycle (legacy gRPC-only)");
606        // Legacy path: presence is not self-managed by an HTTP serve lifecycle,
607        // so run a standalone heartbeat on a child token. (A no-op unless the
608        // instance is registered by an external orchestrator.)
609        let heartbeat_directory = Arc::clone(&directory_api);
610        let heartbeat_gear = opts.gear_name.clone();
611        let heartbeat_instance_id_str = instance_id.to_string();
612        let heartbeat_interval = Duration::from_secs(opts.heartbeat_interval_secs.max(1));
613        let heartbeat_cancel = cancel.child_token();
614        tokio::spawn(async move {
615            info!(interval_secs = ?heartbeat_interval, "Starting legacy heartbeat loop");
616            loop {
617                tokio::select! {
618                    () = heartbeat_cancel.cancelled() => {
619                        info!("Heartbeat loop stopping due to cancellation");
620                        break;
621                    }
622                    () = sleep(heartbeat_interval) => {
623                        if let Err(e) = heartbeat_directory
624                            .send_heartbeat(&heartbeat_gear, &heartbeat_instance_id_str)
625                            .await
626                        {
627                            warn!(error = %e, "Failed to send heartbeat, will retry");
628                        }
629                    }
630                }
631            }
632        });
633        run(run_options).await
634    };
635
636    if let Err(ref e) = result {
637        error!(error = %e, "Gear runtime failed");
638    } else {
639        info!("Gear runtime completed successfully");
640    }
641
642    result
643}
644
645/// Build [`OopServeOptions`] from configuration.
646///
647/// The tenant-plane `BearerAuthenticator` injection point is left unset here;
648/// the app/gear binary supplies that adapter when it needs the tenant-plane
649/// middleware installed. The platform-plane authenticator is constructed from
650/// `oop_http.internal_auth` when the `k8s-auth` feature is enabled.
651async fn build_oop_serve_options(
652    cfg: &super::config::OopHttpConfig,
653    gear_name: &str,
654    instance_id: Uuid,
655    version: Option<String>,
656    heartbeat_interval: Duration,
657    directory: Arc<dyn DirectoryClient>,
658) -> Result<OopServeOptions> {
659    let listen_addr: std::net::SocketAddr = cfg
660        .listen_addr
661        .parse()
662        .with_context(|| format!("invalid oop_http.listen_addr: {}", cfg.listen_addr))?;
663
664    let probe_bind_addr = cfg
665        .probe_bind_addr
666        .as_deref()
667        .map(|s| {
668            s.parse::<std::net::SocketAddr>()
669                .with_context(|| format!("invalid oop_http.probe_bind_addr: {s}"))
670        })
671        .transpose()?;
672
673    let advertise_uri = cfg
674        .advertise_uri
675        .clone()
676        .unwrap_or_else(|| default_advertise_uri(listen_addr));
677
678    // Fail fast on a malformed or unreachable advertise_uri rather than only
679    // when the directory rejects the registration — the same late-failure the
680    // label validation below avoids.
681    validate_advertise_uri(&advertise_uri, cfg.allow_loopback_advertise)?;
682
683    // Fail fast on a mis-configured label (bad charset, over-long, too many)
684    // at start-up, using the same shared rules the directory enforces. Without
685    // this the process would come up, attempt to register, and only then be
686    // permanently rejected by the directory — a confusing late failure for what
687    // is a static configuration error. Checked before authenticator init so a
688    // static config error is reported without any async backend setup.
689    cf_system_sdks::directory::validate_labels(&cfg.labels).with_context(|| {
690        "invalid oop_http.labels: label keys/values must be <=63 chars, <=64 entries, and use \
691         only ASCII alphanumerics plus '-', '_', '.' (starting and ending alphanumeric)"
692    })?;
693
694    let internal_authenticator = build_internal_authenticator(cfg.internal_auth.as_ref()).await?;
695
696    Ok(OopServeOptions {
697        gear_name: gear_name.to_owned(),
698        instance_id: instance_id.to_string(),
699        version,
700        advertise_uri,
701        listen_addr,
702        probe_bind_addr,
703        drain_timeout: Duration::from_secs(cfg.drain_timeout_secs),
704        heartbeat_interval,
705        healthcheck_timeout: Duration::from_millis(cfg.healthcheck_timeout_ms),
706        directory,
707        bearer_authenticator: None,
708        internal_authenticator,
709        labels: cfg.labels.clone(),
710    })
711}
712
713/// Construct the inbound platform-plane authenticator from configuration.
714///
715/// The `shared_secret` provider is built here directly (no external backend).
716/// The `kube` provider initializes the Kubernetes `TokenReview` authenticator
717/// and requires the `k8s-auth` feature; without it, `provider: kube` is an
718/// error rather than a silent enforcement downgrade. `Ok(None)` is returned
719/// only when no `internal_auth` is configured.
720#[cfg_attr(not(feature = "k8s-auth"), allow(clippy::unused_async))]
721async fn build_internal_authenticator(
722    cfg: Option<&toolkit_security::InternalAuthConfig>,
723) -> Result<Option<toolkit_security::DynInternalAuthenticator>> {
724    let Some(cfg) = cfg else {
725        return Ok(None);
726    };
727
728    // Dependency-light providers (shared-secret) build directly here.
729    if let Some(authenticator) = cfg.build_authenticator() {
730        info!("Initializing shared-secret platform-plane authenticator");
731        return Ok(Some(authenticator));
732    }
733
734    #[cfg(feature = "k8s-auth")]
735    {
736        if cfg.is_kube() {
737            info!("Initializing Kubernetes TokenReview platform-plane authenticator");
738            let audiences = cfg.kube_audiences().unwrap_or_default().to_vec();
739            let authenticator =
740                toolkit_k8s_auth::K8sTokenReviewAuthenticator::try_default(audiences)
741                    .await
742                    .context("failed to initialize Kubernetes TokenReview authenticator")?;
743            return Ok(Some(toolkit_security::DynInternalAuthenticator::new(
744                authenticator,
745            )));
746        }
747    }
748    #[cfg(not(feature = "k8s-auth"))]
749    {
750        if cfg.is_kube() {
751            anyhow::bail!("oop_http.internal_auth provider=kube requires the `k8s-auth` feature");
752        }
753    }
754
755    Ok(None)
756}
757
758/// Build the platform-plane credential material from config: the **outbound**
759/// gRPC [`InternalAuthInterceptor`] (for `DirectoryService` calls) and the
760/// transport-agnostic
761/// [`InternalTokenProvider`](toolkit_contract::runtime::config::InternalTokenProvider)
762/// that `#[toolkit::provides]` clients attach as `X-ToolKit-Internal-Token` on
763/// platform-plane methods (`cpt-cf-adr-two-plane-auth`).
764///
765/// Both derive from ONE credential source so they never disagree after a
766/// rotation:
767/// - `shared_secret` → a static token for both.
768/// - `kube` with `token_path` → one [`ServiceAccountTokenReader`] drives both;
769///   its "no token yet" maps to [`CredentialState::Unavailable`] so an attach
770///   site warns rather than silently sending nothing.
771/// - `kube` without `token_path` → inbound-only: disabled interceptor + `None`
772///   provider (warned).
773///
774/// The `match` is exhaustive so a future variant fails to compile here rather
775/// than silently downgrading to no outbound credential.
776async fn build_platform_credentials(
777    cfg: &toolkit_security::InternalAuthConfig,
778    cancel: &CancellationToken,
779) -> Result<(
780    toolkit_transport_grpc::InternalAuthInterceptor,
781    Option<toolkit_contract::runtime::config::InternalTokenProvider>,
782)> {
783    use secrecy::SecretString;
784    use toolkit_contract::runtime::config::{CredentialState, InternalTokenProvider};
785    use toolkit_security::InternalAuthConfig;
786    use toolkit_transport_grpc::{
787        DEFAULT_REFRESH_INTERVAL, InternalAuthInterceptor, ServiceAccountTokenReader,
788    };
789
790    match cfg {
791        InternalAuthConfig::SharedSecret { secret, .. } => {
792            let token = SecretString::from(secret.clone());
793            Ok((
794                InternalAuthInterceptor::from_token(token.clone()),
795                Some(InternalTokenProvider::from_token(token)),
796            ))
797        }
798        InternalAuthConfig::Kube {
799            token_path: Some(path),
800            ..
801        } => {
802            let reader = ServiceAccountTokenReader::with_cancellation(
803                path,
804                DEFAULT_REFRESH_INTERVAL,
805                cancel.child_token(),
806            )
807            .await
808            .context("failed to read projected service-account token for outbound credential")?;
809            let interceptor = reader.interceptor();
810            // The reader yields `None` while the projected file is empty or the
811            // refresh has not populated it; surface that as `Unavailable` (warn)
812            // rather than `NotConfigured` (silent).
813            let token_fn = reader.token_provider();
814            let provider = InternalTokenProvider::new(move || match token_fn() {
815                Some(token) => CredentialState::Available(token),
816                None => CredentialState::Unavailable(
817                    "projected service-account token is currently unavailable \
818                     (file empty or not yet read)"
819                        .into(),
820                ),
821            });
822            Ok((interceptor, Some(provider)))
823        }
824        InternalAuthConfig::Kube {
825            token_path: None, ..
826        } => {
827            warn!(
828                "oop_http.internal_auth: provider=kube without token_path - this participant \
829                 validates inbound platform tokens but will attach NO outbound credential"
830            );
831            Ok((InternalAuthInterceptor::disabled(), None))
832        }
833    }
834}
835
836/// Derive a default advertise URI from the bind address. Unspecified hosts
837/// (`0.0.0.0` / `[::]`) are rewritten to loopback, and every IPv6 host is
838/// enclosed in brackets so the resulting URI is valid when registered as a
839/// REST endpoint.
840fn default_advertise_uri(listen_addr: std::net::SocketAddr) -> String {
841    let host = match listen_addr {
842        std::net::SocketAddr::V4(addr) if addr.ip().is_unspecified() => "127.0.0.1".to_owned(),
843        std::net::SocketAddr::V4(addr) => addr.ip().to_string(),
844        std::net::SocketAddr::V6(addr) if addr.ip().is_unspecified() => "[::1]".to_owned(),
845        std::net::SocketAddr::V6(addr) => format!("[{}]", addr.ip()),
846    };
847    format!("http://{host}:{}", listen_addr.port())
848}
849
850/// Reject a malformed or unreachable `advertise_uri` at start-up (fail fast).
851///
852/// Checks shape (parseable `http`/`https` URL, non-empty host, no userinfo) and,
853/// unless `allow_loopback`, rejects a loopback / unspecified host - a
854/// registered-but-unreachable instance in multi-host Profile 2 / Profile 3
855/// (`cpt-cf-adr-instance-addressable-discovery` section 5). The default derives
856/// loopback from an unspecified bind, so this also covers an *unset* value. The
857/// directory enforces its full endpoint ruleset server-side.
858fn validate_advertise_uri(uri: &str, allow_loopback: bool) -> Result<()> {
859    let parsed = url::Url::parse(uri)
860        .with_context(|| format!("invalid oop_http.advertise_uri: not a valid URL: {uri}"))?;
861    if !matches!(parsed.scheme(), "http" | "https") {
862        anyhow::bail!(
863            "invalid oop_http.advertise_uri: scheme must be http or https (got '{}')",
864            parsed.scheme()
865        );
866    }
867    if parsed.host_str().is_none_or(str::is_empty) {
868        anyhow::bail!("invalid oop_http.advertise_uri: missing host: {uri}");
869    }
870    if !parsed.username().is_empty() || parsed.password().is_some() {
871        anyhow::bail!("invalid oop_http.advertise_uri: must not contain userinfo: {uri}");
872    }
873    let is_loopback = match parsed.host() {
874        Some(url::Host::Ipv4(ip)) => ip.is_loopback() || ip.is_unspecified(),
875        Some(url::Host::Ipv6(ip)) => ip.is_loopback() || ip.is_unspecified(),
876        Some(url::Host::Domain(d)) => d.trim_end_matches('.').eq_ignore_ascii_case("localhost"),
877        None => false,
878    };
879    if !allow_loopback && is_loopback {
880        anyhow::bail!(
881            "invalid oop_http.advertise_uri: '{uri}' is a loopback/unspecified address, which is \
882             unreachable by other gears in multi-host Profile 2 / Profile 3 (a registered-but-\
883             unreachable instance). Set oop_http.advertise_uri to a routable host, or set \
884             oop_http.allow_loopback_advertise = true for single-host / local-dev."
885        );
886    }
887    Ok(())
888}
889
890#[allow(unknown_lints, de1301_no_print_macros)] // direct stdout config print before exit
891fn print_config(config: &AppConfig) {
892    match config.to_yaml() {
893        Ok(yaml) => {
894            println!("{yaml}");
895        }
896        Err(e) => {
897            eprintln!("Failed to render config as YAML: {e}");
898        }
899    }
900}
901
902#[cfg(test)]
903#[cfg_attr(coverage_nightly, coverage(off))]
904#[path = "oop_tests.rs"]
905mod tests;