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, 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 let oop_http = final_config.oop_http.clone();
571
572 let config_provider = Arc::new(final_config);
574
575 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 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 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
645async 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 validate_advertise_uri(&advertise_uri, cfg.allow_loopback_advertise)?;
682
683 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#[cfg_attr(not(feature = "k8s-auth"), allow(clippy::unused_async))]
725async fn build_internal_authenticator(
726 cfg: Option<&toolkit_security::InternalAuthConfig>,
727) -> Result<Option<toolkit_security::DynInternalAuthenticator>> {
728 let Some(cfg) = cfg else {
729 return Ok(None);
730 };
731
732 if let Some(authenticator) = cfg.build_authenticator() {
734 info!("Initializing shared-secret platform-plane authenticator");
735 return Ok(Some(authenticator));
736 }
737
738 #[cfg(feature = "k8s-auth")]
739 {
740 if cfg.is_kube() {
741 info!("Initializing Kubernetes TokenReview platform-plane authenticator");
742 let audiences = cfg.kube_audiences().unwrap_or_default().to_vec();
743 let authenticator = toolkit_k8s_auth::build_cached_k8s_authenticator(
744 audiences,
745 Some(toolkit_security::DEFAULT_TOKEN_REVIEW_CACHE_TTL),
746 )
747 .await
748 .context("failed to initialize Kubernetes TokenReview authenticator")?;
749 return Ok(Some(authenticator));
750 }
751 }
752 #[cfg(not(feature = "k8s-auth"))]
753 {
754 if cfg.is_kube() {
755 anyhow::bail!("oop_http.internal_auth provider=kube requires the `k8s-auth` feature");
756 }
757 }
758
759 anyhow::bail!(
760 "internal_auth is configured but no authenticator could be built for the selected provider"
761 )
762}
763
764async fn build_platform_credentials(
783 cfg: &toolkit_security::InternalAuthConfig,
784 cancel: &CancellationToken,
785) -> Result<(
786 toolkit_transport_grpc::InternalAuthInterceptor,
787 Option<toolkit_contract::runtime::config::InternalTokenProvider>,
788)> {
789 use secrecy::SecretString;
790 use toolkit_contract::runtime::config::{CredentialState, InternalTokenProvider};
791 use toolkit_security::InternalAuthConfig;
792 use toolkit_transport_grpc::{
793 DEFAULT_REFRESH_INTERVAL, InternalAuthInterceptor, ServiceAccountTokenReader,
794 };
795
796 match cfg {
797 InternalAuthConfig::SharedSecret { secret, .. } => {
798 let token = SecretString::from(secret.clone());
799 Ok((
800 InternalAuthInterceptor::from_token(token.clone()),
801 Some(InternalTokenProvider::from_token(token)),
802 ))
803 }
804 InternalAuthConfig::Kube {
805 token_path: Some(path),
806 ..
807 } => {
808 let reader = ServiceAccountTokenReader::with_cancellation(
809 path,
810 DEFAULT_REFRESH_INTERVAL,
811 cancel.child_token(),
812 )
813 .await
814 .context("failed to read projected service-account token for outbound credential")?;
815 let interceptor = reader.interceptor();
816 let token_fn = reader.token_provider();
820 let provider = InternalTokenProvider::new(move || match token_fn() {
821 Some(token) => CredentialState::Available(token),
822 None => CredentialState::Unavailable(
823 "projected service-account token is currently unavailable \
824 (file empty or not yet read)"
825 .into(),
826 ),
827 });
828 Ok((interceptor, Some(provider)))
829 }
830 InternalAuthConfig::Kube {
831 token_path: None, ..
832 } => {
833 warn!(
834 "oop_http.internal_auth: provider=kube without token_path - this participant \
835 validates inbound platform tokens but will attach NO outbound credential"
836 );
837 Ok((InternalAuthInterceptor::disabled(), None))
838 }
839 }
840}
841
842fn default_advertise_uri(listen_addr: std::net::SocketAddr) -> String {
847 let host = match listen_addr {
848 std::net::SocketAddr::V4(addr) if addr.ip().is_unspecified() => "127.0.0.1".to_owned(),
849 std::net::SocketAddr::V4(addr) => addr.ip().to_string(),
850 std::net::SocketAddr::V6(addr) if addr.ip().is_unspecified() => "[::1]".to_owned(),
851 std::net::SocketAddr::V6(addr) => format!("[{}]", addr.ip()),
852 };
853 format!("http://{host}:{}", listen_addr.port())
854}
855
856fn validate_advertise_uri(uri: &str, allow_loopback: bool) -> Result<()> {
865 let parsed = url::Url::parse(uri)
866 .with_context(|| format!("invalid oop_http.advertise_uri: not a valid URL: {uri}"))?;
867 if !matches!(parsed.scheme(), "http" | "https") {
868 anyhow::bail!(
869 "invalid oop_http.advertise_uri: scheme must be http or https (got '{}')",
870 parsed.scheme()
871 );
872 }
873 if parsed.host_str().is_none_or(str::is_empty) {
874 anyhow::bail!("invalid oop_http.advertise_uri: missing host: {uri}");
875 }
876 if !parsed.username().is_empty() || parsed.password().is_some() {
877 anyhow::bail!("invalid oop_http.advertise_uri: must not contain userinfo: {uri}");
878 }
879 let is_loopback = match parsed.host() {
880 Some(url::Host::Ipv4(ip)) => ip.is_loopback() || ip.is_unspecified(),
881 Some(url::Host::Ipv6(ip)) => ip.is_loopback() || ip.is_unspecified(),
882 Some(url::Host::Domain(d)) => d.trim_end_matches('.').eq_ignore_ascii_case("localhost"),
883 None => false,
884 };
885 if !allow_loopback && is_loopback {
886 anyhow::bail!(
887 "invalid oop_http.advertise_uri: '{uri}' is a loopback/unspecified address, which is \
888 unreachable by other gears in multi-host Profile 2 / Profile 3 (a registered-but-\
889 unreachable instance). Set oop_http.advertise_uri to a routable host, or set \
890 oop_http.allow_loopback_advertise = true for single-host / local-dev."
891 );
892 }
893 Ok(())
894}
895
896#[allow(unknown_lints, de1301_no_print_macros)] fn print_config(config: &AppConfig) {
898 match config.to_yaml() {
899 Ok(yaml) => {
900 println!("{yaml}");
901 }
902 Err(e) => {
903 eprintln!("Failed to render config as YAML: {e}");
904 }
905 }
906}
907
908#[cfg(test)]
909#[cfg_attr(coverage_nightly, coverage(off))]
910#[path = "oop_tests.rs"]
911mod tests;