state_admin.rs (61599B)
1 //! Bounded durable idempotency for permissioned Myc admin mutations. 2 3 use core::{fmt, future::Future, pin::Pin}; 4 use std::error::Error; 5 6 use radroots_service_sqlite::{ 7 ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, 8 }; 9 use sha2::{Digest, Sha256}; 10 use sqlx::Row; 11 12 use crate::{ 13 MycAdminRequestDocument, MycAdminResponseDocument, MycAdminRoute, MycStateRepository, 14 state_repository::{PersistedMetadata, RepositoryOperationError, require_expected_metadata}, 15 }; 16 17 /// Maximum encoded length of a durable admin operation identifier. 18 pub const MYC_ADMIN_OPERATION_ID_MAX_BYTES: usize = 128; 19 /// Maximum canonical response-model bytes retained for replay. 20 pub const MYC_ADMIN_OPERATION_RESPONSE_MODEL_MAX_BYTES: usize = 8_192; 21 /// Maximum encoded success envelope for an 8,192-byte model and 128-byte correlation ID. 22 pub const MYC_ADMIN_OPERATION_RESPONSE_ENVELOPE_MAX_UTF8_BYTES: u32 = 8_382; 23 /// Maximum retained completed operations after expiry pruning. 24 pub const MYC_ADMIN_OPERATION_COMPLETED_LIMIT: u16 = 4_096; 25 /// Maximum retained operations whose external outcome is unresolved. 26 pub const MYC_ADMIN_OPERATION_PREPARED_LIMIT: u8 = 128; 27 /// Minimum configurable completed-response retention. 28 pub const MYC_ADMIN_OPERATION_MIN_RETENTION_MS: u64 = 1; 29 /// Maximum configurable completed-response retention. 30 pub const MYC_ADMIN_OPERATION_MAX_RETENTION_MS: u64 = 31_536_000_000; 31 /// Frozen seven-day completed-response retention. 32 pub const MYC_ADMIN_OPERATION_DEFAULT_RETENTION_MS: u64 = 604_800_000; 33 34 const REQUEST_DIGEST_DOMAIN: &[u8] = b"radroots.myc.admin_operation_request.v1\0"; 35 const PRUNE_LIMIT: i64 = 4_096; 36 37 const PRUNE_EXPIRED_SQL: &str = r#"DELETE FROM myc_admin_operations 38 WHERE operation_id IN ( 39 SELECT operation_id FROM myc_admin_operations 40 WHERE state = 'completed' AND expires_at_unix_ms <= ? 41 ORDER BY expires_at_unix_ms, operation_id 42 LIMIT ? 43 )"#; 44 45 const READ_OPERATION_SQL: &str = r#"SELECT 46 CASE WHEN typeof(route) = 'text' AND length(CAST(route AS BLOB)) BETWEEN 1 AND 128 47 THEN route ELSE NULL END AS route, 48 CASE WHEN typeof(request_sha256) = 'blob' AND length(request_sha256) = 32 49 THEN request_sha256 ELSE NULL END AS request_sha256, 50 CASE WHEN typeof(state) = 'text' AND length(CAST(state AS BLOB)) <= 16 51 THEN state ELSE NULL END AS state, 52 CASE WHEN typeof(response_model) = 'blob' AND length(response_model) BETWEEN 1 AND 8192 53 THEN response_model ELSE NULL END AS response_model, 54 typeof(response_model) AS response_model_type, 55 CASE WHEN typeof(response_sha256) = 'blob' AND length(response_sha256) = 32 56 THEN response_sha256 ELSE NULL END AS response_sha256, 57 typeof(response_sha256) AS response_sha256_type, 58 prepared_at_unix_ms, 59 completed_at_unix_ms, 60 typeof(completed_at_unix_ms) AS completed_at_type, 61 expires_at_unix_ms, 62 typeof(expires_at_unix_ms) AS expires_at_type 63 FROM myc_admin_operations 64 WHERE operation_id = ? 65 LIMIT 2"#; 66 67 const READ_COUNTS_SQL: &str = r#"SELECT 68 COUNT(CASE WHEN state = 'completed' THEN 1 END) AS completed_count, 69 COUNT(CASE WHEN state = 'prepared' THEN 1 END) AS prepared_count 70 FROM myc_admin_operations"#; 71 72 const INSERT_PREPARED_SQL: &str = r#"INSERT INTO myc_admin_operations ( 73 operation_id, route, request_sha256, state, prepared_at_unix_ms 74 ) VALUES (?, ?, ?, 'prepared', ?)"#; 75 76 const COMPLETE_OPERATION_SQL: &str = r#"UPDATE myc_admin_operations 77 SET state = 'completed', response_model = ?, response_sha256 = ?, 78 completed_at_unix_ms = ?, expires_at_unix_ms = ? 79 WHERE operation_id = ? AND state = 'prepared'"#; 80 81 /// Stable source-free admin-journal failure classes. 82 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 83 pub enum MycAdminOperationErrorKind { 84 InvalidMode, 85 InvalidInput, 86 OperationConflict, 87 OperationOutcomeUnknown, 88 ResourceExhausted, 89 Binding, 90 Transaction, 91 CommitOutcomeUnknown, 92 } 93 94 impl MycAdminOperationErrorKind { 95 /// Returns the stable machine-readable code. 96 #[must_use] 97 pub const fn code(self) -> &'static str { 98 match self { 99 Self::InvalidMode => "admin_operation_mode_invalid", 100 Self::InvalidInput => "admin_operation_input_invalid", 101 Self::OperationConflict => "operation_conflict", 102 Self::OperationOutcomeUnknown => "operation_outcome_unknown", 103 Self::ResourceExhausted => "resource_exhausted", 104 Self::Binding => "admin_operation_binding_invalid", 105 Self::Transaction => "admin_operation_transaction_failed", 106 Self::CommitOutcomeUnknown => "admin_operation_commit_outcome_unknown", 107 } 108 } 109 } 110 111 /// Redacted admin-journal error. 112 #[derive(Clone, Copy, PartialEq, Eq)] 113 pub struct MycAdminOperationError { 114 kind: MycAdminOperationErrorKind, 115 } 116 117 impl MycAdminOperationError { 118 const fn new(kind: MycAdminOperationErrorKind) -> Self { 119 Self { kind } 120 } 121 122 /// Returns the stable failure class. 123 #[must_use] 124 pub const fn kind(self) -> MycAdminOperationErrorKind { 125 self.kind 126 } 127 128 /// Returns the stable machine-readable code. 129 #[must_use] 130 pub const fn code(self) -> &'static str { 131 self.kind.code() 132 } 133 } 134 135 impl fmt::Display for MycAdminOperationError { 136 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 137 formatter.write_str(match self.kind { 138 MycAdminOperationErrorKind::InvalidMode => { 139 "Myc admin operation requires writable state" 140 } 141 MycAdminOperationErrorKind::InvalidInput => "Myc admin operation input is invalid", 142 MycAdminOperationErrorKind::OperationConflict => { 143 "Myc admin operation identity conflicts with retained evidence" 144 } 145 MycAdminOperationErrorKind::OperationOutcomeUnknown => { 146 "Myc admin operation outcome is unknown" 147 } 148 MycAdminOperationErrorKind::ResourceExhausted => { 149 "Myc admin operation capacity is exhausted" 150 } 151 MycAdminOperationErrorKind::Binding => "Myc admin operation journal binding is invalid", 152 MycAdminOperationErrorKind::Transaction => "Myc admin operation transaction failed", 153 MycAdminOperationErrorKind::CommitOutcomeUnknown => { 154 "Myc admin operation commit outcome is unknown" 155 } 156 }) 157 } 158 } 159 160 impl fmt::Debug for MycAdminOperationError { 161 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 162 formatter 163 .debug_struct("MycAdminOperationError") 164 .field("kind", &self.kind) 165 .finish() 166 } 167 } 168 169 impl Error for MycAdminOperationError {} 170 171 #[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] 172 struct AdminOperationIdBinding(Box<str>); 173 174 impl AdminOperationIdBinding { 175 fn new(value: &str) -> Result<Self, MycAdminOperationError> { 176 let bytes = value.as_bytes(); 177 let valid = !bytes.is_empty() 178 && bytes.len() <= MYC_ADMIN_OPERATION_ID_MAX_BYTES 179 && bytes[0].is_ascii_alphanumeric() 180 && bytes.iter().all(|byte| { 181 byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b':' | b'-') 182 }); 183 valid 184 .then(|| Self(value.into())) 185 .ok_or_else(|| MycAdminOperationError::new(MycAdminOperationErrorKind::InvalidInput)) 186 } 187 188 fn as_str(&self) -> &str { 189 &self.0 190 } 191 } 192 193 impl fmt::Debug for AdminOperationIdBinding { 194 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 195 formatter.write_str("AdminOperationIdBinding([redacted])") 196 } 197 } 198 199 /// Injected UTC millisecond evidence representable by SQLite. 200 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] 201 pub struct MycAdminOperationTimeUnixMs(u64); 202 203 impl MycAdminOperationTimeUnixMs { 204 /// Validates one UTC millisecond instant without reading ambient time. 205 pub fn new(value: u64) -> Result<Self, MycAdminOperationError> { 206 i64::try_from(value) 207 .map(|_| Self(value)) 208 .map_err(|_| MycAdminOperationError::new(MycAdminOperationErrorKind::InvalidInput)) 209 } 210 211 /// Returns the validated instant. 212 #[must_use] 213 pub const fn get(self) -> u64 { 214 self.0 215 } 216 217 fn sqlite_value(self) -> i64 { 218 i64::try_from(self.0).expect("validated admin operation time fits SQLite") 219 } 220 } 221 222 /// Explicit bounded completed-response retention policy. 223 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 224 pub struct MycAdminOperationJournalPolicy { 225 completed_retention_ms: u64, 226 } 227 228 impl MycAdminOperationJournalPolicy { 229 /// Admits the frozen inclusive retention range. 230 pub fn new(completed_retention_ms: u64) -> Result<Self, MycAdminOperationError> { 231 (MYC_ADMIN_OPERATION_MIN_RETENTION_MS..=MYC_ADMIN_OPERATION_MAX_RETENTION_MS) 232 .contains(&completed_retention_ms) 233 .then_some(Self { 234 completed_retention_ms, 235 }) 236 .ok_or_else(|| MycAdminOperationError::new(MycAdminOperationErrorKind::InvalidInput)) 237 } 238 239 /// Returns the exact seven-day policy. 240 #[must_use] 241 pub const fn seven_days() -> Self { 242 Self { 243 completed_retention_ms: MYC_ADMIN_OPERATION_DEFAULT_RETENTION_MS, 244 } 245 } 246 247 /// Returns the admitted retention duration. 248 #[must_use] 249 pub const fn completed_retention_ms(self) -> u64 { 250 self.completed_retention_ms 251 } 252 } 253 254 /// Sealed evidence that one external or cross-resource mutation is unresolved. 255 pub struct MycPreparedAdminOperation { 256 operation_id: AdminOperationIdBinding, 257 route: MycAdminRoute, 258 request_sha256: [u8; 32], 259 prepared_at: MycAdminOperationTimeUnixMs, 260 } 261 262 impl fmt::Debug for MycPreparedAdminOperation { 263 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 264 formatter 265 .debug_struct("MycPreparedAdminOperation") 266 .field("route", &self.route) 267 .field("identity", &"[redacted]") 268 .finish() 269 } 270 } 271 272 /// Result of mutation admission after bounded expiry pruning. 273 pub enum MycAdminOperationAdmission { 274 Prepared(MycPreparedAdminOperation), 275 ExactReplay(MycAdminResponseDocument), 276 } 277 278 impl fmt::Debug for MycAdminOperationAdmission { 279 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 280 formatter.write_str(match self { 281 Self::Prepared(_) => "MycAdminOperationAdmission::Prepared([redacted])", 282 Self::ExactReplay(_) => "MycAdminOperationAdmission::ExactReplay([redacted])", 283 }) 284 } 285 } 286 287 /// Result of completing a previously prepared operation. 288 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 289 pub enum MycAdminOperationCompletion { 290 Completed, 291 ExactReplay, 292 } 293 294 impl MycStateRepository<'_> { 295 /// Atomically commits one SQLite-only admin mutation and its replay receipt. 296 /// 297 /// The supplied operation runs inside the same governed transaction that 298 /// admits the operation identifier and records the canonical response. A 299 /// domain failure therefore leaves neither a prepared journal row nor a 300 /// partial domain effect. 301 pub(crate) async fn execute_database_admin_operation<F>( 302 &self, 303 request: &MycAdminRequestDocument, 304 completed_at: MycAdminOperationTimeUnixMs, 305 policy: MycAdminOperationJournalPolicy, 306 operation: F, 307 ) -> Result<MycAdminResponseDocument, MycAdminOperationError> 308 where 309 F: for<'a, 'b> FnOnce( 310 &'a mut ServiceSqliteTransaction<'b>, 311 ) -> AdminDatabaseOperationFuture<'a> 312 + Send 313 + 'static, 314 { 315 if !self.is_writable() { 316 return Err(MycAdminOperationError::new( 317 MycAdminOperationErrorKind::InvalidMode, 318 )); 319 } 320 let binding = AdminRequestBinding::from_document(request)?; 321 let expires_at = completed_at 322 .get() 323 .checked_add(policy.completed_retention_ms()) 324 .filter(|value| i64::try_from(*value).is_ok()) 325 .ok_or_else(|| MycAdminOperationError::new(MycAdminOperationErrorKind::InvalidInput))?; 326 let expected = PersistedMetadata::from(self.expected()); 327 self.host() 328 .transaction(move |transaction| { 329 Box::pin(async move { 330 require_expected_metadata(transaction, &expected) 331 .await 332 .map_err(AdminJournalOperationError::from)?; 333 let prepared = match prepare_operation(transaction, &binding, completed_at) 334 .await? 335 { 336 MycAdminOperationAdmission::ExactReplay(response) => return Ok(response), 337 MycAdminOperationAdmission::Prepared(prepared) => prepared, 338 }; 339 let response = operation(transaction).await?; 340 if response.route() != binding.route 341 || response.canonical_bytes().is_empty() 342 || response.canonical_bytes().len() 343 > MYC_ADMIN_OPERATION_RESPONSE_MODEL_MAX_BYTES 344 { 345 return Err(AdminJournalOperationError::InvalidInput); 346 } 347 complete_operation( 348 transaction, 349 &PreparedBinding::from_prepared(&prepared), 350 response.canonical_bytes(), 351 completed_at, 352 expires_at, 353 ) 354 .await?; 355 Ok(response) 356 }) 357 }) 358 .await 359 .map_err(map_transaction_error) 360 } 361 362 pub(crate) async fn complete_prepared_database_admin_operation<F>( 363 &self, 364 prepared: &MycPreparedAdminOperation, 365 completed_at: MycAdminOperationTimeUnixMs, 366 policy: MycAdminOperationJournalPolicy, 367 operation: F, 368 ) -> Result<MycAdminResponseDocument, MycAdminOperationError> 369 where 370 F: for<'a, 'b> FnOnce( 371 &'a mut ServiceSqliteTransaction<'b>, 372 ) -> AdminDatabaseOperationFuture<'a> 373 + Send 374 + 'static, 375 { 376 if !self.is_writable() { 377 return Err(MycAdminOperationError::new( 378 MycAdminOperationErrorKind::InvalidMode, 379 )); 380 } 381 if completed_at < prepared.prepared_at { 382 return Err(MycAdminOperationError::new( 383 MycAdminOperationErrorKind::InvalidInput, 384 )); 385 } 386 let expires_at = completed_at 387 .get() 388 .checked_add(policy.completed_retention_ms()) 389 .filter(|value| i64::try_from(*value).is_ok()) 390 .ok_or_else(|| MycAdminOperationError::new(MycAdminOperationErrorKind::InvalidInput))?; 391 let binding = PreparedBinding::from_prepared(prepared); 392 let expected = PersistedMetadata::from(self.expected()); 393 self.host() 394 .transaction(move |transaction| { 395 Box::pin(async move { 396 require_expected_metadata(transaction, &expected) 397 .await 398 .map_err(AdminJournalOperationError::from)?; 399 let response = operation(transaction).await?; 400 if response.route() != binding.route 401 || response.canonical_bytes().is_empty() 402 || response.canonical_bytes().len() 403 > MYC_ADMIN_OPERATION_RESPONSE_MODEL_MAX_BYTES 404 { 405 return Err(AdminJournalOperationError::InvalidInput); 406 } 407 complete_operation( 408 transaction, 409 &binding, 410 response.canonical_bytes(), 411 completed_at, 412 expires_at, 413 ) 414 .await?; 415 Ok(response) 416 }) 417 }) 418 .await 419 .map_err(map_transaction_error) 420 } 421 422 /// Prunes a bounded expired prefix and admits or replays one mutation. 423 pub async fn prepare_admin_operation( 424 &self, 425 request: &MycAdminRequestDocument, 426 observed_at: MycAdminOperationTimeUnixMs, 427 ) -> Result<MycAdminOperationAdmission, MycAdminOperationError> { 428 if !self.is_writable() { 429 return Err(MycAdminOperationError::new( 430 MycAdminOperationErrorKind::InvalidMode, 431 )); 432 } 433 let binding = AdminRequestBinding::from_document(request)?; 434 let expected = PersistedMetadata::from(self.expected()); 435 self.host() 436 .transaction(move |transaction| { 437 Box::pin(async move { 438 require_expected_metadata(transaction, &expected) 439 .await 440 .map_err(AdminJournalOperationError::from)?; 441 prepare_operation(transaction, &binding, observed_at).await 442 }) 443 }) 444 .await 445 .map_err(map_transaction_error) 446 } 447 448 /// Completes one external mutation only after its durable effect exists. 449 pub async fn complete_admin_operation( 450 &self, 451 prepared: &MycPreparedAdminOperation, 452 response: &MycAdminResponseDocument, 453 completed_at: MycAdminOperationTimeUnixMs, 454 policy: MycAdminOperationJournalPolicy, 455 ) -> Result<MycAdminOperationCompletion, MycAdminOperationError> { 456 if !self.is_writable() { 457 return Err(MycAdminOperationError::new( 458 MycAdminOperationErrorKind::InvalidMode, 459 )); 460 } 461 if response.route() != prepared.route 462 || response.canonical_bytes().is_empty() 463 || response.canonical_bytes().len() > MYC_ADMIN_OPERATION_RESPONSE_MODEL_MAX_BYTES 464 || completed_at < prepared.prepared_at 465 { 466 return Err(MycAdminOperationError::new( 467 MycAdminOperationErrorKind::InvalidInput, 468 )); 469 } 470 let expires_at = completed_at 471 .get() 472 .checked_add(policy.completed_retention_ms()) 473 .filter(|value| i64::try_from(*value).is_ok()) 474 .ok_or_else(|| MycAdminOperationError::new(MycAdminOperationErrorKind::InvalidInput))?; 475 let binding = PreparedBinding::from_prepared(prepared); 476 let response = response.canonical_bytes().to_vec().into_boxed_slice(); 477 let expected = PersistedMetadata::from(self.expected()); 478 self.host() 479 .transaction(move |transaction| { 480 Box::pin(async move { 481 require_expected_metadata(transaction, &expected) 482 .await 483 .map_err(AdminJournalOperationError::from)?; 484 complete_operation(transaction, &binding, &response, completed_at, expires_at) 485 .await 486 }) 487 }) 488 .await 489 .map_err(map_transaction_error) 490 } 491 } 492 493 struct AdminRequestBinding { 494 operation_id: AdminOperationIdBinding, 495 route: MycAdminRoute, 496 request_sha256: [u8; 32], 497 } 498 499 impl AdminRequestBinding { 500 fn from_document(request: &MycAdminRequestDocument) -> Result<Self, MycAdminOperationError> { 501 if !request.route().is_mutation() { 502 return Err(MycAdminOperationError::new( 503 MycAdminOperationErrorKind::InvalidInput, 504 )); 505 } 506 let operation_id = request 507 .operation_id() 508 .ok_or_else(|| MycAdminOperationError::new(MycAdminOperationErrorKind::InvalidInput))?; 509 Ok(Self { 510 operation_id: AdminOperationIdBinding::new(operation_id)?, 511 route: request.route(), 512 request_sha256: request_digest(request), 513 }) 514 } 515 } 516 517 struct PreparedBinding { 518 operation_id: AdminOperationIdBinding, 519 route: MycAdminRoute, 520 request_sha256: [u8; 32], 521 prepared_at: MycAdminOperationTimeUnixMs, 522 } 523 524 impl PreparedBinding { 525 fn from_prepared(prepared: &MycPreparedAdminOperation) -> Self { 526 Self { 527 operation_id: prepared.operation_id.clone(), 528 route: prepared.route, 529 request_sha256: prepared.request_sha256, 530 prepared_at: prepared.prepared_at, 531 } 532 } 533 } 534 535 enum StoredOperation { 536 Prepared { 537 route: MycAdminRoute, 538 request_sha256: [u8; 32], 539 prepared_at: MycAdminOperationTimeUnixMs, 540 }, 541 Completed { 542 route: MycAdminRoute, 543 request_sha256: [u8; 32], 544 response: Box<[u8]>, 545 response_sha256: [u8; 32], 546 prepared_at: MycAdminOperationTimeUnixMs, 547 completed_at: MycAdminOperationTimeUnixMs, 548 expires_at: MycAdminOperationTimeUnixMs, 549 }, 550 } 551 552 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 553 pub(crate) enum AdminJournalOperationError { 554 InvalidInput, 555 Conflict, 556 OutcomeUnknown, 557 ResourceExhausted, 558 Binding, 559 Storage, 560 } 561 562 pub(crate) type AdminDatabaseOperationFuture<'a> = Pin< 563 Box< 564 dyn Future<Output = Result<MycAdminResponseDocument, AdminJournalOperationError>> 565 + Send 566 + 'a, 567 >, 568 >; 569 570 impl From<RepositoryOperationError> for AdminJournalOperationError { 571 fn from(error: RepositoryOperationError) -> Self { 572 match error { 573 RepositoryOperationError::Binding => Self::Binding, 574 RepositoryOperationError::Storage => Self::Storage, 575 } 576 } 577 } 578 579 async fn prepare_operation( 580 transaction: &mut ServiceSqliteTransaction<'_>, 581 binding: &AdminRequestBinding, 582 observed_at: MycAdminOperationTimeUnixMs, 583 ) -> Result<MycAdminOperationAdmission, AdminJournalOperationError> { 584 prune_expired(transaction, observed_at).await?; 585 if let Some(existing) = read_operation(transaction, &binding.operation_id).await? { 586 return match existing { 587 StoredOperation::Prepared { 588 route, 589 request_sha256, 590 .. 591 } if route == binding.route && request_sha256 == binding.request_sha256 => { 592 Err(AdminJournalOperationError::OutcomeUnknown) 593 } 594 StoredOperation::Completed { 595 route, 596 request_sha256, 597 response, 598 response_sha256, 599 .. 600 } if route == binding.route && request_sha256 == binding.request_sha256 => { 601 if sha256(&response) != response_sha256 { 602 return Err(AdminJournalOperationError::Binding); 603 } 604 MycAdminResponseDocument::from_canonical_bytes(route, &response) 605 .map(MycAdminOperationAdmission::ExactReplay) 606 .map_err(|_| AdminJournalOperationError::Binding) 607 } 608 StoredOperation::Prepared { .. } | StoredOperation::Completed { .. } => { 609 Err(AdminJournalOperationError::Conflict) 610 } 611 }; 612 } 613 let (completed, prepared) = read_counts(transaction).await?; 614 let reserved = completed 615 .checked_add(prepared) 616 .ok_or(AdminJournalOperationError::Binding)?; 617 if reserved >= u64::from(MYC_ADMIN_OPERATION_COMPLETED_LIMIT) 618 || prepared >= u64::from(MYC_ADMIN_OPERATION_PREPARED_LIMIT) 619 { 620 return Err(AdminJournalOperationError::ResourceExhausted); 621 } 622 let result = sqlx::query(INSERT_PREPARED_SQL) 623 .bind(binding.operation_id.as_str()) 624 .bind(binding.route.operation_id()) 625 .bind(binding.request_sha256.as_slice()) 626 .bind(observed_at.sqlite_value()) 627 .execute(&mut *transaction) 628 .await 629 .map_err(|_| AdminJournalOperationError::Storage)?; 630 require_one(result.rows_affected())?; 631 match read_operation(transaction, &binding.operation_id).await? { 632 Some(StoredOperation::Prepared { 633 route, 634 request_sha256, 635 prepared_at, 636 }) if route == binding.route 637 && request_sha256 == binding.request_sha256 638 && prepared_at == observed_at => 639 { 640 Ok(MycAdminOperationAdmission::Prepared( 641 MycPreparedAdminOperation { 642 operation_id: binding.operation_id.clone(), 643 route, 644 request_sha256, 645 prepared_at, 646 }, 647 )) 648 } 649 Some(_) | None => Err(AdminJournalOperationError::Binding), 650 } 651 } 652 653 async fn complete_operation( 654 transaction: &mut ServiceSqliteTransaction<'_>, 655 binding: &PreparedBinding, 656 response: &[u8], 657 completed_at: MycAdminOperationTimeUnixMs, 658 expires_at: u64, 659 ) -> Result<MycAdminOperationCompletion, AdminJournalOperationError> { 660 let response_sha256 = sha256(response); 661 match read_operation(transaction, &binding.operation_id).await? { 662 Some(StoredOperation::Completed { 663 route, 664 request_sha256, 665 response: existing_response, 666 response_sha256: existing_sha256, 667 .. 668 }) if route == binding.route 669 && request_sha256 == binding.request_sha256 670 && existing_response.as_ref() == response 671 && existing_sha256 == response_sha256 => 672 { 673 return Ok(MycAdminOperationCompletion::ExactReplay); 674 } 675 Some(StoredOperation::Completed { .. }) => { 676 return Err(AdminJournalOperationError::Conflict); 677 } 678 Some(StoredOperation::Prepared { 679 route, 680 request_sha256, 681 prepared_at, 682 }) if route == binding.route 683 && request_sha256 == binding.request_sha256 684 && prepared_at == binding.prepared_at => {} 685 Some(StoredOperation::Prepared { .. }) => { 686 return Err(AdminJournalOperationError::Conflict); 687 } 688 None => return Err(AdminJournalOperationError::Binding), 689 } 690 let (completed, _) = read_counts(transaction).await?; 691 if completed >= u64::from(MYC_ADMIN_OPERATION_COMPLETED_LIMIT) { 692 return Err(AdminJournalOperationError::ResourceExhausted); 693 } 694 let result = sqlx::query(COMPLETE_OPERATION_SQL) 695 .bind(response) 696 .bind(response_sha256.as_slice()) 697 .bind(completed_at.sqlite_value()) 698 .bind(i64::try_from(expires_at).map_err(|_| AdminJournalOperationError::InvalidInput)?) 699 .bind(binding.operation_id.as_str()) 700 .execute(&mut *transaction) 701 .await 702 .map_err(|_| AdminJournalOperationError::Storage)?; 703 require_one(result.rows_affected())?; 704 match read_operation(transaction, &binding.operation_id).await? { 705 Some(StoredOperation::Completed { 706 route, 707 request_sha256, 708 response: actual_response, 709 response_sha256: actual_sha256, 710 prepared_at, 711 completed_at: actual_completed_at, 712 expires_at: actual_expires_at, 713 }) if route == binding.route 714 && request_sha256 == binding.request_sha256 715 && actual_response.as_ref() == response 716 && actual_sha256 == response_sha256 717 && prepared_at == binding.prepared_at 718 && actual_completed_at == completed_at 719 && actual_expires_at.get() == expires_at => 720 { 721 Ok(MycAdminOperationCompletion::Completed) 722 } 723 Some(_) | None => Err(AdminJournalOperationError::Binding), 724 } 725 } 726 727 async fn prune_expired( 728 transaction: &mut ServiceSqliteTransaction<'_>, 729 observed_at: MycAdminOperationTimeUnixMs, 730 ) -> Result<(), AdminJournalOperationError> { 731 sqlx::query(PRUNE_EXPIRED_SQL) 732 .bind(observed_at.sqlite_value()) 733 .bind(PRUNE_LIMIT) 734 .execute(&mut *transaction) 735 .await 736 .map(|_| ()) 737 .map_err(|_| AdminJournalOperationError::Storage) 738 } 739 740 async fn read_counts( 741 transaction: &mut ServiceSqliteTransaction<'_>, 742 ) -> Result<(u64, u64), AdminJournalOperationError> { 743 let rows = sqlx::query(READ_COUNTS_SQL) 744 .fetch_all(&mut *transaction) 745 .await 746 .map_err(|_| AdminJournalOperationError::Storage)?; 747 if rows.len() != 1 { 748 return Err(AdminJournalOperationError::Binding); 749 } 750 let completed = rows[0] 751 .try_get::<i64, _>("completed_count") 752 .map_err(|_| AdminJournalOperationError::Binding)?; 753 let prepared = rows[0] 754 .try_get::<i64, _>("prepared_count") 755 .map_err(|_| AdminJournalOperationError::Binding)?; 756 Ok(( 757 u64::try_from(completed).map_err(|_| AdminJournalOperationError::Binding)?, 758 u64::try_from(prepared).map_err(|_| AdminJournalOperationError::Binding)?, 759 )) 760 } 761 762 async fn read_operation( 763 transaction: &mut ServiceSqliteTransaction<'_>, 764 operation_id: &AdminOperationIdBinding, 765 ) -> Result<Option<StoredOperation>, AdminJournalOperationError> { 766 let rows = sqlx::query(READ_OPERATION_SQL) 767 .bind(operation_id.as_str()) 768 .fetch_all(&mut *transaction) 769 .await 770 .map_err(|_| AdminJournalOperationError::Storage)?; 771 if rows.len() > 1 { 772 return Err(AdminJournalOperationError::Binding); 773 } 774 rows.first().map(decode_operation).transpose() 775 } 776 777 fn decode_operation( 778 row: &sqlx::sqlite::SqliteRow, 779 ) -> Result<StoredOperation, AdminJournalOperationError> { 780 let route = row 781 .try_get::<Option<&str>, _>("route") 782 .map_err(|_| AdminJournalOperationError::Binding)? 783 .and_then(parse_route) 784 .ok_or(AdminJournalOperationError::Binding)?; 785 let request_sha256 = exact_digest(row, "request_sha256")?; 786 let state = row 787 .try_get::<Option<&str>, _>("state") 788 .map_err(|_| AdminJournalOperationError::Binding)? 789 .ok_or(AdminJournalOperationError::Binding)?; 790 let prepared_at = time(row, "prepared_at_unix_ms")?; 791 match state { 792 "prepared" => { 793 require_null(row, "response_model_type")?; 794 require_null(row, "response_sha256_type")?; 795 require_null(row, "completed_at_type")?; 796 require_null(row, "expires_at_type")?; 797 Ok(StoredOperation::Prepared { 798 route, 799 request_sha256, 800 prepared_at, 801 }) 802 } 803 "completed" => { 804 require_type(row, "response_model_type", "blob")?; 805 require_type(row, "response_sha256_type", "blob")?; 806 require_type(row, "completed_at_type", "integer")?; 807 require_type(row, "expires_at_type", "integer")?; 808 let response = row 809 .try_get::<Option<Vec<u8>>, _>("response_model") 810 .map_err(|_| AdminJournalOperationError::Binding)? 811 .ok_or(AdminJournalOperationError::Binding)? 812 .into_boxed_slice(); 813 Ok(StoredOperation::Completed { 814 route, 815 request_sha256, 816 response, 817 response_sha256: exact_digest(row, "response_sha256")?, 818 prepared_at, 819 completed_at: time(row, "completed_at_unix_ms")?, 820 expires_at: time(row, "expires_at_unix_ms")?, 821 }) 822 } 823 _ => Err(AdminJournalOperationError::Binding), 824 } 825 } 826 827 fn request_digest(request: &MycAdminRequestDocument) -> [u8; 32] { 828 let mut hasher = Sha256::new(); 829 hasher.update(REQUEST_DIGEST_DOMAIN); 830 hash_field(&mut hasher, request.route().operation_id().as_bytes()); 831 match request.parameter_binding() { 832 Some((name, value)) => { 833 hasher.update([1]); 834 hash_field(&mut hasher, name.as_bytes()); 835 hash_field(&mut hasher, value.as_bytes()); 836 } 837 None => hasher.update([0]), 838 } 839 hash_field(&mut hasher, request.model_bytes()); 840 hasher.finalize().into() 841 } 842 843 fn hash_field(hasher: &mut Sha256, bytes: &[u8]) { 844 hasher.update( 845 u64::try_from(bytes.len()) 846 .expect("bounded field length") 847 .to_be_bytes(), 848 ); 849 hasher.update(bytes); 850 } 851 852 fn sha256(bytes: &[u8]) -> [u8; 32] { 853 Sha256::digest(bytes).into() 854 } 855 856 fn parse_route(value: &str) -> Option<MycAdminRoute> { 857 MycAdminRoute::ALL 858 .into_iter() 859 .find(|route| route.is_mutation() && route.operation_id() == value) 860 } 861 862 fn exact_digest( 863 row: &sqlx::sqlite::SqliteRow, 864 column: &str, 865 ) -> Result<[u8; 32], AdminJournalOperationError> { 866 row.try_get::<Option<Vec<u8>>, _>(column) 867 .map_err(|_| AdminJournalOperationError::Binding)? 868 .ok_or(AdminJournalOperationError::Binding)? 869 .try_into() 870 .map_err(|_| AdminJournalOperationError::Binding) 871 } 872 873 fn time( 874 row: &sqlx::sqlite::SqliteRow, 875 column: &str, 876 ) -> Result<MycAdminOperationTimeUnixMs, AdminJournalOperationError> { 877 let value = row 878 .try_get::<i64, _>(column) 879 .map_err(|_| AdminJournalOperationError::Binding)?; 880 MycAdminOperationTimeUnixMs::new( 881 u64::try_from(value).map_err(|_| AdminJournalOperationError::Binding)?, 882 ) 883 .map_err(|_| AdminJournalOperationError::Binding) 884 } 885 886 fn require_null( 887 row: &sqlx::sqlite::SqliteRow, 888 column: &str, 889 ) -> Result<(), AdminJournalOperationError> { 890 require_type(row, column, "null") 891 } 892 893 fn require_type( 894 row: &sqlx::sqlite::SqliteRow, 895 column: &str, 896 expected: &str, 897 ) -> Result<(), AdminJournalOperationError> { 898 (row.try_get::<&str, _>(column) 899 .map_err(|_| AdminJournalOperationError::Binding)? 900 == expected) 901 .then_some(()) 902 .ok_or(AdminJournalOperationError::Binding) 903 } 904 905 fn require_one(rows: u64) -> Result<(), AdminJournalOperationError> { 906 (rows == 1) 907 .then_some(()) 908 .ok_or(AdminJournalOperationError::Storage) 909 } 910 911 fn map_transaction_error( 912 error: ServiceSqliteTransactionError<AdminJournalOperationError>, 913 ) -> MycAdminOperationError { 914 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 915 return MycAdminOperationError::new(MycAdminOperationErrorKind::CommitOutcomeUnknown); 916 } 917 let kind = match error.operation_error() { 918 Some(AdminJournalOperationError::InvalidInput) => MycAdminOperationErrorKind::InvalidInput, 919 Some(AdminJournalOperationError::Conflict) => MycAdminOperationErrorKind::OperationConflict, 920 Some(AdminJournalOperationError::OutcomeUnknown) => { 921 MycAdminOperationErrorKind::OperationOutcomeUnknown 922 } 923 Some(AdminJournalOperationError::ResourceExhausted) => { 924 MycAdminOperationErrorKind::ResourceExhausted 925 } 926 Some(AdminJournalOperationError::Binding) => MycAdminOperationErrorKind::Binding, 927 Some(AdminJournalOperationError::Storage) | None => MycAdminOperationErrorKind::Transaction, 928 }; 929 MycAdminOperationError::new(kind) 930 } 931 932 #[cfg(test)] 933 mod tests { 934 use std::{fs, os::unix::fs::PermissionsExt, path::Path}; 935 936 use radroots_service_sqlite::{MigrationAppliedAtUnixSeconds, MigrationBuildIdentity}; 937 use radroots_storage::event::SourceGeneration; 938 939 use super::*; 940 use crate::{ 941 MycConfigProfile, MycRuntimeContext, MycStateHost, MycStateMetadata, 942 RadrootsHostEnvironment, RadrootsPathResolver, RadrootsPlatform, initialize_myc_state, 943 open_myc_state_read_write, parse_myc_cli_v1_from, parse_myc_config_v1, 944 resolve_myc_runtime_context, 945 }; 946 947 const CONFIG: &[u8] = include_bytes!("../contracts/services_hardening/config.v1.example.toml"); 948 949 fn runtime(root: &Path) -> MycRuntimeContext { 950 let invocation = parse_myc_cli_v1_from([ 951 "myc", 952 "--profile", 953 "repo-local", 954 "--instance", 955 "primary", 956 "--repo-local-root", 957 root.to_str().expect("UTF-8 root"), 958 "run", 959 ]) 960 .expect("invocation"); 961 resolve_myc_runtime_context( 962 &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()), 963 &invocation, 964 ) 965 .expect("runtime") 966 } 967 968 fn migration_build() -> MigrationBuildIdentity { 969 MigrationBuildIdentity::new( 970 env!("CARGO_PKG_VERSION"), 971 "1111111111111111111111111111111111111111", 972 "053d0c750bf9cd683c6ea37cefe7e79617ba629f", 973 "rustc-test", 974 "test-target", 975 "service-host", 976 1, 977 crate::MYC_STATE_SCHEMA_VERSION, 978 1, 979 1, 980 1, 981 ) 982 .expect("build") 983 } 984 985 async fn fixture() -> ( 986 tempfile::TempDir, 987 MycRuntimeContext, 988 MycStateMetadata, 989 MycStateHost, 990 ) { 991 let directory = tempfile::tempdir().expect("root"); 992 let runtime = runtime(directory.path()); 993 fs::create_dir_all(runtime.context().paths().state()).expect("state directory"); 994 fs::set_permissions( 995 runtime.context().paths().state(), 996 fs::Permissions::from_mode(0o700), 997 ) 998 .expect("state mode"); 999 let configuration = 1000 parse_myc_config_v1(CONFIG, MycConfigProfile::RepoLocal).expect("configuration"); 1001 let metadata = MycStateMetadata::new( 1002 &runtime, 1003 &configuration, 1004 SourceGeneration::new([0x5a; 32]).expect("generation"), 1005 1_725_000_000_000, 1006 ) 1007 .expect("metadata"); 1008 let applied_at = MigrationAppliedAtUnixSeconds::new(1_725_000_000).expect("time"); 1009 let build = migration_build(); 1010 initialize_myc_state(&runtime, &metadata, applied_at, &build) 1011 .await 1012 .expect("initialize"); 1013 let host = open_myc_state_read_write(&runtime, &metadata, applied_at, &build) 1014 .await 1015 .expect("open"); 1016 (directory, runtime, metadata, host) 1017 } 1018 1019 fn request( 1020 operation_id: &str, 1021 connection_id: &str, 1022 generation: u64, 1023 ) -> MycAdminRequestDocument { 1024 let model = format!( 1025 "{{\"confirmation\":\"approve\",\"expected_generation\":{generation},\"permissions\":\"nip04_decrypt\"}}" 1026 ); 1027 MycAdminRequestDocument::mutation_for_test( 1028 MycAdminRoute::ConnectionApprove, 1029 operation_id, 1030 Some(("connection_id", connection_id)), 1031 model.as_bytes(), 1032 ) 1033 } 1034 1035 fn response_document(operation_id: &str, generation: u64) -> MycAdminResponseDocument { 1036 let model = format!( 1037 "{{\"connection_id\":\"connection-1\",\"current_state\":\"approved\",\"generation\":{generation},\"operation_id\":\"{operation_id}\",\"previous_state\":\"pending\"}}" 1038 ); 1039 MycAdminResponseDocument::from_canonical_bytes( 1040 MycAdminRoute::ConnectionApprove, 1041 model.as_bytes(), 1042 ) 1043 .expect("response") 1044 } 1045 1046 #[test] 1047 fn identifier_policy_time_and_public_diagnostics_are_bounded_and_redacted() { 1048 for valid in [ 1049 "a", 1050 "A0._:-z", 1051 &"x".repeat(MYC_ADMIN_OPERATION_ID_MAX_BYTES), 1052 ] { 1053 assert!(AdminOperationIdBinding::new(valid).is_ok(), "{valid}"); 1054 } 1055 for invalid in ["", "-first", "space value", "slash/value", "é"] { 1056 assert!(AdminOperationIdBinding::new(invalid).is_err(), "{invalid}"); 1057 } 1058 assert!( 1059 AdminOperationIdBinding::new(&"x".repeat(MYC_ADMIN_OPERATION_ID_MAX_BYTES + 1)) 1060 .is_err() 1061 ); 1062 let identifier = AdminOperationIdBinding::new("protected-operation").expect("ID"); 1063 assert_eq!( 1064 format!("{identifier:?}"), 1065 "AdminOperationIdBinding([redacted])" 1066 ); 1067 1068 assert!(MycAdminOperationJournalPolicy::new(0).is_err()); 1069 assert!(MycAdminOperationJournalPolicy::new(MYC_ADMIN_OPERATION_MIN_RETENTION_MS).is_ok()); 1070 assert!(MycAdminOperationJournalPolicy::new(MYC_ADMIN_OPERATION_MAX_RETENTION_MS).is_ok()); 1071 assert!( 1072 MycAdminOperationJournalPolicy::new(MYC_ADMIN_OPERATION_MAX_RETENTION_MS + 1).is_err() 1073 ); 1074 assert_eq!( 1075 MycAdminOperationJournalPolicy::seven_days().completed_retention_ms(), 1076 MYC_ADMIN_OPERATION_DEFAULT_RETENTION_MS 1077 ); 1078 assert!(MycAdminOperationTimeUnixMs::new(i64::MAX as u64).is_ok()); 1079 assert!(MycAdminOperationTimeUnixMs::new(i64::MAX as u64 + 1).is_err()); 1080 assert_eq!(PRUNE_LIMIT, i64::from(MYC_ADMIN_OPERATION_COMPLETED_LIMIT)); 1081 1082 for kind in [ 1083 MycAdminOperationErrorKind::InvalidMode, 1084 MycAdminOperationErrorKind::InvalidInput, 1085 MycAdminOperationErrorKind::OperationConflict, 1086 MycAdminOperationErrorKind::OperationOutcomeUnknown, 1087 MycAdminOperationErrorKind::ResourceExhausted, 1088 MycAdminOperationErrorKind::Binding, 1089 MycAdminOperationErrorKind::Transaction, 1090 MycAdminOperationErrorKind::CommitOutcomeUnknown, 1091 ] { 1092 let error = MycAdminOperationError::new(kind); 1093 let rendered = format!("{error} {error:?}"); 1094 assert!(!rendered.contains("protected-operation")); 1095 assert!(!rendered.contains("/tmp/secret")); 1096 assert!(Error::source(&error).is_none()); 1097 } 1098 } 1099 1100 #[test] 1101 fn request_digest_binds_route_parameter_and_canonical_model_without_retaining_them() { 1102 let first = request("digest-1", "connection-1", 1); 1103 let same = request("digest-1", "connection-1", 1); 1104 let changed_parameter = request("digest-1", "connection-2", 1); 1105 let changed_model = request("digest-1", "connection-1", 2); 1106 assert_eq!(request_digest(&first), request_digest(&same)); 1107 assert_ne!(request_digest(&first), request_digest(&changed_parameter)); 1108 assert_ne!(request_digest(&first), request_digest(&changed_model)); 1109 } 1110 1111 #[tokio::test] 1112 async fn prepare_complete_replay_conflict_and_expiry_are_exact() { 1113 let (_directory, _runtime, _metadata, host) = fixture().await; 1114 let repository = host.repository(); 1115 let first = request("journal-1", "connection-1", 1); 1116 let prepared = match repository 1117 .prepare_admin_operation(&first, MycAdminOperationTimeUnixMs::new(10).unwrap()) 1118 .await 1119 .expect("prepare") 1120 { 1121 MycAdminOperationAdmission::Prepared(prepared) => prepared, 1122 MycAdminOperationAdmission::ExactReplay(_) => panic!("unexpected replay"), 1123 }; 1124 let same_prepared = repository 1125 .prepare_admin_operation(&first, MycAdminOperationTimeUnixMs::new(10).unwrap()) 1126 .await 1127 .expect_err("retained Prepared is ambiguous"); 1128 assert_eq!( 1129 same_prepared.kind(), 1130 MycAdminOperationErrorKind::OperationOutcomeUnknown 1131 ); 1132 let changed = request("journal-1", "connection-2", 1); 1133 assert_eq!( 1134 repository 1135 .prepare_admin_operation(&changed, MycAdminOperationTimeUnixMs::new(10).unwrap()) 1136 .await 1137 .expect_err("path binding conflict") 1138 .kind(), 1139 MycAdminOperationErrorKind::OperationConflict 1140 ); 1141 1142 let response = response_document("journal-1", 1); 1143 assert_eq!( 1144 repository 1145 .complete_admin_operation( 1146 &prepared, 1147 &response, 1148 MycAdminOperationTimeUnixMs::new(20).unwrap(), 1149 MycAdminOperationJournalPolicy::new(2).unwrap(), 1150 ) 1151 .await 1152 .expect("complete"), 1153 MycAdminOperationCompletion::Completed 1154 ); 1155 assert_eq!( 1156 repository 1157 .complete_admin_operation( 1158 &prepared, 1159 &response, 1160 MycAdminOperationTimeUnixMs::new(21).unwrap(), 1161 MycAdminOperationJournalPolicy::new(2).unwrap(), 1162 ) 1163 .await 1164 .expect("idempotent completion"), 1165 MycAdminOperationCompletion::ExactReplay 1166 ); 1167 let different_response = response_document("journal-1", 2); 1168 assert_eq!( 1169 repository 1170 .complete_admin_operation( 1171 &prepared, 1172 &different_response, 1173 MycAdminOperationTimeUnixMs::new(21).unwrap(), 1174 MycAdminOperationJournalPolicy::new(2).unwrap(), 1175 ) 1176 .await 1177 .expect_err("different completion conflicts") 1178 .kind(), 1179 MycAdminOperationErrorKind::OperationConflict 1180 ); 1181 1182 match repository 1183 .prepare_admin_operation(&first, MycAdminOperationTimeUnixMs::new(21).unwrap()) 1184 .await 1185 .expect("replay before expiry") 1186 { 1187 MycAdminOperationAdmission::ExactReplay(replayed) => { 1188 assert_eq!(replayed.canonical_bytes(), response.canonical_bytes()); 1189 } 1190 MycAdminOperationAdmission::Prepared(_) => panic!("unexpected prepare"), 1191 } 1192 assert!(matches!( 1193 repository 1194 .prepare_admin_operation(&first, MycAdminOperationTimeUnixMs::new(22).unwrap()) 1195 .await 1196 .expect("exact expiry prunes before admission"), 1197 MycAdminOperationAdmission::Prepared(_) 1198 )); 1199 host.close().await.expect("close"); 1200 } 1201 1202 #[tokio::test] 1203 async fn database_only_admin_effect_and_receipt_share_one_transaction() { 1204 let (_directory, _runtime, _metadata, host) = fixture().await; 1205 let repository = host.repository(); 1206 let request = request("atomic-admin-1", "connection-1", 1); 1207 let error = repository 1208 .execute_database_admin_operation( 1209 &request, 1210 MycAdminOperationTimeUnixMs::new(100).expect("time"), 1211 MycAdminOperationJournalPolicy::seven_days(), 1212 |transaction| { 1213 Box::pin(async move { 1214 sqlx::query( 1215 r#"INSERT INTO connection_rate_windows ( 1216 rate_kind, subject_scope, subject_sha256, 1217 window_started_at_unix_ms, window_ends_at_unix_ms, 1218 accepted_count, rejected_count, lifetime_accepted_count, 1219 lifetime_rejected_count, last_observed_at_unix_ms, 1220 retention_expires_at_unix_ms 1221 ) VALUES ('connection_admission', 'global', ?, 1, 2, 0, 0, 0, 0, 1, 2)"#, 1222 ) 1223 .bind([0xaa_u8; 32].as_slice()) 1224 .execute(&mut *transaction) 1225 .await 1226 .map_err(|_| AdminJournalOperationError::Storage)?; 1227 Err(AdminJournalOperationError::Binding) 1228 }) 1229 }, 1230 ) 1231 .await 1232 .expect_err("domain failure rolls back"); 1233 assert_eq!(error.kind(), MycAdminOperationErrorKind::Binding); 1234 assert_eq!( 1235 repository 1236 .host() 1237 .transaction(|transaction| { 1238 Box::pin(async move { 1239 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM connection_rate_windows") 1240 .fetch_one(&mut *transaction) 1241 .await 1242 }) 1243 }) 1244 .await 1245 .expect("rolled-back effect"), 1246 0 1247 ); 1248 1249 let expected = response_document("atomic-admin-1", 1); 1250 let expected_bytes = expected.canonical_bytes().to_vec(); 1251 let committed = repository 1252 .execute_database_admin_operation( 1253 &request, 1254 MycAdminOperationTimeUnixMs::new(101).expect("time"), 1255 MycAdminOperationJournalPolicy::seven_days(), 1256 move |transaction| { 1257 let expected_bytes = expected_bytes.clone(); 1258 Box::pin(async move { 1259 sqlx::query( 1260 r#"INSERT INTO connection_rate_windows ( 1261 rate_kind, subject_scope, subject_sha256, 1262 window_started_at_unix_ms, window_ends_at_unix_ms, 1263 accepted_count, rejected_count, lifetime_accepted_count, 1264 lifetime_rejected_count, last_observed_at_unix_ms, 1265 retention_expires_at_unix_ms 1266 ) VALUES ('connection_admission', 'global', ?, 3, 4, 0, 0, 0, 0, 3, 4)"#, 1267 ) 1268 .bind([0xbb_u8; 32].as_slice()) 1269 .execute(&mut *transaction) 1270 .await 1271 .map_err(|_| AdminJournalOperationError::Storage)?; 1272 MycAdminResponseDocument::from_canonical_bytes( 1273 MycAdminRoute::ConnectionApprove, 1274 &expected_bytes, 1275 ) 1276 .map_err(|_| AdminJournalOperationError::Binding) 1277 }) 1278 }, 1279 ) 1280 .await 1281 .expect("atomic commit"); 1282 let replay = repository 1283 .execute_database_admin_operation( 1284 &request, 1285 MycAdminOperationTimeUnixMs::new(102).expect("time"), 1286 MycAdminOperationJournalPolicy::seven_days(), 1287 |_transaction| { 1288 Box::pin( 1289 async move { panic!("exact replay must not execute the domain operation") }, 1290 ) 1291 }, 1292 ) 1293 .await 1294 .expect("exact replay"); 1295 assert_eq!(replay.canonical_bytes(), committed.canonical_bytes()); 1296 assert_eq!( 1297 repository 1298 .host() 1299 .transaction(|transaction| { 1300 Box::pin(async move { 1301 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM connection_rate_windows") 1302 .fetch_one(&mut *transaction) 1303 .await 1304 }) 1305 }) 1306 .await 1307 .expect("single committed effect"), 1308 1 1309 ); 1310 host.close().await.expect("close"); 1311 } 1312 1313 #[tokio::test] 1314 async fn prepared_backup_shape_is_ambiguous_and_persists_no_path_or_request_content() { 1315 let (_directory, _runtime, _metadata, host) = fixture().await; 1316 let request = MycAdminRequestDocument::mutation_for_test( 1317 MycAdminRoute::StateBackup, 1318 "backup-1", 1319 None, 1320 br#"{"destination":"/tmp/never-store-this"}"#, 1321 ); 1322 let repository = host.repository(); 1323 assert!(matches!( 1324 repository 1325 .prepare_admin_operation(&request, MycAdminOperationTimeUnixMs::new(30).unwrap()) 1326 .await 1327 .expect("prepare backup"), 1328 MycAdminOperationAdmission::Prepared(_) 1329 )); 1330 assert_eq!( 1331 repository 1332 .prepare_admin_operation(&request, MycAdminOperationTimeUnixMs::new(31).unwrap()) 1333 .await 1334 .expect_err("backup outcome remains unknown") 1335 .kind(), 1336 MycAdminOperationErrorKind::OperationOutcomeUnknown 1337 ); 1338 let row = repository 1339 .host() 1340 .transaction(|transaction| { 1341 Box::pin(async move { 1342 sqlx::query( 1343 "SELECT route, state, response_model, request_sha256, \ 1344 (SELECT group_concat(name, ',') FROM pragma_table_info('myc_admin_operations')) \ 1345 AS columns FROM myc_admin_operations WHERE operation_id = 'backup-1'", 1346 ) 1347 .fetch_one(&mut *transaction) 1348 .await 1349 }) 1350 }) 1351 .await 1352 .expect("inspect journal"); 1353 assert_eq!( 1354 row.get::<String, _>("route"), 1355 MycAdminRoute::StateBackup.operation_id() 1356 ); 1357 assert_eq!(row.get::<String, _>("state"), "prepared"); 1358 assert!(row.get::<Option<Vec<u8>>, _>("response_model").is_none()); 1359 assert_eq!(row.get::<Vec<u8>, _>("request_sha256").len(), 32); 1360 let columns = row.get::<String, _>("columns"); 1361 for forbidden in [ 1362 "path", 1363 "body", 1364 "correlation", 1365 "credential", 1366 "secret", 1367 "bundle", 1368 ] { 1369 assert!(!columns.contains(forbidden), "{columns}"); 1370 } 1371 host.close().await.expect("close"); 1372 } 1373 1374 #[tokio::test] 1375 async fn exact_completed_and_prepared_caps_fail_closed_after_bounded_pruning() { 1376 let (_directory, _runtime, _metadata, host) = fixture().await; 1377 let repository = host.repository(); 1378 repository 1379 .host() 1380 .transaction(|transaction| { 1381 Box::pin(async move { 1382 let response = b"{}"; 1383 let digest = sha256(response); 1384 for index in 0..(MYC_ADMIN_OPERATION_COMPLETED_LIMIT - 1) { 1385 sqlx::query( 1386 "INSERT INTO myc_admin_operations (operation_id, route, \ 1387 request_sha256, state, response_model, response_sha256, \ 1388 prepared_at_unix_ms, completed_at_unix_ms, expires_at_unix_ms) \ 1389 VALUES (?, ?, ?, 'completed', ?, ?, 1, 1, ?)", 1390 ) 1391 .bind(format!("completed-{index}")) 1392 .bind(MycAdminRoute::ConnectionApprove.operation_id()) 1393 .bind([0x11; 32].as_slice()) 1394 .bind(response.as_slice()) 1395 .bind(digest.as_slice()) 1396 .bind(i64::MAX) 1397 .execute(&mut *transaction) 1398 .await?; 1399 } 1400 Ok::<_, sqlx::Error>(()) 1401 }) 1402 }) 1403 .await 1404 .expect("seed completed capacity"); 1405 let reserved = match repository 1406 .prepare_admin_operation( 1407 &request("reserved-completion", "connection-1", 1), 1408 MycAdminOperationTimeUnixMs::new(2).unwrap(), 1409 ) 1410 .await 1411 .expect("reserve final completed slot") 1412 { 1413 MycAdminOperationAdmission::Prepared(prepared) => prepared, 1414 MycAdminOperationAdmission::ExactReplay(_) => panic!("unexpected replay"), 1415 }; 1416 assert_eq!( 1417 repository 1418 .prepare_admin_operation( 1419 &request("new-completed", "connection-1", 1), 1420 MycAdminOperationTimeUnixMs::new(2).unwrap(), 1421 ) 1422 .await 1423 .expect_err("completed cap") 1424 .kind(), 1425 MycAdminOperationErrorKind::ResourceExhausted 1426 ); 1427 assert_eq!( 1428 repository 1429 .complete_admin_operation( 1430 &reserved, 1431 &response_document("reserved-completion", 1), 1432 MycAdminOperationTimeUnixMs::new(3).unwrap(), 1433 MycAdminOperationJournalPolicy::seven_days(), 1434 ) 1435 .await 1436 .expect("consume reserved completed slot"), 1437 MycAdminOperationCompletion::Completed 1438 ); 1439 assert_eq!( 1440 repository 1441 .prepare_admin_operation( 1442 &request("completed-cap", "connection-1", 1), 1443 MycAdminOperationTimeUnixMs::new(4).unwrap(), 1444 ) 1445 .await 1446 .expect_err("completed cap remains closed") 1447 .kind(), 1448 MycAdminOperationErrorKind::ResourceExhausted 1449 ); 1450 host.close().await.expect("close completed fixture"); 1451 1452 let (_directory, _runtime, _metadata, host) = fixture().await; 1453 let repository = host.repository(); 1454 repository 1455 .host() 1456 .transaction(|transaction| { 1457 Box::pin(async move { 1458 for index in 0..MYC_ADMIN_OPERATION_PREPARED_LIMIT { 1459 sqlx::query( 1460 "INSERT INTO myc_admin_operations (operation_id, route, \ 1461 request_sha256, state, prepared_at_unix_ms) \ 1462 VALUES (?, ?, ?, 'prepared', 1)", 1463 ) 1464 .bind(format!("prepared-{index}")) 1465 .bind(MycAdminRoute::StateBackup.operation_id()) 1466 .bind([0x22; 32].as_slice()) 1467 .execute(&mut *transaction) 1468 .await?; 1469 } 1470 Ok::<_, sqlx::Error>(()) 1471 }) 1472 }) 1473 .await 1474 .expect("seed prepared capacity"); 1475 assert_eq!( 1476 repository 1477 .prepare_admin_operation( 1478 &request("new-prepared", "connection-1", 1), 1479 MycAdminOperationTimeUnixMs::new(2).unwrap(), 1480 ) 1481 .await 1482 .expect_err("prepared cap") 1483 .kind(), 1484 MycAdminOperationErrorKind::ResourceExhausted 1485 ); 1486 host.close().await.expect("close prepared fixture"); 1487 } 1488 1489 #[tokio::test] 1490 async fn schema_admits_exact_response_model_cap_and_rejects_one_byte_over() { 1491 let (_directory, _runtime, _metadata, host) = fixture().await; 1492 let repository = host.repository(); 1493 let exact = vec![b'x'; MYC_ADMIN_OPERATION_RESPONSE_MODEL_MAX_BYTES]; 1494 let exact_digest = sha256(&exact); 1495 repository 1496 .host() 1497 .transaction(|transaction| { 1498 Box::pin(async move { 1499 sqlx::query( 1500 "INSERT INTO myc_admin_operations (operation_id, route, \ 1501 request_sha256, state, response_model, response_sha256, \ 1502 prepared_at_unix_ms, completed_at_unix_ms, expires_at_unix_ms) \ 1503 VALUES ('response-exact', ?, ?, 'completed', ?, ?, 1, 1, 2)", 1504 ) 1505 .bind(MycAdminRoute::ConnectionApprove.operation_id()) 1506 .bind([0x33; 32].as_slice()) 1507 .bind(exact) 1508 .bind(exact_digest.as_slice()) 1509 .execute(&mut *transaction) 1510 .await 1511 .map(|_| ()) 1512 }) 1513 }) 1514 .await 1515 .expect("exact response cap"); 1516 1517 let excessive = vec![b'x'; MYC_ADMIN_OPERATION_RESPONSE_MODEL_MAX_BYTES + 1]; 1518 let excessive_digest = sha256(&excessive); 1519 assert!( 1520 repository 1521 .host() 1522 .transaction(|transaction| { 1523 Box::pin(async move { 1524 sqlx::query( 1525 "INSERT INTO myc_admin_operations (operation_id, route, \ 1526 request_sha256, state, response_model, response_sha256, \ 1527 prepared_at_unix_ms, completed_at_unix_ms, expires_at_unix_ms) \ 1528 VALUES ('response-excessive', ?, ?, 'completed', ?, ?, 1, 1, 2)", 1529 ) 1530 .bind(MycAdminRoute::ConnectionApprove.operation_id()) 1531 .bind([0x44; 32].as_slice()) 1532 .bind(excessive) 1533 .bind(excessive_digest.as_slice()) 1534 .execute(&mut *transaction) 1535 .await 1536 .map(|_| ()) 1537 }) 1538 }) 1539 .await 1540 .is_err() 1541 ); 1542 host.close().await.expect("close"); 1543 } 1544 1545 #[test] 1546 fn schema_and_source_freeze_response_and_pruning_bounds() { 1547 const SOURCE: &str = include_str!("state_admin.rs"); 1548 const CATALOG: &str = include_str!("state_catalog.rs"); 1549 let production = SOURCE 1550 .split("#[cfg(test)]") 1551 .next() 1552 .expect("production source"); 1553 assert!(SOURCE.contains("length(response_model) BETWEEN 1 AND 8192")); 1554 assert!(SOURCE.contains("LIMIT ?")); 1555 assert!(CATALOG.contains("length(response_model) BETWEEN 1 AND 8192")); 1556 assert!(CATALOG.contains("CHECK (length(request_sha256) = 32)")); 1557 assert!(!CATALOG.contains("bundle_path")); 1558 assert!(!production.contains("correlation_id")); 1559 assert!(parse_route(MycAdminRoute::Status.operation_id()).is_none()); 1560 assert_eq!( 1561 parse_route(MycAdminRoute::StateBackup.operation_id()), 1562 Some(MycAdminRoute::StateBackup) 1563 ); 1564 let maximum_envelope = br#"{"contract_version":1,"ok":true,"correlation_id":""#.len() 1565 + radroots_service_host::ADMIN_CORRELATION_ID_MAX_UTF8_BYTES 1566 + br#"","result":"#.len() 1567 + MYC_ADMIN_OPERATION_RESPONSE_MODEL_MAX_BYTES 1568 + 1; 1569 assert_eq!( 1570 maximum_envelope, 1571 MYC_ADMIN_OPERATION_RESPONSE_ENVELOPE_MAX_UTF8_BYTES as usize 1572 ); 1573 } 1574 }