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!(
539 "Connecting to directory service at {}",
540 opts.directory_endpoint
541 );
542 let directory_client = DirectoryGrpcClient::connect(&opts.directory_endpoint).await?;
543 let directory_api: Arc<dyn DirectoryClient> = Arc::new(directory_client);
544
545 info!("Successfully connected to directory service");
546
547 let oop_http = final_config.oop_http.clone();
549
550 let config_provider = Arc::new(final_config);
552
553 let run_options = RunOptions {
555 gears_cfg: config_provider,
556 db: db_options,
557 shutdown: ShutdownOptions::Token(cancel.clone()),
558 clients: vec![ClientRegistration::new::<dyn DirectoryClient>(Arc::clone(
559 &directory_api,
560 ))],
561 instance_id,
562 oop: None, shutdown_deadline: None,
564 };
565
566 let result = if let Some(http_cfg) = oop_http {
571 info!("Starting OoP HTTP-serving lifecycle");
572 let serve = build_oop_serve_options(
573 &http_cfg,
574 &opts.gear_name,
575 instance_id,
576 opts.version.clone(),
577 Duration::from_secs(opts.heartbeat_interval_secs),
578 Arc::clone(&directory_api),
579 )
580 .await?;
581 run_oop_serving(run_options, serve).await
582 } else {
583 info!("Starting gear lifecycle (legacy gRPC-only)");
584 let heartbeat_directory = Arc::clone(&directory_api);
588 let heartbeat_gear = opts.gear_name.clone();
589 let heartbeat_instance_id_str = instance_id.to_string();
590 let heartbeat_interval = Duration::from_secs(opts.heartbeat_interval_secs.max(1));
591 let heartbeat_cancel = cancel.child_token();
592 tokio::spawn(async move {
593 info!(interval_secs = ?heartbeat_interval, "Starting legacy heartbeat loop");
594 loop {
595 tokio::select! {
596 () = heartbeat_cancel.cancelled() => {
597 info!("Heartbeat loop stopping due to cancellation");
598 break;
599 }
600 () = sleep(heartbeat_interval) => {
601 if let Err(e) = heartbeat_directory
602 .send_heartbeat(&heartbeat_gear, &heartbeat_instance_id_str)
603 .await
604 {
605 warn!(error = %e, "Failed to send heartbeat, will retry");
606 }
607 }
608 }
609 }
610 });
611 run(run_options).await
612 };
613
614 if let Err(ref e) = result {
615 error!(error = %e, "Gear runtime failed");
616 } else {
617 info!("Gear runtime completed successfully");
618 }
619
620 result
621}
622
623async fn build_oop_serve_options(
630 cfg: &super::config::OopHttpConfig,
631 gear_name: &str,
632 instance_id: Uuid,
633 version: Option<String>,
634 heartbeat_interval: Duration,
635 directory: Arc<dyn DirectoryClient>,
636) -> Result<OopServeOptions> {
637 let listen_addr: std::net::SocketAddr = cfg
638 .listen_addr
639 .parse()
640 .with_context(|| format!("invalid oop_http.listen_addr: {}", cfg.listen_addr))?;
641
642 let probe_bind_addr = cfg
643 .probe_bind_addr
644 .as_deref()
645 .map(|s| {
646 s.parse::<std::net::SocketAddr>()
647 .with_context(|| format!("invalid oop_http.probe_bind_addr: {s}"))
648 })
649 .transpose()?;
650
651 let advertise_uri = cfg
652 .advertise_uri
653 .clone()
654 .unwrap_or_else(|| default_advertise_uri(listen_addr));
655
656 let internal_authenticator = build_internal_authenticator(cfg.internal_auth.as_ref()).await?;
657
658 Ok(OopServeOptions {
659 gear_name: gear_name.to_owned(),
660 instance_id: instance_id.to_string(),
661 version,
662 advertise_uri,
663 listen_addr,
664 probe_bind_addr,
665 drain_timeout: Duration::from_secs(cfg.drain_timeout_secs),
666 heartbeat_interval,
667 healthcheck_timeout: Duration::from_millis(cfg.healthcheck_timeout_ms),
668 directory,
669 bearer_authenticator: None,
670 internal_authenticator,
671 })
672}
673
674#[cfg_attr(not(feature = "k8s-auth"), allow(clippy::unused_async))]
681async fn build_internal_authenticator(
682 cfg: Option<&super::config::InternalAuthConfig>,
683) -> Result<Option<crate::runtime::DynInternalAuthenticator>> {
684 #[cfg(feature = "k8s-auth")]
685 {
686 if let Some(internal_auth) = cfg {
687 info!("Initializing Kubernetes TokenReview platform-plane authenticator");
688 let authenticator = toolkit_k8s_auth::K8sTokenReviewAuthenticator::try_default(
689 internal_auth.audiences.clone(),
690 )
691 .await
692 .context("failed to initialize Kubernetes TokenReview authenticator")?;
693 return Ok(Some(crate::runtime::DynInternalAuthenticator::new(
694 authenticator,
695 )));
696 }
697 Ok(None)
698 }
699 #[cfg(not(feature = "k8s-auth"))]
700 {
701 if cfg.is_some() {
702 anyhow::bail!(
703 "oop_http.internal_auth is configured but the `k8s-auth` feature is disabled"
704 );
705 }
706 Ok(None)
707 }
708}
709
710fn default_advertise_uri(listen_addr: std::net::SocketAddr) -> String {
715 let host = match listen_addr {
716 std::net::SocketAddr::V4(addr) if addr.ip().is_unspecified() => "127.0.0.1".to_owned(),
717 std::net::SocketAddr::V4(addr) => addr.ip().to_string(),
718 std::net::SocketAddr::V6(addr) if addr.ip().is_unspecified() => "[::1]".to_owned(),
719 std::net::SocketAddr::V6(addr) => format!("[{}]", addr.ip()),
720 };
721 format!("http://{host}:{}", listen_addr.port())
722}
723
724#[allow(unknown_lints, de1301_no_print_macros)] fn print_config(config: &AppConfig) {
726 match config.to_yaml() {
727 Ok(yaml) => {
728 println!("{yaml}");
729 }
730 Err(e) => {
731 eprintln!("Failed to render config as YAML: {e}");
732 }
733 }
734}
735
736#[cfg(test)]
737#[cfg_attr(coverage_nightly, coverage(off))]
738#[path = "oop_tests.rs"]
739mod tests;