state_admin.rs (29466B)
1 //! Bounded durable idempotency for permissioned Rhi 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 RhiAdminRequestDocument, RhiAdminResponseDocument, RhiAdminRoute, RhiStateHost, 14 RhiStateHostMode, 15 }; 16 17 /// Maximum encoded length of a durable admin operation identifier. 18 pub const RHI_ADMIN_OPERATION_ID_MAX_BYTES: usize = 128; 19 /// Maximum canonical response-model bytes retained for replay. 20 pub const RHI_ADMIN_OPERATION_RESPONSE_MODEL_MAX_BYTES: usize = 8_192; 21 /// Maximum retained completed operations after expiry pruning. 22 pub const RHI_ADMIN_OPERATION_COMPLETED_LIMIT: u16 = 4_096; 23 /// Maximum retained operations whose external outcome is unresolved. 24 pub const RHI_ADMIN_OPERATION_PREPARED_LIMIT: u8 = 128; 25 /// Frozen seven-day completed-response retention. 26 pub const RHI_ADMIN_OPERATION_DEFAULT_RETENTION_MS: u64 = 604_800_000; 27 28 const REQUEST_DIGEST_DOMAIN: &[u8] = b"radroots.rhi.admin_operation_request.v1\0"; 29 const PRUNE_LIMIT: i64 = 4_096; 30 31 const PRUNE_EXPIRED_SQL: &str = r#"DELETE FROM rhi_admin_operations 32 WHERE operation_id IN ( 33 SELECT operation_id FROM rhi_admin_operations 34 WHERE state = 'completed' AND expires_at_unix_ms <= ? 35 ORDER BY expires_at_unix_ms, operation_id 36 LIMIT ? 37 )"#; 38 39 const READ_OPERATION_SQL: &str = r#"SELECT 40 CASE WHEN typeof(route) = 'text' AND length(CAST(route AS BLOB)) BETWEEN 1 AND 128 41 THEN route ELSE NULL END AS route, 42 CASE WHEN typeof(request_sha256) = 'blob' AND length(request_sha256) = 32 43 THEN request_sha256 ELSE NULL END AS request_sha256, 44 CASE WHEN typeof(state) = 'text' AND length(CAST(state AS BLOB)) <= 16 45 THEN state ELSE NULL END AS state, 46 CASE WHEN typeof(response_model) = 'blob' AND length(response_model) BETWEEN 1 AND 8192 47 THEN response_model ELSE NULL END AS response_model, 48 typeof(response_model) AS response_model_type, 49 CASE WHEN typeof(response_sha256) = 'blob' AND length(response_sha256) = 32 50 THEN response_sha256 ELSE NULL END AS response_sha256, 51 typeof(response_sha256) AS response_sha256_type, 52 prepared_at_unix_ms, 53 completed_at_unix_ms, 54 typeof(completed_at_unix_ms) AS completed_at_type, 55 expires_at_unix_ms, 56 typeof(expires_at_unix_ms) AS expires_at_type 57 FROM rhi_admin_operations 58 WHERE operation_id = ? 59 LIMIT 2"#; 60 61 const READ_COUNTS_SQL: &str = r#"SELECT 62 COUNT(CASE WHEN state = 'completed' THEN 1 END) AS completed_count, 63 COUNT(CASE WHEN state = 'prepared' THEN 1 END) AS prepared_count 64 FROM rhi_admin_operations"#; 65 66 const INSERT_PREPARED_SQL: &str = r#"INSERT INTO rhi_admin_operations ( 67 operation_id, route, request_sha256, state, prepared_at_unix_ms 68 ) VALUES (?, ?, ?, 'prepared', ?)"#; 69 70 const COMPLETE_OPERATION_SQL: &str = r#"UPDATE rhi_admin_operations 71 SET state = 'completed', response_model = ?, response_sha256 = ?, 72 completed_at_unix_ms = ?, expires_at_unix_ms = ? 73 WHERE operation_id = ? AND state = 'prepared'"#; 74 75 /// Stable source-free admin-journal failure classes. 76 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 77 pub enum RhiAdminOperationErrorKind { 78 InvalidMode, 79 InvalidInput, 80 OperationConflict, 81 OperationOutcomeUnknown, 82 ResourceExhausted, 83 Binding, 84 Transaction, 85 CommitOutcomeUnknown, 86 } 87 88 /// Redacted admin-journal error. 89 #[derive(Clone, Copy, PartialEq, Eq)] 90 pub struct RhiAdminOperationError { 91 kind: RhiAdminOperationErrorKind, 92 } 93 94 impl RhiAdminOperationError { 95 const fn new(kind: RhiAdminOperationErrorKind) -> Self { 96 Self { kind } 97 } 98 99 /// Returns the stable failure class. 100 #[must_use] 101 pub const fn kind(self) -> RhiAdminOperationErrorKind { 102 self.kind 103 } 104 } 105 106 impl fmt::Display for RhiAdminOperationError { 107 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 108 formatter.write_str(match self.kind { 109 RhiAdminOperationErrorKind::InvalidMode => { 110 "RHI admin operation requires writable state" 111 } 112 RhiAdminOperationErrorKind::InvalidInput => "RHI admin operation input is invalid", 113 RhiAdminOperationErrorKind::OperationConflict => { 114 "RHI admin operation identity conflicts with retained evidence" 115 } 116 RhiAdminOperationErrorKind::OperationOutcomeUnknown => { 117 "RHI admin operation outcome is unknown" 118 } 119 RhiAdminOperationErrorKind::ResourceExhausted => { 120 "RHI admin operation capacity is exhausted" 121 } 122 RhiAdminOperationErrorKind::Binding => "RHI admin operation journal binding is invalid", 123 RhiAdminOperationErrorKind::Transaction => "RHI admin operation transaction failed", 124 RhiAdminOperationErrorKind::CommitOutcomeUnknown => { 125 "RHI admin operation commit outcome is unknown" 126 } 127 }) 128 } 129 } 130 131 impl fmt::Debug for RhiAdminOperationError { 132 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 133 formatter 134 .debug_struct("RhiAdminOperationError") 135 .field("kind", &self.kind) 136 .finish() 137 } 138 } 139 140 impl Error for RhiAdminOperationError {} 141 142 #[derive(Clone, PartialEq, Eq, PartialOrd, Ord, Hash)] 143 struct AdminOperationIdBinding(Box<str>); 144 145 impl AdminOperationIdBinding { 146 fn new(value: &str) -> Result<Self, RhiAdminOperationError> { 147 let bytes = value.as_bytes(); 148 let valid = !bytes.is_empty() 149 && bytes.len() <= RHI_ADMIN_OPERATION_ID_MAX_BYTES 150 && bytes[0].is_ascii_alphanumeric() 151 && bytes.iter().all(|byte| { 152 byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b':' | b'-') 153 }); 154 valid 155 .then(|| Self(value.into())) 156 .ok_or_else(|| RhiAdminOperationError::new(RhiAdminOperationErrorKind::InvalidInput)) 157 } 158 159 fn as_str(&self) -> &str { 160 &self.0 161 } 162 } 163 164 impl fmt::Debug for AdminOperationIdBinding { 165 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 166 formatter.write_str("AdminOperationIdBinding([redacted])") 167 } 168 } 169 170 /// Injected UTC millisecond evidence representable by SQLite. 171 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] 172 pub struct RhiAdminOperationTimeUnixMs(u64); 173 174 impl RhiAdminOperationTimeUnixMs { 175 /// Validates one UTC millisecond instant without reading ambient time. 176 pub fn new(value: u64) -> Result<Self, RhiAdminOperationError> { 177 i64::try_from(value) 178 .map(|_| Self(value)) 179 .map_err(|_| RhiAdminOperationError::new(RhiAdminOperationErrorKind::InvalidInput)) 180 } 181 182 /// Returns the validated instant. 183 #[must_use] 184 pub const fn get(self) -> u64 { 185 self.0 186 } 187 188 fn sqlite_value(self) -> i64 { 189 i64::try_from(self.0).expect("validated admin operation time fits SQLite") 190 } 191 } 192 193 /// Explicit bounded completed-response retention policy. 194 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 195 pub struct RhiAdminOperationJournalPolicy { 196 completed_retention_ms: u64, 197 } 198 199 impl RhiAdminOperationJournalPolicy { 200 /// Returns the exact seven-day policy. 201 #[must_use] 202 pub const fn seven_days() -> Self { 203 Self { 204 completed_retention_ms: RHI_ADMIN_OPERATION_DEFAULT_RETENTION_MS, 205 } 206 } 207 208 /// Returns the admitted retention duration. 209 #[must_use] 210 pub const fn completed_retention_ms(self) -> u64 { 211 self.completed_retention_ms 212 } 213 } 214 215 /// Sealed evidence that one external or cross-resource mutation is unresolved. 216 pub struct RhiPreparedAdminOperation { 217 operation_id: AdminOperationIdBinding, 218 route: RhiAdminRoute, 219 request_sha256: [u8; 32], 220 prepared_at: RhiAdminOperationTimeUnixMs, 221 } 222 223 impl fmt::Debug for RhiPreparedAdminOperation { 224 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 225 formatter 226 .debug_struct("RhiPreparedAdminOperation") 227 .field("route", &self.route) 228 .field("identity", &"[redacted]") 229 .finish() 230 } 231 } 232 233 /// Result of mutation admission after bounded expiry pruning. 234 pub enum RhiAdminOperationAdmission { 235 Prepared(RhiPreparedAdminOperation), 236 ExactReplay(RhiAdminResponseDocument), 237 } 238 239 impl fmt::Debug for RhiAdminOperationAdmission { 240 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 241 formatter.write_str(match self { 242 Self::Prepared(_) => "RhiAdminOperationAdmission::Prepared([redacted])", 243 Self::ExactReplay(_) => "RhiAdminOperationAdmission::ExactReplay([redacted])", 244 }) 245 } 246 } 247 248 /// Result of completing a previously prepared operation. 249 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 250 pub enum RhiAdminOperationCompletion { 251 Completed, 252 ExactReplay, 253 } 254 255 pub(crate) struct RhiAdminOperationRepository<'host> { 256 host: &'host RhiStateHost, 257 } 258 259 impl<'host> RhiAdminOperationRepository<'host> { 260 pub(crate) const fn new(host: &'host RhiStateHost) -> Self { 261 Self { host } 262 } 263 264 /// Atomically commits one SQLite-only admin mutation and its replay receipt. 265 pub(crate) async fn execute_database_admin_operation<F>( 266 &self, 267 request: &RhiAdminRequestDocument, 268 completed_at: RhiAdminOperationTimeUnixMs, 269 policy: RhiAdminOperationJournalPolicy, 270 operation: F, 271 ) -> Result<RhiAdminResponseDocument, RhiAdminOperationError> 272 where 273 F: for<'a, 'b> FnOnce( 274 &'a mut ServiceSqliteTransaction<'b>, 275 ) -> AdminDatabaseOperationFuture<'a> 276 + Send 277 + 'static, 278 { 279 if self.host.mode() != RhiStateHostMode::ReadWriteExisting { 280 return Err(RhiAdminOperationError::new( 281 RhiAdminOperationErrorKind::InvalidMode, 282 )); 283 } 284 let binding = AdminRequestBinding::from_document(request)?; 285 let expires_at = completed_at 286 .get() 287 .checked_add(policy.completed_retention_ms()) 288 .filter(|value| i64::try_from(*value).is_ok()) 289 .ok_or_else(|| RhiAdminOperationError::new(RhiAdminOperationErrorKind::InvalidInput))?; 290 self.host 291 .sqlite_host() 292 .transaction(move |transaction| { 293 Box::pin(async move { 294 let prepared = match prepare_operation(transaction, &binding, completed_at) 295 .await? 296 { 297 RhiAdminOperationAdmission::ExactReplay(response) => return Ok(response), 298 RhiAdminOperationAdmission::Prepared(prepared) => prepared, 299 }; 300 let response = operation(transaction).await?; 301 if response.route() != binding.route 302 || response.canonical_bytes().is_empty() 303 || response.canonical_bytes().len() 304 > RHI_ADMIN_OPERATION_RESPONSE_MODEL_MAX_BYTES 305 { 306 return Err(AdminJournalOperationError::InvalidInput); 307 } 308 complete_operation( 309 transaction, 310 &PreparedBinding::from_prepared(&prepared), 311 response.canonical_bytes(), 312 completed_at, 313 expires_at, 314 ) 315 .await?; 316 Ok(response) 317 }) 318 }) 319 .await 320 .map_err(map_transaction_error) 321 } 322 323 /// Prunes a bounded expired prefix and admits or replays one mutation. 324 pub async fn prepare_admin_operation( 325 &self, 326 request: &RhiAdminRequestDocument, 327 observed_at: RhiAdminOperationTimeUnixMs, 328 ) -> Result<RhiAdminOperationAdmission, RhiAdminOperationError> { 329 if self.host.mode() != RhiStateHostMode::ReadWriteExisting { 330 return Err(RhiAdminOperationError::new( 331 RhiAdminOperationErrorKind::InvalidMode, 332 )); 333 } 334 let binding = AdminRequestBinding::from_document(request)?; 335 self.host 336 .sqlite_host() 337 .transaction(move |transaction| { 338 Box::pin(async move { prepare_operation(transaction, &binding, observed_at).await }) 339 }) 340 .await 341 .map_err(map_transaction_error) 342 } 343 344 /// Completes one external mutation only after its durable effect exists. 345 pub async fn complete_admin_operation( 346 &self, 347 prepared: &RhiPreparedAdminOperation, 348 response: &RhiAdminResponseDocument, 349 completed_at: RhiAdminOperationTimeUnixMs, 350 policy: RhiAdminOperationJournalPolicy, 351 ) -> Result<RhiAdminOperationCompletion, RhiAdminOperationError> { 352 if self.host.mode() != RhiStateHostMode::ReadWriteExisting { 353 return Err(RhiAdminOperationError::new( 354 RhiAdminOperationErrorKind::InvalidMode, 355 )); 356 } 357 if response.route() != prepared.route 358 || response.canonical_bytes().is_empty() 359 || response.canonical_bytes().len() > RHI_ADMIN_OPERATION_RESPONSE_MODEL_MAX_BYTES 360 || completed_at < prepared.prepared_at 361 { 362 return Err(RhiAdminOperationError::new( 363 RhiAdminOperationErrorKind::InvalidInput, 364 )); 365 } 366 let expires_at = completed_at 367 .get() 368 .checked_add(policy.completed_retention_ms()) 369 .filter(|value| i64::try_from(*value).is_ok()) 370 .ok_or_else(|| RhiAdminOperationError::new(RhiAdminOperationErrorKind::InvalidInput))?; 371 let binding = PreparedBinding::from_prepared(prepared); 372 let response = response.canonical_bytes().to_vec().into_boxed_slice(); 373 self.host 374 .sqlite_host() 375 .transaction(move |transaction| { 376 Box::pin(async move { 377 complete_operation(transaction, &binding, &response, completed_at, expires_at) 378 .await 379 }) 380 }) 381 .await 382 .map_err(map_transaction_error) 383 } 384 } 385 386 struct AdminRequestBinding { 387 operation_id: AdminOperationIdBinding, 388 route: RhiAdminRoute, 389 request_sha256: [u8; 32], 390 } 391 392 impl AdminRequestBinding { 393 fn from_document(request: &RhiAdminRequestDocument) -> Result<Self, RhiAdminOperationError> { 394 if !request.route().is_mutation() { 395 return Err(RhiAdminOperationError::new( 396 RhiAdminOperationErrorKind::InvalidInput, 397 )); 398 } 399 let operation_id = request 400 .operation_id() 401 .ok_or_else(|| RhiAdminOperationError::new(RhiAdminOperationErrorKind::InvalidInput))?; 402 Ok(Self { 403 operation_id: AdminOperationIdBinding::new(operation_id)?, 404 route: request.route(), 405 request_sha256: request_digest(request), 406 }) 407 } 408 } 409 410 struct PreparedBinding { 411 operation_id: AdminOperationIdBinding, 412 route: RhiAdminRoute, 413 request_sha256: [u8; 32], 414 prepared_at: RhiAdminOperationTimeUnixMs, 415 } 416 417 impl PreparedBinding { 418 fn from_prepared(prepared: &RhiPreparedAdminOperation) -> Self { 419 Self { 420 operation_id: prepared.operation_id.clone(), 421 route: prepared.route, 422 request_sha256: prepared.request_sha256, 423 prepared_at: prepared.prepared_at, 424 } 425 } 426 } 427 428 enum StoredOperation { 429 Prepared { 430 route: RhiAdminRoute, 431 request_sha256: [u8; 32], 432 prepared_at: RhiAdminOperationTimeUnixMs, 433 }, 434 Completed { 435 route: RhiAdminRoute, 436 request_sha256: [u8; 32], 437 response: Box<[u8]>, 438 response_sha256: [u8; 32], 439 prepared_at: RhiAdminOperationTimeUnixMs, 440 completed_at: RhiAdminOperationTimeUnixMs, 441 expires_at: RhiAdminOperationTimeUnixMs, 442 }, 443 } 444 445 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 446 pub(crate) enum AdminJournalOperationError { 447 InvalidInput, 448 Conflict, 449 OutcomeUnknown, 450 ResourceExhausted, 451 Binding, 452 Storage, 453 } 454 455 pub(crate) type AdminDatabaseOperationFuture<'a> = Pin< 456 Box< 457 dyn Future<Output = Result<RhiAdminResponseDocument, AdminJournalOperationError>> 458 + Send 459 + 'a, 460 >, 461 >; 462 463 async fn prepare_operation( 464 transaction: &mut ServiceSqliteTransaction<'_>, 465 binding: &AdminRequestBinding, 466 observed_at: RhiAdminOperationTimeUnixMs, 467 ) -> Result<RhiAdminOperationAdmission, AdminJournalOperationError> { 468 prune_expired(transaction, observed_at).await?; 469 if let Some(existing) = read_operation(transaction, &binding.operation_id).await? { 470 return match existing { 471 StoredOperation::Prepared { 472 route, 473 request_sha256, 474 .. 475 } if route == binding.route && request_sha256 == binding.request_sha256 => { 476 Err(AdminJournalOperationError::OutcomeUnknown) 477 } 478 StoredOperation::Completed { 479 route, 480 request_sha256, 481 response, 482 response_sha256, 483 .. 484 } if route == binding.route && request_sha256 == binding.request_sha256 => { 485 if sha256(&response) != response_sha256 { 486 return Err(AdminJournalOperationError::Binding); 487 } 488 RhiAdminResponseDocument::from_canonical_bytes(route, &response) 489 .map(RhiAdminOperationAdmission::ExactReplay) 490 .map_err(|_| AdminJournalOperationError::Binding) 491 } 492 StoredOperation::Prepared { .. } | StoredOperation::Completed { .. } => { 493 Err(AdminJournalOperationError::Conflict) 494 } 495 }; 496 } 497 let (completed, prepared) = read_counts(transaction).await?; 498 let reserved = completed 499 .checked_add(prepared) 500 .ok_or(AdminJournalOperationError::Binding)?; 501 if reserved >= u64::from(RHI_ADMIN_OPERATION_COMPLETED_LIMIT) 502 || prepared >= u64::from(RHI_ADMIN_OPERATION_PREPARED_LIMIT) 503 { 504 return Err(AdminJournalOperationError::ResourceExhausted); 505 } 506 let result = sqlx::query(INSERT_PREPARED_SQL) 507 .bind(binding.operation_id.as_str()) 508 .bind(binding.route.operation_id()) 509 .bind(binding.request_sha256.as_slice()) 510 .bind(observed_at.sqlite_value()) 511 .execute(&mut *transaction) 512 .await 513 .map_err(|_| AdminJournalOperationError::Storage)?; 514 require_one(result.rows_affected())?; 515 match read_operation(transaction, &binding.operation_id).await? { 516 Some(StoredOperation::Prepared { 517 route, 518 request_sha256, 519 prepared_at, 520 }) if route == binding.route 521 && request_sha256 == binding.request_sha256 522 && prepared_at == observed_at => 523 { 524 Ok(RhiAdminOperationAdmission::Prepared( 525 RhiPreparedAdminOperation { 526 operation_id: binding.operation_id.clone(), 527 route, 528 request_sha256, 529 prepared_at, 530 }, 531 )) 532 } 533 Some(_) | None => Err(AdminJournalOperationError::Binding), 534 } 535 } 536 537 async fn complete_operation( 538 transaction: &mut ServiceSqliteTransaction<'_>, 539 binding: &PreparedBinding, 540 response: &[u8], 541 completed_at: RhiAdminOperationTimeUnixMs, 542 expires_at: u64, 543 ) -> Result<RhiAdminOperationCompletion, AdminJournalOperationError> { 544 let response_sha256 = sha256(response); 545 match read_operation(transaction, &binding.operation_id).await? { 546 Some(StoredOperation::Completed { 547 route, 548 request_sha256, 549 response: existing_response, 550 response_sha256: existing_sha256, 551 .. 552 }) if route == binding.route 553 && request_sha256 == binding.request_sha256 554 && existing_response.as_ref() == response 555 && existing_sha256 == response_sha256 => 556 { 557 return Ok(RhiAdminOperationCompletion::ExactReplay); 558 } 559 Some(StoredOperation::Completed { .. }) => { 560 return Err(AdminJournalOperationError::Conflict); 561 } 562 Some(StoredOperation::Prepared { 563 route, 564 request_sha256, 565 prepared_at, 566 }) if route == binding.route 567 && request_sha256 == binding.request_sha256 568 && prepared_at == binding.prepared_at => {} 569 Some(StoredOperation::Prepared { .. }) => { 570 return Err(AdminJournalOperationError::Conflict); 571 } 572 None => return Err(AdminJournalOperationError::Binding), 573 } 574 let (completed, _) = read_counts(transaction).await?; 575 if completed >= u64::from(RHI_ADMIN_OPERATION_COMPLETED_LIMIT) { 576 return Err(AdminJournalOperationError::ResourceExhausted); 577 } 578 let result = sqlx::query(COMPLETE_OPERATION_SQL) 579 .bind(response) 580 .bind(response_sha256.as_slice()) 581 .bind(completed_at.sqlite_value()) 582 .bind(i64::try_from(expires_at).map_err(|_| AdminJournalOperationError::InvalidInput)?) 583 .bind(binding.operation_id.as_str()) 584 .execute(&mut *transaction) 585 .await 586 .map_err(|_| AdminJournalOperationError::Storage)?; 587 require_one(result.rows_affected())?; 588 match read_operation(transaction, &binding.operation_id).await? { 589 Some(StoredOperation::Completed { 590 route, 591 request_sha256, 592 response: actual_response, 593 response_sha256: actual_sha256, 594 prepared_at, 595 completed_at: actual_completed_at, 596 expires_at: actual_expires_at, 597 }) if route == binding.route 598 && request_sha256 == binding.request_sha256 599 && actual_response.as_ref() == response 600 && actual_sha256 == response_sha256 601 && prepared_at == binding.prepared_at 602 && actual_completed_at == completed_at 603 && actual_expires_at.get() == expires_at => 604 { 605 Ok(RhiAdminOperationCompletion::Completed) 606 } 607 Some(_) | None => Err(AdminJournalOperationError::Binding), 608 } 609 } 610 611 async fn prune_expired( 612 transaction: &mut ServiceSqliteTransaction<'_>, 613 observed_at: RhiAdminOperationTimeUnixMs, 614 ) -> Result<(), AdminJournalOperationError> { 615 sqlx::query(PRUNE_EXPIRED_SQL) 616 .bind(observed_at.sqlite_value()) 617 .bind(PRUNE_LIMIT) 618 .execute(&mut *transaction) 619 .await 620 .map(|_| ()) 621 .map_err(|_| AdminJournalOperationError::Storage) 622 } 623 624 async fn read_counts( 625 transaction: &mut ServiceSqliteTransaction<'_>, 626 ) -> Result<(u64, u64), AdminJournalOperationError> { 627 let rows = sqlx::query(READ_COUNTS_SQL) 628 .fetch_all(&mut *transaction) 629 .await 630 .map_err(|_| AdminJournalOperationError::Storage)?; 631 if rows.len() != 1 { 632 return Err(AdminJournalOperationError::Binding); 633 } 634 let completed = rows[0] 635 .try_get::<i64, _>("completed_count") 636 .map_err(|_| AdminJournalOperationError::Binding)?; 637 let prepared = rows[0] 638 .try_get::<i64, _>("prepared_count") 639 .map_err(|_| AdminJournalOperationError::Binding)?; 640 Ok(( 641 u64::try_from(completed).map_err(|_| AdminJournalOperationError::Binding)?, 642 u64::try_from(prepared).map_err(|_| AdminJournalOperationError::Binding)?, 643 )) 644 } 645 646 async fn read_operation( 647 transaction: &mut ServiceSqliteTransaction<'_>, 648 operation_id: &AdminOperationIdBinding, 649 ) -> Result<Option<StoredOperation>, AdminJournalOperationError> { 650 let rows = sqlx::query(READ_OPERATION_SQL) 651 .bind(operation_id.as_str()) 652 .fetch_all(&mut *transaction) 653 .await 654 .map_err(|_| AdminJournalOperationError::Storage)?; 655 if rows.len() > 1 { 656 return Err(AdminJournalOperationError::Binding); 657 } 658 rows.first().map(decode_operation).transpose() 659 } 660 661 fn decode_operation( 662 row: &sqlx::sqlite::SqliteRow, 663 ) -> Result<StoredOperation, AdminJournalOperationError> { 664 let route = row 665 .try_get::<Option<&str>, _>("route") 666 .map_err(|_| AdminJournalOperationError::Binding)? 667 .and_then(parse_route) 668 .ok_or(AdminJournalOperationError::Binding)?; 669 let request_sha256 = exact_digest(row, "request_sha256")?; 670 let state = row 671 .try_get::<Option<&str>, _>("state") 672 .map_err(|_| AdminJournalOperationError::Binding)? 673 .ok_or(AdminJournalOperationError::Binding)?; 674 let prepared_at = time(row, "prepared_at_unix_ms")?; 675 match state { 676 "prepared" => { 677 require_null(row, "response_model_type")?; 678 require_null(row, "response_sha256_type")?; 679 require_null(row, "completed_at_type")?; 680 require_null(row, "expires_at_type")?; 681 Ok(StoredOperation::Prepared { 682 route, 683 request_sha256, 684 prepared_at, 685 }) 686 } 687 "completed" => { 688 require_type(row, "response_model_type", "blob")?; 689 require_type(row, "response_sha256_type", "blob")?; 690 require_type(row, "completed_at_type", "integer")?; 691 require_type(row, "expires_at_type", "integer")?; 692 let response = row 693 .try_get::<Option<Vec<u8>>, _>("response_model") 694 .map_err(|_| AdminJournalOperationError::Binding)? 695 .ok_or(AdminJournalOperationError::Binding)? 696 .into_boxed_slice(); 697 Ok(StoredOperation::Completed { 698 route, 699 request_sha256, 700 response, 701 response_sha256: exact_digest(row, "response_sha256")?, 702 prepared_at, 703 completed_at: time(row, "completed_at_unix_ms")?, 704 expires_at: time(row, "expires_at_unix_ms")?, 705 }) 706 } 707 _ => Err(AdminJournalOperationError::Binding), 708 } 709 } 710 711 fn request_digest(request: &RhiAdminRequestDocument) -> [u8; 32] { 712 let mut hasher = Sha256::new(); 713 hasher.update(REQUEST_DIGEST_DOMAIN); 714 hash_field(&mut hasher, request.route().operation_id().as_bytes()); 715 hasher.update([0]); 716 hash_field(&mut hasher, request.model_bytes()); 717 hasher.finalize().into() 718 } 719 720 fn hash_field(hasher: &mut Sha256, bytes: &[u8]) { 721 hasher.update( 722 u64::try_from(bytes.len()) 723 .expect("bounded field length") 724 .to_be_bytes(), 725 ); 726 hasher.update(bytes); 727 } 728 729 fn sha256(bytes: &[u8]) -> [u8; 32] { 730 Sha256::digest(bytes).into() 731 } 732 733 fn parse_route(value: &str) -> Option<RhiAdminRoute> { 734 RhiAdminRoute::ALL 735 .into_iter() 736 .find(|route| route.is_mutation() && route.operation_id() == value) 737 } 738 739 fn exact_digest( 740 row: &sqlx::sqlite::SqliteRow, 741 column: &str, 742 ) -> Result<[u8; 32], AdminJournalOperationError> { 743 row.try_get::<Option<Vec<u8>>, _>(column) 744 .map_err(|_| AdminJournalOperationError::Binding)? 745 .ok_or(AdminJournalOperationError::Binding)? 746 .try_into() 747 .map_err(|_| AdminJournalOperationError::Binding) 748 } 749 750 fn time( 751 row: &sqlx::sqlite::SqliteRow, 752 column: &str, 753 ) -> Result<RhiAdminOperationTimeUnixMs, AdminJournalOperationError> { 754 let value = row 755 .try_get::<i64, _>(column) 756 .map_err(|_| AdminJournalOperationError::Binding)?; 757 RhiAdminOperationTimeUnixMs::new( 758 u64::try_from(value).map_err(|_| AdminJournalOperationError::Binding)?, 759 ) 760 .map_err(|_| AdminJournalOperationError::Binding) 761 } 762 763 fn require_null( 764 row: &sqlx::sqlite::SqliteRow, 765 column: &str, 766 ) -> Result<(), AdminJournalOperationError> { 767 require_type(row, column, "null") 768 } 769 770 fn require_type( 771 row: &sqlx::sqlite::SqliteRow, 772 column: &str, 773 expected: &str, 774 ) -> Result<(), AdminJournalOperationError> { 775 (row.try_get::<&str, _>(column) 776 .map_err(|_| AdminJournalOperationError::Binding)? 777 == expected) 778 .then_some(()) 779 .ok_or(AdminJournalOperationError::Binding) 780 } 781 782 fn require_one(rows: u64) -> Result<(), AdminJournalOperationError> { 783 (rows == 1) 784 .then_some(()) 785 .ok_or(AdminJournalOperationError::Storage) 786 } 787 788 fn map_transaction_error( 789 error: ServiceSqliteTransactionError<AdminJournalOperationError>, 790 ) -> RhiAdminOperationError { 791 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 792 return RhiAdminOperationError::new(RhiAdminOperationErrorKind::CommitOutcomeUnknown); 793 } 794 let kind = match error.operation_error() { 795 Some(AdminJournalOperationError::InvalidInput) => RhiAdminOperationErrorKind::InvalidInput, 796 Some(AdminJournalOperationError::Conflict) => RhiAdminOperationErrorKind::OperationConflict, 797 Some(AdminJournalOperationError::OutcomeUnknown) => { 798 RhiAdminOperationErrorKind::OperationOutcomeUnknown 799 } 800 Some(AdminJournalOperationError::ResourceExhausted) => { 801 RhiAdminOperationErrorKind::ResourceExhausted 802 } 803 Some(AdminJournalOperationError::Binding) => RhiAdminOperationErrorKind::Binding, 804 Some(AdminJournalOperationError::Storage) | None => RhiAdminOperationErrorKind::Transaction, 805 }; 806 RhiAdminOperationError::new(kind) 807 }