myc

Self-custodial remote signer for Radroots apps
git clone https://radroots.dev/git/myc.git
Log | Files | Refs | README | LICENSE

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 }