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 match cfg.build_authenticator()? {
756 toolkit_security::BuiltAuthenticator::Built(authenticator) => {
757 info!("Initializing shared-secret platform-plane authenticator");
758 return Ok(Some(authenticator));
759 }
760 toolkit_security::BuiltAuthenticator::RequiresExternalBackend => {}
761 }
762
763 #[cfg(feature = "k8s-auth")]
764 {
765 if cfg.is_kube() {
766 info!("Initializing Kubernetes TokenReview platform-plane authenticator");
767 let audiences = cfg.kube_audiences().unwrap_or_default().to_vec();
768 let authenticator = toolkit_k8s_auth::build_cached_k8s_authenticator(
769 audiences,
770 Some(toolkit_security::DEFAULT_TOKEN_REVIEW_CACHE_TTL),
771 None,
772 )
773 .await
774 .context("failed to initialize Kubernetes TokenReview authenticator")?;
775 return Ok(Some(authenticator));
776 }
777 }
778 #[cfg(not(feature = "k8s-auth"))]
779 {
780 if cfg.is_kube() {
781 anyhow::bail!("oop_http.internal_auth provider=kube requires the `k8s-auth` feature");
782 }
783 }
784
785 anyhow::bail!(
786 "internal_auth is configured but no authenticator could be built for the selected provider"
787 )
788}
789
790async fn build_directory_client(
806 directory_endpoint: &str,
807 internal_auth_cfg: Option<&toolkit_security::InternalAuthConfig>,
808 cancel: &CancellationToken,
809) -> Result<(
810 DirectoryGrpcClient,
811 Option<toolkit_contract::runtime::config::InternalTokenProvider>,
812)> {
813 let client = DirectoryGrpcClient::connect_lazy(directory_endpoint)?;
818
819 let Some(cfg) = internal_auth_cfg else {
820 return Ok((client, None));
821 };
822
823 let (interceptor, provider) = build_platform_credentials(cfg, cancel).await?;
824 let client =
827 DirectoryGrpcClient::connect_lazy_with_interceptor(directory_endpoint, interceptor)?;
828 Ok((client, provider))
829}
830
831async fn build_platform_credentials(
850 cfg: &toolkit_security::InternalAuthConfig,
851 cancel: &CancellationToken,
852) -> Result<(
853 toolkit_transport_grpc::InternalAuthInterceptor,
854 Option<toolkit_contract::runtime::config::InternalTokenProvider>,
855)> {
856 use secrecy::SecretString;
857 use toolkit_contract::runtime::config::{CredentialState, InternalTokenProvider};
858 use toolkit_security::InternalAuthConfig;
859 use toolkit_transport_grpc::{
860 DEFAULT_REFRESH_INTERVAL, InternalAuthInterceptor, ServiceAccountTokenReader,
861 };
862
863 match cfg {
864 InternalAuthConfig::SharedSecret { secret, .. } => {
865 let token = SecretString::from(secret.clone());
866 Ok((
867 InternalAuthInterceptor::from_token(token.clone()),
868 Some(InternalTokenProvider::from_token(token)),
869 ))
870 }
871 InternalAuthConfig::Kube {
872 token_path: Some(path),
873 ..
874 } => {
875 let reader = ServiceAccountTokenReader::with_cancellation(
876 path,
877 DEFAULT_REFRESH_INTERVAL,
878 cancel.child_token(),
879 )
880 .await
881 .context("failed to read projected service-account token for outbound credential")?;
882 let interceptor = reader.interceptor();
883 let token_fn = reader.token_provider();
887 let provider = InternalTokenProvider::new(move || match token_fn() {
888 Some(token) => CredentialState::Available(token),
889 None => CredentialState::Unavailable(
890 "projected service-account token is currently unavailable \
891 (file empty or not yet read)"
892 .into(),
893 ),
894 });
895 Ok((interceptor, Some(provider)))
896 }
897 InternalAuthConfig::Kube {
898 token_path: None, ..
899 } => {
900 warn!(
901 "oop_http.internal_auth: provider=kube without token_path - this participant \
902 validates inbound platform tokens but will attach NO outbound credential"
903 );
904 Ok((InternalAuthInterceptor::disabled(), None))
905 }
906 }
907}
908
909fn default_advertise_uri(listen_addr: std::net::SocketAddr) -> String {
914 let host = match listen_addr {
915 std::net::SocketAddr::V4(addr) if addr.ip().is_unspecified() => "127.0.0.1".to_owned(),
916 std::net::SocketAddr::V4(addr) => addr.ip().to_string(),
917 std::net::SocketAddr::V6(addr) if addr.ip().is_unspecified() => "[::1]".to_owned(),
918 std::net::SocketAddr::V6(addr) => format!("[{}]", addr.ip()),
919 };
920 format!("http://{host}:{}", listen_addr.port())
921}
922
923fn validate_advertise_uri(uri: &str, allow_loopback: bool) -> Result<()> {
932 let parsed = url::Url::parse(uri)
933 .with_context(|| format!("invalid oop_http.advertise_uri: not a valid URL: {uri}"))?;
934 if !matches!(parsed.scheme(), "http" | "https") {
935 anyhow::bail!(
936 "invalid oop_http.advertise_uri: scheme must be http or https (got '{}')",
937 parsed.scheme()
938 );
939 }
940 if parsed.host_str().is_none_or(str::is_empty) {
941 anyhow::bail!("invalid oop_http.advertise_uri: missing host: {uri}");
942 }
943 if !parsed.username().is_empty() || parsed.password().is_some() {
944 anyhow::bail!("invalid oop_http.advertise_uri: must not contain userinfo: {uri}");
945 }
946 let is_loopback = match parsed.host() {
947 Some(url::Host::Ipv4(ip)) => ip.is_loopback() || ip.is_unspecified(),
948 Some(url::Host::Ipv6(ip)) => ip.is_loopback() || ip.is_unspecified(),
949 Some(url::Host::Domain(d)) => d.trim_end_matches('.').eq_ignore_ascii_case("localhost"),
950 None => false,
951 };
952 if !allow_loopback && is_loopback {
953 anyhow::bail!(
954 "invalid oop_http.advertise_uri: '{uri}' is a loopback/unspecified address, which is \
955 unreachable by other gears in multi-host Profile 2 / Profile 3 (a registered-but-\
956 unreachable instance). Set oop_http.advertise_uri to a routable host, or set \
957 oop_http.allow_loopback_advertise = true for single-host / local-dev."
958 );
959 }
960 Ok(())
961}
962
963#[allow(unknown_lints, de1301_no_print_macros)] fn print_config(config: &AppConfig) {
965 match config.to_yaml() {
966 Ok(yaml) => {
967 println!("{yaml}");
968 }
969 Err(e) => {
970 eprintln!("Failed to render config as YAML: {e}");
971 }
972 }
973}
974
975#[cfg(test)]
976#[cfg_attr(coverage_nightly, coverage(off))]
977#[path = "oop_tests.rs"]
978mod tests;