rhi

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

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 }