1use anyhow::{Context, Result};
46use figment::{Figment, providers::Serialized};
47use std::path::{Path, PathBuf};
48use std::sync::Arc;
49use std::time::Duration;
50use tokio::time::sleep;
51use tokio_util::sync::CancellationToken;
52use tracing::{debug, error, info, warn};
53use uuid::Uuid;
54
55use super::config::{
56 AppConfig, CliArgs, LoggingConfig, RenderedDbConfig, RenderedGearConfig,
57 TOOLKIT_MODULE_CONFIG_ENV,
58};
59use crate::bootstrap::host::{init_logging_unified, init_panic_tracing};
60use crate::runtime::{
61 ClientRegistration, DbOptions, RunOptions, ShutdownOptions, TOOLKIT_DIRECTORY_ENDPOINT_ENV,
62 run, shutdown,
63};
64use cf_system_sdks::directory::{DirectoryClient, DirectoryGrpcClient};
65
66#[derive(Debug, Clone)]
68pub struct OopRunOptions {
69 pub gear_name: String,
71
72 pub instance_id: Option<Uuid>,
74
75 pub directory_endpoint: String,
77
78 pub config_path: Option<PathBuf>,
80
81 pub verbose: u8,
83
84 pub print_config: bool,
86
87 pub heartbeat_interval_secs: u64,
89}
90
91impl Default for OopRunOptions {
92 fn default() -> Self {
93 let config_path = std::env::var("TOOLKIT_CONFIG_PATH").ok().map(PathBuf::from);
95
96 let directory_endpoint = std::env::var(TOOLKIT_DIRECTORY_ENDPOINT_ENV)
99 .unwrap_or_else(|_| "http://127.0.0.1:50051".to_owned());
100
101 Self {
102 gear_name: String::new(),
103 instance_id: None,
104 directory_endpoint,
105 config_path,
106 verbose: 0,
107 print_config: false,
108 heartbeat_interval_secs: 5,
109 }
110 }
111}
112
113#[tracing::instrument(
127 level = "debug",
128 skip(local_config, rendered_config),
129 fields(
130 has_rendered = rendered_config.is_some(),
131 has_local_db = local_config.database.is_some()
132 )
133)]
134fn build_oop_config_and_db(
135 local_config: &AppConfig,
136 gear_name: &str,
137 rendered_config: Option<&RenderedGearConfig>,
138) -> Result<(AppConfig, LoggingConfig, DbOptions)> {
139 let home_dir = PathBuf::from(&local_config.server.home_dir);
140
141 let final_config = if let Some(rendered) = rendered_config {
143 let mut config = local_config.clone();
145
146 let gear_entry = config
148 .gears
149 .entry(gear_name.to_owned())
150 .or_insert_with(|| serde_json::json!({}));
151
152 if let Some(obj) = gear_entry.as_object_mut() {
154 if !obj.contains_key("config") || obj["config"].is_null() {
157 obj.insert("config".to_owned(), rendered.config.clone());
158 }
159 }
161
162 debug!(
163 gear = %gear_name,
164 has_rendered_db = %rendered.database.is_some(),
165 has_rendered_logging = %rendered.logging.is_some(),
166 "Using rendered config from master as base, local config as override"
167 );
168
169 config
170 } else {
171 debug!(
173 gear = %gear_name,
174 "No rendered config from master, using local config entirely (standalone mode)"
175 );
176 local_config.clone()
177 };
178
179 let final_logging = merge_logging_configs(
181 rendered_config.as_ref().and_then(|r| r.logging.as_ref()),
182 &local_config.logging,
183 );
184
185 let db_options = build_merged_db_options(
188 &home_dir,
189 gear_name,
190 rendered_config.as_ref().and_then(|r| r.database.as_ref()),
191 local_config,
192 )?;
193
194 Ok((final_config, final_logging, db_options))
195}
196
197fn merge_logging_configs(master: Option<&LoggingConfig>, local: &LoggingConfig) -> LoggingConfig {
202 master
203 .cloned()
204 .unwrap_or_default()
205 .into_iter()
206 .chain(local.clone())
207 .collect()
208}
209
210fn build_merged_db_options(
215 home_dir: &Path,
216 gear_name: &str,
217 rendered_db: Option<&RenderedDbConfig>,
218 local_config: &AppConfig,
219) -> Result<DbOptions> {
220 let has_rendered_db = rendered_db.is_some_and(|db| db.gear.is_some() || db.global.is_some());
222 let has_local_db = local_config.database.is_some()
223 || local_config
224 .gears
225 .get(gear_name)
226 .and_then(|m| m.get("database"))
227 .is_some();
228
229 if !has_rendered_db && !has_local_db {
230 debug!(
231 gear = %gear_name,
232 "No database config available"
233 );
234 return Ok(DbOptions::None);
235 }
236
237 let mut merged_config = serde_json::Map::new();
242
243 if let Some(rendered) = rendered_db {
245 if let Some(ref global) = rendered.global {
247 let global_json = serde_json::to_value(global)
248 .context("Failed to serialize rendered global db config")?;
249 merged_config.insert("database".to_owned(), global_json);
250 }
251
252 if let Some(ref gear_db) = rendered.gear {
254 let gear_db_json = serde_json::to_value(gear_db)
255 .context("Failed to serialize rendered gear db config")?;
256
257 let mut gears = serde_json::Map::new();
258 let mut gear_entry = serde_json::Map::new();
259 gear_entry.insert("database".to_owned(), gear_db_json);
260 gears.insert(gear_name.to_owned(), serde_json::Value::Object(gear_entry));
261 merged_config.insert("gears".to_owned(), serde_json::Value::Object(gears));
262 }
263 }
264
265 if let Some(ref local_db) = local_config.database {
268 let local_db_json =
269 serde_json::to_value(local_db).context("Failed to serialize local global db config")?;
270
271 if let Some(existing) = merged_config.get_mut("database") {
273 merge_json_objects(existing, &local_db_json);
274 } else {
275 merged_config.insert("database".to_owned(), local_db_json);
276 }
277 }
278
279 if let Some(local_gear) = local_config.gears.get(gear_name)
281 && let Some(local_gear_db) = local_gear.get("database")
282 {
283 let gears = merged_config
284 .entry("gears".to_owned())
285 .or_insert_with(|| serde_json::Value::Object(serde_json::Map::new()));
286
287 if let Some(gears_obj) = gears.as_object_mut() {
288 let gear_entry = gears_obj
289 .entry(gear_name.to_owned())
290 .or_insert_with(|| serde_json::Value::Object(serde_json::Map::new()));
291
292 if let Some(gear_obj) = gear_entry.as_object_mut() {
293 if let Some(existing_db) = gear_obj.get_mut("database") {
294 merge_json_objects(existing_db, local_gear_db);
295 } else {
296 gear_obj.insert("database".to_owned(), local_gear_db.clone());
297 }
298 }
299 }
300 }
301
302 debug!(
303 gear = %gear_name,
304 has_rendered = %rendered_db.is_some(),
305 has_local_global = %local_config.database.is_some(),
306 "Building DbManager with merged config"
307 );
308
309 let figment = Figment::new().merge(Serialized::defaults(serde_json::Value::Object(
311 merged_config,
312 )));
313 let db_manager = Arc::new(
314 toolkit_db::DbManager::from_figment(figment, home_dir.to_path_buf())
315 .context("Failed to create DbManager from merged config")?,
316 );
317
318 Ok(DbOptions::Manager(db_manager))
319}
320
321fn merge_json_objects(target: &mut serde_json::Value, source: &serde_json::Value) {
324 if let (Some(target_obj), Some(source_obj)) = (target.as_object_mut(), source.as_object()) {
325 for (key, value) in source_obj {
326 if let Some(target_value) = target_obj.get_mut(key) {
327 if target_value.is_object() && value.is_object() {
329 merge_json_objects(target_value, value);
330 } else {
331 *target_value = value.clone();
332 }
333 } else {
334 target_obj.insert(key.clone(), value.clone());
335 }
336 }
337 } else {
338 *target = source.clone();
340 }
341}
342
343#[tracing::instrument(
394 level = "info",
395 name = "oop_bootstrap",
396 skip(opts),
397 fields(
398 gear = %opts.gear_name,
399 directory = %opts.directory_endpoint
400 )
401)]
402pub async fn run_oop_with_options(opts: OopRunOptions) -> Result<()> {
403 let instance_id = opts.instance_id.unwrap_or_else(Uuid::new_v4);
405
406 let cancel = CancellationToken::new();
409
410 let cancel_for_signals = cancel.clone();
413 tokio::spawn(async move {
414 match shutdown::wait_for_shutdown().await {
415 Ok(()) => {
416 info!(target: "", "------------------");
417 info!("shutdown: signal received in OoP bootstrap");
418 }
419 Err(e) => {
420 warn!(
421 error = %e,
422 "shutdown: primary waiter failed in OoP bootstrap, falling back to ctrl_c()"
423 );
424 _ = tokio::signal::ctrl_c().await;
425 }
426 }
427 cancel_for_signals.cancel();
428 });
429
430 let args = CliArgs {
432 config: opts
433 .config_path
434 .as_ref()
435 .map(|p| p.to_string_lossy().to_string()),
436 print_config: opts.print_config,
437 verbose: opts.verbose,
438 mock: false,
439 };
440
441 let mut config = AppConfig::load_or_default(opts.config_path.as_ref())?;
443 config.apply_cli_overrides(args.verbose);
444
445 let rendered_config = match std::env::var(TOOLKIT_MODULE_CONFIG_ENV) {
448 Ok(json) => RenderedGearConfig::from_json(&json).ok(),
449 Err(_) => None,
450 };
451
452 let (final_config, merged_logging, db_options) =
457 build_oop_config_and_db(&config, &opts.gear_name, rendered_config.as_ref())?;
458
459 #[cfg(feature = "otel")]
463 let otel_cfg = rendered_config
464 .as_ref()
465 .and_then(|rc| rc.opentelemetry.as_ref());
466
467 #[cfg(feature = "otel")]
469 let otel_layer = otel_cfg
470 .filter(|cfg| cfg.tracing.enabled)
471 .map(crate::telemetry::init_tracing)
472 .transpose()?;
473 #[cfg(not(feature = "otel"))]
474 let otel_layer = None;
475
476 #[cfg(feature = "otel")]
479 let metrics_init_error = otel_cfg
480 .filter(|cfg| cfg.metrics.enabled)
481 .and_then(|cfg| crate::telemetry::init::init_metrics_provider(cfg).err());
482
483 init_logging_unified(&merged_logging, &config.server.home_dir, otel_layer);
485
486 #[cfg(feature = "otel")]
488 if let Some(e) = metrics_init_error {
489 tracing::error!(error = %e, "OpenTelemetry metrics not initialized (OoP)");
490 }
491
492 init_panic_tracing();
494
495 if let Some(ref rc) = rendered_config {
497 info!(
498 env_var = TOOLKIT_MODULE_CONFIG_ENV,
499 has_database = rc.database.is_some(),
500 has_config = !rc.config.is_null(),
501 has_logging = rc.logging.is_some(),
502 has_opentelemetry = rc.opentelemetry.is_some(),
503 "Received rendered config from master host"
504 );
505 } else if std::env::var(TOOLKIT_MODULE_CONFIG_ENV).is_ok() {
506 warn!(
507 env_var = TOOLKIT_MODULE_CONFIG_ENV,
508 "Failed to parse rendered config from master host, using local config only"
509 );
510 } else {
511 debug!(
512 env_var = TOOLKIT_MODULE_CONFIG_ENV,
513 "No rendered config from master host, using local config only"
514 );
515 }
516
517 info!(
518 gear = %opts.gear_name,
519 instance_id = %instance_id,
520 directory_endpoint = %opts.directory_endpoint,
521 "OoP gear bootstrap starting"
522 );
523
524 if opts.print_config {
526 print_config(&config);
527 return Ok(());
528 }
529
530 info!(
532 "Connecting to directory service at {}",
533 opts.directory_endpoint
534 );
535 let directory_client = DirectoryGrpcClient::connect(&opts.directory_endpoint).await?;
536 let directory_api: Arc<dyn DirectoryClient> = Arc::new(directory_client);
537
538 info!("Successfully connected to directory service");
539
540 let heartbeat_directory = Arc::clone(&directory_api);
543 let heartbeat_gear = opts.gear_name.clone();
544 let heartbeat_instance_id_str = instance_id.to_string();
545 let heartbeat_interval = Duration::from_secs(opts.heartbeat_interval_secs);
546 let heartbeat_cancel = cancel.child_token();
547
548 tokio::spawn(async move {
549 info!(
550 interval_secs = opts.heartbeat_interval_secs,
551 "Starting heartbeat loop"
552 );
553
554 loop {
555 tokio::select! {
556 () = heartbeat_cancel.cancelled() => {
557 info!("Heartbeat loop stopping due to cancellation");
558 break;
559 }
560 () = sleep(heartbeat_interval) => {
561 match heartbeat_directory
562 .send_heartbeat(&heartbeat_gear, &heartbeat_instance_id_str)
563 .await
564 {
565 Ok(()) => {
566 tracing::debug!("Heartbeat sent successfully");
567 }
568 Err(e) => {
569 warn!(error = %e, "Failed to send heartbeat, will retry");
570 }
571 }
572 }
573 }
574 }
575 });
576
577 let config_provider = Arc::new(final_config);
579
580 info!("Starting gear lifecycle");
585 let run_options = RunOptions {
586 gears_cfg: config_provider,
587 db: db_options,
588 shutdown: ShutdownOptions::Token(cancel.clone()),
589 clients: vec![ClientRegistration::new::<dyn DirectoryClient>(
590 directory_api,
591 )],
592 instance_id,
593 oop: None, shutdown_deadline: None,
595 };
596
597 let result = run(run_options).await;
598
599 if let Err(ref e) = result {
600 error!(error = %e, "Gear runtime failed");
601 } else {
602 info!("Gear runtime completed successfully");
603 }
604
605 result
606}
607
608#[allow(unknown_lints, de1301_no_print_macros)] fn print_config(config: &AppConfig) {
610 match config.to_yaml() {
611 Ok(yaml) => {
612 println!("{yaml}");
613 }
614 Err(e) => {
615 eprintln!("Failed to render config as YAML: {e}");
616 }
617 }
618}
619
620#[cfg(test)]
621#[cfg_attr(coverage_nightly, coverage(off))]
622#[path = "oop_tests.rs"]
623mod tests;