process_v1.rs (38455B)
1 //! Binary-owned execution for one already-admitted command invocation. 2 3 use std::env; 4 use std::io::{Read, Write}; 5 use std::path::{Path, PathBuf}; 6 7 use radroots_service_host::{ 8 AdminClient, AdminClientErrorKind, AdminClientTarget, ContractVersions, EntropySource, 9 SystemEntropy, SystemWallClock, WallClock, 10 }; 11 use radroots_service_sqlite::{ 12 BACKUP_MANIFEST_CANONICAL_MAX_BYTES, BackupCreatedAtUnixMs, IntegrityCheckOutcome, 13 IntegrityCheckedAtUnixMs, MigrationAppliedAtUnixSeconds, MigrationBuildIdentity, OpenMode, 14 ServiceBackupManifest, WriterAuthority, 15 }; 16 use radroots_storage::event::SourceGeneration; 17 use serde_json::{Value, json}; 18 19 use crate::admin_v1::{admin_transport_limits, admit_admin_response_value}; 20 use crate::cli_bootstrap::read_identity_provisioning_document; 21 use crate::config_v1::config_schema_document; 22 use crate::{ 23 MYC_CONFIG_DOCUMENT_MAX_UTF8_BYTES, MYC_CONFIG_SCHEMA, MYC_CONFIG_SCHEMA_VERSION, 24 MYC_OPERATOR_CONTRACT_VERSION, MYC_PROVIDER_CONTRACT_VERSION, 25 MYC_SIGNER_STATUS_CONTRACT_VERSION, MYC_STATE_SCHEMA_VERSION, MycBootstrapProfileV1, 26 MycCliInvocationV1, MycCliOutputModeV1, MycCliPrimaryAuthorityV1, MycCommandV1, 27 MycConfigCommandV1, MycConfigDocumentV1, MycConnectionCountsV1, MycIdentityCommandArgsV1, 28 MycIdentityCommandV1, MycLocalSignerClient, MycOutboxStatusV1, MycProcessResult, 29 MycProviderKind, MycProviderRole, MycRuntimeContext, MycStateBackupArgsV1, MycStateCommandV1, 30 MycStateMetadata, MycStateRestoreArgsV1, RadrootsHostEnvironment, RadrootsPathResolver, 31 RadrootsPlatform, finalize_myc_state_restore, initialize_myc_config_document, 32 initialize_myc_state, load_myc_config_candidate, load_myc_config_document, 33 open_myc_encrypted_identity, open_myc_state_inspection_from_config, 34 open_myc_state_read_write_from_config, plan_myc_cli_v1, provision_myc_encrypted_identity, 35 resolve_myc_runtime_context, resolve_myc_wrapping_credential, stage_myc_state_restore, 36 verify_myc_state_backup, 37 }; 38 39 #[derive(Clone, Copy)] 40 struct ProcessFailure(MycProcessResult); 41 42 type ProcessResult<T> = Result<T, ProcessFailure>; 43 44 struct OfflineStateSnapshot { 45 value: Value, 46 connection_counts: MycConnectionCountsV1, 47 outbox: MycOutboxStatusV1, 48 } 49 50 /// Executes one admitted Myc invocation without reparsing process arguments. 51 /// 52 /// Result bytes are written to stdout only after the governed operation 53 /// succeeds. The binary retains responsibility for emitting the final safe 54 /// structured diagnostic and choosing the process exit code. 55 #[must_use] 56 pub fn execute_myc_cli_v1(invocation: MycCliInvocationV1) -> MycProcessResult { 57 execute(invocation).unwrap_or_else(|failure| failure.0) 58 } 59 60 /// Executes one admitted invocation with a binary-owned process-signal source. 61 /// 62 /// The signal source factory is consulted only for `run` and only after the 63 /// governed Tokio runtime has entered. Non-daemon commands retain the exact 64 /// one-pass execution path used by [`execute_myc_cli_v1`]. 65 #[must_use] 66 pub fn execute_myc_cli_v1_with_signal_source<F, S>( 67 invocation: MycCliInvocationV1, 68 make_signal_source: F, 69 ) -> MycProcessResult 70 where 71 F: FnOnce() -> Option<S>, 72 S: crate::MycProcessSignalSource + 'static, 73 { 74 if !matches!(invocation.command(), MycCommandV1::Run) { 75 return execute_myc_cli_v1(invocation); 76 } 77 execute_run(invocation, make_signal_source).unwrap_or_else(|failure| failure.0) 78 } 79 80 fn execute_run<F, S>( 81 invocation: MycCliInvocationV1, 82 make_signal_source: F, 83 ) -> ProcessResult<MycProcessResult> 84 where 85 F: FnOnce() -> Option<S>, 86 S: crate::MycProcessSignalSource + 'static, 87 { 88 let resolver = RadrootsPathResolver::new(RadrootsPlatform::current(), host_environment()); 89 let runtime = 90 resolve_myc_runtime_context(&resolver, &invocation).map_err(|_| input_failure())?; 91 let configuration = load_myc_config_document(&runtime).map_err(|_| input_failure())?; 92 let applied_at = migration_time()?; 93 let build = migration_build_identity()?; 94 let tokio = build_tokio_runtime(configuration.runtime_thread_limits())?; 95 tokio.block_on(async move { 96 let signals = make_signal_source().ok_or(ProcessFailure( 97 MycProcessResult::ServiceOrDependencyUnavailable, 98 ))?; 99 Ok(crate::runtime_graph::run_myc_daemon( 100 runtime, 101 configuration, 102 applied_at, 103 &build, 104 signals, 105 ) 106 .await) 107 }) 108 } 109 110 fn execute(invocation: MycCliInvocationV1) -> ProcessResult<MycProcessResult> { 111 let plan = plan_myc_cli_v1(&invocation); 112 if matches!( 113 (plan.primary_authority(), invocation.command()), 114 (MycCliPrimaryAuthorityV1::Daemon, MycCommandV1::Run) 115 ) { 116 return Err(ProcessFailure( 117 MycProcessResult::ServiceOrDependencyUnavailable, 118 )); 119 } 120 121 let resolver = RadrootsPathResolver::new(RadrootsPlatform::current(), host_environment()); 122 let runtime = 123 resolve_myc_runtime_context(&resolver, &invocation).map_err(|_| input_failure())?; 124 let output = invocation.output_mode(); 125 126 match (plan.primary_authority(), invocation.command()) { 127 (MycCliPrimaryAuthorityV1::Offline, MycCommandV1::Config(command)) => { 128 execute_config(output, &runtime, command) 129 } 130 ( 131 MycCliPrimaryAuthorityV1::Offline | MycCliPrimaryAuthorityV1::LiveUnixAdmin, 132 MycCommandV1::State(command), 133 ) => execute_state(output, &runtime, command), 134 ( 135 MycCliPrimaryAuthorityV1::Offline | MycCliPrimaryAuthorityV1::LiveUnixAdmin, 136 MycCommandV1::Identity(command), 137 ) => execute_identity(output, &runtime, command), 138 (MycCliPrimaryAuthorityV1::LiveUnixAdmin, MycCommandV1::Status) => { 139 execute_live_or_offline_status(output, &runtime) 140 } 141 (MycCliPrimaryAuthorityV1::Offline, MycCommandV1::Doctor) => { 142 execute_doctor(output, &runtime) 143 } 144 _ => Err(ProcessFailure(MycProcessResult::UnexpectedInternal)), 145 } 146 } 147 148 fn execute_doctor( 149 _output: MycCliOutputModeV1, 150 runtime: &MycRuntimeContext, 151 ) -> ProcessResult<MycProcessResult> { 152 let configuration = load_myc_config_document(runtime).map_err(|_| input_failure())?; 153 let tokio = build_tokio_runtime(configuration.runtime_thread_limits())?; 154 let report = tokio 155 .block_on(crate::run_myc_doctor( 156 runtime, 157 &crate::system_doctor::MycSystemDoctorProbe::new(runtime, &configuration), 158 )) 159 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal))?; 160 emit_bytes(report.canonical_json())?; 161 if report.exit_code() == 0 { 162 Ok(MycProcessResult::Success) 163 } else { 164 Err(ProcessFailure(MycProcessResult::DoctorRequiredCheckFailed)) 165 } 166 } 167 168 fn execute_config( 169 output: MycCliOutputModeV1, 170 runtime: &MycRuntimeContext, 171 command: &MycConfigCommandV1, 172 ) -> ProcessResult<MycProcessResult> { 173 match command { 174 MycConfigCommandV1::Init => { 175 let bytes = read_bounded_stdin(MYC_CONFIG_DOCUMENT_MAX_UTF8_BYTES)?; 176 initialize_myc_config_document(runtime, &bytes).map_err(|_| input_failure())?; 177 emit_simple_success(output, "config_initialized") 178 } 179 MycConfigCommandV1::Validate => { 180 load_myc_config_document(runtime).map_err(|_| input_failure())?; 181 emit_simple_success(output, "config_valid") 182 } 183 MycConfigCommandV1::Show => { 184 let configuration = load_myc_config_document(runtime).map_err(|_| input_failure())?; 185 emit_bytes(configuration.effective().canonical_json().as_bytes()) 186 } 187 MycConfigCommandV1::Schema => emit_bytes(config_schema_document().as_bytes()), 188 MycConfigCommandV1::Apply(arguments) => { 189 let current = load_myc_config_document(runtime).map_err(|_| input_failure())?; 190 let candidate = load_myc_config_candidate(runtime, arguments.candidate_config()) 191 .map_err(|_| input_failure())?; 192 let build = migration_build_identity()?; 193 let applied_at = migration_time()?; 194 let tokio = build_tokio_runtime(current.runtime_thread_limits())?; 195 let outcome = tokio.block_on(async { 196 validate_candidate_providers(runtime, &candidate).await?; 197 let state = 198 open_myc_state_read_write_from_config(runtime, ¤t, applied_at, &build) 199 .await 200 .map_err(|_| state_failure())?; 201 let applied = state 202 .repository() 203 .apply_configuration(¤t, &candidate, applied_at, &build) 204 .await 205 .map_err(|error| match error.kind() { 206 crate::MycConfigApplyErrorKind::PolicyConflict 207 | crate::MycConfigApplyErrorKind::ResourceExhausted => conflict_failure(), 208 _ => state_failure(), 209 }); 210 let closed = state.close().await.map_err(|_| state_failure()); 211 closed.and(applied) 212 })?; 213 emit_value( 214 output, 215 "config_applied", 216 json!({ 217 "generation": outcome.generation(), 218 "revoked_challenges": outcome.revoked_challenge_count(), 219 "revoked_connections": outcome.revoked_connection_count(), 220 }), 221 ) 222 } 223 } 224 } 225 226 fn execute_state( 227 output: MycCliOutputModeV1, 228 runtime: &MycRuntimeContext, 229 command: &MycStateCommandV1, 230 ) -> ProcessResult<MycProcessResult> { 231 let configuration = load_myc_config_document(runtime).map_err(|_| input_failure())?; 232 let tokio = build_tokio_runtime(configuration.runtime_thread_limits())?; 233 match command { 234 MycStateCommandV1::Init => { 235 let generation = source_generation()?; 236 let created_at = wall_time_millis()?; 237 let metadata = MycStateMetadata::new(runtime, &configuration, generation, created_at) 238 .map_err(|_| state_failure())?; 239 let applied_at = migration_time()?; 240 let build = migration_build_identity()?; 241 tokio 242 .block_on(initialize_myc_state(runtime, &metadata, applied_at, &build)) 243 .map_err(|_| state_failure())?; 244 emit_simple_success(output, "state_initialized") 245 } 246 MycStateCommandV1::Status => { 247 if let Some(bytes) = tokio.block_on(live_get( 248 runtime, 249 &configuration, 250 crate::MycAdminRoute::StateStatus, 251 None, 252 ))? { 253 return emit_bytes(&bytes); 254 } 255 let value = tokio.block_on(offline_state_status(runtime, &configuration))?; 256 emit_admin_value(output, crate::MycAdminRoute::StateStatus, value) 257 } 258 MycStateCommandV1::Backup(arguments) => { 259 if let Some(bytes) = tokio.block_on(live_backup(runtime, &configuration, arguments))? { 260 return emit_bytes(&bytes); 261 } 262 let manifest = tokio.block_on(offline_backup(runtime, &configuration, arguments))?; 263 emit_exact_bytes(manifest.canonical_bytes()) 264 } 265 MycStateCommandV1::Restore(arguments) => { 266 tokio.block_on(offline_restore(runtime, &configuration, arguments))?; 267 emit_simple_success(output, "state_restore_finalized") 268 } 269 MycStateCommandV1::Verify => { 270 tokio.block_on(offline_verify(runtime, &configuration))?; 271 emit_simple_success(output, "state_verified") 272 } 273 MycStateCommandV1::Migrate => { 274 let build = migration_build_identity()?; 275 let applied_at = migration_time()?; 276 tokio.block_on(async { 277 let state = open_myc_state_read_write_from_config( 278 runtime, 279 &configuration, 280 applied_at, 281 &build, 282 ) 283 .await 284 .map_err(|_| state_failure())?; 285 state.close().await.map_err(|_| state_failure()) 286 })?; 287 emit_simple_success(output, "state_migrated") 288 } 289 } 290 } 291 292 fn execute_identity( 293 output: MycCliOutputModeV1, 294 runtime: &MycRuntimeContext, 295 command: &MycIdentityCommandV1, 296 ) -> ProcessResult<MycProcessResult> { 297 let configuration = load_myc_config_document(runtime).map_err(|_| input_failure())?; 298 match command { 299 MycIdentityCommandV1::Init(arguments) => { 300 let binding = identity_binding(&configuration, *arguments)?; 301 if binding.kind() != MycProviderKind::EncryptedFile { 302 return Err(conflict_failure()); 303 } 304 let material = read_identity_provisioning_document(std::io::stdin().lock()) 305 .map_err(|_| input_failure())?; 306 let credential = 307 resolve_myc_wrapping_credential(runtime, binding).map_err(|_| state_failure())?; 308 let paths = crate::state_host::state_paths(runtime).map_err(|_| state_failure())?; 309 let mut authority = WriterAuthority::acquire(&paths, OpenMode::Initialize) 310 .map_err(|_| state_failure())? 311 .ok_or_else(state_failure)?; 312 let identity = provision_myc_encrypted_identity(binding, &credential, material) 313 .map_err(|_| state_failure()); 314 let released = authority.release().map_err(|_| state_failure()); 315 let identity = released.and(identity)?; 316 emit_identity( 317 output, 318 arguments.role(), 319 identity.public_identity().as_hex(), 320 0, 321 ) 322 } 323 MycIdentityCommandV1::Status(arguments) | MycIdentityCommandV1::ExportPublic(arguments) => { 324 let tokio = build_tokio_runtime(configuration.runtime_thread_limits())?; 325 let route = match command { 326 MycIdentityCommandV1::Status(_) => crate::MycAdminRoute::IdentityStatus, 327 MycIdentityCommandV1::ExportPublic(_) => crate::MycAdminRoute::IdentityPublic, 328 MycIdentityCommandV1::Init(_) => { 329 return Err(ProcessFailure(MycProcessResult::UnexpectedInternal)); 330 } 331 }; 332 if let Some(bytes) = tokio.block_on(live_get( 333 runtime, 334 &configuration, 335 route, 336 Some(arguments.role()), 337 ))? { 338 return emit_bytes(&bytes); 339 } 340 let (public_key, generation, provider, available) = 341 tokio.block_on(offline_identity(runtime, &configuration, arguments.role()))?; 342 if matches!(command, MycIdentityCommandV1::ExportPublic(_)) { 343 let public_key = public_key.ok_or_else(state_failure)?; 344 emit_identity(output, arguments.role(), &public_key, generation) 345 } else { 346 emit_admin_value( 347 output, 348 crate::MycAdminRoute::IdentityStatus, 349 json!({ 350 "available": available, 351 "configured": true, 352 "generation": generation, 353 "provider": provider.as_str(), 354 "public_key": public_key, 355 "reason_codes": if available { json!([]) } else { json!(["provider_unavailable"]) }, 356 "role": arguments.role().as_str(), 357 }), 358 ) 359 } 360 } 361 } 362 } 363 364 fn execute_live_or_offline_status( 365 output: MycCliOutputModeV1, 366 runtime: &MycRuntimeContext, 367 ) -> ProcessResult<MycProcessResult> { 368 let configuration = load_myc_config_document(runtime).map_err(|_| input_failure())?; 369 let tokio = build_tokio_runtime(configuration.runtime_thread_limits())?; 370 if let Some(bytes) = tokio.block_on(live_get( 371 runtime, 372 &configuration, 373 crate::MycAdminRoute::Status, 374 None, 375 ))? { 376 return emit_bytes(&bytes); 377 } 378 let state = tokio.block_on(offline_service_status(runtime, &configuration))?; 379 emit_admin_value(output, crate::MycAdminRoute::Status, state) 380 } 381 382 async fn live_get( 383 runtime: &MycRuntimeContext, 384 configuration: &MycConfigDocumentV1, 385 route: crate::MycAdminRoute, 386 role: Option<MycProviderRole>, 387 ) -> ProcessResult<Option<Box<[u8]>>> { 388 let client = AdminClient::new( 389 runtime.artifacts().admin_socket(), 390 admin_transport_limits(configuration).map_err(|_| input_failure())?, 391 ) 392 .map_err(|_| input_failure())?; 393 let target = if let Some(role) = role { 394 AdminClientTarget::new(format!("{}?role={}", route.path(), role.as_str())) 395 } else { 396 AdminClientTarget::new(route.path()) 397 } 398 .map_err(|_| input_failure())?; 399 match client.get::<Value>(&target).await { 400 Ok(response) => admit_admin_response_value(route, response.result()) 401 .map(Some) 402 .map_err(|_| ProcessFailure(MycProcessResult::ServiceOrDependencyUnavailable)), 403 Err(error) if error.kind() == AdminClientErrorKind::Connect => Ok(None), 404 Err(error) if error.kind() == AdminClientErrorKind::ServerFailure => { 405 Err(conflict_failure()) 406 } 407 Err(_) => Err(ProcessFailure( 408 MycProcessResult::ServiceOrDependencyUnavailable, 409 )), 410 } 411 } 412 413 async fn live_backup( 414 runtime: &MycRuntimeContext, 415 configuration: &MycConfigDocumentV1, 416 arguments: &MycStateBackupArgsV1, 417 ) -> ProcessResult<Option<Box<[u8]>>> { 418 let client = AdminClient::new( 419 runtime.artifacts().admin_socket(), 420 admin_transport_limits(configuration).map_err(|_| input_failure())?, 421 ) 422 .map_err(|_| input_failure())?; 423 let target = AdminClientTarget::new(crate::MycAdminRoute::StateBackup.path()) 424 .map_err(|_| input_failure())?; 425 let target_path = arguments.target().to_str().ok_or_else(input_failure)?; 426 let request = json!({ 427 "confirmation": "confirm", 428 "expected_generation": arguments.expected_generation(), 429 "target_path": target_path, 430 }); 431 let operation_id = radroots_service_host::AdminOperationId::new(arguments.operation_id()) 432 .map_err(|_| input_failure())?; 433 match client 434 .mutate::<_, Value>(&target, operation_id, None, request) 435 .await 436 { 437 Ok(response) => { 438 admit_admin_response_value(crate::MycAdminRoute::StateBackup, response.result()) 439 .map(Some) 440 .map_err(|_| ProcessFailure(MycProcessResult::ServiceOrDependencyUnavailable)) 441 } 442 Err(error) if error.kind() == AdminClientErrorKind::Connect => Ok(None), 443 Err(error) if error.kind() == AdminClientErrorKind::ServerFailure => { 444 Err(conflict_failure()) 445 } 446 Err(_) => Err(ProcessFailure( 447 MycProcessResult::ServiceOrDependencyUnavailable, 448 )), 449 } 450 } 451 452 async fn offline_state_status( 453 runtime: &MycRuntimeContext, 454 configuration: &MycConfigDocumentV1, 455 ) -> ProcessResult<Value> { 456 inspect_offline_state(runtime, configuration) 457 .await 458 .map(|snapshot| snapshot.value) 459 } 460 461 async fn inspect_offline_state( 462 runtime: &MycRuntimeContext, 463 configuration: &MycConfigDocumentV1, 464 ) -> ProcessResult<OfflineStateSnapshot> { 465 let state = open_myc_state_inspection_from_config(runtime, configuration) 466 .await 467 .map_err(|_| state_failure())?; 468 let inspected = async { 469 let generation = state 470 .repository() 471 .current_configuration_generation() 472 .await 473 .map_err(|_| state_failure())?; 474 let connection_counts = state 475 .repository() 476 .read_runtime_connection_counts() 477 .await 478 .map_err(|_| state_failure())?; 479 let outbox = state 480 .repository() 481 .read_runtime_outbox_status() 482 .await 483 .map_err(|_| state_failure())?; 484 let checked_at = integrity_time()?; 485 let report = state 486 .inspect_integrity(checked_at) 487 .await 488 .map_err(|_| state_failure())?; 489 let schema = state 490 .metadata() 491 .database_identity() 492 .supported_state_schema_version() 493 .get(); 494 let verified = report.sqlite() == IntegrityCheckOutcome::Verified 495 && report.foreign_keys() == IntegrityCheckOutcome::Verified; 496 Ok(OfflineStateSnapshot { 497 value: json!({ 498 "backup_eligible": verified, 499 "generation": generation, 500 "integrity": if verified { "verified" } else { "failed" }, 501 "reason_codes": if verified { json!([]) } else { json!(["database_integrity_failed"]) }, 502 "schema_version": schema, 503 "writer_lock": "free", 504 }), 505 connection_counts, 506 outbox, 507 }) 508 } 509 .await; 510 let closed = state.close().await.map_err(|_| state_failure()); 511 closed.and(inspected) 512 } 513 514 async fn offline_service_status( 515 runtime: &MycRuntimeContext, 516 configuration: &MycConfigDocumentV1, 517 ) -> ProcessResult<Value> { 518 let snapshot = inspect_offline_state(runtime, configuration).await?; 519 let state = &snapshot.value; 520 let generation = state 521 .get("generation") 522 .and_then(Value::as_u64) 523 .ok_or(ProcessFailure(MycProcessResult::UnexpectedInternal))?; 524 let schema_version = state 525 .get("schema_version") 526 .and_then(Value::as_u64) 527 .ok_or(ProcessFailure(MycProcessResult::UnexpectedInternal))?; 528 let integrity = state 529 .get("integrity") 530 .and_then(Value::as_str) 531 .ok_or(ProcessFailure(MycProcessResult::UnexpectedInternal))?; 532 let build = build_info()?; 533 let versions = build.contract_versions(); 534 let digest = crate::state_metadata::normalized_config_digest( 535 configuration.profile(), 536 configuration.normalized(), 537 ) 538 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal))?; 539 let discovery_configured = configuration 540 .provider_contract() 541 .binding(MycProviderRole::Discovery) 542 .is_some(); 543 let unavailable = |configured: bool| { 544 json!({ 545 "available": false, 546 "configured": configured, 547 "reason_codes": ["signer_provider_unavailable"], 548 }) 549 }; 550 Ok(json!({ 551 "build_info": { 552 "contract_versions": { 553 "admin": versions.admin(), 554 "config": versions.config(), 555 "provider": versions.provider(), 556 "state": versions.state(), 557 "status": versions.status(), 558 }, 559 "revision": build.service_commit(), 560 "toolchain": build.rust_version(), 561 "version": build.service_version(), 562 }, 563 "configuration": { 564 "digest": hex::encode(digest.as_bytes()), 565 "schema": MYC_CONFIG_SCHEMA, 566 "schema_version": MYC_CONFIG_SCHEMA_VERSION, 567 "source": if runtime.profile() == MycBootstrapProfileV1::RepoLocal { 568 "derived_repo_local" 569 } else { 570 "explicit_config" 571 }, 572 }, 573 "contract_version": MYC_SIGNER_STATUS_CONTRACT_VERSION, 574 "instance": runtime.context().instance().as_str(), 575 "myc": { 576 "connection_counts": snapshot.connection_counts, 577 "discovery": unavailable(discovery_configured), 578 "outbox": snapshot.outbox, 579 "transport": unavailable(true), 580 "user": unavailable(true), 581 }, 582 "persistence": { 583 "generation": generation, 584 "health": "read_only", 585 "integrity": integrity, 586 "reason_codes": [], 587 "schema_version": schema_version, 588 }, 589 "phase": "unready", 590 "provider": { 591 "discovery": unavailable(discovery_configured), 592 "health": "unavailable", 593 "reason_codes": ["signer_provider_unavailable"], 594 "transport": unavailable(true), 595 "user": unavailable(true), 596 }, 597 "ready": false, 598 "reason_codes": [ 599 "required_relay_unavailable", 600 "signer_provider_unavailable", 601 "subscriber_not_active" 602 ], 603 "service": "myc", 604 "transport": { 605 "connected_relay_count": 0, 606 "health": "unavailable", 607 "reason_codes": ["required_relay_unavailable", "subscriber_not_active"], 608 "required_relays_ready": false, 609 }, 610 "uptime_millis": 0, 611 })) 612 } 613 614 async fn offline_backup( 615 runtime: &MycRuntimeContext, 616 configuration: &MycConfigDocumentV1, 617 arguments: &MycStateBackupArgsV1, 618 ) -> ProcessResult<radroots_service_sqlite::ServiceBackupManifest> { 619 let build = migration_build_identity()?; 620 let applied_at = migration_time()?; 621 let state = open_myc_state_read_write_from_config(runtime, configuration, applied_at, &build) 622 .await 623 .map_err(|_| state_failure())?; 624 let captured = async { 625 let generation = state 626 .repository() 627 .current_configuration_generation() 628 .await 629 .map_err(|_| state_failure())?; 630 if u64::from(generation) != arguments.expected_generation() { 631 return Err(conflict_failure()); 632 } 633 let created_at = backup_time()?; 634 state 635 .capture_online_backup(arguments.target(), created_at) 636 .await 637 .map_err(|_| state_failure()) 638 } 639 .await; 640 let closed = state.close().await.map_err(|_| state_failure()); 641 closed.and(captured) 642 } 643 644 async fn offline_restore( 645 runtime: &MycRuntimeContext, 646 configuration: &MycConfigDocumentV1, 647 arguments: &MycStateRestoreArgsV1, 648 ) -> ProcessResult<()> { 649 let manifest = read_bounded_file(arguments.manifest(), BACKUP_MANIFEST_CANONICAL_MAX_BYTES)?; 650 let parsed = 651 ServiceBackupManifest::from_canonical_bytes(&manifest).map_err(|_| state_failure())?; 652 if parsed.digest() != arguments.manifest_sha256() { 653 return Err(state_failure()); 654 } 655 let expected = MycStateMetadata::new( 656 runtime, 657 configuration, 658 parsed.source_generation(), 659 parsed.created_at_unix_ms().get(), 660 ) 661 .map_err(|_| state_failure())?; 662 let verified = verify_myc_state_backup( 663 &manifest, 664 arguments.manifest_sha256(), 665 arguments.bundle(), 666 &expected, 667 arguments.maximum_state_bytes(), 668 ) 669 .map_err(|_| state_failure())?; 670 let staged = stage_myc_state_restore(runtime, &expected, verified) 671 .await 672 .map_err(|_| state_failure())?; 673 finalize_myc_state_restore(staged) 674 .await 675 .map_err(|_| state_failure()) 676 } 677 678 async fn offline_verify( 679 runtime: &MycRuntimeContext, 680 configuration: &MycConfigDocumentV1, 681 ) -> ProcessResult<()> { 682 let build = migration_build_identity()?; 683 let applied_at = migration_time()?; 684 let state = open_myc_state_read_write_from_config(runtime, configuration, applied_at, &build) 685 .await 686 .map_err(|_| state_failure())?; 687 let report = state 688 .inspect_integrity(integrity_time()?) 689 .await 690 .map_err(|_| state_failure()); 691 let closed = state.close().await.map_err(|_| state_failure()); 692 match closed.and(report) { 693 Ok(report) 694 if report.sqlite() == IntegrityCheckOutcome::Verified 695 && report.foreign_keys() == IntegrityCheckOutcome::Verified => 696 { 697 Ok(()) 698 } 699 Err(error) => Err(error), 700 _ => Err(state_failure()), 701 } 702 } 703 704 async fn offline_identity( 705 runtime: &MycRuntimeContext, 706 configuration: &MycConfigDocumentV1, 707 role: MycProviderRole, 708 ) -> ProcessResult<(Option<String>, u16, MycProviderKind, bool)> { 709 let state = open_myc_state_inspection_from_config(runtime, configuration) 710 .await 711 .map_err(|_| state_failure())?; 712 let generation = state 713 .repository() 714 .current_configuration_generation() 715 .await 716 .map_err(|_| state_failure()); 717 let closed = state.close().await.map_err(|_| state_failure()); 718 let generation = closed.and(generation)?; 719 let binding = configuration 720 .provider_contract() 721 .binding(role) 722 .ok_or_else(input_failure)?; 723 match binding.kind() { 724 MycProviderKind::EncryptedFile => { 725 let credential = 726 resolve_myc_wrapping_credential(runtime, binding).map_err(|_| state_failure())?; 727 let identity = 728 open_myc_encrypted_identity(binding, &credential).map_err(|_| state_failure())?; 729 Ok(( 730 Some(identity.public_identity().as_hex().to_owned()), 731 generation, 732 binding.kind(), 733 true, 734 )) 735 } 736 MycProviderKind::LocalSigner => { 737 let _client = MycLocalSignerClient::new(binding).map_err(|_| state_failure())?; 738 Ok((None, generation, binding.kind(), false)) 739 } 740 } 741 } 742 743 fn identity_binding( 744 configuration: &MycConfigDocumentV1, 745 arguments: MycIdentityCommandArgsV1, 746 ) -> ProcessResult<&crate::MycProviderBinding> { 747 configuration 748 .provider_contract() 749 .binding(arguments.role()) 750 .ok_or_else(input_failure) 751 } 752 753 async fn validate_candidate_providers( 754 runtime: &MycRuntimeContext, 755 candidate: &MycConfigDocumentV1, 756 ) -> ProcessResult<()> { 757 let cancellation = crate::MycTaskCancellation::uncancelled(); 758 let executor = 759 crate::provider_executor::MycProviderExecutor::open(runtime, candidate, &cancellation) 760 .await 761 .map_err(|_| state_failure())?; 762 executor 763 .probe_all(wall_time_millis()?, provider_probe_seed()?, &cancellation) 764 .await 765 .map_err(|_| state_failure()) 766 } 767 768 fn provider_probe_seed() -> ProcessResult<[u8; 32]> { 769 let mut seed = [0_u8; 32]; 770 SystemEntropy 771 .fill_bytes(&mut seed) 772 .map_err(|_| state_failure())?; 773 Ok(seed) 774 } 775 776 fn migration_build_identity() -> ProcessResult<MigrationBuildIdentity> { 777 let build = build_info()?; 778 let versions = build.contract_versions(); 779 MigrationBuildIdentity::new( 780 build.service_version(), 781 build.service_commit(), 782 build.lib_revision(), 783 build.rust_version(), 784 build.target(), 785 build.feature_profile(), 786 versions.config(), 787 versions.state(), 788 versions.admin(), 789 versions.status(), 790 versions.provider(), 791 ) 792 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal)) 793 } 794 795 fn build_info() -> ProcessResult<radroots_service_host::BuildInfo> { 796 let versions = ContractVersions::new( 797 MYC_CONFIG_SCHEMA_VERSION, 798 MYC_STATE_SCHEMA_VERSION, 799 MYC_OPERATOR_CONTRACT_VERSION, 800 MYC_SIGNER_STATUS_CONTRACT_VERSION, 801 MYC_PROVIDER_CONTRACT_VERSION, 802 ) 803 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal))?; 804 radroots_service_host::compile_time_build_info!( 805 feature_profile: "service-host", 806 contract_versions: versions, 807 ) 808 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal)) 809 } 810 811 fn source_generation() -> ProcessResult<SourceGeneration> { 812 for _ in 0..4 { 813 let mut bytes = [0_u8; 32]; 814 SystemEntropy 815 .fill_bytes(&mut bytes) 816 .map_err(|_| state_failure())?; 817 if let Ok(generation) = SourceGeneration::new(bytes) { 818 return Ok(generation); 819 } 820 } 821 Err(state_failure()) 822 } 823 824 fn wall_time_seconds() -> ProcessResult<u64> { 825 SystemWallClock 826 .now_utc() 827 .map(|time| time.get()) 828 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal)) 829 } 830 831 fn wall_time_millis() -> ProcessResult<u64> { 832 wall_time_seconds()? 833 .checked_mul(1_000) 834 .filter(|value| *value <= i64::MAX as u64) 835 .ok_or(ProcessFailure(MycProcessResult::UnexpectedInternal)) 836 } 837 838 fn migration_time() -> ProcessResult<MigrationAppliedAtUnixSeconds> { 839 MigrationAppliedAtUnixSeconds::new(wall_time_seconds()?) 840 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal)) 841 } 842 843 fn backup_time() -> ProcessResult<BackupCreatedAtUnixMs> { 844 BackupCreatedAtUnixMs::new(wall_time_millis()?) 845 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal)) 846 } 847 848 fn integrity_time() -> ProcessResult<IntegrityCheckedAtUnixMs> { 849 IntegrityCheckedAtUnixMs::new(wall_time_millis()?) 850 .ok_or(ProcessFailure(MycProcessResult::UnexpectedInternal)) 851 } 852 853 fn build_tokio_runtime( 854 limits: crate::MycRuntimeThreadLimitsV1, 855 ) -> ProcessResult<tokio::runtime::Runtime> { 856 tokio::runtime::Builder::new_multi_thread() 857 .worker_threads(limits.worker_threads()) 858 .max_blocking_threads(limits.blocking_threads()) 859 .enable_all() 860 .build() 861 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal)) 862 } 863 864 fn host_environment() -> RadrootsHostEnvironment { 865 let path = |name| { 866 env::var_os(name) 867 .filter(|value| !value.is_empty()) 868 .map(PathBuf::from) 869 }; 870 RadrootsHostEnvironment { 871 home_dir: path("HOME"), 872 xdg_config_home: path("XDG_CONFIG_HOME"), 873 xdg_data_home: path("XDG_DATA_HOME"), 874 xdg_state_home: path("XDG_STATE_HOME"), 875 xdg_cache_home: path("XDG_CACHE_HOME"), 876 xdg_runtime_dir: path("XDG_RUNTIME_DIR"), 877 appdata_dir: path("APPDATA"), 878 localappdata_dir: path("LOCALAPPDATA"), 879 } 880 } 881 882 fn read_bounded_stdin(maximum: usize) -> ProcessResult<Vec<u8>> { 883 let mut reader = std::io::stdin().lock(); 884 let mut bytes = Vec::with_capacity(maximum.min(64 * 1_024).saturating_add(1)); 885 Read::by_ref(&mut reader) 886 .take(u64::try_from(maximum).unwrap_or(u64::MAX).saturating_add(1)) 887 .read_to_end(&mut bytes) 888 .map_err(|_| input_failure())?; 889 if bytes.len() > maximum { 890 return Err(input_failure()); 891 } 892 Ok(bytes) 893 } 894 895 #[cfg(any(target_os = "linux", target_os = "macos"))] 896 fn read_bounded_file(path: &Path, maximum: usize) -> ProcessResult<Vec<u8>> { 897 crate::config_loader::read_secure_bounded_file(path, maximum).map_err(|_| state_failure()) 898 } 899 900 #[cfg(not(any(target_os = "linux", target_os = "macos")))] 901 fn read_bounded_file(_path: &Path, _maximum: usize) -> ProcessResult<Vec<u8>> { 902 Err(state_failure()) 903 } 904 905 fn emit_identity( 906 output: MycCliOutputModeV1, 907 role: MycProviderRole, 908 public_key: &str, 909 generation: u16, 910 ) -> ProcessResult<MycProcessResult> { 911 emit_admin_value( 912 output, 913 crate::MycAdminRoute::IdentityPublic, 914 json!({ 915 "generation": generation, 916 "public_key": public_key, 917 "role": role.as_str(), 918 }), 919 ) 920 } 921 922 fn emit_simple_success( 923 output: MycCliOutputModeV1, 924 code: &'static str, 925 ) -> ProcessResult<MycProcessResult> { 926 emit_value(output, code, json!({"ok": true})) 927 } 928 929 fn emit_admin_value( 930 _output: MycCliOutputModeV1, 931 route: crate::MycAdminRoute, 932 value: Value, 933 ) -> ProcessResult<MycProcessResult> { 934 let bytes = admit_admin_response_value(route, &value) 935 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal))?; 936 emit_bytes(&bytes) 937 } 938 939 fn emit_value( 940 output: MycCliOutputModeV1, 941 code: &'static str, 942 value: Value, 943 ) -> ProcessResult<MycProcessResult> { 944 let bytes = match output { 945 MycCliOutputModeV1::Json => serde_json::to_vec(&value) 946 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal))?, 947 MycCliOutputModeV1::Human => { 948 if value.is_object() && value.as_object().is_some_and(|object| object.len() > 1) { 949 serde_json::to_vec(&value) 950 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal))? 951 } else { 952 code.as_bytes().to_vec() 953 } 954 } 955 }; 956 emit_bytes(&bytes) 957 } 958 959 fn emit_bytes(bytes: &[u8]) -> ProcessResult<MycProcessResult> { 960 let mut stdout = std::io::stdout().lock(); 961 stdout 962 .write_all(bytes) 963 .and_then(|()| { 964 if bytes.ends_with(b"\n") { 965 Ok(()) 966 } else { 967 stdout.write_all(b"\n") 968 } 969 }) 970 .and_then(|()| stdout.flush()) 971 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal))?; 972 Ok(MycProcessResult::Success) 973 } 974 975 fn emit_exact_bytes(bytes: &[u8]) -> ProcessResult<MycProcessResult> { 976 let mut stdout = std::io::stdout().lock(); 977 stdout 978 .write_all(bytes) 979 .and_then(|()| stdout.flush()) 980 .map_err(|_| ProcessFailure(MycProcessResult::UnexpectedInternal))?; 981 Ok(MycProcessResult::Success) 982 } 983 984 const fn input_failure() -> ProcessFailure { 985 ProcessFailure(MycProcessResult::InputOrConfiguration) 986 } 987 988 const fn state_failure() -> ProcessFailure { 989 ProcessFailure(MycProcessResult::StateOrIdentityUnavailable) 990 } 991 992 const fn conflict_failure() -> ProcessFailure { 993 ProcessFailure(MycProcessResult::OperationRejectedOrConflict) 994 } 995 996 #[cfg(test)] 997 mod tests { 998 #[test] 999 fn runtime_limits_are_used_without_cpu_derived_defaults() { 1000 let source = include_str!("process_v1.rs") 1001 .split("#[cfg(test)]") 1002 .next() 1003 .expect("production source"); 1004 assert_eq!(source.matches("Builder::new_multi_thread()").count(), 1); 1005 assert!(!source.contains("available_parallelism")); 1006 assert!(source.contains("worker_threads(limits.worker_threads())")); 1007 assert!(source.contains("max_blocking_threads(limits.blocking_threads())")); 1008 } 1009 1010 #[test] 1011 fn host_environment_uses_only_standard_path_inputs() { 1012 let source = include_str!("process_v1.rs"); 1013 assert!(!source.contains("path(\"MYC_")); 1014 for name in [ 1015 "HOME", 1016 "XDG_CONFIG_HOME", 1017 "XDG_DATA_HOME", 1018 "XDG_STATE_HOME", 1019 "XDG_CACHE_HOME", 1020 "XDG_RUNTIME_DIR", 1021 "APPDATA", 1022 "LOCALAPPDATA", 1023 ] { 1024 assert!(source.contains(&format!("path(\"{name}\")"))); 1025 } 1026 } 1027 }