source_ingest.rs (44095B)
1 //! Bounded relay-source ingestion with generation-fenced checkpoint commit. 2 3 use core::{cmp::Ordering, fmt}; 4 use std::{collections::BTreeSet, error::Error}; 5 6 use radroots_event::{SignedEvent, id::TradeId}; 7 use radroots_service_host::UnixTimeSeconds; 8 use radroots_service_sqlite::{ 9 ServiceSqliteTransaction, ServiceSqliteTransactionError, ServiceSqliteTransactionErrorKind, 10 }; 11 use radroots_transport::{ 12 FetchRequest, Target, TargetSet, 13 outcome::FetchTargetState, 14 source::{FetchBounds, FetchCursor, FetchSelector, NextPage}, 15 }; 16 use serde_json::Value; 17 use sqlx::Row; 18 19 use crate::{ 20 RhiAdmittedTradeMutationEvent, RhiConfigDocumentV1, RhiEvidencePolicyDigest, RhiStateHostMode, 21 RhiStateRepositories, RhiTradeMutationAdmissionLimits, RhiTradeMutationAuthoredTimePolicy, 22 RhiTradeMutationObservedAtUnixSeconds, RhiTradeSourceObservation, RhiTransportAdapters, 23 admit_rhi_trade_mutation_event, state_metadata, 24 state_trade::{PersistenceOperationError, PersistenceRecord, persist}, 25 }; 26 27 /// Exact version of the RHI relay-source ingestion contract. 28 pub const RHI_TRADE_SOURCE_INGEST_CONTRACT_VERSION: u32 = 1; 29 30 /// Maximum distinct event identities admitted from one source attempt. 31 pub const RHI_TRADE_SOURCE_RESULT_MAX_EVENTS: usize = 4_096; 32 33 /// Maximum aggregate original event bytes admitted from one source attempt. 34 pub const RHI_TRADE_SOURCE_RESULT_MAX_BYTES: usize = 8 * 1024 * 1024; 35 36 const SOURCE_SELECTOR: &str = "trade_mutation_lineage_v1"; 37 const FETCH_REQUEST_ID_MAX_BYTES: usize = 256; 38 const FETCH_PAGE_MAX_EVENTS: u16 = 1_000; 39 const EVENT_KINDS: [u32; 5] = [3470, 3471, 3472, 3473, 3474]; 40 41 const READ_CHECKPOINT_SQL: &str = r#"SELECT 42 cursor_created_at_unix_s, 43 length(cursor_event_id) AS cursor_event_id_bytes, 44 substr(cursor_event_id, 1, 33) AS cursor_event_id, 45 revision, 46 completed_at_unix_s 47 FROM relay_checkpoints 48 WHERE source_id = ? AND selector_id = ? AND evidence_policy_sha256 = ? AND trade_id = ? 49 LIMIT 1"#; 50 const INSERT_CHECKPOINT_SQL: &str = r#"INSERT INTO relay_checkpoints ( 51 source_id, selector_id, evidence_policy_sha256, trade_id, 52 cursor_created_at_unix_s, cursor_event_id, revision, completed_at_unix_s 53 ) VALUES (?, ?, ?, ?, ?, ?, 1, ?)"#; 54 const UPDATE_CHECKPOINT_SQL: &str = r#"UPDATE relay_checkpoints 55 SET cursor_created_at_unix_s = ?, cursor_event_id = ?, 56 revision = revision + 1, completed_at_unix_s = ? 57 WHERE source_id = ? AND selector_id = ? AND evidence_policy_sha256 = ? AND trade_id = ? 58 AND revision = ?"#; 59 const READ_DIRTY_SQL: &str = r#"SELECT generation, 60 length(evidence_policy_sha256) AS evidence_policy_bytes, 61 substr(evidence_policy_sha256, 1, 33) AS evidence_policy_sha256, 62 updated_at_unix_s 63 FROM trade_dirty_generations 64 WHERE trade_id = ? 65 LIMIT 1"#; 66 const INSERT_DIRTY_SQL: &str = r#"INSERT INTO trade_dirty_generations ( 67 trade_id, generation, evidence_policy_sha256, updated_at_unix_s 68 ) VALUES (?, 1, ?, ?)"#; 69 const UPDATE_DIRTY_SQL: &str = r#"UPDATE trade_dirty_generations 70 SET generation = generation + 1, evidence_policy_sha256 = ?, updated_at_unix_s = ? 71 WHERE trade_id = ? AND generation = ?"#; 72 73 /// Stable terminal classification for one exact relay-source attempt. 74 #[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)] 75 pub enum RhiTradeSourceCompletion { 76 Complete, 77 IncompleteTimeout, 78 IncompleteUnavailable, 79 IncompleteResourceLimit, 80 IncompleteUnknown, 81 Unsupported, 82 } 83 84 impl RhiTradeSourceCompletion { 85 /// Returns the exact machine-contract spelling. 86 #[must_use] 87 pub const fn code(self) -> &'static str { 88 match self { 89 Self::Complete => "complete", 90 Self::IncompleteTimeout => "incomplete_timeout", 91 Self::IncompleteUnavailable => "incomplete_unavailable", 92 Self::IncompleteResourceLimit => "incomplete_resource_limit", 93 Self::IncompleteUnknown => "incomplete_unknown", 94 Self::Unsupported => "unsupported", 95 } 96 } 97 98 #[must_use] 99 pub(crate) const fn allows_checkpoint(self) -> bool { 100 matches!(self, Self::Complete) 101 } 102 } 103 104 /// Monotonic per-trade invalidation generation. 105 #[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] 106 pub struct RhiTradeDirtyGeneration(u64); 107 108 impl RhiTradeDirtyGeneration { 109 /// Returns the positive durable generation. 110 #[must_use] 111 pub const fn get(self) -> u64 { 112 self.0 113 } 114 } 115 116 /// Exact canonically admitted cursor tuple for one source scope. 117 #[derive(Clone, Copy, PartialEq, Eq)] 118 pub struct RhiTradeSourceCursor { 119 created_at_unix_seconds: u64, 120 event_id: [u8; 32], 121 } 122 123 impl RhiTradeSourceCursor { 124 pub(crate) const fn from_verified_parts( 125 created_at_unix_seconds: u64, 126 event_id: [u8; 32], 127 ) -> Self { 128 Self { 129 created_at_unix_seconds, 130 event_id, 131 } 132 } 133 134 /// Returns the inclusive event-authored UTC second. 135 #[must_use] 136 pub const fn created_at_unix_seconds(self) -> u64 { 137 self.created_at_unix_seconds 138 } 139 140 /// Returns the exact verified Nostr event identifier bytes. 141 #[must_use] 142 pub const fn event_id(self) -> [u8; 32] { 143 self.event_id 144 } 145 } 146 147 impl fmt::Debug for RhiTradeSourceCursor { 148 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 149 formatter 150 .debug_struct("RhiTradeSourceCursor") 151 .field("created_at_unix_seconds", &self.created_at_unix_seconds) 152 .field("event_id", &"[redacted]") 153 .finish() 154 } 155 } 156 157 /// Caller-owned, injected timing and request evidence for one fetch attempt. 158 pub struct RhiTradeSourceAttempt { 159 request_id: Box<str>, 160 attempt_started_at: UnixTimeSeconds, 161 observed_at: RhiTradeMutationObservedAtUnixSeconds, 162 authored_time_policy: RhiTradeMutationAuthoredTimePolicy, 163 } 164 165 impl RhiTradeSourceAttempt { 166 /// Validates a bounded request identity and explicit, ordered timestamps. 167 pub fn new( 168 request_id: impl AsRef<str>, 169 attempt_started_at: UnixTimeSeconds, 170 observed_at: RhiTradeMutationObservedAtUnixSeconds, 171 authored_time_policy: RhiTradeMutationAuthoredTimePolicy, 172 ) -> Result<Self, RhiTradeSourceIngestError> { 173 let request_id = request_id.as_ref(); 174 if request_id.is_empty() 175 || request_id.len() > FETCH_REQUEST_ID_MAX_BYTES 176 || request_id != request_id.trim() 177 || request_id.chars().any(char::is_control) 178 || attempt_started_at.get() == 0 179 || i64::try_from(attempt_started_at.get()).is_err() 180 || observed_at.get() < attempt_started_at.get() 181 { 182 return Err(failure(RhiTradeSourceIngestErrorKind::InvalidInput)); 183 } 184 Ok(Self { 185 request_id: request_id.into(), 186 attempt_started_at, 187 observed_at, 188 authored_time_policy, 189 }) 190 } 191 } 192 193 impl fmt::Debug for RhiTradeSourceAttempt { 194 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 195 formatter 196 .debug_struct("RhiTradeSourceAttempt") 197 .field("request_id", &"[redacted]") 198 .field("attempt_started_at", &self.attempt_started_at.get()) 199 .field("observed_at", &self.observed_at.get()) 200 .field("authored_time_policy", &self.authored_time_policy) 201 .finish() 202 } 203 } 204 205 /// Stable source-free relay-ingest failure classification. 206 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 207 pub enum RhiTradeSourceIngestErrorKind { 208 InvalidMode, 209 InvalidInput, 210 InvalidConfiguration, 211 GenerationConflict, 212 Storage, 213 CommitOutcomeUnknown, 214 } 215 216 impl RhiTradeSourceIngestErrorKind { 217 /// Returns the stable machine-readable code. 218 #[must_use] 219 pub const fn code(self) -> &'static str { 220 match self { 221 Self::InvalidMode => "trade_source_mode_invalid", 222 Self::InvalidInput => "trade_source_input_invalid", 223 Self::InvalidConfiguration => "trade_source_configuration_invalid", 224 Self::GenerationConflict => "trade_source_generation_conflict", 225 Self::Storage => "trade_source_storage_failed", 226 Self::CommitOutcomeUnknown => "trade_source_commit_outcome_unknown", 227 } 228 } 229 } 230 231 /// Redacted source-free relay-ingest failure. 232 #[derive(Clone, Copy, PartialEq, Eq)] 233 pub struct RhiTradeSourceIngestError { 234 kind: RhiTradeSourceIngestErrorKind, 235 } 236 237 impl RhiTradeSourceIngestError { 238 /// Returns the stable failure kind. 239 #[must_use] 240 pub const fn kind(self) -> RhiTradeSourceIngestErrorKind { 241 self.kind 242 } 243 244 /// Returns the stable machine-readable code. 245 #[must_use] 246 pub const fn code(self) -> &'static str { 247 self.kind.code() 248 } 249 } 250 251 impl fmt::Debug for RhiTradeSourceIngestError { 252 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 253 formatter 254 .debug_struct("RhiTradeSourceIngestError") 255 .field("kind", &self.kind) 256 .finish() 257 } 258 } 259 260 impl fmt::Display for RhiTradeSourceIngestError { 261 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 262 formatter.write_str(match self.kind { 263 RhiTradeSourceIngestErrorKind::InvalidMode => { 264 "RHI trade-source ingestion requires writable state" 265 } 266 RhiTradeSourceIngestErrorKind::InvalidInput => { 267 "RHI trade-source attempt input is invalid" 268 } 269 RhiTradeSourceIngestErrorKind::InvalidConfiguration => { 270 "RHI trade-source configuration is invalid" 271 } 272 RhiTradeSourceIngestErrorKind::GenerationConflict => { 273 "RHI trade-source generation changed during the attempt" 274 } 275 RhiTradeSourceIngestErrorKind::Storage => "RHI trade-source state transaction failed", 276 RhiTradeSourceIngestErrorKind::CommitOutcomeUnknown => { 277 "RHI trade-source commit outcome is unknown" 278 } 279 }) 280 } 281 } 282 283 impl Error for RhiTradeSourceIngestError {} 284 285 /// Durable outcome of one bounded source fetch and atomic evidence commit. 286 #[derive(Clone, Copy, PartialEq, Eq)] 287 pub struct RhiTradeSourceIngestOutcome { 288 completion: RhiTradeSourceCompletion, 289 received_events: u32, 290 admitted_events: u32, 291 rejected_events: u32, 292 duplicate_events: u32, 293 inserted_mutations: u32, 294 inserted_signed_events: u32, 295 inserted_observations: u32, 296 checkpoint: Option<RhiTradeSourceCursor>, 297 checkpoint_advanced: bool, 298 dirty_generation: Option<RhiTradeDirtyGeneration>, 299 dirty_generation_advanced: bool, 300 } 301 302 impl RhiTradeSourceIngestOutcome { 303 #[must_use] 304 pub const fn completion(self) -> RhiTradeSourceCompletion { 305 self.completion 306 } 307 308 #[must_use] 309 pub const fn received_events(self) -> u32 { 310 self.received_events 311 } 312 313 #[must_use] 314 pub const fn admitted_events(self) -> u32 { 315 self.admitted_events 316 } 317 318 #[must_use] 319 pub const fn rejected_events(self) -> u32 { 320 self.rejected_events 321 } 322 323 #[must_use] 324 pub const fn duplicate_events(self) -> u32 { 325 self.duplicate_events 326 } 327 328 #[must_use] 329 pub const fn inserted_mutations(self) -> u32 { 330 self.inserted_mutations 331 } 332 333 #[must_use] 334 pub const fn inserted_signed_events(self) -> u32 { 335 self.inserted_signed_events 336 } 337 338 #[must_use] 339 pub const fn inserted_observations(self) -> u32 { 340 self.inserted_observations 341 } 342 343 #[must_use] 344 pub const fn checkpoint(self) -> Option<RhiTradeSourceCursor> { 345 self.checkpoint 346 } 347 348 #[must_use] 349 pub const fn checkpoint_advanced(self) -> bool { 350 self.checkpoint_advanced 351 } 352 353 #[must_use] 354 pub const fn dirty_generation(self) -> Option<RhiTradeDirtyGeneration> { 355 self.dirty_generation 356 } 357 358 #[must_use] 359 pub const fn dirty_generation_advanced(self) -> bool { 360 self.dirty_generation_advanced 361 } 362 } 363 364 impl fmt::Debug for RhiTradeSourceIngestOutcome { 365 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 366 formatter 367 .debug_struct("RhiTradeSourceIngestOutcome") 368 .field("completion", &self.completion) 369 .field("received_events", &self.received_events) 370 .field("admitted_events", &self.admitted_events) 371 .field("rejected_events", &self.rejected_events) 372 .field("duplicate_events", &self.duplicate_events) 373 .field("inserted_mutations", &self.inserted_mutations) 374 .field("inserted_signed_events", &self.inserted_signed_events) 375 .field("inserted_observations", &self.inserted_observations) 376 .field("checkpoint", &self.checkpoint) 377 .field("checkpoint_advanced", &self.checkpoint_advanced) 378 .field("dirty_generation", &self.dirty_generation) 379 .field("dirty_generation_advanced", &self.dirty_generation_advanced) 380 .finish() 381 } 382 } 383 384 /// Fetches one exact configured relay source and atomically commits admitted evidence. 385 pub async fn ingest_rhi_trade_source( 386 repositories: &RhiStateRepositories<'_>, 387 transports: &RhiTransportAdapters, 388 configuration: &RhiConfigDocumentV1, 389 source_id: &str, 390 trade_id: TradeId, 391 attempt: RhiTradeSourceAttempt, 392 ) -> Result<RhiTradeSourceIngestOutcome, RhiTradeSourceIngestError> { 393 let host = repositories.host(); 394 if host.mode() != RhiStateHostMode::ReadWriteExisting { 395 return Err(failure(RhiTradeSourceIngestErrorKind::InvalidMode)); 396 } 397 let source = ConfiguredSource::new(host, configuration, source_id)?; 398 let initial = read_initial_state(repositories, &source, trade_id).await?; 399 let fetched = fetch_source(transports, &source, trade_id, &attempt, initial.checkpoint).await?; 400 commit_source_result(repositories, source, trade_id, attempt, initial, fetched).await 401 } 402 403 /// Atomically commits one event that was already admitted from a governed 404 /// subscription. This keeps subscription delivery on the same checkpoint, 405 /// provenance, dirty-generation, and evidence transaction used by paged 406 /// source ingestion without issuing a second network request. 407 #[cfg(any(target_os = "linux", target_os = "macos"))] 408 pub(crate) async fn ingest_rhi_subscribed_trade_event( 409 repositories: &RhiStateRepositories<'_>, 410 configuration: &RhiConfigDocumentV1, 411 source_id: &str, 412 admitted: RhiAdmittedTradeMutationEvent, 413 attempt: RhiTradeSourceAttempt, 414 ) -> Result<RhiTradeSourceIngestOutcome, RhiTradeSourceIngestError> { 415 let host = repositories.host(); 416 if host.mode() != RhiStateHostMode::ReadWriteExisting { 417 return Err(failure(RhiTradeSourceIngestErrorKind::InvalidMode)); 418 } 419 let source = ConfiguredSource::new(host, configuration, source_id)?; 420 let trade_id = admitted.mutation().trade_id; 421 if admitted.original_bytes().len() > source.maximum_bytes || source.maximum_events == 0 { 422 return Err(failure(RhiTradeSourceIngestErrorKind::InvalidInput)); 423 } 424 let cursor = RhiTradeSourceCursor { 425 created_at_unix_seconds: admitted.authored_at_unix_seconds(), 426 event_id: *admitted.event_id().as_bytes(), 427 }; 428 let initial = read_initial_state(repositories, &source, trade_id).await?; 429 let fetched = FetchedSource { 430 completion: RhiTradeSourceCompletion::Complete, 431 received_events: 1, 432 duplicate_events: 0, 433 admitted: vec![admitted], 434 rejected_events: 0, 435 cursor_candidate: Some(cursor), 436 }; 437 commit_source_result(repositories, source, trade_id, attempt, initial, fetched).await 438 } 439 440 struct ConfiguredSource { 441 source_id: Box<str>, 442 relay_url: Box<str>, 443 policy: RhiEvidencePolicyDigest, 444 deadline_ms: u64, 445 lookback_seconds: u64, 446 overlap_seconds: u64, 447 maximum_events: usize, 448 maximum_bytes: usize, 449 admission_limits: RhiTradeMutationAdmissionLimits, 450 } 451 452 impl ConfiguredSource { 453 fn new( 454 host: &crate::RhiStateHost, 455 configuration: &RhiConfigDocumentV1, 456 source_id: &str, 457 ) -> Result<Self, RhiTradeSourceIngestError> { 458 let normalized = configuration.normalized(); 459 let source = configured_source(normalized, source_id) 460 .ok_or_else(|| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; 461 if source.pointer("/kind").and_then(Value::as_str) != Some("nostr_relay") 462 || source.pointer("/selector").and_then(Value::as_str) != Some(SOURCE_SELECTOR) 463 { 464 return Err(failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration)); 465 } 466 let relay_id = source 467 .pointer("/relay_id") 468 .and_then(Value::as_str) 469 .ok_or_else(|| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; 470 let relay = normalized 471 .pointer("/relays") 472 .and_then(Value::as_array) 473 .and_then(|relays| { 474 relays 475 .iter() 476 .find(|relay| relay.pointer("/id").and_then(Value::as_str) == Some(relay_id)) 477 }) 478 .filter(|relay| relay.pointer("/read").and_then(Value::as_bool) == Some(true)) 479 .ok_or_else(|| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; 480 let relay_url = relay 481 .pointer("/url") 482 .and_then(Value::as_str) 483 .ok_or_else(|| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; 484 let policy = state_metadata::evidence_policy_digest(normalized) 485 .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; 486 if policy != host.metadata().evidence_policy_digest() { 487 return Err(failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration)); 488 } 489 let deadline_ms = exact_u64(source, "/deadline_ms", 100, 30_000)?; 490 let lookback_seconds = exact_u64(source, "/lookback_seconds", 60, 2_678_400)?; 491 let overlap_seconds = exact_u64(source, "/overlap_seconds", 1, 86_400)?; 492 if overlap_seconds > lookback_seconds { 493 return Err(failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration)); 494 } 495 let maximum_events = exact_usize( 496 normalized, 497 "/resource_limits/source_results/events", 498 1, 499 RHI_TRADE_SOURCE_RESULT_MAX_EVENTS, 500 )?; 501 let maximum_bytes = exact_usize( 502 normalized, 503 "/resource_limits/source_results/bytes", 504 1, 505 RHI_TRADE_SOURCE_RESULT_MAX_BYTES, 506 )?; 507 Target::nostr_relay(relay_url) 508 .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; 509 let admission_limits = RhiTradeMutationAdmissionLimits::from_config(configuration) 510 .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; 511 Ok(Self { 512 source_id: source_id.into(), 513 relay_url: relay_url.into(), 514 policy, 515 deadline_ms, 516 lookback_seconds, 517 overlap_seconds, 518 maximum_events, 519 maximum_bytes, 520 admission_limits, 521 }) 522 } 523 } 524 525 #[derive(Clone, Copy, PartialEq, Eq)] 526 pub(crate) struct Checkpoint { 527 pub(crate) cursor: RhiTradeSourceCursor, 528 pub(crate) revision: u64, 529 pub(crate) completed_at_unix_s: u64, 530 } 531 532 #[derive(Clone, Copy, PartialEq, Eq)] 533 pub(crate) struct DirtyState { 534 pub(crate) generation: RhiTradeDirtyGeneration, 535 pub(crate) policy: RhiEvidencePolicyDigest, 536 pub(crate) updated_at_unix_s: u64, 537 } 538 539 #[derive(Clone, Copy)] 540 struct InitialState { 541 checkpoint: Option<Checkpoint>, 542 dirty: Option<DirtyState>, 543 } 544 545 struct Candidate { 546 event: SignedEvent, 547 } 548 549 struct FetchedSource { 550 completion: RhiTradeSourceCompletion, 551 received_events: usize, 552 duplicate_events: usize, 553 admitted: Vec<RhiAdmittedTradeMutationEvent>, 554 rejected_events: usize, 555 cursor_candidate: Option<RhiTradeSourceCursor>, 556 } 557 558 async fn fetch_source( 559 transports: &RhiTransportAdapters, 560 source: &ConfiguredSource, 561 trade_id: TradeId, 562 attempt: &RhiTradeSourceAttempt, 563 checkpoint: Option<Checkpoint>, 564 ) -> Result<FetchedSource, RhiTradeSourceIngestError> { 565 let deadline_unix_ms = attempt 566 .attempt_started_at 567 .get() 568 .checked_mul(1_000) 569 .and_then(|value| value.checked_add(source.deadline_ms)) 570 .ok_or_else(|| failure(RhiTradeSourceIngestErrorKind::InvalidInput))?; 571 let since = checkpoint.map_or_else( 572 || { 573 attempt 574 .attempt_started_at 575 .get() 576 .saturating_sub(source.lookback_seconds) 577 }, 578 |value| { 579 value 580 .cursor 581 .created_at_unix_seconds 582 .saturating_sub(source.overlap_seconds) 583 }, 584 ); 585 let selector = FetchSelector::all() 586 .with_kinds(EVENT_KINDS.to_vec()) 587 .and_then(|selector| selector.with_exact_tag_value('d', trade_id.to_hex())) 588 .and_then(|selector| selector.with_since_unix_seconds(since)) 589 .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; 590 let target = Target::nostr_relay(source.relay_url.as_ref()) 591 .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; 592 let fingerprint = target.fingerprint().clone(); 593 let targets = TargetSet::new(vec![target]) 594 .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?; 595 let mut adapter_cursor = None::<FetchCursor>; 596 let mut seen_adapter_cursors = BTreeSet::new(); 597 let mut candidates = Vec::new(); 598 let mut received_events = 0_usize; 599 let mut received_bytes = 0_usize; 600 let completion = loop { 601 let remaining = source.maximum_events.saturating_sub(received_events); 602 let request_limit = usize::min(usize::from(FETCH_PAGE_MAX_EVENTS), remaining + 1); 603 let bounds = FetchBounds::new( 604 u16::try_from(request_limit) 605 .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration))?, 606 deadline_unix_ms, 607 ) 608 .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidInput))?; 609 let mut request = FetchRequest::new(attempt.request_id.as_ref(), targets.clone(), bounds) 610 .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidInput))? 611 .with_selector(selector.clone()); 612 if let Some(cursor) = adapter_cursor.take() { 613 request = request.with_cursor(cursor); 614 } 615 let page = match transports.evidence_source().fetch(request.clone()).await { 616 Ok(page) => page, 617 Err(radroots_transport::Error::UnsupportedOperation) => { 618 break RhiTradeSourceCompletion::Unsupported; 619 } 620 Err(_) => break RhiTradeSourceCompletion::IncompleteUnknown, 621 }; 622 if page.validate_for_request(&request).is_err() { 623 break RhiTradeSourceCompletion::IncompleteUnknown; 624 } 625 let outcome = page 626 .target_outcomes() 627 .iter() 628 .find(|outcome| outcome.target() == &fingerprint); 629 let Some(outcome) = outcome.filter(|_| page.target_outcomes().len() == 1) else { 630 break RhiTradeSourceCompletion::IncompleteUnknown; 631 }; 632 let target_state = outcome.state(); 633 match target_state { 634 FetchTargetState::Complete | FetchTargetState::Partial => {} 635 FetchTargetState::Unavailable | FetchTargetState::FailedRetryable => { 636 break RhiTradeSourceCompletion::IncompleteUnavailable; 637 } 638 FetchTargetState::FailedTerminal => { 639 break RhiTradeSourceCompletion::IncompleteUnknown; 640 } 641 FetchTargetState::Cancelled => break RhiTradeSourceCompletion::IncompleteTimeout, 642 } 643 for observed in page.events() { 644 received_events = received_events.saturating_add(1); 645 if received_events > source.maximum_events { 646 break; 647 } 648 received_bytes = match received_bytes.checked_add(observed.event().raw_json().len()) { 649 Some(value) if value <= source.maximum_bytes => value, 650 _ => { 651 received_events = source.maximum_events.saturating_add(1); 652 break; 653 } 654 }; 655 candidates.push(Candidate { 656 event: observed.event().clone(), 657 }); 658 } 659 if received_events > source.maximum_events { 660 break RhiTradeSourceCompletion::IncompleteResourceLimit; 661 } 662 if target_state == FetchTargetState::Partial { 663 break RhiTradeSourceCompletion::IncompleteUnknown; 664 } 665 match page.next_page() { 666 NextPage::Complete => break RhiTradeSourceCompletion::Complete, 667 NextPage::Cancelled { .. } => break RhiTradeSourceCompletion::IncompleteTimeout, 668 NextPage::Cursor(cursor) => { 669 if page.events().is_empty() 670 || !seen_adapter_cursors.insert(cursor.as_str().to_owned()) 671 { 672 break RhiTradeSourceCompletion::IncompleteUnknown; 673 } 674 adapter_cursor = Some(cursor.clone()); 675 } 676 } 677 }; 678 679 candidates.sort_by(compare_candidate); 680 let mut duplicate_events = 0_usize; 681 let mut admitted_signed_event_ids = BTreeSet::new(); 682 let mut admitted = Vec::with_capacity(candidates.len()); 683 let mut rejected_events = 0_usize; 684 let mut cursor_candidate = None; 685 for candidate in candidates { 686 match admit_rhi_trade_mutation_event( 687 source.admission_limits, 688 candidate.event.raw_json().as_bytes(), 689 attempt.observed_at, 690 attempt.authored_time_policy, 691 ) { 692 Ok(event) if event.mutation().trade_id == trade_id => { 693 if !admitted_signed_event_ids 694 .insert((*event.event_id().as_bytes(), event.event_signature_bytes())) 695 { 696 duplicate_events = duplicate_events.saturating_add(1); 697 continue; 698 } 699 let cursor = RhiTradeSourceCursor { 700 created_at_unix_seconds: event.authored_at_unix_seconds(), 701 event_id: *event.event_id().as_bytes(), 702 }; 703 cursor_candidate = Some(cursor_candidate.map_or(cursor, |current| { 704 if compare_cursor(current, cursor).is_lt() { 705 cursor 706 } else { 707 current 708 } 709 })); 710 admitted.push(event); 711 } 712 Ok(_) | Err(_) => rejected_events = rejected_events.saturating_add(1), 713 } 714 } 715 Ok(FetchedSource { 716 completion, 717 received_events: received_events.min(source.maximum_events), 718 duplicate_events, 719 admitted, 720 rejected_events, 721 cursor_candidate, 722 }) 723 } 724 725 async fn read_initial_state( 726 repositories: &RhiStateRepositories<'_>, 727 source: &ConfiguredSource, 728 trade_id: TradeId, 729 ) -> Result<InitialState, RhiTradeSourceIngestError> { 730 let source_id = source.source_id.clone(); 731 let policy = source.policy; 732 repositories 733 .host() 734 .sqlite_host() 735 .transaction(move |transaction| { 736 Box::pin(async move { 737 Ok(InitialState { 738 checkpoint: read_checkpoint(transaction, source_id.as_ref(), policy, trade_id) 739 .await?, 740 dirty: read_dirty(transaction, trade_id).await?, 741 }) 742 }) 743 }) 744 .await 745 .map_err(map_transaction_error) 746 } 747 748 async fn commit_source_result( 749 repositories: &RhiStateRepositories<'_>, 750 source: ConfiguredSource, 751 trade_id: TradeId, 752 attempt: RhiTradeSourceAttempt, 753 initial: InitialState, 754 fetched: FetchedSource, 755 ) -> Result<RhiTradeSourceIngestOutcome, RhiTradeSourceIngestError> { 756 let source_id = source.source_id; 757 let policy = source.policy; 758 repositories 759 .host() 760 .sqlite_host() 761 .transaction(move |transaction| { 762 Box::pin(async move { 763 if read_checkpoint(transaction, source_id.as_ref(), policy, trade_id).await? 764 != initial.checkpoint 765 || read_dirty(transaction, trade_id).await? != initial.dirty 766 { 767 return Err(SourceOperationError::GenerationConflict); 768 } 769 let mut inserted_mutations = 0_u32; 770 let mut inserted_signed_events = 0_u32; 771 let mut inserted_observations = 0_u32; 772 let admitted_events = u32::try_from(fetched.admitted.len()) 773 .map_err(|_| SourceOperationError::Storage)?; 774 for admitted in fetched.admitted { 775 let observation = 776 RhiTradeSourceObservation::from_parts(source_id.clone(), policy, &admitted); 777 let record = PersistenceRecord::from_admitted(admitted) 778 .map_err(|_| SourceOperationError::Storage)?; 779 let persisted = persist(transaction, &record, &observation) 780 .await 781 .map_err(SourceOperationError::Persistence)?; 782 inserted_mutations = inserted_mutations 783 .checked_add(u32::from(persisted.mutation_inserted())) 784 .ok_or(SourceOperationError::Storage)?; 785 inserted_signed_events = inserted_signed_events 786 .checked_add(u32::from(persisted.signed_event_inserted())) 787 .ok_or(SourceOperationError::Storage)?; 788 inserted_observations = inserted_observations 789 .checked_add(u32::from(persisted.observation_inserted())) 790 .ok_or(SourceOperationError::Storage)?; 791 } 792 let new_relevant_evidence = inserted_mutations != 0 || inserted_signed_events != 0; 793 let (dirty_generation, dirty_generation_advanced) = if new_relevant_evidence { 794 let generation = advance_dirty_generation( 795 transaction, 796 trade_id, 797 policy, 798 attempt.observed_at.get(), 799 initial.dirty, 800 ) 801 .await?; 802 (Some(generation), true) 803 } else { 804 (initial.dirty.map(|dirty| dirty.generation), false) 805 }; 806 let mut checkpoint = initial.checkpoint.map(|value| value.cursor); 807 let mut checkpoint_advanced = false; 808 if fetched.completion.allows_checkpoint() 809 && fetched.cursor_candidate.is_some_and(|candidate| { 810 initial 811 .checkpoint 812 .is_none_or(|current| compare_cursor(current.cursor, candidate).is_lt()) 813 }) 814 { 815 let Some(candidate) = fetched.cursor_candidate else { 816 return Err(SourceOperationError::Storage); 817 }; 818 write_checkpoint( 819 transaction, 820 source_id.as_ref(), 821 policy, 822 trade_id, 823 initial.checkpoint, 824 candidate, 825 attempt.observed_at.get(), 826 ) 827 .await?; 828 checkpoint = Some(candidate); 829 checkpoint_advanced = true; 830 } 831 Ok(RhiTradeSourceIngestOutcome { 832 completion: fetched.completion, 833 received_events: u32::try_from(fetched.received_events) 834 .map_err(|_| SourceOperationError::Storage)?, 835 admitted_events, 836 rejected_events: u32::try_from(fetched.rejected_events) 837 .map_err(|_| SourceOperationError::Storage)?, 838 duplicate_events: u32::try_from(fetched.duplicate_events) 839 .map_err(|_| SourceOperationError::Storage)?, 840 inserted_mutations, 841 inserted_signed_events, 842 inserted_observations, 843 checkpoint, 844 checkpoint_advanced, 845 dirty_generation, 846 dirty_generation_advanced, 847 }) 848 }) 849 }) 850 .await 851 .map_err(map_transaction_error) 852 } 853 854 pub(crate) async fn advance_dirty_generation( 855 transaction: &mut ServiceSqliteTransaction<'_>, 856 trade_id: TradeId, 857 policy: RhiEvidencePolicyDigest, 858 updated_at_unix_s: u64, 859 expected: Option<DirtyState>, 860 ) -> Result<RhiTradeDirtyGeneration, SourceOperationError> { 861 let updated_at = i64::try_from(updated_at_unix_s).map_err(|_| SourceOperationError::Storage)?; 862 match expected { 863 None => { 864 let result = sqlx::query(INSERT_DIRTY_SQL) 865 .bind(trade_id.as_bytes().as_slice()) 866 .bind(policy.as_bytes().as_slice()) 867 .bind(updated_at) 868 .execute(&mut *transaction) 869 .await 870 .map_err(|_| SourceOperationError::Storage)?; 871 if result.rows_affected() != 1 { 872 return Err(SourceOperationError::GenerationConflict); 873 } 874 Ok(RhiTradeDirtyGeneration(1)) 875 } 876 Some(current) 877 if current.policy == policy && updated_at_unix_s >= current.updated_at_unix_s => 878 { 879 let next = current 880 .generation 881 .get() 882 .checked_add(1) 883 .filter(|value| i64::try_from(*value).is_ok()) 884 .ok_or(SourceOperationError::Storage)?; 885 let result = sqlx::query(UPDATE_DIRTY_SQL) 886 .bind(policy.as_bytes().as_slice()) 887 .bind(updated_at) 888 .bind(trade_id.as_bytes().as_slice()) 889 .bind( 890 i64::try_from(current.generation.get()) 891 .map_err(|_| SourceOperationError::Storage)?, 892 ) 893 .execute(&mut *transaction) 894 .await 895 .map_err(|_| SourceOperationError::Storage)?; 896 if result.rows_affected() != 1 { 897 return Err(SourceOperationError::GenerationConflict); 898 } 899 Ok(RhiTradeDirtyGeneration(next)) 900 } 901 Some(_) => Err(SourceOperationError::GenerationConflict), 902 } 903 } 904 905 pub(crate) async fn read_dirty( 906 transaction: &mut ServiceSqliteTransaction<'_>, 907 trade_id: TradeId, 908 ) -> Result<Option<DirtyState>, SourceOperationError> { 909 let Some(row) = sqlx::query(READ_DIRTY_SQL) 910 .bind(trade_id.as_bytes().as_slice()) 911 .fetch_optional(&mut *transaction) 912 .await 913 .map_err(|_| SourceOperationError::Storage)? 914 else { 915 return Ok(None); 916 }; 917 let generation = positive_i64_u64(&row, "generation")?; 918 let policy = exact_digest(&row, "evidence_policy_sha256", "evidence_policy_bytes")?; 919 let updated_at_unix_s = nonnegative_i64_u64(&row, "updated_at_unix_s")?; 920 Ok(Some(DirtyState { 921 generation: RhiTradeDirtyGeneration(generation), 922 policy: RhiEvidencePolicyDigest::from_bytes(policy), 923 updated_at_unix_s, 924 })) 925 } 926 927 pub(crate) async fn read_checkpoint( 928 transaction: &mut ServiceSqliteTransaction<'_>, 929 source_id: &str, 930 policy: RhiEvidencePolicyDigest, 931 trade_id: TradeId, 932 ) -> Result<Option<Checkpoint>, SourceOperationError> { 933 let Some(row) = sqlx::query(READ_CHECKPOINT_SQL) 934 .bind(source_id) 935 .bind(SOURCE_SELECTOR) 936 .bind(policy.as_bytes().as_slice()) 937 .bind(trade_id.as_bytes().as_slice()) 938 .fetch_optional(&mut *transaction) 939 .await 940 .map_err(|_| SourceOperationError::Storage)? 941 else { 942 return Ok(None); 943 }; 944 Ok(Some(Checkpoint { 945 cursor: RhiTradeSourceCursor { 946 created_at_unix_seconds: nonnegative_i64_u64(&row, "cursor_created_at_unix_s")?, 947 event_id: exact_digest(&row, "cursor_event_id", "cursor_event_id_bytes")?, 948 }, 949 revision: positive_i64_u64(&row, "revision")?, 950 completed_at_unix_s: positive_i64_u64(&row, "completed_at_unix_s")?, 951 })) 952 } 953 954 pub(crate) async fn write_checkpoint( 955 transaction: &mut ServiceSqliteTransaction<'_>, 956 source_id: &str, 957 policy: RhiEvidencePolicyDigest, 958 trade_id: TradeId, 959 current: Option<Checkpoint>, 960 next: RhiTradeSourceCursor, 961 completed_at_unix_s: u64, 962 ) -> Result<(), SourceOperationError> { 963 let created_at = 964 i64::try_from(next.created_at_unix_seconds).map_err(|_| SourceOperationError::Storage)?; 965 let completed_at = 966 i64::try_from(completed_at_unix_s).map_err(|_| SourceOperationError::Storage)?; 967 let result = match current { 968 None => { 969 sqlx::query(INSERT_CHECKPOINT_SQL) 970 .bind(source_id) 971 .bind(SOURCE_SELECTOR) 972 .bind(policy.as_bytes().as_slice()) 973 .bind(trade_id.as_bytes().as_slice()) 974 .bind(created_at) 975 .bind(next.event_id.as_slice()) 976 .bind(completed_at) 977 .execute(&mut *transaction) 978 .await 979 } 980 Some(current) => { 981 sqlx::query(UPDATE_CHECKPOINT_SQL) 982 .bind(created_at) 983 .bind(next.event_id.as_slice()) 984 .bind(completed_at) 985 .bind(source_id) 986 .bind(SOURCE_SELECTOR) 987 .bind(policy.as_bytes().as_slice()) 988 .bind(trade_id.as_bytes().as_slice()) 989 .bind(i64::try_from(current.revision).map_err(|_| SourceOperationError::Storage)?) 990 .execute(&mut *transaction) 991 .await 992 } 993 } 994 .map_err(|_| SourceOperationError::Storage)?; 995 if result.rows_affected() == 1 { 996 Ok(()) 997 } else { 998 Err(SourceOperationError::GenerationConflict) 999 } 1000 } 1001 1002 fn compare_candidate(left: &Candidate, right: &Candidate) -> Ordering { 1003 left.event 1004 .created_at() 1005 .cmp(&right.event.created_at()) 1006 .then_with(|| left.event.id().as_bytes().cmp(right.event.id().as_bytes())) 1007 .then_with(|| { 1008 left.event 1009 .sig() 1010 .as_bytes() 1011 .cmp(right.event.sig().as_bytes()) 1012 }) 1013 } 1014 1015 pub(crate) fn compare_cursor(left: RhiTradeSourceCursor, right: RhiTradeSourceCursor) -> Ordering { 1016 (left.created_at_unix_seconds, left.event_id) 1017 .cmp(&(right.created_at_unix_seconds, right.event_id)) 1018 } 1019 1020 fn configured_source<'a>(configuration: &'a Value, source_id: &str) -> Option<&'a Value> { 1021 if source_id.is_empty() 1022 || source_id.len() > 64 1023 || !source_id.bytes().enumerate().all(|(index, byte)| { 1024 if index == 0 { 1025 byte.is_ascii_lowercase() 1026 } else { 1027 byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-') 1028 } 1029 }) 1030 { 1031 return None; 1032 } 1033 configuration 1034 .pointer("/evidence/sources")? 1035 .as_array()? 1036 .iter() 1037 .find(|source| source.pointer("/source_id").and_then(Value::as_str) == Some(source_id)) 1038 } 1039 1040 fn exact_u64( 1041 value: &Value, 1042 pointer: &str, 1043 minimum: u64, 1044 maximum: u64, 1045 ) -> Result<u64, RhiTradeSourceIngestError> { 1046 value 1047 .pointer(pointer) 1048 .and_then(Value::as_u64) 1049 .filter(|value| (minimum..=maximum).contains(value)) 1050 .ok_or_else(|| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration)) 1051 } 1052 1053 fn exact_usize( 1054 value: &Value, 1055 pointer: &str, 1056 minimum: usize, 1057 maximum: usize, 1058 ) -> Result<usize, RhiTradeSourceIngestError> { 1059 exact_u64( 1060 value, 1061 pointer, 1062 u64::try_from(minimum).unwrap_or(u64::MAX), 1063 u64::try_from(maximum).unwrap_or(u64::MAX), 1064 ) 1065 .and_then(|value| { 1066 usize::try_from(value) 1067 .map_err(|_| failure(RhiTradeSourceIngestErrorKind::InvalidConfiguration)) 1068 }) 1069 } 1070 1071 fn exact_digest( 1072 row: &sqlx::sqlite::SqliteRow, 1073 field: &str, 1074 length_field: &str, 1075 ) -> Result<[u8; 32], SourceOperationError> { 1076 if row.try_get::<i64, _>(length_field).ok() != Some(32) { 1077 return Err(SourceOperationError::Storage); 1078 } 1079 row.try_get::<Vec<u8>, _>(field) 1080 .map_err(|_| SourceOperationError::Storage)? 1081 .try_into() 1082 .map_err(|_| SourceOperationError::Storage) 1083 } 1084 1085 fn positive_i64_u64( 1086 row: &sqlx::sqlite::SqliteRow, 1087 field: &str, 1088 ) -> Result<u64, SourceOperationError> { 1089 row.try_get::<i64, _>(field) 1090 .ok() 1091 .filter(|value| *value > 0) 1092 .and_then(|value| u64::try_from(value).ok()) 1093 .ok_or(SourceOperationError::Storage) 1094 } 1095 1096 fn nonnegative_i64_u64( 1097 row: &sqlx::sqlite::SqliteRow, 1098 field: &str, 1099 ) -> Result<u64, SourceOperationError> { 1100 row.try_get::<i64, _>(field) 1101 .ok() 1102 .filter(|value| *value >= 0) 1103 .and_then(|value| u64::try_from(value).ok()) 1104 .ok_or(SourceOperationError::Storage) 1105 } 1106 1107 #[derive(Clone, Copy, Debug, PartialEq, Eq)] 1108 pub(crate) enum SourceOperationError { 1109 GenerationConflict, 1110 Persistence(PersistenceOperationError), 1111 Storage, 1112 } 1113 1114 fn map_transaction_error( 1115 error: ServiceSqliteTransactionError<SourceOperationError>, 1116 ) -> RhiTradeSourceIngestError { 1117 if error.kind() == ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown { 1118 return failure(RhiTradeSourceIngestErrorKind::CommitOutcomeUnknown); 1119 } 1120 failure(match error.operation_error().copied() { 1121 Some(SourceOperationError::GenerationConflict) => { 1122 RhiTradeSourceIngestErrorKind::GenerationConflict 1123 } 1124 Some(SourceOperationError::Persistence(_)) | Some(SourceOperationError::Storage) | None => { 1125 RhiTradeSourceIngestErrorKind::Storage 1126 } 1127 }) 1128 } 1129 1130 const fn failure(kind: RhiTradeSourceIngestErrorKind) -> RhiTradeSourceIngestError { 1131 RhiTradeSourceIngestError { kind } 1132 } 1133 1134 #[cfg(test)] 1135 mod tests { 1136 use super::*; 1137 1138 #[test] 1139 fn equal_timestamp_cursor_order_uses_verified_event_id() { 1140 let lower = RhiTradeSourceCursor { 1141 created_at_unix_seconds: 100, 1142 event_id: [0x11; 32], 1143 }; 1144 let higher = RhiTradeSourceCursor { 1145 created_at_unix_seconds: 100, 1146 event_id: [0x22; 32], 1147 }; 1148 assert_eq!(compare_cursor(lower, higher), Ordering::Less); 1149 assert_eq!(compare_cursor(higher, lower), Ordering::Greater); 1150 assert_eq!(compare_cursor(lower, lower), Ordering::Equal); 1151 } 1152 1153 #[test] 1154 fn completion_codes_and_checkpoint_policy_are_closed() { 1155 let vectors = [ 1156 (RhiTradeSourceCompletion::Complete, "complete", true), 1157 ( 1158 RhiTradeSourceCompletion::IncompleteTimeout, 1159 "incomplete_timeout", 1160 false, 1161 ), 1162 ( 1163 RhiTradeSourceCompletion::IncompleteUnavailable, 1164 "incomplete_unavailable", 1165 false, 1166 ), 1167 ( 1168 RhiTradeSourceCompletion::IncompleteResourceLimit, 1169 "incomplete_resource_limit", 1170 false, 1171 ), 1172 ( 1173 RhiTradeSourceCompletion::IncompleteUnknown, 1174 "incomplete_unknown", 1175 false, 1176 ), 1177 (RhiTradeSourceCompletion::Unsupported, "unsupported", false), 1178 ]; 1179 for (completion, code, checkpoint) in vectors { 1180 assert_eq!(completion.code(), code); 1181 assert_eq!(completion.allows_checkpoint(), checkpoint); 1182 } 1183 } 1184 }