process_v1.rs (28797B)
1 //! Binary-owned execution for one already-admitted command invocation. 2 3 use std::env; 4 use std::io::{Read, Write}; 5 use std::num::NonZeroU64; 6 use std::path::{Path, PathBuf}; 7 8 #[cfg(any(target_os = "linux", target_os = "macos"))] 9 use radroots_service_host::{ 10 AdminClient, AdminClientErrorKind, AdminClientTarget, AdminOperationId, 11 }; 12 use radroots_service_host::{ 13 ContractVersions, EntropySource, SystemEntropy, SystemWallClock, WallClock, 14 }; 15 use radroots_service_sqlite::{ 16 BACKUP_MANIFEST_CANONICAL_MAX_BYTES, BackupManifestSha256, IntegrityCheckedAtUnixMs, 17 MigrationAppliedAtUnixSeconds, MigrationBuildIdentity, 18 }; 19 use radroots_storage::event::SourceGeneration; 20 use serde_json::{Value, json}; 21 22 #[cfg(any(target_os = "linux", target_os = "macos"))] 23 use crate::admin_v1::{admin_transport_limits, admit_admin_response_value}; 24 use crate::cli_bootstrap::read_identity_provisioning_document; 25 use crate::{ 26 RHI_ADMIN_CONTRACT_VERSION, RHI_CONFIG_DOCUMENT_MAX_UTF8_BYTES, RHI_CONFIG_SCHEMA_VERSION, 27 RHI_PROVIDER_CONTRACT_VERSION, RHI_STATE_SCHEMA_VERSION, RHI_STATUS_CONTRACT_VERSION, 28 RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, RhiCliInvocationV1, 29 RhiCliOutputModeV1, RhiCliPrimaryAuthorityV1, RhiCommandV1, RhiConfigCommandV1, 30 RhiIdentityCommandV1, RhiProcessResult, RhiRuntimeContext, RhiStateCommandV1, RhiStateMetadata, 31 apply_rhi_configuration, finalize_rhi_state_restore, initialize_rhi_config_document, 32 initialize_rhi_state, load_rhi_config_candidate, load_rhi_config_document, 33 open_rhi_state_read_write_from_config, plan_rhi_cli_v1, provision_rhi_encrypted_identity, 34 resolve_rhi_runtime_context, resolve_rhi_wrapping_credential, stage_rhi_state_restore, 35 verify_rhi_state_backup, 36 }; 37 #[cfg(any(target_os = "linux", target_os = "macos"))] 38 use crate::{ 39 RhiMetricsCommandV1, RhiPresenceCommandV1, RhiPublicationCommandV1, RhiReconciliationCommandV1, 40 RhiSourcesCommandV1, RhiTradeCommandV1, 41 }; 42 43 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 44 struct ProcessFailure(RhiProcessResult); 45 46 type ProcessResult<T> = Result<T, ProcessFailure>; 47 48 /// Executes one admitted RHI invocation without reparsing process arguments. 49 #[must_use] 50 pub fn execute_rhi_cli_v1(invocation: RhiCliInvocationV1) -> RhiProcessResult { 51 execute(invocation).unwrap_or_else(|failure| failure.0) 52 } 53 54 /// Executes one admitted invocation with a binary-owned process-signal source. 55 #[must_use] 56 pub fn execute_rhi_cli_v1_with_signal_source<F, S>( 57 invocation: RhiCliInvocationV1, 58 make_signal_source: F, 59 ) -> RhiProcessResult 60 where 61 F: FnOnce() -> Option<S>, 62 S: crate::RhiProcessSignalSource + 'static, 63 { 64 if !matches!(invocation.command(), RhiCommandV1::Run) { 65 return execute_rhi_cli_v1(invocation); 66 } 67 execute_run(invocation, make_signal_source).unwrap_or_else(|failure| failure.0) 68 } 69 70 #[cfg(any(target_os = "linux", target_os = "macos"))] 71 fn execute_run<F, S>( 72 invocation: RhiCliInvocationV1, 73 make_signal_source: F, 74 ) -> ProcessResult<RhiProcessResult> 75 where 76 F: FnOnce() -> Option<S>, 77 S: crate::RhiProcessSignalSource + 'static, 78 { 79 let runtime = resolve_runtime(&invocation)?; 80 let configuration = load_rhi_config_document(&runtime).map_err(|_| input_failure())?; 81 let applied_at = migration_time()?; 82 let build = migration_build_identity()?; 83 let tokio = build_tokio_runtime(configuration.runtime_thread_limits())?; 84 tokio.block_on(async move { 85 let signals = make_signal_source().ok_or(ProcessFailure( 86 RhiProcessResult::ServiceOrDependencyUnavailable, 87 ))?; 88 Ok(crate::runtime_graph::run_rhi_daemon( 89 runtime, 90 configuration, 91 applied_at, 92 &build, 93 signals, 94 ) 95 .await) 96 }) 97 } 98 99 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 100 fn execute_run<F, S>( 101 _invocation: RhiCliInvocationV1, 102 _make_signal_source: F, 103 ) -> ProcessResult<RhiProcessResult> 104 where 105 F: FnOnce() -> Option<S>, 106 S: crate::RhiProcessSignalSource + 'static, 107 { 108 Err(ProcessFailure( 109 RhiProcessResult::ServiceOrDependencyUnavailable, 110 )) 111 } 112 113 fn execute(invocation: RhiCliInvocationV1) -> ProcessResult<RhiProcessResult> { 114 let plan = plan_rhi_cli_v1(&invocation); 115 if plan.primary_authority() == RhiCliPrimaryAuthorityV1::Daemon { 116 return Err(ProcessFailure( 117 RhiProcessResult::ServiceOrDependencyUnavailable, 118 )); 119 } 120 let runtime = resolve_runtime(&invocation)?; 121 let output = invocation.output_mode(); 122 match (plan.primary_authority(), invocation.command()) { 123 (RhiCliPrimaryAuthorityV1::Offline, RhiCommandV1::Config(command)) => { 124 execute_config(output, &runtime, command) 125 } 126 (RhiCliPrimaryAuthorityV1::Offline, RhiCommandV1::State(command)) => { 127 execute_state(output, &runtime, command) 128 } 129 (RhiCliPrimaryAuthorityV1::Offline, RhiCommandV1::Identity(command)) => { 130 execute_identity(output, &runtime, *command) 131 } 132 (RhiCliPrimaryAuthorityV1::Offline, RhiCommandV1::Doctor) => { 133 execute_doctor(output, &runtime) 134 } 135 (RhiCliPrimaryAuthorityV1::LiveUnixAdmin, command) => { 136 execute_live(output, &runtime, command) 137 } 138 _ => Err(ProcessFailure(RhiProcessResult::UnexpectedInternal)), 139 } 140 } 141 142 fn execute_config( 143 output: RhiCliOutputModeV1, 144 runtime: &RhiRuntimeContext, 145 command: &RhiConfigCommandV1, 146 ) -> ProcessResult<RhiProcessResult> { 147 match command { 148 RhiConfigCommandV1::Init => { 149 let bytes = read_bounded_stdin(RHI_CONFIG_DOCUMENT_MAX_UTF8_BYTES)?; 150 initialize_rhi_config_document(runtime, &bytes).map_err(|_| input_failure())?; 151 emit_simple_success(output, "config_initialized") 152 } 153 RhiConfigCommandV1::Validate => { 154 load_rhi_config_document(runtime).map_err(|_| input_failure())?; 155 emit_simple_success(output, "config_valid") 156 } 157 RhiConfigCommandV1::Schema => emit_bytes(include_bytes!( 158 "../contracts/services_hardening/config.v1.schema.json" 159 )), 160 RhiConfigCommandV1::Apply(arguments) => { 161 let current = load_rhi_config_document(runtime).map_err(|_| input_failure())?; 162 let candidate = load_rhi_config_candidate(runtime, arguments.candidate_config()) 163 .map_err(|_| input_failure())?; 164 let tokio = build_tokio_runtime(current.runtime_thread_limits())?; 165 let outcome = tokio 166 .block_on(apply_rhi_configuration( 167 runtime, 168 ¤t, 169 &candidate, 170 migration_time()?, 171 &migration_build_identity()?, 172 )) 173 .map_err(|_| conflict_failure())?; 174 emit_value( 175 output, 176 "config_applied", 177 json!({"generation": outcome.generation()}), 178 ) 179 } 180 RhiConfigCommandV1::Show => Err(ProcessFailure(RhiProcessResult::UnexpectedInternal)), 181 } 182 } 183 184 fn execute_state( 185 output: RhiCliOutputModeV1, 186 runtime: &RhiRuntimeContext, 187 command: &RhiStateCommandV1, 188 ) -> ProcessResult<RhiProcessResult> { 189 let configuration = load_rhi_config_document(runtime).map_err(|_| input_failure())?; 190 let tokio = build_tokio_runtime(configuration.runtime_thread_limits())?; 191 match command { 192 RhiStateCommandV1::Init => { 193 let metadata = RhiStateMetadata::new( 194 runtime, 195 &configuration, 196 source_generation()?, 197 wall_time_millis()?, 198 ) 199 .map_err(|_| state_failure())?; 200 tokio 201 .block_on(initialize_rhi_state( 202 runtime, 203 &metadata, 204 migration_time()?, 205 &migration_build_identity()?, 206 )) 207 .map_err(|_| state_failure())?; 208 emit_simple_success(output, "state_initialized") 209 } 210 RhiStateCommandV1::Restore(arguments) => { 211 let state = tokio 212 .block_on(open_rhi_state_read_write_from_config( 213 runtime, 214 &configuration, 215 migration_time()?, 216 &migration_build_identity()?, 217 )) 218 .map_err(|_| state_failure())?; 219 let metadata = state.metadata().clone(); 220 tokio.block_on(state.close()).map_err(|_| state_failure())?; 221 let manifest = 222 read_bounded_file(arguments.manifest(), BACKUP_MANIFEST_CANONICAL_MAX_BYTES)?; 223 let digest = decode_hex_32(arguments.manifest_sha256())?; 224 let verified = verify_rhi_state_backup( 225 &manifest, 226 BackupManifestSha256::from_bytes(digest), 227 arguments.bundle(), 228 &metadata, 229 NonZeroU64::new(arguments.maximum_state_bytes()).ok_or_else(input_failure)?, 230 ) 231 .map_err(|_| state_failure())?; 232 let staged = tokio 233 .block_on(stage_rhi_state_restore(runtime, &metadata, verified)) 234 .map_err(|_| state_failure())?; 235 tokio 236 .block_on(finalize_rhi_state_restore(staged)) 237 .map_err(|_| state_failure())?; 238 emit_simple_success(output, "state_restore_finalized") 239 } 240 RhiStateCommandV1::Verify | RhiStateCommandV1::Migrate => { 241 let state = tokio 242 .block_on(open_rhi_state_read_write_from_config( 243 runtime, 244 &configuration, 245 migration_time()?, 246 &migration_build_identity()?, 247 )) 248 .map_err(|_| state_failure())?; 249 if matches!(command, RhiStateCommandV1::Verify) { 250 let checked_at = 251 IntegrityCheckedAtUnixMs::new(wall_time_millis()?).ok_or_else(state_failure)?; 252 tokio 253 .block_on(state.inspect_integrity(checked_at)) 254 .map_err(|_| state_failure())?; 255 } 256 tokio.block_on(state.close()).map_err(|_| state_failure())?; 257 emit_simple_success( 258 output, 259 if matches!(command, RhiStateCommandV1::Verify) { 260 "state_verified" 261 } else { 262 "state_migrated" 263 }, 264 ) 265 } 266 RhiStateCommandV1::Status | RhiStateCommandV1::Backup(_) => { 267 Err(ProcessFailure(RhiProcessResult::UnexpectedInternal)) 268 } 269 } 270 } 271 272 fn execute_identity( 273 output: RhiCliOutputModeV1, 274 runtime: &RhiRuntimeContext, 275 command: RhiIdentityCommandV1, 276 ) -> ProcessResult<RhiProcessResult> { 277 if command != RhiIdentityCommandV1::Init { 278 return Err(ProcessFailure(RhiProcessResult::UnexpectedInternal)); 279 } 280 let configuration = load_rhi_config_document(runtime).map_err(|_| input_failure())?; 281 let tokio = build_tokio_runtime(configuration.runtime_thread_limits())?; 282 let state = tokio 283 .block_on(open_rhi_state_read_write_from_config( 284 runtime, 285 &configuration, 286 migration_time()?, 287 &migration_build_identity()?, 288 )) 289 .map_err(|_| state_failure())?; 290 let metadata = state.metadata().clone(); 291 tokio.block_on(state.close()).map_err(|_| state_failure())?; 292 let binding = crate::RhiIdentityEnvelopeBinding::from_configuration(&configuration, &metadata) 293 .map_err(|_| state_failure())?; 294 let credential = 295 resolve_rhi_wrapping_credential(runtime, &binding).map_err(|_| state_failure())?; 296 let material = read_identity_provisioning_document(std::io::stdin().lock()) 297 .map_err(|_| input_failure())?; 298 let identity = provision_rhi_encrypted_identity(&binding, &credential, material) 299 .map_err(|_| state_failure())?; 300 emit_value( 301 output, 302 "identity_initialized", 303 json!({"generation": 0, "public_key": identity.public_identity().as_hex(), "role": "service"}), 304 ) 305 } 306 307 #[cfg(any(target_os = "linux", target_os = "macos"))] 308 fn execute_doctor( 309 _output: RhiCliOutputModeV1, 310 runtime: &RhiRuntimeContext, 311 ) -> ProcessResult<RhiProcessResult> { 312 let configuration = load_rhi_config_document(runtime).map_err(|_| input_failure())?; 313 let tokio = build_tokio_runtime(configuration.runtime_thread_limits())?; 314 let report = tokio 315 .block_on(crate::run_rhi_doctor( 316 runtime, 317 &crate::system_doctor::RhiSystemDoctorProbe::new(runtime, &configuration), 318 )) 319 .map_err(|_| ProcessFailure(RhiProcessResult::UnexpectedInternal))?; 320 emit_bytes(report.canonical_json())?; 321 if report.exit_code() == 0 { 322 Ok(RhiProcessResult::Success) 323 } else { 324 Err(ProcessFailure(RhiProcessResult::DoctorRequiredCheckFailed)) 325 } 326 } 327 328 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 329 fn execute_doctor( 330 _output: RhiCliOutputModeV1, 331 _runtime: &RhiRuntimeContext, 332 ) -> ProcessResult<RhiProcessResult> { 333 Err(ProcessFailure( 334 RhiProcessResult::ServiceOrDependencyUnavailable, 335 )) 336 } 337 338 #[cfg(any(target_os = "linux", target_os = "macos"))] 339 fn execute_live( 340 _output: RhiCliOutputModeV1, 341 runtime: &RhiRuntimeContext, 342 command: &RhiCommandV1, 343 ) -> ProcessResult<RhiProcessResult> { 344 let configuration = load_rhi_config_document(runtime).map_err(|_| input_failure())?; 345 let tokio = build_tokio_runtime(configuration.runtime_thread_limits())?; 346 let bytes = tokio.block_on(live_command(runtime, &configuration, command))?; 347 emit_bytes(&bytes) 348 } 349 350 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 351 fn execute_live( 352 _output: RhiCliOutputModeV1, 353 _runtime: &RhiRuntimeContext, 354 _command: &RhiCommandV1, 355 ) -> ProcessResult<RhiProcessResult> { 356 Err(ProcessFailure( 357 RhiProcessResult::ServiceOrDependencyUnavailable, 358 )) 359 } 360 361 #[cfg(any(target_os = "linux", target_os = "macos"))] 362 async fn live_command( 363 runtime: &RhiRuntimeContext, 364 configuration: &crate::RhiConfigDocumentV1, 365 command: &RhiCommandV1, 366 ) -> ProcessResult<Box<[u8]>> { 367 use crate::RhiAdminRoute as Route; 368 let (route, target, mutation) = match command { 369 RhiCommandV1::Config(RhiConfigCommandV1::Show) => ( 370 Route::EffectiveConfig, 371 Route::EffectiveConfig.path().to_owned(), 372 None, 373 ), 374 RhiCommandV1::State(RhiStateCommandV1::Status) => ( 375 Route::StateStatus, 376 Route::StateStatus.path().to_owned(), 377 None, 378 ), 379 RhiCommandV1::State(RhiStateCommandV1::Backup(arguments)) => ( 380 Route::StateBackup, 381 Route::StateBackup.path().to_owned(), 382 Some(( 383 arguments.operation_id(), 384 json!({ 385 "confirmation": "confirm", 386 "expected_generation": arguments.expected_generation(), 387 "target_path": arguments.target().to_str().ok_or_else(input_failure)?, 388 }), 389 )), 390 ), 391 RhiCommandV1::Identity(RhiIdentityCommandV1::Status) => ( 392 Route::IdentityStatus, 393 format!("{}?role=service", Route::IdentityStatus.path()), 394 None, 395 ), 396 RhiCommandV1::Identity(RhiIdentityCommandV1::ExportPublic) => ( 397 Route::IdentityPublic, 398 format!("{}?role=service", Route::IdentityPublic.path()), 399 None, 400 ), 401 RhiCommandV1::Status => (Route::Status, Route::Status.path().to_owned(), None), 402 RhiCommandV1::Metrics(RhiMetricsCommandV1::Snapshot) => ( 403 Route::MetricsSnapshot, 404 Route::MetricsSnapshot.path().to_owned(), 405 None, 406 ), 407 RhiCommandV1::Reconciliation(RhiReconciliationCommandV1::Status) => ( 408 Route::ReconciliationStatus, 409 Route::ReconciliationStatus.path().to_owned(), 410 None, 411 ), 412 RhiCommandV1::Reconciliation(RhiReconciliationCommandV1::Jobs(page)) => ( 413 Route::ReconciliationJobs, 414 paged_target(Route::ReconciliationJobs.path(), page), 415 None, 416 ), 417 RhiCommandV1::Reconciliation(RhiReconciliationCommandV1::Refresh(arguments)) => ( 418 Route::ReconciliationRefresh, 419 Route::ReconciliationRefresh.path().to_owned(), 420 Some(( 421 arguments.operation_id(), 422 json!({ 423 "expected_dirty_generation": arguments.expected_dirty_generation(), 424 "trade_id": arguments.trade_id(), 425 }), 426 )), 427 ), 428 RhiCommandV1::Sources(RhiSourcesCommandV1::List(page)) => ( 429 Route::Sources, 430 paged_target(Route::Sources.path(), page), 431 None, 432 ), 433 RhiCommandV1::Trade(RhiTradeCommandV1::Projection(arguments)) => ( 434 Route::TradeProjection, 435 Route::TradeProjection 436 .path() 437 .replace("{trade_id}", arguments.trade_id()), 438 None, 439 ), 440 RhiCommandV1::Trade(RhiTradeCommandV1::ReportCurrent(arguments)) => ( 441 Route::TradeReportCurrent, 442 Route::TradeReportCurrent 443 .path() 444 .replace("{trade_id}", arguments.trade_id()), 445 None, 446 ), 447 RhiCommandV1::Trade(RhiTradeCommandV1::Reports(arguments)) => ( 448 Route::TradeReports, 449 paged_target( 450 &Route::TradeReports 451 .path() 452 .replace("{trade_id}", arguments.trade_id()), 453 arguments.page(), 454 ), 455 None, 456 ), 457 RhiCommandV1::Publication(RhiPublicationCommandV1::Backlog(page)) => ( 458 Route::PublicationBacklog, 459 paged_target(Route::PublicationBacklog.path(), page), 460 None, 461 ), 462 RhiCommandV1::Publication(RhiPublicationCommandV1::Targets(page)) => ( 463 Route::PublicationTargets, 464 paged_target(Route::PublicationTargets.path(), page), 465 None, 466 ), 467 RhiCommandV1::Publication(RhiPublicationCommandV1::Retry(arguments)) => ( 468 Route::PublicationRetry, 469 Route::PublicationRetry.path().to_owned(), 470 Some(( 471 arguments.operation_id(), 472 json!({ 473 "expected_generation": arguments.expected_generation(), 474 "workflow_id": arguments.workflow_id(), 475 }), 476 )), 477 ), 478 RhiCommandV1::Presence(RhiPresenceCommandV1::Desired) => ( 479 Route::PresenceDesired, 480 Route::PresenceDesired.path().to_owned(), 481 None, 482 ), 483 RhiCommandV1::Presence(RhiPresenceCommandV1::Render(arguments)) => ( 484 Route::PresenceRender, 485 Route::PresenceRender.path().to_owned(), 486 Some(( 487 arguments.operation_id(), 488 json!({"expected_generation": arguments.expected_generation()}), 489 )), 490 ), 491 RhiCommandV1::Presence(RhiPresenceCommandV1::Refresh(arguments)) => ( 492 Route::PresenceRefresh, 493 Route::PresenceRefresh.path().to_owned(), 494 Some(( 495 arguments.operation_id(), 496 json!({"expected_generation": arguments.expected_generation()}), 497 )), 498 ), 499 _ => return Err(ProcessFailure(RhiProcessResult::UnexpectedInternal)), 500 }; 501 let client = AdminClient::new( 502 runtime.artifacts().admin_socket(), 503 admin_transport_limits(configuration).map_err(|_| input_failure())?, 504 ) 505 .map_err(|_| input_failure())?; 506 let target = AdminClientTarget::new(target).map_err(|_| input_failure())?; 507 let result = if let Some((operation_id, request)) = mutation { 508 let operation_id = AdminOperationId::new(operation_id).map_err(|_| input_failure())?; 509 client 510 .mutate::<_, Value>(&target, operation_id, None, request) 511 .await 512 .map(|response| response.result().clone()) 513 } else { 514 client 515 .get::<Value>(&target) 516 .await 517 .map(|response| response.result().clone()) 518 }; 519 match result { 520 Ok(value) => admit_admin_response_value(route, &value) 521 .map_err(|_| ProcessFailure(RhiProcessResult::ServiceOrDependencyUnavailable)), 522 Err(error) if error.kind() == AdminClientErrorKind::ServerFailure => { 523 Err(conflict_failure()) 524 } 525 Err(_) => Err(ProcessFailure( 526 RhiProcessResult::ServiceOrDependencyUnavailable, 527 )), 528 } 529 } 530 531 #[cfg(any(target_os = "linux", target_os = "macos"))] 532 fn paged_target(path: &str, page: &crate::RhiPageQueryArgsV1) -> String { 533 let mut target = format!("{path}?limit={}", page.limit()); 534 if let Some(cursor) = page.cursor() { 535 target.push_str("&cursor="); 536 target.push_str(cursor); 537 } 538 target 539 } 540 541 fn resolve_runtime(invocation: &RhiCliInvocationV1) -> ProcessResult<RhiRuntimeContext> { 542 let resolver = RadrootsPathResolver::new(RadrootsPlatform::current(), host_environment()); 543 resolve_rhi_runtime_context(&resolver, invocation).map_err(|_| input_failure()) 544 } 545 546 fn migration_build_identity() -> ProcessResult<MigrationBuildIdentity> { 547 let versions = ContractVersions::new( 548 RHI_CONFIG_SCHEMA_VERSION, 549 RHI_STATE_SCHEMA_VERSION, 550 RHI_ADMIN_CONTRACT_VERSION, 551 RHI_STATUS_CONTRACT_VERSION, 552 RHI_PROVIDER_CONTRACT_VERSION, 553 ) 554 .map_err(|_| ProcessFailure(RhiProcessResult::UnexpectedInternal))?; 555 let build = radroots_service_host::compile_time_build_info!( 556 feature_profile: "service-host", 557 contract_versions: versions, 558 ) 559 .map_err(|_| ProcessFailure(RhiProcessResult::UnexpectedInternal))?; 560 MigrationBuildIdentity::new( 561 build.service_version(), 562 build.service_commit(), 563 build.lib_revision(), 564 build.rust_version(), 565 build.target(), 566 build.feature_profile(), 567 versions.config(), 568 versions.state(), 569 versions.admin(), 570 versions.status(), 571 versions.provider(), 572 ) 573 .map_err(|_| ProcessFailure(RhiProcessResult::UnexpectedInternal)) 574 } 575 576 fn source_generation() -> ProcessResult<SourceGeneration> { 577 for _ in 0..4 { 578 let mut bytes = [0_u8; 32]; 579 SystemEntropy 580 .fill_bytes(&mut bytes) 581 .map_err(|_| state_failure())?; 582 if let Ok(generation) = SourceGeneration::new(bytes) { 583 return Ok(generation); 584 } 585 } 586 Err(state_failure()) 587 } 588 589 fn wall_time_seconds() -> ProcessResult<u64> { 590 SystemWallClock 591 .now_utc() 592 .map(|time| time.get()) 593 .map_err(|_| ProcessFailure(RhiProcessResult::UnexpectedInternal)) 594 } 595 596 fn wall_time_millis() -> ProcessResult<u64> { 597 wall_time_seconds()? 598 .checked_mul(1_000) 599 .filter(|value| *value <= i64::MAX as u64) 600 .ok_or(ProcessFailure(RhiProcessResult::UnexpectedInternal)) 601 } 602 603 fn migration_time() -> ProcessResult<MigrationAppliedAtUnixSeconds> { 604 MigrationAppliedAtUnixSeconds::new(wall_time_seconds()?) 605 .map_err(|_| ProcessFailure(RhiProcessResult::UnexpectedInternal)) 606 } 607 608 fn build_tokio_runtime( 609 limits: crate::RhiRuntimeThreadLimitsV1, 610 ) -> ProcessResult<tokio::runtime::Runtime> { 611 if tokio::runtime::Handle::try_current().is_ok() { 612 return Err(ProcessFailure(RhiProcessResult::UnexpectedInternal)); 613 } 614 tokio::runtime::Builder::new_multi_thread() 615 .worker_threads(limits.worker_threads()) 616 .max_blocking_threads(limits.blocking_threads()) 617 .enable_all() 618 .build() 619 .map_err(|_| ProcessFailure(RhiProcessResult::UnexpectedInternal)) 620 } 621 622 fn host_environment() -> RadrootsHostEnvironment { 623 let path = |name| { 624 env::var_os(name) 625 .filter(|value| !value.is_empty()) 626 .map(PathBuf::from) 627 }; 628 RadrootsHostEnvironment { 629 home_dir: path("HOME"), 630 xdg_config_home: path("XDG_CONFIG_HOME"), 631 xdg_data_home: path("XDG_DATA_HOME"), 632 xdg_state_home: path("XDG_STATE_HOME"), 633 xdg_cache_home: path("XDG_CACHE_HOME"), 634 xdg_runtime_dir: path("XDG_RUNTIME_DIR"), 635 appdata_dir: path("APPDATA"), 636 localappdata_dir: path("LOCALAPPDATA"), 637 } 638 } 639 640 fn read_bounded_stdin(maximum: usize) -> ProcessResult<Vec<u8>> { 641 let mut reader = std::io::stdin().lock(); 642 let mut bytes = Vec::with_capacity(maximum.min(64 * 1_024).saturating_add(1)); 643 Read::by_ref(&mut reader) 644 .take(u64::try_from(maximum).unwrap_or(u64::MAX).saturating_add(1)) 645 .read_to_end(&mut bytes) 646 .map_err(|_| input_failure())?; 647 if bytes.len() > maximum { 648 return Err(input_failure()); 649 } 650 Ok(bytes) 651 } 652 653 #[cfg(any(target_os = "linux", target_os = "macos"))] 654 fn read_bounded_file(path: &Path, maximum: usize) -> ProcessResult<Vec<u8>> { 655 crate::config_loader::read_secure_bounded_file(path, maximum).map_err(|_| state_failure()) 656 } 657 658 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 659 fn read_bounded_file(_path: &Path, _maximum: usize) -> ProcessResult<Vec<u8>> { 660 Err(state_failure()) 661 } 662 663 fn decode_hex_32(value: &str) -> ProcessResult<[u8; 32]> { 664 if value.len() != 64 { 665 return Err(input_failure()); 666 } 667 let mut output = [0_u8; 32]; 668 for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() { 669 output[index] = (hex_nibble(pair[0])? << 4) | hex_nibble(pair[1])?; 670 } 671 Ok(output) 672 } 673 674 const fn hex_nibble(value: u8) -> ProcessResult<u8> { 675 match value { 676 b'0'..=b'9' => Ok(value - b'0'), 677 b'a'..=b'f' => Ok(value - b'a' + 10), 678 _ => Err(input_failure()), 679 } 680 } 681 682 fn emit_simple_success( 683 output: RhiCliOutputModeV1, 684 code: &'static str, 685 ) -> ProcessResult<RhiProcessResult> { 686 emit_value(output, code, json!({"ok": true})) 687 } 688 689 fn emit_value( 690 output: RhiCliOutputModeV1, 691 code: &'static str, 692 value: Value, 693 ) -> ProcessResult<RhiProcessResult> { 694 let bytes = match output { 695 RhiCliOutputModeV1::Json => serde_json::to_vec(&value) 696 .map_err(|_| ProcessFailure(RhiProcessResult::UnexpectedInternal))?, 697 RhiCliOutputModeV1::Human => code.as_bytes().to_vec(), 698 }; 699 emit_bytes(&bytes) 700 } 701 702 fn emit_bytes(bytes: &[u8]) -> ProcessResult<RhiProcessResult> { 703 let mut stdout = std::io::stdout().lock(); 704 stdout 705 .write_all(bytes) 706 .and_then(|()| { 707 if bytes.ends_with(b"\n") { 708 Ok(()) 709 } else { 710 stdout.write_all(b"\n") 711 } 712 }) 713 .and_then(|()| stdout.flush()) 714 .map_err(|_| ProcessFailure(RhiProcessResult::UnexpectedInternal))?; 715 Ok(RhiProcessResult::Success) 716 } 717 718 const fn input_failure() -> ProcessFailure { 719 ProcessFailure(RhiProcessResult::InputOrConfiguration) 720 } 721 722 const fn state_failure() -> ProcessFailure { 723 ProcessFailure(RhiProcessResult::StateOrIdentityUnavailable) 724 } 725 726 const fn conflict_failure() -> ProcessFailure { 727 ProcessFailure(RhiProcessResult::OperationRejectedOrConflict) 728 } 729 730 #[cfg(test)] 731 mod tests { 732 use super::*; 733 734 #[test] 735 fn runtime_limits_are_explicit_and_digest_decoding_is_strict() { 736 let source = include_str!("process_v1.rs") 737 .split("#[cfg(test)]") 738 .next() 739 .expect("production source"); 740 assert!(source.contains("worker_threads(limits.worker_threads())")); 741 assert!(source.contains("max_blocking_threads(limits.blocking_threads())")); 742 assert!(!source.contains("available_parallelism")); 743 assert_eq!( 744 decode_hex_32(&"ab".repeat(32)).expect("lowercase digest"), 745 [0xab; 32] 746 ); 747 assert!(decode_hex_32(&"AB".repeat(32)).is_err()); 748 } 749 750 #[test] 751 fn nested_tokio_runtime_creation_fails_closed_without_panicking() { 752 let configuration = crate::parse_rhi_config_v1( 753 include_bytes!("../contracts/services_hardening/config.v1.example.toml"), 754 crate::RhiConfigProfile::Production, 755 ) 756 .expect("configuration fixture"); 757 let outer = tokio::runtime::Builder::new_current_thread() 758 .enable_all() 759 .build() 760 .expect("outer runtime"); 761 let result = 762 outer.block_on(async { build_tokio_runtime(configuration.runtime_thread_limits()) }); 763 assert_eq!( 764 result.expect_err("nested runtime must be rejected"), 765 ProcessFailure(RhiProcessResult::UnexpectedInternal) 766 ); 767 } 768 }