1use 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#[derive(Debug, Clone)]
69pub struct OopRunOptions {
70 pub gear_name: String,
72
73 pub instance_id: Option<Uuid>,
75
76 pub directory_endpoint: String,
78
79 pub config_path: Option<PathBuf>,
81
82 pub verbose: u8,
84
85 pub print_config: bool,
87
88 pub heartbeat_interval_secs: u64,
90
91 pub version: Option<String>,
94}
95
96impl Default for OopRunOptions {
97 fn default() -> Self {
98 let config_path = std::env::var("TOOLKIT_CONFIG_PATH").ok().map(PathBuf::from);
100
101 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#[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 let final_config = if let Some(rendered) = rendered_config {
149 let mut config = local_config.clone();
151
152 let gear_entry = config
154 .gears
155 .entry(gear_name.to_owned())
156 .or_insert_with(|| serde_json::json!({}));
157
158 if let Some(obj) = gear_entry.as_object_mut() {
160 if !obj.contains_key("config") || obj["config"].is_null() {
163 obj.insert("config".to_owned(), rendered.config.clone());
164 }
165 }
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 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 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 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
203fn 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
216fn 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 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 let mut merged_config = serde_json::Map::new();
248
249 if let Some(rendered) = rendered_db {
251 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 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 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 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 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 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
327fn 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 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 *target = source.clone();
346 }
347}
348
349#[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 let instance_id = opts.instance_id.unwrap_or_else(Uuid::new_v4);
412
413 let cancel = CancellationToken::new();
416
417 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 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 let mut config = AppConfig::load_or_default(opts.config_path.as_ref())?;
450 config.apply_cli_overrides(args.verbose);
451
452 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 let (final_config, merged_logging, db_options) =
464 build_oop_config_and_db(&config, &opts.gear_name, rendered_config.as_ref())?;
465
466 #[cfg(feature = "otel")]
470 let otel_cfg = rendered_config
471 .as_ref()
472 .and_then(|rc| rc.opentelemetry.as_ref());
473
474 #[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 #[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 init_logging_unified(&merged_logging, &config.server.home_dir, otel_layer);
492
493 #[cfg(feature = "otel")]
495 if let Some(e) = metrics_init_error {
496 tracing::error!(error = %e, "OpenTelemetry metrics not initialized (OoP)");
497 }
498
499 init_panic_tracing();
501
502 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 if opts.print_config {
533 print_config(&config);
534 return Ok(());
535 }
536
537 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 let directory_client = if let Some(cfg) = internal_auth_cfg {
549 let interceptor = toolkit_transport_grpc::build_internal_auth_interceptor(cfg).await?;
550 info!("Attaching platform-plane credential to outbound DirectoryService calls");
551 DirectoryGrpcClient::connect_with_interceptor(&opts.directory_endpoint, interceptor).await?
552 } else {
553 DirectoryGrpcClient::connect(&opts.directory_endpoint).await?
554 };
555 let directory_api: Arc<dyn DirectoryClient> = Arc::new(directory_client);
556
557 info!("Successfully connected to directory service");
558
559 let oop_http = final_config.oop_http.clone();
561
562 let config_provider = Arc::new(final_config);
564
565 let run_options = RunOptions {
567 gears_cfg: config_provider,
568 db: db_options,
569 shutdown: ShutdownOptions::Token(cancel.clone()),
570 clients: vec![ClientRegistration::new::<dyn DirectoryClient>(Arc::clone(
571 &directory_api,
572 ))],
573 instance_id,
574 oop: None, shutdown_deadline: None,
576 };
577
578 let result = if let Some(http_cfg) = oop_http {
583 info!("Starting OoP HTTP-serving lifecycle");
584 let serve = build_oop_serve_options(
585 &http_cfg,
586 &opts.gear_name,
587 instance_id,
588 opts.version.clone(),
589 Duration::from_secs(opts.heartbeat_interval_secs),
590 Arc::clone(&directory_api),
591 )
592 .await?;
593 run_oop_serving(run_options, serve).await
594 } else {
595 info!("Starting gear lifecycle (legacy gRPC-only)");
596 let heartbeat_directory = Arc::clone(&directory_api);
600 let heartbeat_gear = opts.gear_name.clone();
601 let heartbeat_instance_id_str = instance_id.to_string();
602 let heartbeat_interval = Duration::from_secs(opts.heartbeat_interval_secs.max(1));
603 let heartbeat_cancel = cancel.child_token();
604 tokio::spawn(async move {
605 info!(interval_secs = ?heartbeat_interval, "Starting legacy heartbeat loop");
606 loop {
607 tokio::select! {
608 () = heartbeat_cancel.cancelled() => {
609 info!("Heartbeat loop stopping due to cancellation");
610 break;
611 }
612 () = sleep(heartbeat_interval) => {
613 if let Err(e) = heartbeat_directory
614 .send_heartbeat(&heartbeat_gear, &heartbeat_instance_id_str)
615 .await
616 {
617 warn!(error = %e, "Failed to send heartbeat, will retry");
618 }
619 }
620 }
621 }
622 });
623 run(run_options).await
624 };
625
626 if let Err(ref e) = result {
627 error!(error = %e, "Gear runtime failed");
628 } else {
629 info!("Gear runtime completed successfully");
630 }
631
632 result
633}
634
635async fn build_oop_serve_options(
642 cfg: &super::config::OopHttpConfig,
643 gear_name: &str,
644 instance_id: Uuid,
645 version: Option<String>,
646 heartbeat_interval: Duration,
647 directory: Arc<dyn DirectoryClient>,
648) -> Result<OopServeOptions> {
649 let listen_addr: std::net::SocketAddr = cfg
650 .listen_addr
651 .parse()
652 .with_context(|| format!("invalid oop_http.listen_addr: {}", cfg.listen_addr))?;
653
654 let probe_bind_addr = cfg
655 .probe_bind_addr
656 .as_deref()
657 .map(|s| {
658 s.parse::<std::net::SocketAddr>()
659 .with_context(|| format!("invalid oop_http.probe_bind_addr: {s}"))
660 })
661 .transpose()?;
662
663 let advertise_uri = cfg
664 .advertise_uri
665 .clone()
666 .unwrap_or_else(|| default_advertise_uri(listen_addr));
667
668 let internal_authenticator = build_internal_authenticator(cfg.internal_auth.as_ref()).await?;
669
670 Ok(OopServeOptions {
671 gear_name: gear_name.to_owned(),
672 instance_id: instance_id.to_string(),
673 version,
674 advertise_uri,
675 listen_addr,
676 probe_bind_addr,
677 drain_timeout: Duration::from_secs(cfg.drain_timeout_secs),
678 heartbeat_interval,
679 healthcheck_timeout: Duration::from_millis(cfg.healthcheck_timeout_ms),
680 directory,
681 bearer_authenticator: None,
682 internal_authenticator,
683 })
684}
685
686#[cfg_attr(not(feature = "k8s-auth"), allow(clippy::unused_async))]
694async fn build_internal_authenticator(
695 cfg: Option<&toolkit_security::InternalAuthConfig>,
696) -> Result<Option<toolkit_security::DynInternalAuthenticator>> {
697 let Some(cfg) = cfg else {
698 return Ok(None);
699 };
700
701 if let Some(authenticator) = cfg.build_authenticator() {
703 info!("Initializing shared-secret platform-plane authenticator");
704 return Ok(Some(authenticator));
705 }
706
707 #[cfg(feature = "k8s-auth")]
708 {
709 if cfg.is_kube() {
710 info!("Initializing Kubernetes TokenReview platform-plane authenticator");
711 let audiences = cfg.kube_audiences().unwrap_or_default().to_vec();
712 let authenticator =
713 toolkit_k8s_auth::K8sTokenReviewAuthenticator::try_default(audiences)
714 .await
715 .context("failed to initialize Kubernetes TokenReview authenticator")?;
716 return Ok(Some(toolkit_security::DynInternalAuthenticator::new(
717 authenticator,
718 )));
719 }
720 }
721 #[cfg(not(feature = "k8s-auth"))]
722 {
723 if cfg.is_kube() {
724 anyhow::bail!("oop_http.internal_auth provider=kube requires the `k8s-auth` feature");
725 }
726 }
727
728 Ok(None)
729}
730
731fn default_advertise_uri(listen_addr: std::net::SocketAddr) -> String {
736 let host = match listen_addr {
737 std::net::SocketAddr::V4(addr) if addr.ip().is_unspecified() => "127.0.0.1".to_owned(),
738 std::net::SocketAddr::V4(addr) => addr.ip().to_string(),
739 std::net::SocketAddr::V6(addr) if addr.ip().is_unspecified() => "[::1]".to_owned(),
740 std::net::SocketAddr::V6(addr) => format!("[{}]", addr.ip()),
741 };
742 format!("http://{host}:{}", listen_addr.port())
743}
744
745#[allow(unknown_lints, de1301_no_print_macros)] fn print_config(config: &AppConfig) {
747 match config.to_yaml() {
748 Ok(yaml) => {
749 println!("{yaml}");
750 }
751 Err(e) => {
752 eprintln!("Failed to render config as YAML: {e}");
753 }
754 }
755}
756
757#[cfg(test)]
758#[cfg_attr(coverage_nightly, coverage(off))]
759#[path = "oop_tests.rs"]
760mod tests;