rhi

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

runtime_admin.rs (85181B)


      1 //! State-backed implementation of the final RHI Unix-admin contract.
      2 
      3 use core::sync::atomic::{AtomicBool, Ordering};
      4 use std::{path::PathBuf, sync::Arc};
      5 
      6 use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD};
      7 use hmac::{Hmac, Mac as _};
      8 use radroots_event::id::TradeId;
      9 use radroots_service_sqlite::BackupCreatedAtUnixMs;
     10 use serde_json::{Map, Value, json};
     11 use sha2::{Digest as _, Sha256};
     12 use sqlx::Row;
     13 
     14 use crate::{
     15     RhiAdminFuture, RhiAdminHandler, RhiAdminHandlerError, RhiAdminHandlerErrorKind,
     16     RhiAdminRequestDocument, RhiAdminResponseDocument, RhiAdminRoute, RhiConfigDocumentV1,
     17     RhiDecryptedIdentity, RhiPresenceDesiredAuthority, RhiPresenceUnixMilliseconds,
     18     RhiPublicationAuthority, RhiReconciliationJobErrorKind, RhiReconciliationJobPolicy,
     19     RhiReconciliationUnixMilliseconds, RhiStateHost, RhiStatusReader, RhiTimeEntropyAdapters,
     20     build_rhi_signed_presence_documents,
     21     state_admin::{
     22         AdminJournalOperationError, RhiAdminOperationAdmission, RhiAdminOperationError,
     23         RhiAdminOperationErrorKind, RhiAdminOperationJournalPolicy, RhiAdminOperationRepository,
     24         RhiAdminOperationTimeUnixMs,
     25     },
     26     state_config,
     27 };
     28 
     29 const REPORT_SELECT: &str = r#"SELECT
     30     length(report.statement_sha256) AS statement_bytes,
     31     substr(report.statement_sha256, 1, 33) AS statement_sha256,
     32     length(report.manifest_sha256) AS manifest_bytes,
     33     substr(report.manifest_sha256, 1, 33) AS manifest_sha256,
     34     length(report.projection_sha256) AS projection_bytes,
     35     substr(report.projection_sha256, 1, 33) AS projection_sha256,
     36     length(report.trade_id) AS trade_id_bytes,
     37     substr(report.trade_id, 1, 17) AS trade_id,
     38     length(report.claim_mutation_id) AS claim_bytes,
     39     substr(report.claim_mutation_id, 1, 33) AS claim_mutation_id,
     40     length(report.issuer_public_key) AS issuer_bytes,
     41     substr(report.issuer_public_key, 1, 33) AS issuer_public_key,
     42     report.outcome,
     43     length(report.canonical_report) AS canonical_report_bytes,
     44     substr(report.canonical_report, 1, 16385) AS canonical_report,
     45     report.observed_at_unix_s,
     46     CASE WHEN report.supersedes_statement_sha256 IS NULL THEN NULL
     47          ELSE length(report.supersedes_statement_sha256) END AS supersedes_bytes,
     48     CASE WHEN report.supersedes_statement_sha256 IS NULL THEN NULL
     49          ELSE substr(report.supersedes_statement_sha256, 1, 33) END AS supersedes_statement_sha256,
     50     length(event.event_id) AS event_id_bytes,
     51     substr(event.event_id, 1, 33) AS event_id,
     52     manifest.trade_generation,
     53     length(manifest.evidence_policy_sha256) AS policy_bytes,
     54     substr(manifest.evidence_policy_sha256, 1, 33) AS evidence_policy_sha256,
     55     projection.reducer_contract,
     56     projection.reducer_contract_version,
     57     (SELECT COUNT(*) FROM evidence_reconciliation_sources AS source
     58         WHERE source.attempt_id = manifest.attempt_id) AS source_count,
     59     (SELECT COUNT(*) FROM evidence_reconciliation_sources AS source
     60         WHERE source.attempt_id = manifest.attempt_id AND source.required = 1
     61           AND source.completion != 'complete') AS incomplete_required,
     62     (SELECT COUNT(*) FROM evidence_reconciliation_sources AS source
     63         WHERE source.attempt_id = manifest.attempt_id AND source.required = 1
     64           AND source.completion = 'unsupported') AS unsupported_required,
     65     (SELECT COALESCE(SUM(source.accepted_event_count), 0)
     66         FROM evidence_reconciliation_sources AS source
     67         WHERE source.attempt_id = manifest.attempt_id) AS accepted_event_count
     68 FROM attestation_reports AS report
     69 JOIN signed_attestation_events AS event
     70     ON event.statement_sha256 = report.statement_sha256
     71 JOIN evidence_manifests AS manifest
     72     ON manifest.manifest_sha256 = report.manifest_sha256
     73 JOIN trade_projections AS projection
     74     ON projection.projection_sha256 = report.projection_sha256
     75 "#;
     76 
     77 const REPORT_ROWS_SUFFIX: &str = r#"WHERE report.trade_id = ? AND report.observed_at_unix_s <= ?
     78         AND (? IS NULL OR report.outcome = ?)
     79     ORDER BY report.observed_at_unix_s DESC, report.statement_sha256 DESC
     80     LIMIT ? OFFSET ?"#;
     81 
     82 const CURRENT_REPORT_SUFFIX: &str = r#"WHERE report.trade_id = ?
     83         AND NOT EXISTS (
     84             SELECT 1 FROM attestation_reports AS successor
     85             WHERE successor.supersedes_statement_sha256 = report.statement_sha256
     86         )
     87     ORDER BY report.observed_at_unix_s DESC, report.statement_sha256 DESC
     88     LIMIT 2"#;
     89 
     90 const PUBLICATION_BACKLOG_SQL: &str = r#"SELECT
     91     length(outbox.outbox_id) AS outbox_id_bytes,
     92     substr(outbox.outbox_id, 1, 33) AS outbox_id,
     93     length(event.statement_sha256) AS statement_bytes,
     94     substr(event.statement_sha256, 1, 33) AS statement_sha256,
     95     length(outbox.event_id) AS event_id_bytes,
     96     substr(outbox.event_id, 1, 33) AS event_id,
     97     length(outbox.event_sha256) AS event_sha256_bytes,
     98     substr(outbox.event_sha256, 1, 33) AS event_sha256,
     99     outbox.state AS outbox_state,
    100     outbox.next_attempt_unix_ms,
    101     COALESCE(MAX(target.attempt_count), 0) AS attempt_count,
    102     CASE
    103         WHEN outbox.state = 'pending' THEN 'pending'
    104         WHEN outbox.state = 'leased' THEN 'submitted'
    105         WHEN outbox.state = 'complete' THEN 'accepted'
    106         WHEN SUM(CASE WHEN target.state = 'auth_required' THEN 1 ELSE 0 END) > 0 THEN 'auth_required'
    107         WHEN SUM(CASE WHEN target.state = 'rate_limited' THEN 1 ELSE 0 END) > 0 THEN 'rate_limited'
    108         WHEN SUM(CASE WHEN target.state = 'rejected' THEN 1 ELSE 0 END) > 0 THEN 'rejected'
    109         WHEN SUM(CASE WHEN target.state = 'failed' THEN 1 ELSE 0 END) > 0 THEN 'failed'
    110         ELSE 'unknown'
    111     END AS public_state
    112 FROM publication_outbox AS outbox
    113 JOIN signed_attestation_events AS event ON event.event_id = outbox.event_id
    114 JOIN publication_targets AS target ON target.outbox_id = outbox.outbox_id
    115 WHERE outbox.updated_at_unix_ms <= ?
    116 GROUP BY outbox.outbox_id
    117 HAVING (? IS NULL OR public_state = ?)
    118 ORDER BY outbox.updated_at_unix_ms DESC, outbox.outbox_id DESC
    119 LIMIT ? OFFSET ?"#;
    120 
    121 const PUBLICATION_TARGETS_SQL: &str = r#"SELECT
    122     length(target.outbox_id) AS outbox_id_bytes,
    123     substr(target.outbox_id, 1, 33) AS outbox_id,
    124     target.relay_id, target.state, target.attempt_count,
    125     (SELECT MAX(attempt.finished_at_unix_ms)
    126         FROM publication_attempts AS attempt
    127         WHERE attempt.outbox_id = target.outbox_id
    128           AND attempt.target_ordinal = target.target_ordinal) AS last_attempt_unix_ms
    129 FROM publication_targets AS target
    130 WHERE target.updated_at_unix_ms <= ?
    131   AND (? IS NULL OR target.outbox_id = ?)
    132   AND (? IS NULL OR target.state = ?)
    133 ORDER BY target.updated_at_unix_ms DESC, target.outbox_id DESC, target.target_ordinal
    134 LIMIT ? OFFSET ?"#;
    135 
    136 pub(crate) struct RuntimeAdminHandler {
    137     pub(crate) state: Arc<RhiStateHost>,
    138     pub(crate) configuration: Arc<RhiConfigDocumentV1>,
    139     pub(crate) identity: Arc<RhiDecryptedIdentity>,
    140     pub(crate) publication: Arc<RhiPublicationAuthority>,
    141     pub(crate) presence: Arc<RhiPresenceDesiredAuthority>,
    142     pub(crate) status: RhiStatusReader,
    143     pub(crate) accepting_mutations: Arc<AtomicBool>,
    144     pub(crate) cursor_key: [u8; 32],
    145     pub(crate) time_entropy: RhiTimeEntropyAdapters,
    146 }
    147 
    148 impl RuntimeAdminHandler {
    149     async fn handle_inner(
    150         &self,
    151         request: RhiAdminRequestDocument,
    152     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
    153         if request.route().is_mutation() && !self.accepting_mutations.load(Ordering::Acquire) {
    154             return Err(failure(RhiAdminHandlerErrorKind::Unavailable));
    155         }
    156         match request.route() {
    157             RhiAdminRoute::Status => {
    158                 let snapshot = self.status.snapshot();
    159                 let value = serde_json::from_slice(snapshot.detailed_status_json())
    160                     .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    161                 response(request.route(), value)
    162             }
    163             RhiAdminRoute::EffectiveConfig => RhiAdminResponseDocument::from_canonical_bytes(
    164                 request.route(),
    165                 self.configuration.effective().canonical_json().as_bytes(),
    166             )
    167             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal)),
    168             RhiAdminRoute::IdentityStatus | RhiAdminRoute::IdentityPublic => {
    169                 self.identity(&request).await
    170             }
    171             RhiAdminRoute::StateStatus => self.state_status(request.route()).await,
    172             RhiAdminRoute::StateBackup => self.backup(request).await,
    173             RhiAdminRoute::MetricsSnapshot => self.metrics_snapshot(request.route()),
    174             RhiAdminRoute::ReconciliationStatus => self.reconciliation_status(request).await,
    175             RhiAdminRoute::ReconciliationJobs => self.reconciliation_jobs(request).await,
    176             RhiAdminRoute::ReconciliationRefresh => self.reconciliation_refresh(request).await,
    177             RhiAdminRoute::Sources => self.sources(request).await,
    178             RhiAdminRoute::TradeProjection => self.trade_projection(request).await,
    179             RhiAdminRoute::TradeReportCurrent => self.trade_report_current(request).await,
    180             RhiAdminRoute::TradeReports => self.trade_reports(request).await,
    181             RhiAdminRoute::PublicationBacklog => self.publication_backlog(request).await,
    182             RhiAdminRoute::PublicationTargets => self.publication_targets(request).await,
    183             RhiAdminRoute::PublicationRetry => self.publication_retry(request).await,
    184             RhiAdminRoute::PresenceDesired => self.presence_desired(request.route()).await,
    185             RhiAdminRoute::PresenceRender => self.presence_render(request).await,
    186             RhiAdminRoute::PresenceRefresh => self.presence_refresh(request).await,
    187         }
    188     }
    189 
    190     async fn identity(
    191         &self,
    192         request: &RhiAdminRequestDocument,
    193     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
    194         let model = request_model(request)?;
    195         if model.pointer("/role").and_then(Value::as_str) != Some("service") {
    196             return Err(failure(RhiAdminHandlerErrorKind::Internal));
    197         }
    198         let generation = current_generation(&self.state).await?;
    199         let value = if request.route() == RhiAdminRoute::IdentityPublic {
    200             json!({
    201                 "generation": generation,
    202                 "public_key": self.identity.public_identity().as_hex(),
    203                 "role": "service",
    204             })
    205         } else {
    206             json!({
    207                 "available": true,
    208                 "configured": true,
    209                 "generation": generation,
    210                 "provider": "encrypted_file",
    211                 "public_key": self.identity.public_identity().as_hex(),
    212                 "reason_codes": [],
    213                 "role": "service",
    214             })
    215         };
    216         response(request.route(), value)
    217     }
    218 
    219     async fn state_status(
    220         &self,
    221         route: RhiAdminRoute,
    222     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
    223         response(
    224             route,
    225             json!({
    226                 "backup_eligible": true,
    227                 "generation": current_generation(&self.state).await?,
    228                 "integrity": "verified",
    229                 "reason_codes": [],
    230                 "schema_version": self.state.metadata().initial_database_metadata().state_schema_version().get(),
    231                 "writer_lock": "held_by_daemon",
    232             }),
    233         )
    234     }
    235 
    236     fn metrics_snapshot(
    237         &self,
    238         route: RhiAdminRoute,
    239     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
    240         let snapshot = self.status.operations_cache().snapshot();
    241         let mut metrics = Map::new();
    242         for sample in snapshot.metrics().samples() {
    243             let mut key = sample.name().as_str().to_owned();
    244             for label in sample.labels() {
    245                 key.push('_');
    246                 key.push_str(label.value());
    247             }
    248             let value = match sample.value() {
    249                 radroots_service_host::MetricValue::Counter(value) => value,
    250                 radroots_service_host::MetricValue::Gauge(value) => {
    251                     u64::try_from(value).map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?
    252                 }
    253             };
    254             if metrics.insert(key, Value::from(value)).is_some() {
    255                 return Err(failure(RhiAdminHandlerErrorKind::Internal));
    256             }
    257         }
    258         response(
    259             route,
    260             json!({
    261                 "captured_at_utc": self.now_seconds()?,
    262                 "metrics": metrics,
    263             }),
    264         )
    265     }
    266 
    267     async fn backup(
    268         &self,
    269         request: RhiAdminRequestDocument,
    270     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
    271         let model = request_model(&request)?;
    272         let expected_generation = required_u64(&model, "/expected_generation")?;
    273         let generation = current_generation(&self.state).await?;
    274         if expected_generation != generation {
    275             return Err(failure(RhiAdminHandlerErrorKind::Conflict));
    276         }
    277         let now_ms = self.now_millis()?;
    278         let prepared = match RhiAdminOperationRepository::new(&self.state)
    279             .prepare_admin_operation(&request, operation_time(now_ms)?)
    280             .await
    281             .map_err(map_journal_error)?
    282         {
    283             RhiAdminOperationAdmission::ExactReplay(response) => return Ok(response),
    284             RhiAdminOperationAdmission::Prepared(prepared) => prepared,
    285         };
    286         let target = model
    287             .pointer("/target_path")
    288             .and_then(Value::as_str)
    289             .map(PathBuf::from)
    290             .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?;
    291         let manifest = self
    292             .state
    293             .capture_online_backup(
    294                 &target,
    295                 BackupCreatedAtUnixMs::new(now_ms)
    296                     .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?,
    297             )
    298             .await
    299             .map_err(|_| failure(RhiAdminHandlerErrorKind::Unavailable))?;
    300         let completed_ms = self.now_millis()?;
    301         let completed_seconds = completed_ms / 1_000;
    302         let completed = response(
    303             request.route(),
    304             json!({
    305                 "completed_at_utc": completed_seconds,
    306                 "manifest_digest": lower_hex(manifest.digest().as_bytes()),
    307                 "operation_id": required_operation_id(&request)?,
    308                 "snapshot_generation": generation,
    309             }),
    310         )?;
    311         RhiAdminOperationRepository::new(&self.state)
    312             .complete_admin_operation(
    313                 &prepared,
    314                 &completed,
    315                 operation_time(completed_ms)?,
    316                 RhiAdminOperationJournalPolicy::seven_days(),
    317             )
    318             .await
    319             .map_err(map_journal_error)?;
    320         Ok(completed)
    321     }
    322 
    323     async fn presence_desired(
    324         &self,
    325         route: RhiAdminRoute,
    326     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
    327         let desired = self
    328             .state
    329             .repositories()
    330             .desired_presence()
    331             .current()
    332             .await
    333             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?
    334             .ok_or_else(|| failure(RhiAdminHandlerErrorKind::NotFound))?;
    335         let (state, digests) = self.presence_state(desired.generation()).await?;
    336         response(
    337             route,
    338             json!({
    339                 "document_digests": digests,
    340                 "generation": desired.generation(),
    341                 "reason_codes": [],
    342                 "state": if desired.mode().code() == "disabled" { "disabled" } else { state },
    343             }),
    344         )
    345     }
    346 
    347     async fn presence_render(
    348         &self,
    349         request: RhiAdminRequestDocument,
    350     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
    351         let model = request_model(&request)?;
    352         let expected = required_u64(&model, "/expected_generation")?;
    353         let desired = self
    354             .state
    355             .repositories()
    356             .desired_presence()
    357             .commit(&self.presence)
    358             .await
    359             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    360         if desired.state().generation() != expected {
    361             return Err(failure(RhiAdminHandlerErrorKind::Conflict));
    362         }
    363         let now_ms = self.now_millis()?;
    364         let prepared = match RhiAdminOperationRepository::new(&self.state)
    365             .prepare_admin_operation(&request, operation_time(now_ms)?)
    366             .await
    367             .map_err(map_journal_error)?
    368         {
    369             RhiAdminOperationAdmission::ExactReplay(response) => return Ok(response),
    370             RhiAdminOperationAdmission::Prepared(prepared) => prepared,
    371         };
    372         let documents = build_rhi_signed_presence_documents(
    373             desired,
    374             &self.presence,
    375             &self.identity,
    376             self.time_entropy
    377                 .now_utc()
    378                 .map_err(|_| failure(RhiAdminHandlerErrorKind::Unavailable))?,
    379             self.time_entropy.entropy(),
    380         )
    381         .map_err(|_| failure(RhiAdminHandlerErrorKind::Unavailable))?;
    382         let digests = presence_document_digests(&documents);
    383         self.state
    384             .repositories()
    385             .presence_outbox()
    386             .commit_signed_presence(
    387                 &documents,
    388                 RhiPresenceUnixMilliseconds::new(now_ms)
    389                     .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?,
    390             )
    391             .await
    392             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    393         let completed_ms = self.now_millis()?;
    394         let completed = response(
    395             request.route(),
    396             json!({
    397                 "document_digests": digests,
    398                 "generation": expected,
    399                 "operation_id": required_operation_id(&request)?,
    400             }),
    401         )?;
    402         RhiAdminOperationRepository::new(&self.state)
    403             .complete_admin_operation(
    404                 &prepared,
    405                 &completed,
    406                 operation_time(completed_ms)?,
    407                 RhiAdminOperationJournalPolicy::seven_days(),
    408             )
    409             .await
    410             .map_err(map_journal_error)?;
    411         Ok(completed)
    412     }
    413 
    414     fn now_millis(&self) -> Result<u64, RhiAdminHandlerError> {
    415         self.time_entropy
    416             .now_utc_milliseconds()
    417             .map_err(|_| failure(RhiAdminHandlerErrorKind::Unavailable))
    418     }
    419 
    420     fn now_seconds(&self) -> Result<u64, RhiAdminHandlerError> {
    421         self.time_entropy
    422             .now_utc()
    423             .map(|value| value.get())
    424             .map_err(|_| failure(RhiAdminHandlerErrorKind::Unavailable))
    425     }
    426 
    427     async fn read_dirty_generation(&self, trade: TradeId) -> Result<u64, RhiAdminHandlerError> {
    428         self.state
    429             .sqlite_host()
    430             .transaction(move |transaction| {
    431                 Box::pin(async move {
    432                     let rows = sqlx::query(
    433                         "SELECT generation FROM trade_dirty_generations WHERE trade_id = ? LIMIT 2",
    434                     )
    435                     .bind(trade.as_bytes().as_slice())
    436                     .fetch_all(&mut *transaction)
    437                     .await
    438                     .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    439                     if rows.len() != 1 {
    440                         return Err(if rows.is_empty() {
    441                             failure(RhiAdminHandlerErrorKind::NotFound)
    442                         } else {
    443                             failure(RhiAdminHandlerErrorKind::Internal)
    444                         });
    445                     }
    446                     row_u64(&rows[0], "generation")
    447                 })
    448             })
    449             .await
    450             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
    451     }
    452 
    453     async fn current_report_row(
    454         &self,
    455         trade: TradeId,
    456     ) -> Result<sqlx::sqlite::SqliteRow, RhiAdminHandlerError> {
    457         self.state
    458             .sqlite_host()
    459             .transaction(move |transaction| {
    460                 Box::pin(async move {
    461                     let sql = format!("{REPORT_SELECT}{CURRENT_REPORT_SUFFIX}");
    462                     // Both fragments are private compile-time constants; no caller data enters SQL.
    463                     let rows = sqlx::query(sqlx::AssertSqlSafe(sql.as_str()))
    464                         .bind(trade.as_bytes().as_slice())
    465                         .fetch_all(&mut *transaction)
    466                         .await
    467                         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    468                     if rows.len() != 1 {
    469                         return Err(if rows.is_empty() {
    470                             failure(RhiAdminHandlerErrorKind::NotFound)
    471                         } else {
    472                             failure(RhiAdminHandlerErrorKind::Internal)
    473                         });
    474                     }
    475                     Ok(rows.into_iter().next().expect("one row was established"))
    476                 })
    477             })
    478             .await
    479             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
    480     }
    481 
    482     fn cursor_or_new(
    483         &self,
    484         model: &Value,
    485         route: RhiAdminRoute,
    486         query_digest: [u8; 32],
    487         new_snapshot: u64,
    488     ) -> Result<(u64, u32), RhiAdminHandlerError> {
    489         model
    490             .pointer("/cursor")
    491             .and_then(Value::as_str)
    492             .map(|cursor| self.decode_cursor(route, cursor, query_digest))
    493             .transpose()
    494             .map(|cursor| cursor.unwrap_or((new_snapshot, 0)))
    495     }
    496 
    497     fn encode_cursor(
    498         &self,
    499         route: RhiAdminRoute,
    500         snapshot: u64,
    501         offset: u32,
    502         query_digest: [u8; 32],
    503     ) -> String {
    504         encode_runtime_cursor(&self.cursor_key, route, snapshot, offset, query_digest)
    505     }
    506 
    507     fn decode_cursor(
    508         &self,
    509         route: RhiAdminRoute,
    510         encoded: &str,
    511         query_digest: [u8; 32],
    512     ) -> Result<(u64, u32), RhiAdminHandlerError> {
    513         decode_runtime_cursor(&self.cursor_key, route, encoded, query_digest)
    514     }
    515 }
    516 
    517 impl RhiAdminHandler for RuntimeAdminHandler {
    518     fn handle<'a>(&'a self, request: RhiAdminRequestDocument) -> RhiAdminFuture<'a> {
    519         Box::pin(async move { self.handle_inner(request).await })
    520     }
    521 }
    522 
    523 fn presence_document_digests(documents: &crate::RhiSignedPresenceDocuments) -> Map<String, Value> {
    524     documents
    525         .documents()
    526         .iter()
    527         .map(|document| {
    528             (
    529                 document.kind().code().to_owned(),
    530                 Value::String(lower_hex(document.signed_event_sha256())),
    531             )
    532         })
    533         .collect()
    534 }
    535 
    536 fn request_model(request: &RhiAdminRequestDocument) -> Result<Value, RhiAdminHandlerError> {
    537     serde_json::from_slice(request.model_bytes())
    538         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
    539 }
    540 
    541 fn required_u64(value: &Value, pointer: &str) -> Result<u64, RhiAdminHandlerError> {
    542     value
    543         .pointer(pointer)
    544         .and_then(Value::as_u64)
    545         .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))
    546 }
    547 
    548 fn required_operation_id(request: &RhiAdminRequestDocument) -> Result<&str, RhiAdminHandlerError> {
    549     request
    550         .operation_id()
    551         .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))
    552 }
    553 
    554 fn operation_time(value: u64) -> Result<RhiAdminOperationTimeUnixMs, RhiAdminHandlerError> {
    555     RhiAdminOperationTimeUnixMs::new(value).map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
    556 }
    557 
    558 async fn current_generation(state: &RhiStateHost) -> Result<u64, RhiAdminHandlerError> {
    559     state_config::current_generation(state)
    560         .await
    561         .map(u64::from)
    562         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
    563 }
    564 
    565 fn response(
    566     route: RhiAdminRoute,
    567     value: Value,
    568 ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
    569     let bytes =
    570         serde_json::to_vec(&value).map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    571     RhiAdminResponseDocument::from_canonical_bytes(route, &bytes)
    572         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
    573 }
    574 
    575 fn map_journal_error(error: RhiAdminOperationError) -> RhiAdminHandlerError {
    576     failure(match error.kind() {
    577         RhiAdminOperationErrorKind::OperationConflict => {
    578             RhiAdminHandlerErrorKind::OperationIdConflict
    579         }
    580         RhiAdminOperationErrorKind::ResourceExhausted => RhiAdminHandlerErrorKind::Unavailable,
    581         RhiAdminOperationErrorKind::OperationOutcomeUnknown
    582         | RhiAdminOperationErrorKind::CommitOutcomeUnknown => RhiAdminHandlerErrorKind::Unavailable,
    583         RhiAdminOperationErrorKind::InvalidMode
    584         | RhiAdminOperationErrorKind::InvalidInput
    585         | RhiAdminOperationErrorKind::Binding
    586         | RhiAdminOperationErrorKind::Transaction => RhiAdminHandlerErrorKind::Internal,
    587     })
    588 }
    589 
    590 const fn failure(kind: RhiAdminHandlerErrorKind) -> RhiAdminHandlerError {
    591     RhiAdminHandlerError::new(kind)
    592 }
    593 
    594 fn lower_hex(bytes: &[u8]) -> String {
    595     use core::fmt::Write as _;
    596 
    597     let mut encoded = String::with_capacity(bytes.len().saturating_mul(2));
    598     for byte in bytes {
    599         write!(&mut encoded, "{byte:02x}").expect("writing to String cannot fail");
    600     }
    601     encoded
    602 }
    603 
    604 fn parse_trade_id(value: &str) -> Result<TradeId, RhiAdminHandlerError> {
    605     TradeId::parse(value).map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
    606 }
    607 
    608 fn parse_digest(value: &str) -> Result<[u8; 32], RhiAdminHandlerError> {
    609     if value.len() != 64
    610         || value
    611             .as_bytes()
    612             .iter()
    613             .any(|byte| !byte.is_ascii_hexdigit() || byte.is_ascii_uppercase())
    614     {
    615         return Err(failure(RhiAdminHandlerErrorKind::Internal));
    616     }
    617     let mut bytes = [0_u8; 32];
    618     for (index, chunk) in value.as_bytes().chunks_exact(2).enumerate() {
    619         let high =
    620             hex_nibble(chunk[0]).ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?;
    621         let low =
    622             hex_nibble(chunk[1]).ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?;
    623         bytes[index] = (high << 4) | low;
    624     }
    625     Ok(bytes)
    626 }
    627 
    628 const fn hex_nibble(byte: u8) -> Option<u8> {
    629     match byte {
    630         b'0'..=b'9' => Some(byte - b'0'),
    631         b'a'..=b'f' => Some(byte - b'a' + 10),
    632         _ => None,
    633     }
    634 }
    635 
    636 fn page_limit(model: &Value) -> Result<u16, RhiAdminHandlerError> {
    637     model
    638         .pointer("/limit")
    639         .and_then(Value::as_u64)
    640         .and_then(|value| u16::try_from(value).ok())
    641         .filter(|value| (1..=200).contains(value))
    642         .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))
    643 }
    644 
    645 fn query_digest(model: &Value) -> Result<[u8; 32], RhiAdminHandlerError> {
    646     let mut model = model.clone();
    647     model
    648         .as_object_mut()
    649         .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?
    650         .remove("cursor");
    651     let bytes =
    652         serde_json::to_vec(&model).map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    653     Ok(Sha256::digest(bytes).into())
    654 }
    655 
    656 fn cursor_authenticator(key: &[u8; 32], payload: &[u8]) -> Hmac<Sha256> {
    657     let mut digest = Hmac::<Sha256>::new_from_slice(key).expect("SHA-256 HMAC accepts every key");
    658     digest.update(b"radroots.rhi.admin_cursor.v1\0");
    659     digest.update(
    660         &u64::try_from(payload.len())
    661             .expect("cursor payload length fits u64")
    662             .to_be_bytes(),
    663     );
    664     digest.update(payload);
    665     digest
    666 }
    667 
    668 fn cursor_tag(key: &[u8; 32], payload: &[u8]) -> [u8; 32] {
    669     cursor_authenticator(key, payload)
    670         .finalize()
    671         .into_bytes()
    672         .into()
    673 }
    674 
    675 fn encode_runtime_cursor(
    676     key: &[u8; 32],
    677     route: RhiAdminRoute,
    678     snapshot: u64,
    679     offset: u32,
    680     query_digest: [u8; 32],
    681 ) -> String {
    682     let mut bytes = Vec::with_capacity(78);
    683     bytes.push(1);
    684     bytes.push(route_code(route));
    685     bytes.extend_from_slice(&snapshot.to_be_bytes());
    686     bytes.extend_from_slice(&offset.to_be_bytes());
    687     bytes.extend_from_slice(&query_digest);
    688     let tag = cursor_tag(key, &bytes);
    689     bytes.extend_from_slice(&tag);
    690     URL_SAFE_NO_PAD.encode(bytes)
    691 }
    692 
    693 fn decode_runtime_cursor(
    694     key: &[u8; 32],
    695     route: RhiAdminRoute,
    696     encoded: &str,
    697     query_digest: [u8; 32],
    698 ) -> Result<(u64, u32), RhiAdminHandlerError> {
    699     let bytes = URL_SAFE_NO_PAD
    700         .decode(encoded)
    701         .map_err(|_| failure(RhiAdminHandlerErrorKind::InvalidCursor))?;
    702     if bytes.len() != 78
    703         || bytes[0] != 1
    704         || bytes[1] != route_code(route)
    705         || bytes[14..46] != query_digest
    706     {
    707         return Err(failure(RhiAdminHandlerErrorKind::InvalidCursor));
    708     }
    709     cursor_authenticator(key, &bytes[..46])
    710         .verify_slice(&bytes[46..])
    711         .map_err(|_| failure(RhiAdminHandlerErrorKind::InvalidCursor))?;
    712     let snapshot = u64::from_be_bytes(
    713         bytes[2..10]
    714             .try_into()
    715             .map_err(|_| failure(RhiAdminHandlerErrorKind::InvalidCursor))?,
    716     );
    717     let offset = u32::from_be_bytes(
    718         bytes[10..14]
    719             .try_into()
    720             .map_err(|_| failure(RhiAdminHandlerErrorKind::InvalidCursor))?,
    721     );
    722     Ok((snapshot, offset))
    723 }
    724 
    725 const fn route_code(route: RhiAdminRoute) -> u8 {
    726     match route {
    727         RhiAdminRoute::Status => 0,
    728         RhiAdminRoute::EffectiveConfig => 1,
    729         RhiAdminRoute::IdentityStatus => 2,
    730         RhiAdminRoute::IdentityPublic => 3,
    731         RhiAdminRoute::StateStatus => 4,
    732         RhiAdminRoute::StateBackup => 5,
    733         RhiAdminRoute::MetricsSnapshot => 6,
    734         RhiAdminRoute::ReconciliationStatus => 7,
    735         RhiAdminRoute::ReconciliationJobs => 8,
    736         RhiAdminRoute::ReconciliationRefresh => 9,
    737         RhiAdminRoute::Sources => 10,
    738         RhiAdminRoute::TradeProjection => 11,
    739         RhiAdminRoute::TradeReportCurrent => 12,
    740         RhiAdminRoute::TradeReports => 13,
    741         RhiAdminRoute::PublicationBacklog => 14,
    742         RhiAdminRoute::PublicationTargets => 15,
    743         RhiAdminRoute::PublicationRetry => 16,
    744         RhiAdminRoute::PresenceDesired => 17,
    745         RhiAdminRoute::PresenceRender => 18,
    746         RhiAdminRoute::PresenceRefresh => 19,
    747     }
    748 }
    749 
    750 fn row_u64(row: &sqlx::sqlite::SqliteRow, name: &str) -> Result<u64, RhiAdminHandlerError> {
    751     row.try_get::<i64, _>(name)
    752         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
    753         .and_then(|value| {
    754             u64::try_from(value).map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
    755         })
    756 }
    757 
    758 fn optional_nonnegative_i64(
    759     row: &sqlx::sqlite::SqliteRow,
    760     name: &str,
    761 ) -> Result<Option<u64>, RhiAdminHandlerError> {
    762     row.try_get::<Option<i64>, _>(name)
    763         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?
    764         .map(|value| u64::try_from(value).map_err(|_| failure(RhiAdminHandlerErrorKind::Internal)))
    765         .transpose()
    766 }
    767 
    768 fn exact_blob<const N: usize>(
    769     row: &sqlx::sqlite::SqliteRow,
    770     value: &str,
    771     length: &str,
    772 ) -> Result<[u8; N], RhiAdminHandlerError> {
    773     if row_u64(row, length)? != N as u64 {
    774         return Err(failure(RhiAdminHandlerErrorKind::Internal));
    775     }
    776     row.try_get::<Vec<u8>, _>(value)
    777         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?
    778         .try_into()
    779         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
    780 }
    781 
    782 fn optional_exact_blob<const N: usize>(
    783     row: &sqlx::sqlite::SqliteRow,
    784     value: &str,
    785     length: &str,
    786 ) -> Result<Option<[u8; N]>, RhiAdminHandlerError> {
    787     let length = row
    788         .try_get::<Option<i64>, _>(length)
    789         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    790     match length {
    791         None => {
    792             if row
    793                 .try_get::<Option<Vec<u8>>, _>(value)
    794                 .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?
    795                 .is_some()
    796             {
    797                 return Err(failure(RhiAdminHandlerErrorKind::Internal));
    798             }
    799             Ok(None)
    800         }
    801         Some(length) if length == N as i64 => row
    802             .try_get::<Option<Vec<u8>>, _>(value)
    803             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?
    804             .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?
    805             .try_into()
    806             .map(Some)
    807             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal)),
    808         Some(_) => Err(failure(RhiAdminHandlerErrorKind::Internal)),
    809     }
    810 }
    811 
    812 fn bounded_blob(
    813     row: &sqlx::sqlite::SqliteRow,
    814     value: &str,
    815     length: &str,
    816     maximum: usize,
    817 ) -> Result<Box<[u8]>, RhiAdminHandlerError> {
    818     let length = row_u64(row, length)?;
    819     if length == 0 || length > maximum as u64 {
    820         return Err(failure(RhiAdminHandlerErrorKind::Internal));
    821     }
    822     let bytes = row
    823         .try_get::<Vec<u8>, _>(value)
    824         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    825     if bytes.len() as u64 != length {
    826         return Err(failure(RhiAdminHandlerErrorKind::Internal));
    827     }
    828     Ok(bytes.into_boxed_slice())
    829 }
    830 
    831 fn job_summary(row: &sqlx::sqlite::SqliteRow) -> Result<Value, RhiAdminHandlerError> {
    832     let job_id = exact_blob::<32>(row, "job_id", "job_id_bytes")?;
    833     let trade_id = exact_blob::<16>(row, "trade_id", "trade_id_bytes")?;
    834     let attempt_count = row_u64(row, "attempt_count")?;
    835     let state = row
    836         .try_get::<&str, _>("state")
    837         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    838     let state = match state {
    839         "ready" if attempt_count == 0 => "pending",
    840         "ready" => "retry_scheduled",
    841         "leased" => "leased",
    842         "completed" => "completed",
    843         "exhausted" | "superseded" => "failed",
    844         _ => return Err(failure(RhiAdminHandlerErrorKind::Internal)),
    845     };
    846     let mut value = json!({
    847         "attempt_count": attempt_count,
    848         "dirty_generation": row_u64(row, "input_generation")?,
    849         "job_id": lower_hex(&job_id),
    850         "reason_codes": [],
    851         "scheduled_at_utc": row_u64(row, "created_at_unix_ms")? / 1_000,
    852         "state": state,
    853         "trade_id": lower_hex(&trade_id),
    854     });
    855     if let Some(expires) = optional_nonnegative_i64(row, "lease_expires_unix_ms")? {
    856         value["lease_expires_at_utc"] = Value::from(expires / 1_000);
    857     }
    858     Ok(value)
    859 }
    860 
    861 struct ConfiguredSourceSummary {
    862     source_id: Box<str>,
    863     required: bool,
    864 }
    865 
    866 fn configured_source_summaries(
    867     configuration: &RhiConfigDocumentV1,
    868 ) -> Result<Vec<ConfiguredSourceSummary>, RhiAdminHandlerError> {
    869     let sources = configuration
    870         .normalized()
    871         .pointer("/evidence/sources")
    872         .and_then(Value::as_array)
    873         .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?;
    874     if sources.is_empty() || sources.len() > 16 {
    875         return Err(failure(RhiAdminHandlerErrorKind::Internal));
    876     }
    877     sources
    878         .iter()
    879         .map(|source| {
    880             let source_id = source
    881                 .pointer("/source_id")
    882                 .and_then(Value::as_str)
    883                 .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?;
    884             let required = source
    885                 .pointer("/required")
    886                 .and_then(Value::as_bool)
    887                 .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?;
    888             Ok(ConfiguredSourceSummary {
    889                 source_id: source_id.into(),
    890                 required,
    891             })
    892         })
    893         .collect()
    894 }
    895 
    896 fn request_trade_id(request: &RhiAdminRequestDocument) -> Result<TradeId, RhiAdminHandlerError> {
    897     request
    898         .parameter("trade_id")
    899         .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))
    900         .and_then(parse_trade_id)
    901 }
    902 
    903 fn contract_outcome_to_db(value: &str) -> Result<&'static str, RhiAdminHandlerError> {
    904     match value {
    905         "Valid" => Ok("valid"),
    906         "Invalid" => Ok("invalid"),
    907         "Indeterminate" => Ok("indeterminate"),
    908         _ => Err(failure(RhiAdminHandlerErrorKind::Internal)),
    909     }
    910 }
    911 
    912 fn db_outcome_to_contract(value: &str) -> Result<&'static str, RhiAdminHandlerError> {
    913     match value {
    914         "valid" => Ok("Valid"),
    915         "invalid" => Ok("Invalid"),
    916         "indeterminate" => Ok("Indeterminate"),
    917         _ => Err(failure(RhiAdminHandlerErrorKind::Internal)),
    918     }
    919 }
    920 
    921 fn report_coverage(row: &sqlx::sqlite::SqliteRow) -> Result<&'static str, RhiAdminHandlerError> {
    922     let source_count = row_u64(row, "source_count")?;
    923     let incomplete = row_u64(row, "incomplete_required")?;
    924     let unsupported = row_u64(row, "unsupported_required")?;
    925     let accepted = row_u64(row, "accepted_event_count")?;
    926     if source_count == 0 {
    927         return Err(failure(RhiAdminHandlerErrorKind::Internal));
    928     }
    929     Ok(if unsupported > 0 {
    930         "Unsupported"
    931     } else if incomplete > 0 && accepted == 0 {
    932         "Missing"
    933     } else if incomplete > 0 {
    934         "Partial"
    935     } else {
    936         "ScopeSatisfied"
    937     })
    938 }
    939 
    940 fn report_detail(row: &sqlx::sqlite::SqliteRow) -> Result<Value, RhiAdminHandlerError> {
    941     let canonical = bounded_blob(row, "canonical_report", "canonical_report_bytes", 16_384)?;
    942     let mut report: Value = serde_json::from_slice(&canonical)
    943         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    944     let object = report
    945         .as_object_mut()
    946         .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?;
    947     let statement = lower_hex(&exact_blob::<32>(
    948         row,
    949         "statement_sha256",
    950         "statement_bytes",
    951     )?);
    952     let manifest = lower_hex(&exact_blob::<32>(row, "manifest_sha256", "manifest_bytes")?);
    953     let projection = lower_hex(&exact_blob::<32>(
    954         row,
    955         "projection_sha256",
    956         "projection_bytes",
    957     )?);
    958     let trade = lower_hex(&exact_blob::<16>(row, "trade_id", "trade_id_bytes")?);
    959     let claim = lower_hex(&exact_blob::<32>(row, "claim_mutation_id", "claim_bytes")?);
    960     let issuer = lower_hex(&exact_blob::<32>(row, "issuer_public_key", "issuer_bytes")?);
    961     let policy = lower_hex(&exact_blob::<32>(
    962         row,
    963         "evidence_policy_sha256",
    964         "policy_bytes",
    965     )?);
    966     let observed = row_u64(row, "observed_at_unix_s")?;
    967     let generation = row_u64(row, "trade_generation")?;
    968     let outcome = row
    969         .try_get::<&str, _>("outcome")
    970         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    971     let reducer_contract = row
    972         .try_get::<&str, _>("reducer_contract")
    973         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
    974     let reducer_version = row_u64(row, "reducer_contract_version")?;
    975     let required_matches = [
    976         ("report_id", statement.as_str()),
    977         ("evidence_manifest_digest", manifest.as_str()),
    978         ("projection_digest", projection.as_str()),
    979         ("trade_id", trade.as_str()),
    980         ("claim_mutation_id", claim.as_str()),
    981         ("issuer_pubkey", issuer.as_str()),
    982         ("evidence_policy_digest", policy.as_str()),
    983         ("reducer_contract_id", reducer_contract),
    984     ]
    985     .into_iter()
    986     .all(|(field, expected)| object.get(field).and_then(Value::as_str) == Some(expected));
    987     if !required_matches
    988         || object.get("trade_generation").and_then(Value::as_u64) != Some(generation)
    989         || object.get("observed_at_unix_s").and_then(Value::as_u64) != Some(observed)
    990         || object
    991             .get("reducer_contract_version")
    992             .and_then(Value::as_u64)
    993             != Some(reducer_version)
    994         || object.get("outcome").and_then(Value::as_str) != Some(outcome)
    995         || object.get("statement_digest").and_then(Value::as_str) != Some(statement.as_str())
    996     {
    997         return Err(failure(RhiAdminHandlerErrorKind::Internal));
    998     }
    999     let supersedes =
   1000         optional_exact_blob::<32>(row, "supersedes_statement_sha256", "supersedes_bytes")?;
   1001     match supersedes {
   1002         Some(value)
   1003             if object.get("supersedes_report_id").and_then(Value::as_str)
   1004                 == Some(lower_hex(&value).as_str()) => {}
   1005         None if object
   1006             .get("supersedes_report_id")
   1007             .is_some_and(Value::is_null) =>
   1008         {
   1009             object.remove("supersedes_report_id");
   1010             object.remove("supersedes_event_id");
   1011         }
   1012         _ => return Err(failure(RhiAdminHandlerErrorKind::Internal)),
   1013     }
   1014     object.remove("statement_digest");
   1015     object.insert(
   1016         "attestation_event_id".to_owned(),
   1017         Value::String(lower_hex(&exact_blob::<32>(
   1018             row,
   1019             "event_id",
   1020             "event_id_bytes",
   1021         )?)),
   1022     );
   1023     object.insert(
   1024         "coverage".to_owned(),
   1025         Value::String(report_coverage(row)?.to_owned()),
   1026     );
   1027     object.insert(
   1028         "outcome".to_owned(),
   1029         Value::String(db_outcome_to_contract(outcome)?.to_owned()),
   1030     );
   1031     Ok(report)
   1032 }
   1033 
   1034 fn report_summary(row: &sqlx::sqlite::SqliteRow) -> Result<Value, RhiAdminHandlerError> {
   1035     let supersedes =
   1036         optional_exact_blob::<32>(row, "supersedes_statement_sha256", "supersedes_bytes")?;
   1037     let outcome = row
   1038         .try_get::<&str, _>("outcome")
   1039         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1040     let mut value = json!({
   1041         "attestation_event_id": lower_hex(&exact_blob::<32>(row, "event_id", "event_id_bytes")?),
   1042         "claim_mutation_id": lower_hex(&exact_blob::<32>(row, "claim_mutation_id", "claim_bytes")?),
   1043         "coverage": report_coverage(row)?,
   1044         "observed_at_utc": row_u64(row, "observed_at_unix_s")?,
   1045         "outcome": db_outcome_to_contract(outcome)?,
   1046         "report_id": lower_hex(&exact_blob::<32>(row, "statement_sha256", "statement_bytes")?),
   1047         "trade_id": lower_hex(&exact_blob::<16>(row, "trade_id", "trade_id_bytes")?),
   1048     });
   1049     if let Some(supersedes) = supersedes {
   1050         value["supersedes_report_id"] = Value::String(lower_hex(&supersedes));
   1051     }
   1052     Ok(value)
   1053 }
   1054 
   1055 fn publication_summary(row: &sqlx::sqlite::SqliteRow) -> Result<Value, RhiAdminHandlerError> {
   1056     let state = row
   1057         .try_get::<&str, _>("public_state")
   1058         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1059     if !matches!(
   1060         state,
   1061         "pending"
   1062             | "submitted"
   1063             | "accepted"
   1064             | "rejected"
   1065             | "rate_limited"
   1066             | "auth_required"
   1067             | "failed"
   1068             | "unknown"
   1069     ) {
   1070         return Err(failure(RhiAdminHandlerErrorKind::Internal));
   1071     }
   1072     let mut value = json!({
   1073         "attempt_count": row_u64(row, "attempt_count")?,
   1074         "event_id": lower_hex(&exact_blob::<32>(row, "event_id", "event_id_bytes")?),
   1075         "exact_bytes_digest": lower_hex(&exact_blob::<32>(row, "event_sha256", "event_sha256_bytes")?),
   1076         "reason_codes": if state == "failed" || state == "unknown" { vec!["publication_blocked"] } else { Vec::<&str>::new() },
   1077         "report_id": lower_hex(&exact_blob::<32>(row, "statement_sha256", "statement_bytes")?),
   1078         "state": state,
   1079         "workflow_id": lower_hex(&exact_blob::<32>(row, "outbox_id", "outbox_id_bytes")?),
   1080     });
   1081     if let Some(next) = optional_nonnegative_i64(row, "next_attempt_unix_ms")? {
   1082         value["next_attempt_at_utc"] = Value::from(next / 1_000);
   1083     }
   1084     Ok(value)
   1085 }
   1086 
   1087 fn publication_target_summary(
   1088     row: &sqlx::sqlite::SqliteRow,
   1089 ) -> Result<Value, RhiAdminHandlerError> {
   1090     let target = row
   1091         .try_get::<&str, _>("relay_id")
   1092         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1093     let state = row
   1094         .try_get::<&str, _>("state")
   1095         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1096     if target.is_empty()
   1097         || target.len() > 64
   1098         || !matches!(
   1099             state,
   1100             "pending"
   1101                 | "submitted"
   1102                 | "accepted"
   1103                 | "rejected"
   1104                 | "rate_limited"
   1105                 | "auth_required"
   1106                 | "failed"
   1107                 | "unknown"
   1108         )
   1109     {
   1110         return Err(failure(RhiAdminHandlerErrorKind::Internal));
   1111     }
   1112     let mut value = json!({
   1113         "attempt_count": row_u64(row, "attempt_count")?,
   1114         "reason_codes": if state == "failed" || state == "unknown" { vec!["publication_target_failed"] } else { Vec::<&str>::new() },
   1115         "state": state,
   1116         "target_id": target,
   1117         "workflow_id": lower_hex(&exact_blob::<32>(row, "outbox_id", "outbox_id_bytes")?),
   1118     });
   1119     if let Some(last) = optional_nonnegative_i64(row, "last_attempt_unix_ms")? {
   1120         value["last_attempt_at_utc"] = Value::from(last / 1_000);
   1121     }
   1122     Ok(value)
   1123 }
   1124 
   1125 fn operation_row_u64(
   1126     row: &sqlx::sqlite::SqliteRow,
   1127     name: &str,
   1128 ) -> Result<u64, AdminJournalOperationError> {
   1129     row.try_get::<i64, _>(name)
   1130         .ok()
   1131         .and_then(|value| u64::try_from(value).ok())
   1132         .ok_or(AdminJournalOperationError::Binding)
   1133 }
   1134 
   1135 fn operation_exact_blob<const N: usize>(
   1136     row: &sqlx::sqlite::SqliteRow,
   1137     value: &str,
   1138     length: &str,
   1139 ) -> Result<[u8; N], AdminJournalOperationError> {
   1140     if operation_row_u64(row, length)? != N as u64 {
   1141         return Err(AdminJournalOperationError::Binding);
   1142     }
   1143     row.try_get::<Vec<u8>, _>(value)
   1144         .map_err(|_| AdminJournalOperationError::Binding)?
   1145         .try_into()
   1146         .map_err(|_| AdminJournalOperationError::Binding)
   1147 }
   1148 
   1149 fn presence_kind(row: &sqlx::sqlite::SqliteRow) -> Result<&str, RhiAdminHandlerError> {
   1150     match row
   1151         .try_get::<&str, _>("document_kind")
   1152         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?
   1153     {
   1154         kind @ ("service_profile" | "application_handler") => Ok(kind),
   1155         _ => Err(failure(RhiAdminHandlerErrorKind::Internal)),
   1156     }
   1157 }
   1158 
   1159 fn operation_presence_kind(
   1160     row: &sqlx::sqlite::SqliteRow,
   1161 ) -> Result<&str, AdminJournalOperationError> {
   1162     match row
   1163         .try_get::<&str, _>("document_kind")
   1164         .map_err(|_| AdminJournalOperationError::Binding)?
   1165     {
   1166         kind @ ("service_profile" | "application_handler") => Ok(kind),
   1167         _ => Err(AdminJournalOperationError::Binding),
   1168     }
   1169 }
   1170 
   1171 // Domain-route methods are kept below the transport-independent common boundary.
   1172 impl RuntimeAdminHandler {
   1173     async fn reconciliation_status(
   1174         &self,
   1175         request: RhiAdminRequestDocument,
   1176     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
   1177         let policy_digest = lower_hex(self.state.metadata().evidence_policy_digest().as_bytes());
   1178         let value = self
   1179             .state
   1180             .sqlite_host()
   1181             .transaction(|transaction| {
   1182                 Box::pin(async move {
   1183                     let row = sqlx::query(
   1184                         r#"SELECT
   1185                             (SELECT COUNT(*) FROM trade_dirty_generations) AS dirty_count,
   1186                             COUNT(CASE WHEN state = 'ready' AND attempt_count = 0 THEN 1 END) AS pending_count,
   1187                             COUNT(CASE WHEN state = 'ready' AND attempt_count > 0 THEN 1 END) AS retry_count,
   1188                             COUNT(CASE WHEN state = 'leased' THEN 1 END) AS leased_count,
   1189                             COUNT(CASE WHEN state = 'completed' THEN 1 END) AS completed_count,
   1190                             COUNT(CASE WHEN state IN ('exhausted', 'superseded') THEN 1 END) AS failed_count,
   1191                             MIN(CASE WHEN state = 'ready' THEN created_at_unix_ms END) AS oldest_pending_ms,
   1192                             (SELECT COUNT(*) FROM evidence_reconciliation_sources
   1193                                 WHERE required = 1 AND completion != 'complete') AS required_failures
   1194                         FROM reconciliation_jobs"#,
   1195                     )
   1196                     .fetch_one(&mut *transaction)
   1197                     .await
   1198                     .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1199                     let mut counts = Map::new();
   1200                     for (code, column) in [
   1201                         ("pending", "pending_count"),
   1202                         ("retry_scheduled", "retry_count"),
   1203                         ("leased", "leased_count"),
   1204                         ("completed", "completed_count"),
   1205                         ("failed", "failed_count"),
   1206                     ] {
   1207                         counts.insert(code.to_owned(), Value::from(row_u64(&row, column)?));
   1208                     }
   1209                     let mut value = json!({
   1210                         "dirty_trade_count": row_u64(&row, "dirty_count")?,
   1211                         "job_counts": counts,
   1212                         "policy_digest": policy_digest,
   1213                         "required_source_failures": row_u64(&row, "required_failures")?,
   1214                     });
   1215                     if let Some(value_ms) = optional_nonnegative_i64(&row, "oldest_pending_ms")? {
   1216                         value["oldest_pending_at_utc"] = Value::from(value_ms / 1_000);
   1217                     }
   1218                     Ok::<Value, RhiAdminHandlerError>(value)
   1219                 })
   1220             })
   1221             .await
   1222             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1223         response(request.route(), value)
   1224     }
   1225 
   1226     async fn reconciliation_jobs(
   1227         &self,
   1228         request: RhiAdminRequestDocument,
   1229     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
   1230         let model = request_model(&request)?;
   1231         let limit = page_limit(&model)?;
   1232         let state_filter = model
   1233             .pointer("/state")
   1234             .and_then(Value::as_str)
   1235             .map(str::to_owned);
   1236         let trade_filter = model
   1237             .pointer("/trade_id")
   1238             .and_then(Value::as_str)
   1239             .map(parse_trade_id)
   1240             .transpose()?
   1241             .map(|trade| trade.as_bytes().to_vec());
   1242         let query_digest = query_digest(&model)?;
   1243         let (snapshot, offset) =
   1244             self.cursor_or_new(&model, request.route(), query_digest, self.now_millis()?)?;
   1245         let fetch_limit = i64::from(limit) + 1;
   1246         let offset_i64 = i64::from(offset);
   1247         let rows = self
   1248             .state
   1249             .sqlite_host()
   1250             .transaction(move |transaction| {
   1251                 Box::pin(async move {
   1252                     sqlx::query(
   1253                         r#"SELECT
   1254                             length(job_id) AS job_id_bytes, substr(job_id, 1, 33) AS job_id,
   1255                             length(trade_id) AS trade_id_bytes, substr(trade_id, 1, 17) AS trade_id,
   1256                             state, attempt_count, input_generation, created_at_unix_ms,
   1257                             lease_expires_unix_ms
   1258                         FROM reconciliation_jobs
   1259                         WHERE updated_at_unix_ms <= ?
   1260                           AND (? IS NULL OR
   1261                             (? = 'pending' AND state = 'ready' AND attempt_count = 0) OR
   1262                             (? = 'retry_scheduled' AND state = 'ready' AND attempt_count > 0) OR
   1263                             (? = 'leased' AND state = 'leased') OR
   1264                             (? = 'completed' AND state = 'completed') OR
   1265                             (? = 'failed' AND state IN ('exhausted', 'superseded')))
   1266                           AND (? IS NULL OR trade_id = ?)
   1267                         ORDER BY updated_at_unix_ms DESC, job_id DESC
   1268                         LIMIT ? OFFSET ?"#,
   1269                     )
   1270                     .bind(
   1271                         i64::try_from(snapshot)
   1272                             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?,
   1273                     )
   1274                     .bind(state_filter.as_deref())
   1275                     .bind(state_filter.as_deref())
   1276                     .bind(state_filter.as_deref())
   1277                     .bind(state_filter.as_deref())
   1278                     .bind(state_filter.as_deref())
   1279                     .bind(state_filter.as_deref())
   1280                     .bind(trade_filter.as_deref())
   1281                     .bind(trade_filter.as_deref())
   1282                     .bind(fetch_limit)
   1283                     .bind(offset_i64)
   1284                     .fetch_all(&mut *transaction)
   1285                     .await
   1286                     .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
   1287                 })
   1288             })
   1289             .await
   1290             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1291         let has_more = rows.len() > usize::from(limit);
   1292         let items = rows
   1293             .iter()
   1294             .take(usize::from(limit))
   1295             .map(job_summary)
   1296             .collect::<Result<Vec<_>, _>>()?;
   1297         let mut value = json!({"items":items,"snapshot_generation":snapshot});
   1298         if has_more {
   1299             value["next_cursor"] = Value::String(
   1300                 self.encode_cursor(
   1301                     request.route(),
   1302                     snapshot,
   1303                     offset
   1304                         .checked_add(u32::from(limit))
   1305                         .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?,
   1306                     query_digest,
   1307                 ),
   1308             );
   1309         }
   1310         response(request.route(), value)
   1311     }
   1312 
   1313     async fn reconciliation_refresh(
   1314         &self,
   1315         request: RhiAdminRequestDocument,
   1316     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
   1317         let model = request_model(&request)?;
   1318         let trade = model
   1319             .pointer("/trade_id")
   1320             .and_then(Value::as_str)
   1321             .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))
   1322             .and_then(parse_trade_id)?;
   1323         let expected = required_u64(&model, "/expected_dirty_generation")?;
   1324         let actual = self.read_dirty_generation(trade).await?;
   1325         if actual != expected {
   1326             return Err(failure(RhiAdminHandlerErrorKind::Conflict));
   1327         }
   1328         let now_ms = self.now_millis()?;
   1329         let prepared = match RhiAdminOperationRepository::new(&self.state)
   1330             .prepare_admin_operation(&request, operation_time(now_ms)?)
   1331             .await
   1332             .map_err(map_journal_error)?
   1333         {
   1334             RhiAdminOperationAdmission::ExactReplay(response) => return Ok(response),
   1335             RhiAdminOperationAdmission::Prepared(prepared) => prepared,
   1336         };
   1337         let policy = RhiReconciliationJobPolicy::from_configuration(&self.configuration)
   1338             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1339         let outcome = self
   1340             .state
   1341             .repositories()
   1342             .reconciliation_jobs()
   1343             .schedule_trade(
   1344                 trade,
   1345                 policy,
   1346                 RhiReconciliationUnixMilliseconds::new(now_ms)
   1347                     .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?,
   1348             )
   1349             .await
   1350             .map_err(|error| match error.kind() {
   1351                 RhiReconciliationJobErrorKind::DirtyGenerationConflict => {
   1352                     failure(RhiAdminHandlerErrorKind::Conflict)
   1353                 }
   1354                 RhiReconciliationJobErrorKind::QueueFull => {
   1355                     failure(RhiAdminHandlerErrorKind::Unavailable)
   1356                 }
   1357                 _ => failure(RhiAdminHandlerErrorKind::Internal),
   1358             })?;
   1359         let job = outcome.job();
   1360         let completed_ms = self.now_millis()?;
   1361         let completed = response(
   1362             request.route(),
   1363             json!({
   1364                 "dirty_generation": job.input_generation(),
   1365                 "job_id": lower_hex(job.id().as_bytes()),
   1366                 "operation_id": required_operation_id(&request)?,
   1367                 "trade_id": lower_hex(job.trade_id().as_bytes()),
   1368             }),
   1369         )?;
   1370         RhiAdminOperationRepository::new(&self.state)
   1371             .complete_admin_operation(
   1372                 &prepared,
   1373                 &completed,
   1374                 operation_time(completed_ms)?,
   1375                 RhiAdminOperationJournalPolicy::seven_days(),
   1376             )
   1377             .await
   1378             .map_err(map_journal_error)?;
   1379         Ok(completed)
   1380     }
   1381     async fn sources(
   1382         &self,
   1383         request: RhiAdminRequestDocument,
   1384     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
   1385         let model = request_model(&request)?;
   1386         let limit = page_limit(&model)?;
   1387         let required_filter = model.pointer("/required").and_then(Value::as_bool);
   1388         let completion_filter = model
   1389             .pointer("/completion")
   1390             .and_then(Value::as_str)
   1391             .map(str::to_owned);
   1392         let sources = configured_source_summaries(&self.configuration)?;
   1393         let query_digest = query_digest(&model)?;
   1394         let (snapshot, offset) =
   1395             self.cursor_or_new(&model, request.route(), query_digest, self.now_millis()?)?;
   1396         let snapshot_sql =
   1397             i64::try_from(snapshot).map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1398         let rows = self
   1399             .state
   1400             .sqlite_host()
   1401             .transaction(move |transaction| {
   1402                 Box::pin(async move {
   1403                     let mut items = Vec::with_capacity(sources.len());
   1404                     for source in sources {
   1405                         let attempt = sqlx::query(
   1406                             r#"SELECT completion, finished_unix_ms
   1407                             FROM evidence_reconciliation_sources
   1408                             WHERE source_id = ? AND finished_unix_ms <= ?
   1409                             ORDER BY finished_unix_ms DESC, attempt_id DESC, request_id DESC
   1410                             LIMIT 1"#,
   1411                         )
   1412                         .bind(source.source_id.as_ref())
   1413                         .bind(snapshot_sql)
   1414                         .fetch_optional(&mut *transaction)
   1415                         .await
   1416                         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1417                         let completion = attempt
   1418                             .as_ref()
   1419                             .map(|row| {
   1420                                 row.try_get::<&str, _>("completion")
   1421                                     .map(str::to_owned)
   1422                                     .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
   1423                             })
   1424                             .transpose()?
   1425                             .unwrap_or_else(|| "incomplete_unknown".to_owned());
   1426                         if required_filter.is_some_and(|required| required != source.required)
   1427                             || completion_filter
   1428                                 .as_deref()
   1429                                 .is_some_and(|expected| expected != completion)
   1430                         {
   1431                             continue;
   1432                         }
   1433                         let checkpoint = sqlx::query(
   1434                             r#"SELECT cursor_created_at_unix_s,
   1435                                 length(cursor_event_id) AS event_id_bytes,
   1436                                 substr(cursor_event_id, 1, 33) AS cursor_event_id
   1437                             FROM relay_checkpoints
   1438                             WHERE source_id = ? AND completed_at_unix_s <= ?
   1439                             ORDER BY completed_at_unix_s DESC, trade_id DESC
   1440                             LIMIT 1"#,
   1441                         )
   1442                         .bind(source.source_id.as_ref())
   1443                         .bind(snapshot_sql / 1_000)
   1444                         .fetch_optional(&mut *transaction)
   1445                         .await
   1446                         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1447                         let mut value = json!({
   1448                             "completion": completion,
   1449                             "reason_codes": [],
   1450                             "required": source.required,
   1451                             "source_id": source.source_id,
   1452                         });
   1453                         if let Some(row) = attempt.as_ref() {
   1454                             value["last_attempt_at_utc"] =
   1455                                 Value::from(row_u64(row, "finished_unix_ms")? / 1_000);
   1456                         }
   1457                         if let Some(row) = checkpoint.as_ref() {
   1458                             value["cursor"] = json!({
   1459                                 "created_at_unix_seconds": row_u64(row, "cursor_created_at_unix_s")?,
   1460                                 "event_id_lowercase_hex": lower_hex(&exact_blob::<32>(row, "cursor_event_id", "event_id_bytes")?),
   1461                             });
   1462                         }
   1463                         items.push(value);
   1464                     }
   1465                     Ok::<Vec<Value>, RhiAdminHandlerError>(items)
   1466                 })
   1467             })
   1468             .await
   1469             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1470         let start =
   1471             usize::try_from(offset).map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1472         if start > rows.len() {
   1473             return Err(failure(RhiAdminHandlerErrorKind::InvalidCursor));
   1474         }
   1475         let end = start.saturating_add(usize::from(limit)).min(rows.len());
   1476         let items = rows[start..end].to_vec();
   1477         let mut value = json!({"items":items,"snapshot_generation":snapshot});
   1478         if end < rows.len() {
   1479             value["next_cursor"] = Value::String(self.encode_cursor(
   1480                 request.route(),
   1481                 snapshot,
   1482                 u32::try_from(end).map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?,
   1483                 query_digest,
   1484             ));
   1485         }
   1486         response(request.route(), value)
   1487     }
   1488     async fn trade_projection(
   1489         &self,
   1490         request: RhiAdminRequestDocument,
   1491     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
   1492         let trade = request_trade_id(&request)?;
   1493         let row = self.current_report_row(trade).await?;
   1494         let report = report_detail(&row)?;
   1495         response(
   1496             request.route(),
   1497             json!({
   1498                 "coverage": report["coverage"].clone(),
   1499                 "dirty_generation": report["trade_generation"].clone(),
   1500                 "manifest_digest": report["evidence_manifest_digest"].clone(),
   1501                 "observed_at_utc": report["observed_at_unix_s"].clone(),
   1502                 "outcome": report["outcome"].clone(),
   1503                 "policy_digest": report["evidence_policy_digest"].clone(),
   1504                 "projection_digest": report["projection_digest"].clone(),
   1505                 "reason_codes": report["reason_codes"].clone(),
   1506                 "trade_id": report["trade_id"].clone(),
   1507             }),
   1508         )
   1509     }
   1510     async fn trade_report_current(
   1511         &self,
   1512         request: RhiAdminRequestDocument,
   1513     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
   1514         let row = self.current_report_row(request_trade_id(&request)?).await?;
   1515         response(request.route(), report_detail(&row)?)
   1516     }
   1517     async fn trade_reports(
   1518         &self,
   1519         request: RhiAdminRequestDocument,
   1520     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
   1521         let model = request_model(&request)?;
   1522         let limit = page_limit(&model)?;
   1523         let trade = request_trade_id(&request)?;
   1524         let outcome = model
   1525             .pointer("/outcome")
   1526             .and_then(Value::as_str)
   1527             .map(contract_outcome_to_db)
   1528             .transpose()?
   1529             .map(str::to_owned);
   1530         let query_digest = query_digest(&model)?;
   1531         let (snapshot, offset) =
   1532             self.cursor_or_new(&model, request.route(), query_digest, self.now_seconds()?)?;
   1533         let trade_bytes = trade.as_bytes().to_vec();
   1534         let fetch_limit = i64::from(limit) + 1;
   1535         let offset_i64 = i64::from(offset);
   1536         let snapshot_sql =
   1537             i64::try_from(snapshot).map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1538         let rows = self
   1539             .state
   1540             .sqlite_host()
   1541             .transaction(move |transaction| {
   1542                 Box::pin(async move {
   1543                     let sql = format!("{REPORT_SELECT}{REPORT_ROWS_SUFFIX}");
   1544                     // Both fragments are private compile-time constants; all filters remain bound.
   1545                     sqlx::query(sqlx::AssertSqlSafe(sql.as_str()))
   1546                         .bind(trade_bytes)
   1547                         .bind(snapshot_sql)
   1548                         .bind(outcome.as_deref())
   1549                         .bind(outcome.as_deref())
   1550                         .bind(fetch_limit)
   1551                         .bind(offset_i64)
   1552                         .fetch_all(&mut *transaction)
   1553                         .await
   1554                         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
   1555                 })
   1556             })
   1557             .await
   1558             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1559         let has_more = rows.len() > usize::from(limit);
   1560         let items = rows
   1561             .iter()
   1562             .take(usize::from(limit))
   1563             .map(report_summary)
   1564             .collect::<Result<Vec<_>, _>>()?;
   1565         let mut value = json!({"items":items,"snapshot_generation":snapshot});
   1566         if has_more {
   1567             value["next_cursor"] = Value::String(
   1568                 self.encode_cursor(
   1569                     request.route(),
   1570                     snapshot,
   1571                     offset
   1572                         .checked_add(u32::from(limit))
   1573                         .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?,
   1574                     query_digest,
   1575                 ),
   1576             );
   1577         }
   1578         response(request.route(), value)
   1579     }
   1580     async fn publication_backlog(
   1581         &self,
   1582         request: RhiAdminRequestDocument,
   1583     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
   1584         let model = request_model(&request)?;
   1585         let limit = page_limit(&model)?;
   1586         let state_filter = model
   1587             .pointer("/state")
   1588             .and_then(Value::as_str)
   1589             .map(str::to_owned);
   1590         let query_digest = query_digest(&model)?;
   1591         let (snapshot, offset) =
   1592             self.cursor_or_new(&model, request.route(), query_digest, self.now_millis()?)?;
   1593         let rows = self
   1594             .state
   1595             .sqlite_host()
   1596             .transaction(move |transaction| {
   1597                 Box::pin(async move {
   1598                     sqlx::query(PUBLICATION_BACKLOG_SQL)
   1599                         .bind(
   1600                             i64::try_from(snapshot)
   1601                                 .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?,
   1602                         )
   1603                         .bind(state_filter.as_deref())
   1604                         .bind(state_filter.as_deref())
   1605                         .bind(i64::from(limit) + 1)
   1606                         .bind(i64::from(offset))
   1607                         .fetch_all(&mut *transaction)
   1608                         .await
   1609                         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
   1610                 })
   1611             })
   1612             .await
   1613             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1614         let has_more = rows.len() > usize::from(limit);
   1615         let items = rows
   1616             .iter()
   1617             .take(usize::from(limit))
   1618             .map(publication_summary)
   1619             .collect::<Result<Vec<_>, _>>()?;
   1620         let mut value = json!({"items":items,"snapshot_generation":snapshot});
   1621         if has_more {
   1622             value["next_cursor"] = Value::String(
   1623                 self.encode_cursor(
   1624                     request.route(),
   1625                     snapshot,
   1626                     offset
   1627                         .checked_add(u32::from(limit))
   1628                         .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?,
   1629                     query_digest,
   1630                 ),
   1631             );
   1632         }
   1633         response(request.route(), value)
   1634     }
   1635     async fn publication_targets(
   1636         &self,
   1637         request: RhiAdminRequestDocument,
   1638     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
   1639         let model = request_model(&request)?;
   1640         let limit = page_limit(&model)?;
   1641         let state_filter = model
   1642             .pointer("/state")
   1643             .and_then(Value::as_str)
   1644             .map(str::to_owned);
   1645         let workflow = model
   1646             .pointer("/workflow_id")
   1647             .and_then(Value::as_str)
   1648             .map(parse_digest)
   1649             .transpose()?
   1650             .map(Vec::from);
   1651         let query_digest = query_digest(&model)?;
   1652         let (snapshot, offset) =
   1653             self.cursor_or_new(&model, request.route(), query_digest, self.now_millis()?)?;
   1654         let rows = self
   1655             .state
   1656             .sqlite_host()
   1657             .transaction(move |transaction| {
   1658                 Box::pin(async move {
   1659                     sqlx::query(PUBLICATION_TARGETS_SQL)
   1660                         .bind(
   1661                             i64::try_from(snapshot)
   1662                                 .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?,
   1663                         )
   1664                         .bind(workflow.as_deref())
   1665                         .bind(workflow.as_deref())
   1666                         .bind(state_filter.as_deref())
   1667                         .bind(state_filter.as_deref())
   1668                         .bind(i64::from(limit) + 1)
   1669                         .bind(i64::from(offset))
   1670                         .fetch_all(&mut *transaction)
   1671                         .await
   1672                         .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
   1673                 })
   1674             })
   1675             .await
   1676             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1677         let has_more = rows.len() > usize::from(limit);
   1678         let items = rows
   1679             .iter()
   1680             .take(usize::from(limit))
   1681             .map(publication_target_summary)
   1682             .collect::<Result<Vec<_>, _>>()?;
   1683         let mut value = json!({"items":items,"snapshot_generation":snapshot});
   1684         if has_more {
   1685             value["next_cursor"] = Value::String(
   1686                 self.encode_cursor(
   1687                     request.route(),
   1688                     snapshot,
   1689                     offset
   1690                         .checked_add(u32::from(limit))
   1691                         .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))?,
   1692                     query_digest,
   1693                 ),
   1694             );
   1695         }
   1696         response(request.route(), value)
   1697     }
   1698     async fn publication_retry(
   1699         &self,
   1700         request: RhiAdminRequestDocument,
   1701     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
   1702         let model = request_model(&request)?;
   1703         let workflow = model
   1704             .pointer("/workflow_id")
   1705             .and_then(Value::as_str)
   1706             .ok_or_else(|| failure(RhiAdminHandlerErrorKind::Internal))
   1707             .and_then(parse_digest)?;
   1708         let expected = required_u64(&model, "/expected_generation")?;
   1709         let now_ms = self.now_millis()?;
   1710         let operation_id = required_operation_id(&request)?.to_owned();
   1711         let route = request.route();
   1712         let expected_target_count = self.publication.targets().len();
   1713         RhiAdminOperationRepository::new(&self.state)
   1714             .execute_database_admin_operation(
   1715                 &request,
   1716                 operation_time(now_ms)?,
   1717                 RhiAdminOperationJournalPolicy::seven_days(),
   1718                 move |transaction| {
   1719                     Box::pin(async move {
   1720                         let rows = sqlx::query(
   1721                             r#"SELECT state, revision, target_count,
   1722                                 length(event_sha256) AS event_sha256_bytes,
   1723                                 substr(event_sha256, 1, 33) AS event_sha256
   1724                             FROM publication_outbox WHERE outbox_id = ? LIMIT 2"#,
   1725                         )
   1726                         .bind(workflow.as_slice())
   1727                         .fetch_all(&mut *transaction)
   1728                         .await
   1729                         .map_err(|_| AdminJournalOperationError::Storage)?;
   1730                         if rows.len() != 1 {
   1731                             return Err(if rows.is_empty() {
   1732                                 AdminJournalOperationError::Conflict
   1733                             } else {
   1734                                 AdminJournalOperationError::Binding
   1735                             });
   1736                         }
   1737                         let row = &rows[0];
   1738                         let state = row
   1739                             .try_get::<&str, _>("state")
   1740                             .map_err(|_| AdminJournalOperationError::Binding)?;
   1741                         let revision = operation_row_u64(row, "revision")?;
   1742                         let target_count = operation_row_u64(row, "target_count")?;
   1743                         let event_sha256 = operation_exact_blob::<32>(
   1744                             row,
   1745                             "event_sha256",
   1746                             "event_sha256_bytes",
   1747                         )?;
   1748                         if state != "blocked"
   1749                             || revision != expected
   1750                             || usize::try_from(target_count).ok() != Some(expected_target_count)
   1751                         {
   1752                             return Err(AdminJournalOperationError::Conflict);
   1753                         }
   1754                         sqlx::query(
   1755                             r#"UPDATE publication_targets
   1756                             SET state = 'pending', revision = revision + 1,
   1757                                 next_attempt_unix_ms = ?, updated_at_unix_ms = ?
   1758                             WHERE outbox_id = ? AND state != 'accepted'"#,
   1759                         )
   1760                         .bind(i64::try_from(now_ms).map_err(|_| AdminJournalOperationError::InvalidInput)?)
   1761                         .bind(i64::try_from(now_ms).map_err(|_| AdminJournalOperationError::InvalidInput)?)
   1762                         .bind(workflow.as_slice())
   1763                         .execute(&mut *transaction)
   1764                         .await
   1765                         .map_err(|_| AdminJournalOperationError::Storage)?;
   1766                         let changed = sqlx::query(
   1767                             r#"UPDATE publication_outbox
   1768                             SET state = 'pending', revision = revision + 1,
   1769                                 next_attempt_unix_ms = ?, updated_at_unix_ms = ?
   1770                             WHERE outbox_id = ? AND state = 'blocked' AND revision = ?"#,
   1771                         )
   1772                         .bind(i64::try_from(now_ms).map_err(|_| AdminJournalOperationError::InvalidInput)?)
   1773                         .bind(i64::try_from(now_ms).map_err(|_| AdminJournalOperationError::InvalidInput)?)
   1774                         .bind(workflow.as_slice())
   1775                         .bind(i64::try_from(expected).map_err(|_| AdminJournalOperationError::InvalidInput)?)
   1776                         .execute(&mut *transaction)
   1777                         .await
   1778                         .map_err(|_| AdminJournalOperationError::Storage)?;
   1779                         if changed.rows_affected() != 1 {
   1780                             return Err(AdminJournalOperationError::Conflict);
   1781                         }
   1782                         response(
   1783                             route,
   1784                             json!({
   1785                                 "exact_bytes_digest": lower_hex(&event_sha256),
   1786                                 "generation": expected.checked_add(1).ok_or(AdminJournalOperationError::InvalidInput)?,
   1787                                 "operation_id": operation_id,
   1788                                 "target_count": target_count,
   1789                                 "workflow_id": lower_hex(&workflow),
   1790                             }),
   1791                         )
   1792                         .map_err(|_| AdminJournalOperationError::Binding)
   1793                     })
   1794                 },
   1795             )
   1796             .await
   1797             .map_err(map_journal_error)
   1798     }
   1799     async fn presence_refresh(
   1800         &self,
   1801         request: RhiAdminRequestDocument,
   1802     ) -> Result<RhiAdminResponseDocument, RhiAdminHandlerError> {
   1803         let model = request_model(&request)?;
   1804         let expected = required_u64(&model, "/expected_generation")?;
   1805         let operation_id = required_operation_id(&request)?.to_owned();
   1806         let now_ms = self.now_millis()?;
   1807         let route = request.route();
   1808         RhiAdminOperationRepository::new(&self.state)
   1809             .execute_database_admin_operation(
   1810                 &request,
   1811                 operation_time(now_ms)?,
   1812                 RhiAdminOperationJournalPolicy::seven_days(),
   1813                 move |transaction| {
   1814                     Box::pin(async move {
   1815                         let desired = sqlx::query(
   1816                             r#"SELECT generation, enabled, target_count
   1817                             FROM presence_desired_state WHERE singleton = 1 LIMIT 2"#,
   1818                         )
   1819                         .fetch_all(&mut *transaction)
   1820                         .await
   1821                         .map_err(|_| AdminJournalOperationError::Storage)?;
   1822                         if desired.len() != 1 {
   1823                             return Err(AdminJournalOperationError::Binding);
   1824                         }
   1825                         let generation = operation_row_u64(&desired[0], "generation")?;
   1826                         let enabled = operation_row_u64(&desired[0], "enabled")?;
   1827                         let target_count = operation_row_u64(&desired[0], "target_count")?;
   1828                         if generation != expected || enabled != 1 {
   1829                             return Err(AdminJournalOperationError::Conflict);
   1830                         }
   1831                         let rows = sqlx::query(
   1832                             r#"SELECT document_kind,
   1833                                 length(event_sha256) AS event_sha256_bytes,
   1834                                 substr(event_sha256, 1, 33) AS event_sha256,
   1835                                 state
   1836                             FROM presence_outbox
   1837                             WHERE desired_generation = ?
   1838                             ORDER BY document_kind LIMIT 3"#,
   1839                         )
   1840                         .bind(
   1841                             i64::try_from(expected)
   1842                                 .map_err(|_| AdminJournalOperationError::InvalidInput)?,
   1843                         )
   1844                         .fetch_all(&mut *transaction)
   1845                         .await
   1846                         .map_err(|_| AdminJournalOperationError::Storage)?;
   1847                         if rows.is_empty() || rows.len() > 2 {
   1848                             return Err(AdminJournalOperationError::Conflict);
   1849                         }
   1850                         let mut digests = Map::new();
   1851                         for row in &rows {
   1852                             let kind = operation_presence_kind(row)?;
   1853                             let digest = operation_exact_blob::<32>(
   1854                                 row,
   1855                                 "event_sha256",
   1856                                 "event_sha256_bytes",
   1857                             )?;
   1858                             if digests
   1859                                 .insert(kind.to_owned(), Value::String(lower_hex(&digest)))
   1860                                 .is_some()
   1861                             {
   1862                                 return Err(AdminJournalOperationError::Binding);
   1863                             }
   1864                         }
   1865                         let changed = sqlx::query(
   1866                             r#"UPDATE presence_outbox
   1867                             SET state = 'pending', revision = revision + 1,
   1868                                 next_attempt_unix_ms = ?, updated_at_unix_ms = ?
   1869                             WHERE desired_generation = ? AND state = 'blocked'"#,
   1870                         )
   1871                         .bind(
   1872                             i64::try_from(now_ms)
   1873                                 .map_err(|_| AdminJournalOperationError::InvalidInput)?,
   1874                         )
   1875                         .bind(
   1876                             i64::try_from(now_ms)
   1877                                 .map_err(|_| AdminJournalOperationError::InvalidInput)?,
   1878                         )
   1879                         .bind(
   1880                             i64::try_from(expected)
   1881                                 .map_err(|_| AdminJournalOperationError::InvalidInput)?,
   1882                         )
   1883                         .execute(&mut *transaction)
   1884                         .await
   1885                         .map_err(|_| AdminJournalOperationError::Storage)?;
   1886                         if changed.rows_affected() == 0
   1887                             || changed.rows_affected() as usize > rows.len()
   1888                         {
   1889                             return Err(AdminJournalOperationError::Conflict);
   1890                         }
   1891                         response(
   1892                             route,
   1893                             json!({
   1894                                 "document_digests": digests,
   1895                                 "generation": expected,
   1896                                 "operation_id": operation_id,
   1897                                 "state": "pending",
   1898                                 "target_count": target_count,
   1899                             }),
   1900                         )
   1901                         .map_err(|_| AdminJournalOperationError::Binding)
   1902                     })
   1903                 },
   1904             )
   1905             .await
   1906             .map_err(map_journal_error)
   1907     }
   1908 
   1909     async fn presence_state(
   1910         &self,
   1911         generation: u64,
   1912     ) -> Result<(&'static str, Map<String, Value>), RhiAdminHandlerError> {
   1913         self.state
   1914             .sqlite_host()
   1915             .transaction(move |transaction| {
   1916                 Box::pin(async move {
   1917                     let rows = sqlx::query(
   1918                         r#"SELECT document_kind, state,
   1919                             length(event_sha256) AS event_sha256_bytes,
   1920                             substr(event_sha256, 1, 33) AS event_sha256
   1921                         FROM presence_outbox
   1922                         WHERE desired_generation = ?
   1923                         ORDER BY document_kind LIMIT 3"#,
   1924                     )
   1925                     .bind(
   1926                         i64::try_from(generation)
   1927                             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?,
   1928                     )
   1929                     .fetch_all(&mut *transaction)
   1930                     .await
   1931                     .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?;
   1932                     if rows.len() > 2 {
   1933                         return Err(failure(RhiAdminHandlerErrorKind::Internal));
   1934                     }
   1935                     let mut digests = Map::new();
   1936                     let mut states = Vec::with_capacity(rows.len());
   1937                     for row in &rows {
   1938                         let kind = presence_kind(row)?;
   1939                         let digest = exact_blob::<32>(row, "event_sha256", "event_sha256_bytes")?;
   1940                         if digests
   1941                             .insert(kind.to_owned(), Value::String(lower_hex(&digest)))
   1942                             .is_some()
   1943                         {
   1944                             return Err(failure(RhiAdminHandlerErrorKind::Internal));
   1945                         }
   1946                         states.push(
   1947                             row.try_get::<&str, _>("state")
   1948                                 .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))?,
   1949                         );
   1950                     }
   1951                     let state = if rows.is_empty() {
   1952                         "dirty"
   1953                     } else if states.contains(&"leased") {
   1954                         "submitted"
   1955                     } else if states.contains(&"pending") {
   1956                         "pending"
   1957                     } else if states.contains(&"blocked") {
   1958                         "failed"
   1959                     } else if states.iter().all(|state| *state == "complete") {
   1960                         "accepted"
   1961                     } else if states.iter().all(|state| *state == "superseded") {
   1962                         "dirty"
   1963                     } else {
   1964                         "unknown"
   1965                     };
   1966                     Ok::<(&'static str, Map<String, Value>), RhiAdminHandlerError>((state, digests))
   1967                 })
   1968             })
   1969             .await
   1970             .map_err(|_| failure(RhiAdminHandlerErrorKind::Internal))
   1971     }
   1972 }
   1973 
   1974 #[cfg(test)]
   1975 mod tests {
   1976     use super::*;
   1977 
   1978     #[test]
   1979     fn cursors_are_bounded_authenticated_and_bound_to_route_filters_and_key() {
   1980         let key = [0x11; 32];
   1981         let query = [0x22; 32];
   1982         let cursor = encode_runtime_cursor(
   1983             &key,
   1984             RhiAdminRoute::ReconciliationJobs,
   1985             u64::MAX,
   1986             u32::MAX,
   1987             query,
   1988         );
   1989         assert_eq!(cursor.len(), 104);
   1990         assert!(!cursor.contains('='));
   1991         assert_eq!(
   1992             decode_runtime_cursor(&key, RhiAdminRoute::ReconciliationJobs, &cursor, query)
   1993                 .expect("exact cursor"),
   1994             (u64::MAX, u32::MAX)
   1995         );
   1996 
   1997         for result in [
   1998             decode_runtime_cursor(
   1999                 &[0x12; 32],
   2000                 RhiAdminRoute::ReconciliationJobs,
   2001                 &cursor,
   2002                 query,
   2003             ),
   2004             decode_runtime_cursor(&key, RhiAdminRoute::Sources, &cursor, query),
   2005             decode_runtime_cursor(&key, RhiAdminRoute::ReconciliationJobs, &cursor, [0x23; 32]),
   2006         ] {
   2007             assert_eq!(
   2008                 result.expect_err("mismatched cursor binding").kind(),
   2009                 RhiAdminHandlerErrorKind::InvalidCursor
   2010             );
   2011         }
   2012 
   2013         let mut tampered = cursor.into_bytes();
   2014         tampered[20] = if tampered[20] == b'A' { b'B' } else { b'A' };
   2015         let tampered = String::from_utf8(tampered).expect("ASCII cursor");
   2016         assert_eq!(
   2017             decode_runtime_cursor(&key, RhiAdminRoute::ReconciliationJobs, &tampered, query,)
   2018                 .expect_err("tampered cursor")
   2019                 .kind(),
   2020             RhiAdminHandlerErrorKind::InvalidCursor
   2021         );
   2022     }
   2023 }