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 init_logging_unified(&merged_logging, &config.server.home_dir, otel_layer);
495
496 #[cfg(feature = "otel")]
498 if let Some(e) = metrics_init_error {
499 tracing::error!(error = %e, "OpenTelemetry metrics not initialized (OoP)");
500 }
501
502 init_panic_tracing();
504
505 if let Some(ref rc) = rendered_config {
507 info!(
508 env_var = TOOLKIT_MODULE_CONFIG_ENV,
509 has_database = rc.database.is_some(),
510 has_config = !rc.config.is_null(),
511 has_logging = rc.logging.is_some(),
512 has_opentelemetry = rc.opentelemetry.is_some(),
513 "Received rendered config from master host"
514 );
515 } else if std::env::var(TOOLKIT_MODULE_CONFIG_ENV).is_ok() {
516 warn!(
517 env_var = TOOLKIT_MODULE_CONFIG_ENV,
518 "Failed to parse rendered config from master host, using local config only"
519 );
520 } else {
521 debug!(
522 env_var = TOOLKIT_MODULE_CONFIG_ENV,
523 "No rendered config from master host, using local config only"
524 );
525 }
526
527 info!(
528 gear = %opts.gear_name,
529 instance_id = %instance_id,
530 directory_endpoint = %opts.directory_endpoint,
531 "OoP gear bootstrap starting"
532 );
533
534 if opts.print_config {
536 print_config(&config);
537 return Ok(());
538 }
539
540 info!(
557 "Creating directory service client (lazy connect) for {}",
558 opts.directory_endpoint
559 );
560 let internal_auth_cfg = final_config
561 .oop_http
562 .as_ref()
563 .and_then(|h| h.internal_auth.as_ref());
564 let (directory_client, internal_token_provider) =
565 build_directory_client(&opts.directory_endpoint, internal_auth_cfg, &cancel).await?;
566 let directory_api: Arc<dyn DirectoryClient> = Arc::new(directory_client);
567
568 info!("Directory service client ready (will connect on first use)");
569
570 let oop_http = final_config.oop_http.clone();
572
573 let config_provider = Arc::new(final_config);
575
576 let run_options = RunOptions::new(
579 config_provider,
580 db_options,
581 ShutdownOptions::Token(cancel.clone()),
582 instance_id,
583 )
584 .with_clients(vec![ClientRegistration::new::<dyn DirectoryClient>(
585 Arc::clone(&directory_api),
586 )])
587 .with_internal_token_provider(internal_token_provider);
588
589 let result = if let Some(http_cfg) = oop_http {
594 info!("Starting OoP HTTP-serving lifecycle");
595 let serve = build_oop_serve_options(
596 &http_cfg,
597 &opts.gear_name,
598 instance_id,
599 opts.version.clone(),
600 Duration::from_secs(opts.heartbeat_interval_secs),
601 Arc::clone(&directory_api),
602 )
603 .await?;
604 run_oop_serving(run_options, serve).await
605 } else {
606 info!("Starting gear lifecycle (legacy gRPC-only)");
607 let heartbeat_directory = Arc::clone(&directory_api);
611 let heartbeat_gear = opts.gear_name.clone();
612 let heartbeat_instance_id_str = instance_id.to_string();
613 let heartbeat_interval = Duration::from_secs(opts.heartbeat_interval_secs.max(1));
614 let heartbeat_cancel = cancel.child_token();
615 tokio::spawn(async move {
616 info!(interval_secs = ?heartbeat_interval, "Starting legacy heartbeat loop");
617 loop {
618 tokio::select! {
619 () = heartbeat_cancel.cancelled() => {
620 info!("Heartbeat loop stopping due to cancellation");
621 break;
622 }
623 () = sleep(heartbeat_interval) => {
624 if let Err(e) = heartbeat_directory
625 .send_heartbeat(&heartbeat_gear, &heartbeat_instance_id_str)
626 .await
627 {
628 warn!(error = %e, "Failed to send heartbeat, will retry");
629 }
630 }
631 }
632 }
633 });
634 run(run_options).await
635 };
636
637 if let Err(ref e) = result {
638 error!(error = %e, "Gear runtime failed");
639 } else {
640 info!("Gear runtime completed successfully");
641 }
642
643 result
644}
645
646async fn build_oop_serve_options(
653 cfg: &super::config::OopHttpConfig,
654 gear_name: &str,
655 instance_id: Uuid,
656 version: Option<String>,
657 heartbeat_interval: Duration,
658 directory: Arc<dyn DirectoryClient>,
659) -> Result<OopServeOptions> {
660 let listen_addr: std::net::SocketAddr = cfg
661 .listen_addr
662 .parse()
663 .with_context(|| format!("invalid oop_http.listen_addr: {}", cfg.listen_addr))?;
664
665 let probe_bind_addr = cfg
666 .probe_bind_addr
667 .as_deref()
668 .map(|s| {
669 s.parse::<std::net::SocketAddr>()
670 .with_context(|| format!("invalid oop_http.probe_bind_addr: {s}"))
671 })
672 .transpose()?;
673
674 let advertise_uri = cfg
675 .advertise_uri
676 .clone()
677 .unwrap_or_else(|| default_advertise_uri(listen_addr));
678
679 validate_advertise_uri(&advertise_uri, cfg.allow_loopback_advertise)?;
683
684 cf_system_sdks::directory::validate_labels(&cfg.labels).with_context(|| {
691 "invalid oop_http.labels: label keys/values must be <=63 chars, <=64 entries, and use \
692 only ASCII alphanumerics plus '-', '_', '.' (starting and ending alphanumeric)"
693 })?;
694
695 let internal_authenticator = build_internal_authenticator(cfg.internal_auth.as_ref()).await?;
696
697 Ok(OopServeOptions {
698 gear_name: gear_name.to_owned(),
699 instance_id: instance_id.to_string(),
700 version,
701 advertise_uri,
702 listen_addr,
703 probe_bind_addr,
704 drain_timeout: Duration::from_secs(cfg.drain_timeout_secs),
705 heartbeat_interval,
706 healthcheck_timeout: Duration::from_millis(cfg.healthcheck_timeout_ms),
707 directory,
708 bearer_authenticator: None,
709 internal_authenticator,
710 labels: cfg.labels.clone(),
711 })
712}
713
714#[cfg_attr(not(feature = "k8s-auth"), allow(clippy::unused_async))]
726async fn build_internal_authenticator(
727 cfg: Option<&toolkit_security::InternalAuthConfig>,
728) -> Result<Option<toolkit_security::DynInternalAuthenticator>> {
729 let Some(cfg) = cfg else {
730 return Ok(None);
731 };
732
733 if let Some(authenticator) = cfg.build_authenticator() {
735 info!("Initializing shared-secret platform-plane authenticator");
736 return Ok(Some(authenticator));
737 }
738
739 #[cfg(feature = "k8s-auth")]
740 {
741 if cfg.is_kube() {
742 info!("Initializing Kubernetes TokenReview platform-plane authenticator");
743 let audiences = cfg.kube_audiences().unwrap_or_default().to_vec();
744 let authenticator = toolkit_k8s_auth::build_cached_k8s_authenticator(
745 audiences,
746 Some(toolkit_security::DEFAULT_TOKEN_REVIEW_CACHE_TTL),
747 )
748 .await
749 .context("failed to initialize Kubernetes TokenReview authenticator")?;
750 return Ok(Some(authenticator));
751 }
752 }
753 #[cfg(not(feature = "k8s-auth"))]
754 {
755 if cfg.is_kube() {
756 anyhow::bail!("oop_http.internal_auth provider=kube requires the `k8s-auth` feature");
757 }
758 }
759
760 anyhow::bail!(
761 "internal_auth is configured but no authenticator could be built for the selected provider"
762 )
763}
764
765async fn build_directory_client(
781 directory_endpoint: &str,
782 internal_auth_cfg: Option<&toolkit_security::InternalAuthConfig>,
783 cancel: &CancellationToken,
784) -> Result<(
785 DirectoryGrpcClient,
786 Option<toolkit_contract::runtime::config::InternalTokenProvider>,
787)> {
788 let client = DirectoryGrpcClient::connect_lazy(directory_endpoint)?;
793
794 let Some(cfg) = internal_auth_cfg else {
795 return Ok((client, None));
796 };
797
798 let (interceptor, provider) = build_platform_credentials(cfg, cancel).await?;
799 let client =
802 DirectoryGrpcClient::connect_lazy_with_interceptor(directory_endpoint, interceptor)?;
803 Ok((client, provider))
804}
805
806async fn build_platform_credentials(
825 cfg: &toolkit_security::InternalAuthConfig,
826 cancel: &CancellationToken,
827) -> Result<(
828 toolkit_transport_grpc::InternalAuthInterceptor,
829 Option<toolkit_contract::runtime::config::InternalTokenProvider>,
830)> {
831 use secrecy::SecretString;
832 use toolkit_contract::runtime::config::{CredentialState, InternalTokenProvider};
833 use toolkit_security::InternalAuthConfig;
834 use toolkit_transport_grpc::{
835 DEFAULT_REFRESH_INTERVAL, InternalAuthInterceptor, ServiceAccountTokenReader,
836 };
837
838 match cfg {
839 InternalAuthConfig::SharedSecret { secret, .. } => {
840 let token = SecretString::from(secret.clone());
841 Ok((
842 InternalAuthInterceptor::from_token(token.clone()),
843 Some(InternalTokenProvider::from_token(token)),
844 ))
845 }
846 InternalAuthConfig::Kube {
847 token_path: Some(path),
848 ..
849 } => {
850 let reader = ServiceAccountTokenReader::with_cancellation(
851 path,
852 DEFAULT_REFRESH_INTERVAL,
853 cancel.child_token(),
854 )
855 .await
856 .context("failed to read projected service-account token for outbound credential")?;
857 let interceptor = reader.interceptor();
858 let token_fn = reader.token_provider();
862 let provider = InternalTokenProvider::new(move || match token_fn() {
863 Some(token) => CredentialState::Available(token),
864 None => CredentialState::Unavailable(
865 "projected service-account token is currently unavailable \
866 (file empty or not yet read)"
867 .into(),
868 ),
869 });
870 Ok((interceptor, Some(provider)))
871 }
872 InternalAuthConfig::Kube {
873 token_path: None, ..
874 } => {
875 warn!(
876 "oop_http.internal_auth: provider=kube without token_path - this participant \
877 validates inbound platform tokens but will attach NO outbound credential"
878 );
879 Ok((InternalAuthInterceptor::disabled(), None))
880 }
881 }
882}
883
884fn default_advertise_uri(listen_addr: std::net::SocketAddr) -> String {
889 let host = match listen_addr {
890 std::net::SocketAddr::V4(addr) if addr.ip().is_unspecified() => "127.0.0.1".to_owned(),
891 std::net::SocketAddr::V4(addr) => addr.ip().to_string(),
892 std::net::SocketAddr::V6(addr) if addr.ip().is_unspecified() => "[::1]".to_owned(),
893 std::net::SocketAddr::V6(addr) => format!("[{}]", addr.ip()),
894 };
895 format!("http://{host}:{}", listen_addr.port())
896}
897
898fn validate_advertise_uri(uri: &str, allow_loopback: bool) -> Result<()> {
907 let parsed = url::Url::parse(uri)
908 .with_context(|| format!("invalid oop_http.advertise_uri: not a valid URL: {uri}"))?;
909 if !matches!(parsed.scheme(), "http" | "https") {
910 anyhow::bail!(
911 "invalid oop_http.advertise_uri: scheme must be http or https (got '{}')",
912 parsed.scheme()
913 );
914 }
915 if parsed.host_str().is_none_or(str::is_empty) {
916 anyhow::bail!("invalid oop_http.advertise_uri: missing host: {uri}");
917 }
918 if !parsed.username().is_empty() || parsed.password().is_some() {
919 anyhow::bail!("invalid oop_http.advertise_uri: must not contain userinfo: {uri}");
920 }
921 let is_loopback = match parsed.host() {
922 Some(url::Host::Ipv4(ip)) => ip.is_loopback() || ip.is_unspecified(),
923 Some(url::Host::Ipv6(ip)) => ip.is_loopback() || ip.is_unspecified(),
924 Some(url::Host::Domain(d)) => d.trim_end_matches('.').eq_ignore_ascii_case("localhost"),
925 None => false,
926 };
927 if !allow_loopback && is_loopback {
928 anyhow::bail!(
929 "invalid oop_http.advertise_uri: '{uri}' is a loopback/unspecified address, which is \
930 unreachable by other gears in multi-host Profile 2 / Profile 3 (a registered-but-\
931 unreachable instance). Set oop_http.advertise_uri to a routable host, or set \
932 oop_http.allow_loopback_advertise = true for single-host / local-dev."
933 );
934 }
935 Ok(())
936}
937
938#[allow(unknown_lints, de1301_no_print_macros)] fn print_config(config: &AppConfig) {
940 match config.to_yaml() {
941 Ok(yaml) => {
942 println!("{yaml}");
943 }
944 Err(e) => {
945 eprintln!("Failed to render config as YAML: {e}");
946 }
947 }
948}
949
950#[cfg(test)]
951#[cfg_attr(coverage_nightly, coverage(off))]
952#[path = "oop_tests.rs"]
953mod tests;