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 }