1use anyhow::{Context, Result};
49use figment::{Figment, providers::Serialized};
50use std::path::{Path, PathBuf};
51use std::sync::Arc;
52use std::time::Duration;
53use tokio::time::sleep;
54use tokio_util::sync::CancellationToken;
55use tracing::{debug, error, info, warn};
56use uuid::Uuid;
57
58use super::config::{
59 AppConfig, CliArgs, LoggingConfig, RenderedDbConfig, RenderedGearConfig,
60 TOOLKIT_MODULE_CONFIG_ENV,
61};
62use crate::bootstrap::host::{init_logging_unified, init_panic_tracing};
63use crate::runtime::{
64 ClientRegistration, DbOptions, OopServeOptions, RunOptions, ShutdownOptions,
65 TOOLKIT_DIRECTORY_ENDPOINT_ENV, run, run_oop_serving, shutdown,
66};
67use cf_system_sdks::directory::{DirectoryClient, DirectoryGrpcClient};
68
69#[derive(Debug, Clone)]
71pub struct OopRunOptions {
72 pub gear_name: String,
74
75 pub instance_id: Option<Uuid>,
77
78 pub directory_endpoint: String,
80
81 pub config_path: Option<PathBuf>,
83
84 pub verbose: u8,
86
87 pub print_config: bool,
89
90 pub heartbeat_interval_secs: u64,
92
93 pub version: Option<String>,
96}
97
98impl Default for OopRunOptions {
99 fn default() -> Self {
100 let config_path = std::env::var("TOOLKIT_CONFIG_PATH").ok().map(PathBuf::from);
102
103 let directory_endpoint = std::env::var(TOOLKIT_DIRECTORY_ENDPOINT_ENV)
106 .unwrap_or_else(|_| "http://127.0.0.1:50051".to_owned());
107
108 Self {
109 gear_name: String::new(),
110 instance_id: None,
111 directory_endpoint,
112 config_path,
113 verbose: 0,
114 print_config: false,
115 heartbeat_interval_secs: 5,
116 version: None,
117 }
118 }
119}
120
121#[tracing::instrument(
135 level = "debug",
136 skip(local_config, rendered_config),
137 fields(
138 has_rendered = rendered_config.is_some(),
139 has_local_db = local_config.database.is_some()
140 )
141)]
142fn build_oop_config_and_db(
143 local_config: &AppConfig,
144 gear_name: &str,
145 rendered_config: Option<&RenderedGearConfig>,
146) -> Result<(AppConfig, LoggingConfig, DbOptions)> {
147 let home_dir = PathBuf::from(&local_config.server.home_dir);
148
149 let final_config = if let Some(rendered) = rendered_config {
151 let mut config = local_config.clone();
153
154 let gear_entry = config
156 .gears
157 .entry(gear_name.to_owned())
158 .or_insert_with(|| serde_json::json!({}));
159
160 if let Some(obj) = gear_entry.as_object_mut() {
162 if !obj.contains_key("config") || obj["config"].is_null() {
165 obj.insert("config".to_owned(), rendered.config.clone());
166 }
167 }
169
170 debug!(
171 gear = %gear_name,
172 has_rendered_db = %rendered.database.is_some(),
173 has_rendered_logging = %rendered.logging.is_some(),
174 "Using rendered config from master as base, local config as override"
175 );
176
177 config
178 } else {
179 debug!(
181 gear = %gear_name,
182 "No rendered config from master, using local config entirely (standalone mode)"
183 );
184 local_config.clone()
185 };
186
187 let final_logging = merge_logging_configs(
189 rendered_config.as_ref().and_then(|r| r.logging.as_ref()),
190 &local_config.logging,
191 );
192
193 let db_options = build_merged_db_options(
196 &home_dir,
197 gear_name,
198 rendered_config.as_ref().and_then(|r| r.database.as_ref()),
199 local_config,
200 )?;
201
202 Ok((final_config, final_logging, db_options))
203}
204
205fn merge_logging_configs(master: Option<&LoggingConfig>, local: &LoggingConfig) -> LoggingConfig {
210 master
211 .cloned()
212 .unwrap_or_default()
213 .into_iter()
214 .chain(local.clone())
215 .collect()
216}
217
218fn build_merged_db_options(
223 home_dir: &Path,
224 gear_name: &str,
225 rendered_db: Option<&RenderedDbConfig>,
226 local_config: &AppConfig,
227) -> Result<DbOptions> {
228 let has_rendered_db = rendered_db.is_some_and(|db| db.gear.is_some() || db.global.is_some());
230 let has_local_db = local_config.database.is_some()
231 || local_config
232 .gears
233 .get(gear_name)
234 .and_then(|m| m.get("database"))
235 .is_some();
236
237 if !has_rendered_db && !has_local_db {
238 debug!(
239 gear = %gear_name,
240 "No database config available"
241 );
242 return Ok(DbOptions::None);
243 }
244
245 let mut merged_config = serde_json::Map::new();
250
251 if let Some(rendered) = rendered_db {
253 if let Some(ref global) = rendered.global {
255 let global_json = serde_json::to_value(global)
256 .context("Failed to serialize rendered global db config")?;
257 merged_config.insert("database".to_owned(), global_json);
258 }
259
260 if let Some(ref gear_db) = rendered.gear {
262 let gear_db_json = serde_json::to_value(gear_db)
263 .context("Failed to serialize rendered gear db config")?;
264
265 let mut gears = serde_json::Map::new();
266 let mut gear_entry = serde_json::Map::new();
267 gear_entry.insert("database".to_owned(), gear_db_json);
268 gears.insert(gear_name.to_owned(), serde_json::Value::Object(gear_entry));
269 merged_config.insert("gears".to_owned(), serde_json::Value::Object(gears));
270 }
271 }
272
273 if let Some(ref local_db) = local_config.database {
276 let local_db_json =
277 serde_json::to_value(local_db).context("Failed to serialize local global db config")?;
278
279 if let Some(existing) = merged_config.get_mut("database") {
281 merge_json_objects(existing, &local_db_json);
282 } else {
283 merged_config.insert("database".to_owned(), local_db_json);
284 }
285 }
286
287 if let Some(local_gear) = local_config.gears.get(gear_name)
289 && let Some(local_gear_db) = local_gear.get("database")
290 {
291 let gears = merged_config
292 .entry("gears".to_owned())
293 .or_insert_with(|| serde_json::Value::Object(serde_json::Map::new()));
294
295 if let Some(gears_obj) = gears.as_object_mut() {
296 let gear_entry = gears_obj
297 .entry(gear_name.to_owned())
298 .or_insert_with(|| serde_json::Value::Object(serde_json::Map::new()));
299
300 if let Some(gear_obj) = gear_entry.as_object_mut() {
301 if let Some(existing_db) = gear_obj.get_mut("database") {
302 merge_json_objects(existing_db, local_gear_db);
303 } else {
304 gear_obj.insert("database".to_owned(), local_gear_db.clone());
305 }
306 }
307 }
308 }
309
310 debug!(
311 gear = %gear_name,
312 has_rendered = %rendered_db.is_some(),
313 has_local_global = %local_config.database.is_some(),
314 "Building DbManager with merged config"
315 );
316
317 let figment = Figment::new().merge(Serialized::defaults(serde_json::Value::Object(
319 merged_config,
320 )));
321 let db_manager = Arc::new(
322 toolkit_db::DbManager::from_figment(figment, home_dir.to_path_buf())
323 .context("Failed to create DbManager from merged config")?,
324 );
325
326 Ok(DbOptions::Manager(db_manager))
327}
328
329fn merge_json_objects(target: &mut serde_json::Value, source: &serde_json::Value) {
332 if let (Some(target_obj), Some(source_obj)) = (target.as_object_mut(), source.as_object()) {
333 for (key, value) in source_obj {
334 if let Some(target_value) = target_obj.get_mut(key) {
335 if target_value.is_object() && value.is_object() {
337 merge_json_objects(target_value, value);
338 } else {
339 *target_value = value.clone();
340 }
341 } else {
342 target_obj.insert(key.clone(), value.clone());
343 }
344 }
345 } else {
346 *target = source.clone();
348 }
349}
350
351#[tracing::instrument(
404 level = "info",
405 name = "oop_bootstrap",
406 skip(opts),
407 fields(
408 gear = %opts.gear_name,
409 directory = %opts.directory_endpoint
410 )
411)]
412pub async fn run_oop_with_options(opts: OopRunOptions) -> Result<()> {
413 let instance_id = opts.instance_id.unwrap_or_else(Uuid::new_v4);
415
416 let cancel = CancellationToken::new();
419
420 let cancel_for_signals = cancel.clone();
423 tokio::spawn(async move {
424 match shutdown::wait_for_shutdown().await {
425 Ok(()) => {
426 info!(target: "", "------------------");
427 info!("shutdown: signal received in OoP bootstrap");
428 }
429 Err(e) => {
430 warn!(
431 error = %e,
432 "shutdown: primary waiter failed in OoP bootstrap, falling back to ctrl_c()"
433 );
434 _ = tokio::signal::ctrl_c().await;
435 }
436 }
437 cancel_for_signals.cancel();
438 });
439
440 let args = CliArgs {
442 config: opts
443 .config_path
444 .as_ref()
445 .map(|p| p.to_string_lossy().to_string()),
446 print_config: opts.print_config,
447 verbose: opts.verbose,
448 mock: false,
449 };
450
451 let mut config = AppConfig::load_or_default(opts.config_path.as_ref())?;
453 config.apply_cli_overrides(args.verbose);
454
455 let rendered_config = match std::env::var(TOOLKIT_MODULE_CONFIG_ENV) {
458 Ok(json) => RenderedGearConfig::from_json(&json).ok(),
459 Err(_) => None,
460 };
461
462 let (final_config, merged_logging, db_options) =
467 build_oop_config_and_db(&config, &opts.gear_name, rendered_config.as_ref())?;
468
469 #[cfg(feature = "otel")]
473 let otel_cfg = rendered_config
474 .as_ref()
475 .and_then(|rc| rc.opentelemetry.as_ref());
476
477 #[cfg(feature = "otel")]
479 let otel_layer = otel_cfg
480 .filter(|cfg| cfg.tracing.enabled)
481 .map(crate::telemetry::init_tracing)
482 .transpose()?;
483 #[cfg(not(feature = "otel"))]
484 let otel_layer = None;
485
486 #[cfg(feature = "otel")]
489 let metrics_init_error = otel_cfg
490 .filter(|cfg| cfg.metrics.enabled)
491 .and_then(|cfg| crate::telemetry::init::init_metrics_provider(cfg).err());
492
493 #[cfg(feature = "otel")]
497 let inject_trace_ids =
498 otel_cfg.is_some_and(crate::telemetry::OpenTelemetryConfig::inject_trace_ids_into_logs);
499 #[cfg(not(feature = "otel"))]
500 let inject_trace_ids = false;
501
502 init_logging_unified(
503 &merged_logging,
504 &config.server.home_dir,
505 otel_layer,
506 inject_trace_ids,
507 );
508
509 #[cfg(feature = "otel")]
511 if let Some(e) = metrics_init_error {
512 tracing::error!(error = %e, "OpenTelemetry metrics not initialized (OoP)");
513 }
514
515 init_panic_tracing();
517
518 if let Some(ref rc) = rendered_config {
520 info!(
521 env_var = TOOLKIT_MODULE_CONFIG_ENV,
522 has_database = rc.database.is_some(),
523 has_config = !rc.config.is_null(),
524 has_logging = rc.logging.is_some(),
525 has_opentelemetry = rc.opentelemetry.is_some(),
526 "Received rendered config from master host"
527 );
528 } else if std::env::var(TOOLKIT_MODULE_CONFIG_ENV).is_ok() {
529 warn!(
530 env_var = TOOLKIT_MODULE_CONFIG_ENV,
531 "Failed to parse rendered config from master host, using local config only"
532 );
533 } else {
534 debug!(
535 env_var = TOOLKIT_MODULE_CONFIG_ENV,
536 "No rendered config from master host, using local config only"
537 );
538 }
539
540 info!(
541 gear = %opts.gear_name,
542 instance_id = %instance_id,
543 directory_endpoint = %opts.directory_endpoint,
544 "OoP gear bootstrap starting"
545 );
546
547 if opts.print_config {
549 print_config(&config);
550 return Ok(());
551 }
552
553 info!(
570 "Creating directory service client (lazy connect) for {}",
571 opts.directory_endpoint
572 );
573 let internal_auth_cfg = final_config
574 .oop_http
575 .as_ref()
576 .and_then(|h| h.internal_auth.as_ref());
577 let (directory_client, internal_token_provider) =
578 build_directory_client(&opts.directory_endpoint, internal_auth_cfg, &cancel).await?;
579 let directory_api: Arc<dyn DirectoryClient> = Arc::new(directory_client);
580
581 info!("Directory service client ready (will connect on first use)");
582
583 let oop_http = final_config.oop_http.clone();
585
586 let config_provider = Arc::new(final_config);
588
589 let run_options = RunOptions::new(
592 config_provider,
593 db_options,
594 ShutdownOptions::Token(cancel.clone()),
595 instance_id,
596 )
597 .with_clients(vec![ClientRegistration::new::<dyn DirectoryClient>(
598 Arc::clone(&directory_api),
599 )])
600 .with_internal_token_provider(internal_token_provider);
601
602 let result = if let Some(http_cfg) = oop_http {
607 info!("Starting OoP HTTP-serving lifecycle");
608 let serve = build_oop_serve_options(
609 &http_cfg,
610 &opts.gear_name,
611 instance_id,
612 opts.version.clone(),
613 Duration::from_secs(opts.heartbeat_interval_secs),
614 Arc::clone(&directory_api),
615 )
616 .await?;
617 run_oop_serving(run_options, serve).await
618 } else {
619 info!("Starting gear lifecycle (legacy gRPC-only)");
620 let heartbeat_directory = Arc::clone(&directory_api);
624 let heartbeat_gear = opts.gear_name.clone();
625 let heartbeat_instance_id_str = instance_id.to_string();
626 let heartbeat_interval = Duration::from_secs(opts.heartbeat_interval_secs.max(1));
627 let heartbeat_cancel = cancel.child_token();
628 tokio::spawn(async move {
629 info!(interval_secs = ?heartbeat_interval, "Starting legacy heartbeat loop");
630 loop {
631 tokio::select! {
632 () = heartbeat_cancel.cancelled() => {
633 info!("Heartbeat loop stopping due to cancellation");
634 break;
635 }
636 () = sleep(heartbeat_interval) => {
637 if let Err(e) = heartbeat_directory
638 .send_heartbeat(&heartbeat_gear, &heartbeat_instance_id_str)
639 .await
640 {
641 warn!(error = %e, "Failed to send heartbeat, will retry");
642 }
643 }
644 }
645 }
646 });
647 run(run_options).await
648 };
649
650 if let Err(ref e) = result {
651 error!(error = %e, "Gear runtime failed");
652 } else {
653 info!("Gear runtime completed successfully");
654 }
655
656 #[cfg(feature = "otel")]
660 crate::bootstrap::run::tracing_shutdown().await;
661
662 result
663}
664
665async fn build_oop_serve_options(
672 cfg: &super::config::OopHttpConfig,
673 gear_name: &str,
674 instance_id: Uuid,
675 version: Option<String>,
676 heartbeat_interval: Duration,
677 directory: Arc<dyn DirectoryClient>,
678) -> Result<OopServeOptions> {
679 let listen_addr: std::net::SocketAddr = cfg
680 .listen_addr
681 .parse()
682 .with_context(|| format!("invalid oop_http.listen_addr: {}", cfg.listen_addr))?;
683
684 let probe_bind_addr = cfg
685 .probe_bind_addr
686 .as_deref()
687 .map(|s| {
688 s.parse::<std::net::SocketAddr>()
689 .with_context(|| format!("invalid oop_http.probe_bind_addr: {s}"))
690 })
691 .transpose()?;
692
693 let advertise_uri = cfg
694 .advertise_uri
695 .clone()
696 .unwrap_or_else(|| default_advertise_uri(listen_addr));
697
698 validate_advertise_uri(&advertise_uri, cfg.allow_loopback_advertise)?;
702
703 cf_system_sdks::directory::validate_labels(&cfg.labels).with_context(|| {
710 "invalid oop_http.labels: label keys/values must be <=63 chars, <=64 entries, and use \
711 only ASCII alphanumerics plus '-', '_', '.' (starting and ending alphanumeric)"
712 })?;
713
714 let internal_authenticator = build_internal_authenticator(cfg.internal_auth.as_ref()).await?;
715
716 Ok(OopServeOptions {
717 gear_name: gear_name.to_owned(),
718 instance_id: instance_id.to_string(),
719 version,
720 advertise_uri,
721 listen_addr,
722 probe_bind_addr,
723 drain_timeout: Duration::from_secs(cfg.drain_timeout_secs),
724 heartbeat_interval,
725 healthcheck_timeout: Duration::from_millis(cfg.healthcheck_timeout_ms),
726 directory,
727 bearer_authenticator: None,
728 internal_authenticator,
729 labels: cfg.labels.clone(),
730 })
731}
732
733#[cfg_attr(not(feature = "k8s-auth"), allow(clippy::unused_async))]
745async fn build_internal_authenticator(
746 cfg: Option<&toolkit_security::InternalAuthConfig>,
747) -> Result<Option<toolkit_security::DynInternalAuthenticator>> {
748 let Some(cfg) = cfg else {
749 return Ok(None);
750 };
751
752 if let Some(authenticator) = cfg.build_authenticator() {
754 info!("Initializing shared-secret platform-plane authenticator");
755 return Ok(Some(authenticator));
756 }
757
758 #[cfg(feature = "k8s-auth")]
759 {
760 if cfg.is_kube() {
761 info!("Initializing Kubernetes TokenReview platform-plane authenticator");
762 let audiences = cfg.kube_audiences().unwrap_or_default().to_vec();
763 let authenticator = toolkit_k8s_auth::build_cached_k8s_authenticator(
764 audiences,
765 Some(toolkit_security::DEFAULT_TOKEN_REVIEW_CACHE_TTL),
766 )
767 .await
768 .context("failed to initialize Kubernetes TokenReview authenticator")?;
769 return Ok(Some(authenticator));
770 }
771 }
772 #[cfg(not(feature = "k8s-auth"))]
773 {
774 if cfg.is_kube() {
775 anyhow::bail!("oop_http.internal_auth provider=kube requires the `k8s-auth` feature");
776 }
777 }
778
779 anyhow::bail!(
780 "internal_auth is configured but no authenticator could be built for the selected provider"
781 )
782}
783
784async fn build_directory_client(
800 directory_endpoint: &str,
801 internal_auth_cfg: Option<&toolkit_security::InternalAuthConfig>,
802 cancel: &CancellationToken,
803) -> Result<(
804 DirectoryGrpcClient,
805 Option<toolkit_contract::runtime::config::InternalTokenProvider>,
806)> {
807 let client = DirectoryGrpcClient::connect_lazy(directory_endpoint)?;
812
813 let Some(cfg) = internal_auth_cfg else {
814 return Ok((client, None));
815 };
816
817 let (interceptor, provider) = build_platform_credentials(cfg, cancel).await?;
818 let client =
821 DirectoryGrpcClient::connect_lazy_with_interceptor(directory_endpoint, interceptor)?;
822 Ok((client, provider))
823}
824
825async fn build_platform_credentials(
844 cfg: &toolkit_security::InternalAuthConfig,
845 cancel: &CancellationToken,
846) -> Result<(
847 toolkit_transport_grpc::InternalAuthInterceptor,
848 Option<toolkit_contract::runtime::config::InternalTokenProvider>,
849)> {
850 use secrecy::SecretString;
851 use toolkit_contract::runtime::config::{CredentialState, InternalTokenProvider};
852 use toolkit_security::InternalAuthConfig;
853 use toolkit_transport_grpc::{
854 DEFAULT_REFRESH_INTERVAL, InternalAuthInterceptor, ServiceAccountTokenReader,
855 };
856
857 match cfg {
858 InternalAuthConfig::SharedSecret { secret, .. } => {
859 let token = SecretString::from(secret.clone());
860 Ok((
861 InternalAuthInterceptor::from_token(token.clone()),
862 Some(InternalTokenProvider::from_token(token)),
863 ))
864 }
865 InternalAuthConfig::Kube {
866 token_path: Some(path),
867 ..
868 } => {
869 let reader = ServiceAccountTokenReader::with_cancellation(
870 path,
871 DEFAULT_REFRESH_INTERVAL,
872 cancel.child_token(),
873 )
874 .await
875 .context("failed to read projected service-account token for outbound credential")?;
876 let interceptor = reader.interceptor();
877 let token_fn = reader.token_provider();
881 let provider = InternalTokenProvider::new(move || match token_fn() {
882 Some(token) => CredentialState::Available(token),
883 None => CredentialState::Unavailable(
884 "projected service-account token is currently unavailable \
885 (file empty or not yet read)"
886 .into(),
887 ),
888 });
889 Ok((interceptor, Some(provider)))
890 }
891 InternalAuthConfig::Kube {
892 token_path: None, ..
893 } => {
894 warn!(
895 "oop_http.internal_auth: provider=kube without token_path - this participant \
896 validates inbound platform tokens but will attach NO outbound credential"
897 );
898 Ok((InternalAuthInterceptor::disabled(), None))
899 }
900 }
901}
902
903fn default_advertise_uri(listen_addr: std::net::SocketAddr) -> String {
908 let host = match listen_addr {
909 std::net::SocketAddr::V4(addr) if addr.ip().is_unspecified() => "127.0.0.1".to_owned(),
910 std::net::SocketAddr::V4(addr) => addr.ip().to_string(),
911 std::net::SocketAddr::V6(addr) if addr.ip().is_unspecified() => "[::1]".to_owned(),
912 std::net::SocketAddr::V6(addr) => format!("[{}]", addr.ip()),
913 };
914 format!("http://{host}:{}", listen_addr.port())
915}
916
917fn validate_advertise_uri(uri: &str, allow_loopback: bool) -> Result<()> {
926 let parsed = url::Url::parse(uri)
927 .with_context(|| format!("invalid oop_http.advertise_uri: not a valid URL: {uri}"))?;
928 if !matches!(parsed.scheme(), "http" | "https") {
929 anyhow::bail!(
930 "invalid oop_http.advertise_uri: scheme must be http or https (got '{}')",
931 parsed.scheme()
932 );
933 }
934 if parsed.host_str().is_none_or(str::is_empty) {
935 anyhow::bail!("invalid oop_http.advertise_uri: missing host: {uri}");
936 }
937 if !parsed.username().is_empty() || parsed.password().is_some() {
938 anyhow::bail!("invalid oop_http.advertise_uri: must not contain userinfo: {uri}");
939 }
940 let is_loopback = match parsed.host() {
941 Some(url::Host::Ipv4(ip)) => ip.is_loopback() || ip.is_unspecified(),
942 Some(url::Host::Ipv6(ip)) => ip.is_loopback() || ip.is_unspecified(),
943 Some(url::Host::Domain(d)) => d.trim_end_matches('.').eq_ignore_ascii_case("localhost"),
944 None => false,
945 };
946 if !allow_loopback && is_loopback {
947 anyhow::bail!(
948 "invalid oop_http.advertise_uri: '{uri}' is a loopback/unspecified address, which is \
949 unreachable by other gears in multi-host Profile 2 / Profile 3 (a registered-but-\
950 unreachable instance). Set oop_http.advertise_uri to a routable host, or set \
951 oop_http.allow_loopback_advertise = true for single-host / local-dev."
952 );
953 }
954 Ok(())
955}
956
957#[allow(unknown_lints, de1301_no_print_macros)] fn print_config(config: &AppConfig) {
959 match config.to_yaml() {
960 Ok(yaml) => {
961 println!("{yaml}");
962 }
963 Err(e) => {
964 eprintln!("Failed to render config as YAML: {e}");
965 }
966 }
967}
968
969#[cfg(test)]
970#[cfg_attr(coverage_nightly, coverage(off))]
971#[path = "oop_tests.rs"]
972mod tests;