rhi

Coordinated trade for connected markets
git clone https://radroots.dev/git/rhi.git
Log | Files | Refs | README | LICENSE

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 }