lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit 148daafbdc0aa10fb9cdd8e32a34b2681e0b4a98
parent cc2199e29ed5e42363dad506661c0464c7c7344f
Author: triesap <tyson@radroots.org>
Date:   Mon, 27 Jul 2026 08:43:22 +0000

transport: authenticate persisted fetch receipts

- bind serialized fetch evidence to its canonical requested relay set
- derive every aggregate count from sealed event and outcome records
- enforce raw receipt and complete-request diagnostic budgets
- preserve typed evidence for each item skipped by resource limits

Diffstat:
Mcrates/transport_nostr/src/fetch.rs | 586+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
Mcrates/transport_nostr/tests/transport.rs | 524+++++++++++++++++++++++++++++++++++++++++++++++--------------------------------
2 files changed, 862 insertions(+), 248 deletions(-)

diff --git a/crates/transport_nostr/src/fetch.rs b/crates/transport_nostr/src/fetch.rs @@ -669,11 +669,7 @@ impl RadrootsRelayFetchEventReceipt { "event receipt dispositions are mutually exclusive", )); } - if (self.inserted - || self.duplicate - || self.not_persisted - || self.out_of_filter - || self.skipped_over_limit) + if (self.inserted || self.duplicate || self.not_persisted || self.out_of_filter) && self.event_id.is_none() { return Err(invalid_fetch_receipt( @@ -690,11 +686,7 @@ impl RadrootsRelayFetchEventReceipt { "malformed receipts cannot identify or verify an event", )); } - if (self.inserted - || self.duplicate - || self.not_persisted - || self.out_of_filter - || self.skipped_over_limit) + if (self.inserted || self.duplicate || self.not_persisted || self.out_of_filter) && self.verification != RadrootsRelayFetchEventVerification::Verified { return Err(invalid_fetch_receipt( @@ -1031,28 +1023,91 @@ pub struct RadrootsRelayFetchedEventsReceipt { pub relay_outcomes: Vec<RadrootsRelayFetchRelayOutcome>, } -#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +impl RadrootsRelayFetchedEventsReceipt { + pub fn target_relays(&self) -> &[String] { + &self.target_relays + } + + pub fn connected_relays(&self) -> &[String] { + &self.connected_relays + } + + pub fn failed_relays(&self) -> &[RadrootsRelayFetchFailure] { + &self.failed_relays + } + + pub fn events(&self) -> &[RadrootsRelayFetchedEvent] { + &self.events + } + + pub fn event_receipts(&self) -> &[RadrootsRelayFetchEventReceipt] { + &self.event_receipts + } + + pub fn duplicate_count(&self) -> usize { + self.duplicate_count + } + + pub fn verification_failed_count(&self) -> usize { + self.verification_failed_count + } + + pub fn malformed_count(&self) -> usize { + self.malformed_count + } + + pub fn out_of_filter_count(&self) -> usize { + self.out_of_filter_count + } + + pub fn skipped_over_limit_count(&self) -> usize { + self.skipped_over_limit_count + } + + pub fn eose_count(&self) -> usize { + self.eose_count + } + + pub fn truncated_count(&self) -> usize { + self.truncated_count + } + + pub fn closed_count(&self) -> usize { + self.closed_count + } + + pub fn notice_count(&self) -> usize { + self.notice_count + } + + pub fn relay_outcomes(&self) -> &[RadrootsRelayFetchRelayOutcome] { + &self.relay_outcomes + } +} + +#[derive(Clone, Debug, PartialEq, Eq, Serialize)] pub struct RadrootsRelayFetchReceipt { - pub inserted_count: usize, - pub duplicate_count: usize, - pub not_persisted_count: usize, - pub malformed_count: usize, - pub out_of_filter_count: usize, - pub skipped_over_limit_count: usize, - pub verification_failed_count: usize, - pub admission_unsupported_count: usize, - pub admission_invalid_count: usize, - pub valid_stream_eligible_count: usize, - pub visible_count: usize, - pub not_admitted_count: usize, - pub not_current_count: usize, - pub suppressed_count: usize, - pub eose_count: usize, - pub truncated_count: usize, - pub closed_count: usize, - pub notice_count: usize, - pub events: Vec<RadrootsRelayFetchEventReceipt>, - pub relay_outcomes: Vec<RadrootsRelayFetchRelayOutcome>, + target_relays: Vec<String>, + inserted_count: usize, + duplicate_count: usize, + not_persisted_count: usize, + malformed_count: usize, + out_of_filter_count: usize, + skipped_over_limit_count: usize, + verification_failed_count: usize, + admission_unsupported_count: usize, + admission_invalid_count: usize, + valid_stream_eligible_count: usize, + visible_count: usize, + not_admitted_count: usize, + not_current_count: usize, + suppressed_count: usize, + eose_count: usize, + truncated_count: usize, + closed_count: usize, + notice_count: usize, + events: Vec<RadrootsRelayFetchEventReceipt>, + relay_outcomes: Vec<RadrootsRelayFetchRelayOutcome>, } pub trait RadrootsRelayFetchAdapter: Send + Sync { @@ -1227,7 +1282,7 @@ where for event in &receipt.events { event.clone().checked()?; } - Ok(receipt) + receipt.checked() } fn relay_fetch_admission( @@ -1344,8 +1399,266 @@ impl RadrootsRelayProcessedFetch { } impl RadrootsRelayFetchReceipt { + fn checked(mut self) -> Result<Self, RadrootsRelayTransportError> { + self.target_relays = canonical_fetch_receipt_targets(self.target_relays)?; + let raw_receipt_count = self + .events + .len() + .checked_add(self.relay_outcomes.len()) + .ok_or_else(|| { + invalid_fetch_receipt("raw_receipt_count", "raw receipt count overflowed") + })?; + if raw_receipt_count > RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX { + return Err(RadrootsRelayTransportError::FetchLimitTooLarge { + field: "raw_receipt_count", + max: RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX, + actual: raw_receipt_count, + }); + } + for event in &self.events { + event.clone().checked()?; + ensure_fetch_receipt_relay_requested(&self.target_relays, event.relay_url())?; + } + validate_fetch_relay_outcomes(&self.target_relays, &self.relay_outcomes)?; + validate_fetch_diagnostic_budget(&self.events, &self.relay_outcomes)?; + + ensure_fetch_receipt_count( + "inserted_count", + self.inserted_count, + self.events + .iter() + .filter(|event| event.was_inserted()) + .count(), + )?; + ensure_fetch_receipt_count( + "duplicate_count", + self.duplicate_count, + self.events + .iter() + .filter(|event| event.was_duplicate()) + .count(), + )?; + ensure_fetch_receipt_count( + "not_persisted_count", + self.not_persisted_count, + self.events + .iter() + .filter(|event| event.was_not_persisted()) + .count(), + )?; + ensure_fetch_receipt_count( + "malformed_count", + self.malformed_count, + self.events + .iter() + .filter(|event| event.is_malformed()) + .count(), + )?; + ensure_fetch_receipt_count( + "out_of_filter_count", + self.out_of_filter_count, + self.events + .iter() + .filter(|event| event.is_out_of_filter()) + .count(), + )?; + ensure_fetch_receipt_count( + "skipped_over_limit_count", + self.skipped_over_limit_count, + self.events + .iter() + .filter(|event| event.was_skipped_over_limit()) + .count(), + )?; + ensure_fetch_receipt_count( + "verification_failed_count", + self.verification_failed_count, + self.events + .iter() + .filter(|event| event.verification() == RadrootsRelayFetchEventVerification::Failed) + .count(), + )?; + ensure_fetch_receipt_count( + "admission_unsupported_count", + self.admission_unsupported_count, + self.events + .iter() + .filter(|event| event.admission() == RadrootsRelayFetchEventAdmission::Unsupported) + .count(), + )?; + ensure_fetch_receipt_count( + "admission_invalid_count", + self.admission_invalid_count, + self.events + .iter() + .filter(|event| event.admission() == RadrootsRelayFetchEventAdmission::Invalid) + .count(), + )?; + ensure_fetch_receipt_count( + "valid_stream_eligible_count", + self.valid_stream_eligible_count, + self.events + .iter() + .filter(|event| { + event.valid_stream() == RadrootsRelayFetchEventValidStream::Eligible + }) + .count(), + )?; + ensure_fetch_receipt_count( + "visible_count", + self.visible_count, + self.events + .iter() + .filter(|event| event.visibility() == RadrootsRelayFetchEventVisibility::Visible) + .count(), + )?; + ensure_fetch_receipt_count( + "not_admitted_count", + self.not_admitted_count, + self.events + .iter() + .filter(|event| { + event.visibility() == RadrootsRelayFetchEventVisibility::NotAdmitted + }) + .count(), + )?; + ensure_fetch_receipt_count( + "not_current_count", + self.not_current_count, + self.events + .iter() + .filter(|event| event.visibility() == RadrootsRelayFetchEventVisibility::NotCurrent) + .count(), + )?; + ensure_fetch_receipt_count( + "suppressed_count", + self.suppressed_count, + self.events + .iter() + .filter(|event| event.visibility() == RadrootsRelayFetchEventVisibility::Suppressed) + .count(), + )?; + for (field, supplied, kind) in [ + ( + "eose_count", + self.eose_count, + RadrootsRelayFetchOutcomeKind::Eose, + ), + ( + "truncated_count", + self.truncated_count, + RadrootsRelayFetchOutcomeKind::Truncated, + ), + ( + "closed_count", + self.closed_count, + RadrootsRelayFetchOutcomeKind::Closed, + ), + ( + "notice_count", + self.notice_count, + RadrootsRelayFetchOutcomeKind::Notice, + ), + ] { + ensure_fetch_receipt_count( + field, + supplied, + self.relay_outcomes + .iter() + .filter(|outcome| outcome.kind() == kind) + .count(), + )?; + } + Ok(self) + } + + pub fn target_relays(&self) -> &[String] { + &self.target_relays + } + + pub fn inserted_count(&self) -> usize { + self.inserted_count + } + + pub fn duplicate_count(&self) -> usize { + self.duplicate_count + } + + pub fn not_persisted_count(&self) -> usize { + self.not_persisted_count + } + + pub fn malformed_count(&self) -> usize { + self.malformed_count + } + + pub fn out_of_filter_count(&self) -> usize { + self.out_of_filter_count + } + + pub fn skipped_over_limit_count(&self) -> usize { + self.skipped_over_limit_count + } + + pub fn verification_failed_count(&self) -> usize { + self.verification_failed_count + } + + pub fn admission_unsupported_count(&self) -> usize { + self.admission_unsupported_count + } + + pub fn admission_invalid_count(&self) -> usize { + self.admission_invalid_count + } + + pub fn valid_stream_eligible_count(&self) -> usize { + self.valid_stream_eligible_count + } + + pub fn visible_count(&self) -> usize { + self.visible_count + } + + pub fn not_admitted_count(&self) -> usize { + self.not_admitted_count + } + + pub fn not_current_count(&self) -> usize { + self.not_current_count + } + + pub fn suppressed_count(&self) -> usize { + self.suppressed_count + } + + pub fn eose_count(&self) -> usize { + self.eose_count + } + + pub fn truncated_count(&self) -> usize { + self.truncated_count + } + + pub fn closed_count(&self) -> usize { + self.closed_count + } + + pub fn notice_count(&self) -> usize { + self.notice_count + } + + pub fn events(&self) -> &[RadrootsRelayFetchEventReceipt] { + &self.events + } + + pub fn relay_outcomes(&self) -> &[RadrootsRelayFetchRelayOutcome] { + &self.relay_outcomes + } + fn from_processed_counts(processed: &RadrootsRelayProcessedFetch) -> Self { Self { + target_relays: processed.target_relays.clone(), inserted_count: 0, duplicate_count: 0, not_persisted_count: 0, @@ -1441,6 +1754,175 @@ impl RadrootsRelayFetchReceipt { } } +fn canonical_fetch_receipt_targets( + target_relays: Vec<String>, +) -> Result<Vec<String>, RadrootsRelayTransportError> { + if target_relays.is_empty() { + return Err(RadrootsRelayTransportError::EmptyTargetSet); + } + if target_relays.len() > radroots_transport::RADROOTS_TRANSPORT_TARGET_MAX_COUNT { + return Err(RadrootsRelayTransportError::FetchLimitTooLarge { + field: "target_relay_count", + max: radroots_transport::RADROOTS_TRANSPORT_TARGET_MAX_COUNT, + actual: target_relays.len(), + }); + } + let mut canonical = Vec::with_capacity(target_relays.len()); + let mut seen = BTreeSet::new(); + for relay_url in target_relays { + let relay_url = canonical_fetch_receipt_relay_url(relay_url.as_str())?; + if !seen.insert(relay_url.clone()) { + return Err(invalid_fetch_receipt( + "target_relays", + format!("duplicate relay URL `{relay_url}`"), + )); + } + canonical.push(relay_url); + } + Ok(canonical) +} + +fn ensure_fetch_receipt_relay_requested( + target_relays: &[String], + relay_url: &str, +) -> Result<(), RadrootsRelayTransportError> { + if !target_relays.iter().any(|target| target == relay_url) { + return Err(invalid_fetch_receipt( + "relay_url", + format!("relay URL `{relay_url}` was not requested"), + )); + } + Ok(()) +} + +fn validate_fetch_relay_outcomes( + target_relays: &[String], + relay_outcomes: &[RadrootsRelayFetchRelayOutcome], +) -> Result<(), RadrootsRelayTransportError> { + let mut terminal_relays = BTreeSet::new(); + for outcome in relay_outcomes { + ensure_fetch_receipt_relay_requested(target_relays, outcome.relay_url())?; + if outcome.kind() != RadrootsRelayFetchOutcomeKind::Notice + && !terminal_relays.insert(outcome.relay_url()) + { + return Err(invalid_fetch_receipt( + "relay_outcomes", + format!( + "duplicate or conflicting terminal outcome for `{}`", + outcome.relay_url() + ), + )); + } + } + Ok(()) +} + +fn validate_fetch_diagnostic_budget( + events: &[RadrootsRelayFetchEventReceipt], + relay_outcomes: &[RadrootsRelayFetchRelayOutcome], +) -> Result<(), RadrootsRelayTransportError> { + let mut diagnostic_bytes = 0usize; + for message in events + .iter() + .filter_map(RadrootsRelayFetchEventReceipt::message) + .chain( + relay_outcomes + .iter() + .filter_map(RadrootsRelayFetchRelayOutcome::message), + ) + { + diagnostic_bytes = diagnostic_bytes.checked_add(message.len()).ok_or( + RadrootsRelayTransportError::DiagnosticLimitExceeded { + field: "fetch_request_diagnostics", + max: radroots_transport::RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES, + actual: usize::MAX, + }, + )?; + if diagnostic_bytes > radroots_transport::RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES { + return Err(RadrootsRelayTransportError::DiagnosticLimitExceeded { + field: "fetch_request_diagnostics", + max: radroots_transport::RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES, + actual: diagnostic_bytes, + }); + } + } + Ok(()) +} + +fn ensure_fetch_receipt_count( + field: &'static str, + supplied: usize, + actual: usize, +) -> Result<(), RadrootsRelayTransportError> { + if supplied != actual { + return Err(invalid_fetch_receipt( + field, + format!("supplied count {supplied} does not match derived count {actual}"), + )); + } + Ok(()) +} + +#[derive(Deserialize)] +#[serde(deny_unknown_fields)] +struct RadrootsRelayFetchReceiptWire { + target_relays: Vec<String>, + inserted_count: usize, + duplicate_count: usize, + not_persisted_count: usize, + malformed_count: usize, + out_of_filter_count: usize, + skipped_over_limit_count: usize, + verification_failed_count: usize, + admission_unsupported_count: usize, + admission_invalid_count: usize, + valid_stream_eligible_count: usize, + visible_count: usize, + not_admitted_count: usize, + not_current_count: usize, + suppressed_count: usize, + eose_count: usize, + truncated_count: usize, + closed_count: usize, + notice_count: usize, + events: Vec<RadrootsRelayFetchEventReceipt>, + relay_outcomes: Vec<RadrootsRelayFetchRelayOutcome>, +} + +impl<'de> Deserialize<'de> for RadrootsRelayFetchReceipt { + fn deserialize<D>(deserializer: D) -> Result<Self, D::Error> + where + D: Deserializer<'de>, + { + let wire = RadrootsRelayFetchReceiptWire::deserialize(deserializer)?; + Self { + target_relays: wire.target_relays, + inserted_count: wire.inserted_count, + duplicate_count: wire.duplicate_count, + not_persisted_count: wire.not_persisted_count, + malformed_count: wire.malformed_count, + out_of_filter_count: wire.out_of_filter_count, + skipped_over_limit_count: wire.skipped_over_limit_count, + verification_failed_count: wire.verification_failed_count, + admission_unsupported_count: wire.admission_unsupported_count, + admission_invalid_count: wire.admission_invalid_count, + valid_stream_eligible_count: wire.valid_stream_eligible_count, + visible_count: wire.visible_count, + not_admitted_count: wire.not_admitted_count, + not_current_count: wire.not_current_count, + suppressed_count: wire.suppressed_count, + eose_count: wire.eose_count, + truncated_count: wire.truncated_count, + closed_count: wire.closed_count, + notice_count: wire.notice_count, + events: wire.events, + relay_outcomes: wire.relay_outcomes, + } + .checked() + .map_err(de::Error::custom) + } +} + #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum RadrootsRelayFetchRawBudgetExhaustion { RawEvents, @@ -1528,6 +2010,13 @@ fn process_relay_fetch_items( if target_relays.is_empty() { return Err(RadrootsRelayTransportError::EmptyTargetSet); } + if items.len() > RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX { + return Err(RadrootsRelayTransportError::FetchLimitTooLarge { + field: "raw_item_count", + max: RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX, + actual: items.len(), + }); + } let mut processed = RadrootsRelayProcessedFetch { target_relays, items: Vec::new(), @@ -1572,6 +2061,27 @@ fn process_relay_fetch_items( RadrootsRelayFetchItemBody::Event { raw_json, .. } => { if raw_budget.charge(raw_json.len()).is_err() { processed.skipped_over_limit_count += 1; + processed + .items + .push(RadrootsRelayProcessedFetchItem::Receipt( + RadrootsRelayFetchEventReceipt { + relay_url, + event_id: None, + inserted: false, + duplicate: false, + not_persisted: false, + malformed: false, + out_of_filter: false, + skipped_over_limit: true, + verification: RadrootsRelayFetchEventVerification::NotEvaluated, + admission: RadrootsRelayFetchEventAdmission::NotEvaluated, + admission_code: None, + valid_stream: RadrootsRelayFetchEventValidStream::NotEvaluated, + visibility: RadrootsRelayFetchEventVisibility::NotEvaluated, + message: None, + } + .checked()?, + )); continue; } if raw_json.len() > DEFAULT_RAW_JSON_MAX_BYTES { @@ -1621,7 +2131,7 @@ fn process_relay_fetch_items( admission_code: None, valid_stream: RadrootsRelayFetchEventValidStream::NotEvaluated, visibility: RadrootsRelayFetchEventVisibility::NotEvaluated, - message: Some("event JSON parse failed".to_owned()), + message: None, } .checked()?, )); @@ -1673,7 +2183,7 @@ fn process_relay_fetch_items( admission_code: None, valid_stream: RadrootsRelayFetchEventValidStream::NotEvaluated, visibility: RadrootsRelayFetchEventVisibility::NotEvaluated, - message: Some("event did not match relay fetch filters".to_owned()), + message: None, } .checked()?, )); @@ -1713,9 +2223,7 @@ fn process_relay_fetch_items( admission_code: None, valid_stream: RadrootsRelayFetchEventValidStream::NotEvaluated, visibility: RadrootsRelayFetchEventVisibility::NotEvaluated, - message: Some( - "accepted relay fetch event limit reached".to_owned(), - ), + message: None, } .checked()?, )); @@ -1800,7 +2308,7 @@ fn accepted_fetch_event_receipt( admission_code: None, valid_stream: RadrootsRelayFetchEventValidStream::NotEvaluated, visibility: RadrootsRelayFetchEventVisibility::NotEvaluated, - message: Some("event accepted by relay fetch filters".to_owned()), + message: None, } .checked() } @@ -1822,7 +2330,7 @@ fn duplicate_fetch_event_receipt( admission_code: None, valid_stream: RadrootsRelayFetchEventValidStream::NotEvaluated, visibility: RadrootsRelayFetchEventVisibility::NotEvaluated, - message: Some("event ID was already observed in this relay fetch".to_owned()), + message: None, } .checked() } diff --git a/crates/transport_nostr/tests/transport.rs b/crates/transport_nostr/tests/transport.rs @@ -38,14 +38,14 @@ use radroots_transport_nostr::{ RadrootsRelayFetchEventAdmission, RadrootsRelayFetchEventValidStream, RadrootsRelayFetchEventVerification, RadrootsRelayFetchEventVisibility, RadrootsRelayFetchFailure, RadrootsRelayFetchFilters, RadrootsRelayFetchItem, - RadrootsRelayFetchMode, RadrootsRelayFetchOutcomeKind, RadrootsRelayFetchRelayOutcome, - RadrootsRelayFetchRequest, RadrootsRelayFetchedEvent, RadrootsRelayOutcome, - RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter, RadrootsRelayPublishRelayReceipt, - RadrootsRelayPublishRequest, RadrootsRelayTargetSet, RadrootsRelayTransportError, - RadrootsRelayUrl, RadrootsRelayUrlPolicy, fetch_and_ingest_relay_events, fetch_relay_events, - fetch_relay_events_blocking, publish_claimed_outbox_event, - publish_claimed_outbox_event_with_transport, publish_signed_event, - verified_signed_event_payload, + RadrootsRelayFetchMode, RadrootsRelayFetchOutcomeKind, RadrootsRelayFetchReceipt, + RadrootsRelayFetchRelayOutcome, RadrootsRelayFetchRequest, RadrootsRelayFetchedEvent, + RadrootsRelayOutcome, RadrootsRelayOutcomeKind, RadrootsRelayPublishAdapter, + RadrootsRelayPublishRelayReceipt, RadrootsRelayPublishRequest, RadrootsRelayTargetSet, + RadrootsRelayTransportError, RadrootsRelayUrl, RadrootsRelayUrlPolicy, + fetch_and_ingest_relay_events, fetch_relay_events, fetch_relay_events_blocking, + publish_claimed_outbox_event, publish_claimed_outbox_event_with_transport, + publish_signed_event, verified_signed_event_payload, }; use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; @@ -2395,10 +2395,10 @@ fn fetch_blocking_facade_runs_mock_adapter() { let receipt = fetch_relay_events_blocking(&adapter, post_relay_fetch_request(1_090, 10)) .expect("blocking fetch"); - assert_eq!(receipt.events.len(), 1); - assert_eq!(receipt.events[0].event().id.to_hex(), accepted_id); - assert_eq!(receipt.events[0].observed_at_ms(), 1_090); - assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]); + assert_eq!(receipt.events().len(), 1); + assert_eq!(receipt.events()[0].event().id.to_hex(), accepted_id); + assert_eq!(receipt.events()[0].observed_at_ms(), 1_090); + assert_eq!(receipt.connected_relays(), vec![RELAY_PRIMARY_WSS]); } #[tokio::test] @@ -2416,10 +2416,10 @@ async fn fetch_canonicalizes_adapter_relay_spelling_and_uses_request_observation .await .expect("canonical fetch"); - assert_eq!(receipt.events.len(), 1); - assert_eq!(receipt.events[0].relay_url(), RELAY_PRIMARY_WSS); - assert_eq!(receipt.events[0].observed_at_ms(), 1_091); - assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]); + assert_eq!(receipt.events().len(), 1); + assert_eq!(receipt.events()[0].relay_url(), RELAY_PRIMARY_WSS); + assert_eq!(receipt.events()[0].observed_at_ms(), 1_091); + assert_eq!(receipt.connected_relays(), vec![RELAY_PRIMARY_WSS]); } #[tokio::test] @@ -2435,16 +2435,16 @@ async fn fetch_verifies_events_before_acceptance_budgeting() { .await .expect("verified fetch"); - assert_eq!(receipt.events.len(), 1); - assert_eq!(receipt.events[0].event().id.to_hex(), accepted_id); - assert_eq!(receipt.verification_failed_count, 1); - assert_eq!(receipt.skipped_over_limit_count, 0); + assert_eq!(receipt.events().len(), 1); + assert_eq!(receipt.events()[0].event().id.to_hex(), accepted_id); + assert_eq!(receipt.verification_failed_count(), 1); + assert_eq!(receipt.skipped_over_limit_count(), 0); assert_eq!( - receipt.event_receipts[0].verification(), + receipt.event_receipts()[0].verification(), RadrootsRelayFetchEventVerification::Failed ); assert_eq!( - receipt.event_receipts[1].verification(), + receipt.event_receipts()[1].verification(), RadrootsRelayFetchEventVerification::Verified ); } @@ -2467,20 +2467,20 @@ async fn fetch_deduplicates_event_ids_without_starving_unique_events() { assert_eq!( receipt - .events + .events() .iter() .map(|event| event.event().id.to_hex()) .collect::<Vec<_>>(), vec![first_id.clone(), second_id] ); - assert_eq!(receipt.duplicate_count, 1); - assert_eq!(receipt.skipped_over_limit_count, 0); - assert_eq!(receipt.event_receipts.len(), 3); + assert_eq!(receipt.duplicate_count(), 1); + assert_eq!(receipt.skipped_over_limit_count(), 0); + assert_eq!(receipt.event_receipts().len(), 3); assert_eq!( - receipt.event_receipts[1].event_id(), + receipt.event_receipts()[1].event_id(), Some(first_id.as_str()) ); - assert!(receipt.event_receipts[1].was_duplicate()); + assert!(receipt.event_receipts()[1].was_duplicate()); } #[tokio::test] @@ -2494,13 +2494,13 @@ async fn fetch_reports_local_truncation_without_claiming_eose() { .await .expect("truncated fetch"); - assert!(receipt.events.is_empty()); - assert!(receipt.connected_relays.is_empty()); - assert_eq!(receipt.eose_count, 0); - assert_eq!(receipt.truncated_count, 1); - assert_eq!(receipt.relay_outcomes.len(), 1); + assert!(receipt.events().is_empty()); + assert!(receipt.connected_relays().is_empty()); + assert_eq!(receipt.eose_count(), 0); + assert_eq!(receipt.truncated_count(), 1); + assert_eq!(receipt.relay_outcomes().len(), 1); assert_eq!( - receipt.relay_outcomes[0].kind(), + receipt.relay_outcomes()[0].kind(), RadrootsRelayFetchOutcomeKind::Truncated ); } @@ -2600,78 +2600,78 @@ async fn fetch_ingests_events_and_records_transport_observations() { .await .expect("fetch ingest"); - assert_eq!(receipt.inserted_count, 3); - assert_eq!(receipt.duplicate_count, 1); - assert_eq!(receipt.not_persisted_count, 0); - assert_eq!(receipt.verification_failed_count, 1); - assert_eq!(receipt.admission_unsupported_count, 1); - assert_eq!(receipt.admission_invalid_count, 1); - assert_eq!(receipt.valid_stream_eligible_count, 2); - assert_eq!(receipt.visible_count, 2); - assert_eq!(receipt.not_admitted_count, 2); - assert_eq!(receipt.not_current_count, 0); - assert_eq!(receipt.suppressed_count, 0); - assert_eq!(receipt.malformed_count, 1); - assert_eq!(receipt.eose_count, 1); - assert_eq!(receipt.closed_count, 2); - assert_eq!(receipt.notice_count, 1); + assert_eq!(receipt.inserted_count(), 3); + assert_eq!(receipt.duplicate_count(), 1); + assert_eq!(receipt.not_persisted_count(), 0); + assert_eq!(receipt.verification_failed_count(), 1); + assert_eq!(receipt.admission_unsupported_count(), 1); + assert_eq!(receipt.admission_invalid_count(), 1); + assert_eq!(receipt.valid_stream_eligible_count(), 2); + assert_eq!(receipt.visible_count(), 2); + assert_eq!(receipt.not_admitted_count(), 2); + assert_eq!(receipt.not_current_count(), 0); + assert_eq!(receipt.suppressed_count(), 0); + assert_eq!(receipt.malformed_count(), 1); + assert_eq!(receipt.eose_count(), 1); + assert_eq!(receipt.closed_count(), 2); + assert_eq!(receipt.notice_count(), 1); assert_eq!( - receipt.inserted_count, + receipt.inserted_count(), receipt - .events + .events() .iter() .filter(|event| event.was_inserted()) .count() ); assert_eq!( - receipt.duplicate_count, + receipt.duplicate_count(), receipt - .events + .events() .iter() .filter(|event| event.was_duplicate()) .count() ); assert_eq!( - receipt.not_persisted_count, + receipt.not_persisted_count(), receipt - .events + .events() .iter() .filter(|event| event.was_not_persisted()) .count() ); assert_eq!( - receipt.admission_unsupported_count, + receipt.admission_unsupported_count(), receipt - .events + .events() .iter() .filter(|event| event.admission() == RadrootsRelayFetchEventAdmission::Unsupported) .count() ); assert_eq!( - receipt.admission_invalid_count, + receipt.admission_invalid_count(), receipt - .events + .events() .iter() .filter(|event| event.admission() == RadrootsRelayFetchEventAdmission::Invalid) .count() ); assert_eq!( - receipt.verification_failed_count, + receipt.verification_failed_count(), receipt - .events + .events() .iter() .filter(|event| event.verification() == RadrootsRelayFetchEventVerification::Failed) .count() ); assert_eq!( - receipt.malformed_count, + receipt.malformed_count(), receipt - .events + .events() .iter() .filter(|event| event.is_malformed()) .count() ); - assert!(receipt.events.iter().all(|event| { + assert!(receipt.events().iter().all(|event| { usize::from(event.was_inserted()) + usize::from(event.was_duplicate()) + usize::from(event.was_not_persisted()) @@ -2679,101 +2679,104 @@ async fn fetch_ingests_events_and_records_transport_observations() { && (!event.is_malformed() || event.verification() == RadrootsRelayFetchEventVerification::NotEvaluated) })); - assert_eq!(receipt.relay_outcomes.len(), 4); - assert_eq!(receipt.relay_outcomes[0].relay_url(), RELAY_PRIMARY_WSS); + assert_eq!(receipt.relay_outcomes().len(), 4); + assert_eq!(receipt.relay_outcomes()[0].relay_url(), RELAY_PRIMARY_WSS); assert_eq!( - receipt.relay_outcomes[0].kind(), + receipt.relay_outcomes()[0].kind(), RadrootsRelayFetchOutcomeKind::Eose ); - assert!(receipt.relay_outcomes[0].relay_outcome().is_none()); - assert_eq!(receipt.relay_outcomes[1].relay_url(), RELAY_SECONDARY_WSS); + assert!(receipt.relay_outcomes()[0].relay_outcome().is_none()); + assert_eq!(receipt.relay_outcomes()[1].relay_url(), RELAY_SECONDARY_WSS); assert_eq!( - receipt.relay_outcomes[1] + receipt.relay_outcomes()[1] .relay_outcome() .expect("auth outcome") .kind(), RadrootsRelayOutcomeKind::AuthRequired ); - assert_eq!(receipt.relay_outcomes[2].relay_url(), RELAY_TERTIARY_WSS); + assert_eq!(receipt.relay_outcomes()[2].relay_url(), RELAY_TERTIARY_WSS); assert_eq!( - receipt.relay_outcomes[2] + receipt.relay_outcomes()[2] .relay_outcome() .expect("restricted outcome") .kind(), RadrootsRelayOutcomeKind::Restricted ); assert_eq!( - receipt.relay_outcomes[3].kind(), + receipt.relay_outcomes()[3].kind(), RadrootsRelayFetchOutcomeKind::Notice ); - assert!(receipt.relay_outcomes[3].relay_outcome().is_none()); + assert!(receipt.relay_outcomes()[3].relay_outcome().is_none()); assert_eq!( - receipt.events[0].admission(), + receipt.events()[0].admission(), RadrootsRelayFetchEventAdmission::Admitted ); assert_eq!( - receipt.events[0].valid_stream(), + receipt.events()[0].valid_stream(), RadrootsRelayFetchEventValidStream::Eligible ); assert_eq!( - receipt.events[0].visibility(), + receipt.events()[0].visibility(), RadrootsRelayFetchEventVisibility::Visible ); assert_eq!( - receipt.events[1].admission(), + receipt.events()[1].admission(), RadrootsRelayFetchEventAdmission::Admitted ); assert_eq!( - receipt.events[1].valid_stream(), + receipt.events()[1].valid_stream(), RadrootsRelayFetchEventValidStream::Eligible ); assert_eq!( - receipt.events[1].visibility(), + receipt.events()[1].visibility(), RadrootsRelayFetchEventVisibility::Visible ); assert_eq!( - receipt.events[2].admission(), + receipt.events()[2].admission(), RadrootsRelayFetchEventAdmission::Unsupported ); - assert_eq!(receipt.events[2].admission_code(), Some("unsupported_kind")); assert_eq!( - receipt.events[2].valid_stream(), + receipt.events()[2].admission_code(), + Some("unsupported_kind") + ); + assert_eq!( + receipt.events()[2].valid_stream(), RadrootsRelayFetchEventValidStream::Ineligible ); assert_eq!( - receipt.events[2].visibility(), + receipt.events()[2].visibility(), RadrootsRelayFetchEventVisibility::NotAdmitted ); assert_eq!( - receipt.events[3].admission(), + receipt.events()[3].admission(), RadrootsRelayFetchEventAdmission::Invalid ); assert_eq!( - receipt.events[3].admission_code(), + receipt.events()[3].admission_code(), Some("reply_event_id_invalid") ); assert_eq!( - receipt.events[3].valid_stream(), + receipt.events()[3].valid_stream(), RadrootsRelayFetchEventValidStream::Ineligible ); assert_eq!( - receipt.events[3].visibility(), + receipt.events()[3].visibility(), RadrootsRelayFetchEventVisibility::NotAdmitted ); assert_eq!( - receipt.events[4].verification(), + receipt.events()[4].verification(), RadrootsRelayFetchEventVerification::Failed ); assert_eq!( - receipt.events[4].admission(), + receipt.events()[4].admission(), RadrootsRelayFetchEventAdmission::NotEvaluated ); assert_eq!( - receipt.events[5].verification(), + receipt.events()[5].verification(), RadrootsRelayFetchEventVerification::NotEvaluated ); assert_eq!( - receipt.events[5].admission(), + receipt.events()[5].admission(), RadrootsRelayFetchEventAdmission::NotEvaluated ); @@ -2811,6 +2814,31 @@ async fn fetch_ingests_events_and_records_transport_observations() { assert!(!serialized_event.contains_key("projection_eligible")); assert!(!serialized_event.contains_key("admission_status")); + assert_eq!( + serde_json::from_value::<RadrootsRelayFetchReceipt>(serialized.clone()) + .expect("strict aggregate fetch receipt reload"), + receipt + ); + let mut wrong_count = serialized.clone(); + wrong_count["inserted_count"] = serde_json::json!(0); + assert!(serde_json::from_value::<RadrootsRelayFetchReceipt>(wrong_count).is_err()); + let mut unrequested_event = serialized.clone(); + unrequested_event["events"][0]["relay_url"] = + serde_json::json!("wss://unrequested.example.com"); + assert!(serde_json::from_value::<RadrootsRelayFetchReceipt>(unrequested_event).is_err()); + let mut duplicate_targets = serialized.clone(); + duplicate_targets["target_relays"] = serde_json::json!([RELAY_PRIMARY_WSS, RELAY_PRIMARY_WSS]); + assert!(serde_json::from_value::<RadrootsRelayFetchReceipt>(duplicate_targets).is_err()); + let mut diagnostic_over = serialized.clone(); + diagnostic_over["events"][0]["message"] = + serde_json::json!("a".repeat(RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES / 2)); + diagnostic_over["events"][1]["message"] = + serde_json::json!("b".repeat(RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES / 2 + 1)); + assert!(serde_json::from_value::<RadrootsRelayFetchReceipt>(diagnostic_over).is_err()); + let mut unknown = serialized; + unknown["extra"] = serde_json::json!(true); + assert!(serde_json::from_value::<RadrootsRelayFetchReceipt>(unknown).is_err()); + let observations = store .observations_for_event(signed.id_str()) .await @@ -2852,24 +2880,24 @@ async fn fetch_reports_final_replaceable_visibility_when_newer_arrives_first() { .await .expect("fetch ingest"); - assert_eq!(receipt.inserted_count, 2); - assert_eq!(receipt.valid_stream_eligible_count, 2); - assert_eq!(receipt.visible_count, 1); - assert_eq!(receipt.not_current_count, 1); + assert_eq!(receipt.inserted_count(), 2); + assert_eq!(receipt.valid_stream_eligible_count(), 2); + assert_eq!(receipt.visible_count(), 1); + assert_eq!(receipt.not_current_count(), 1); assert_eq!( - receipt.events[0].visibility(), + receipt.events()[0].visibility(), RadrootsRelayFetchEventVisibility::Visible ); assert_eq!( - receipt.events[1].admission(), + receipt.events()[1].admission(), RadrootsRelayFetchEventAdmission::Admitted ); assert_eq!( - receipt.events[1].valid_stream(), + receipt.events()[1].valid_stream(), RadrootsRelayFetchEventValidStream::Eligible ); assert_eq!( - receipt.events[1].visibility(), + receipt.events()[1].visibility(), RadrootsRelayFetchEventVisibility::NotCurrent ); } @@ -2901,16 +2929,16 @@ async fn fetch_reports_final_replaceable_visibility_when_older_arrives_first() { .await .expect("fetch ingest"); - assert_eq!(receipt.inserted_count, 2); - assert_eq!(receipt.valid_stream_eligible_count, 2); - assert_eq!(receipt.visible_count, 1); - assert_eq!(receipt.not_current_count, 1); + assert_eq!(receipt.inserted_count(), 2); + assert_eq!(receipt.valid_stream_eligible_count(), 2); + assert_eq!(receipt.visible_count(), 1); + assert_eq!(receipt.not_current_count(), 1); assert_eq!( - receipt.events[0].visibility(), + receipt.events()[0].visibility(), RadrootsRelayFetchEventVisibility::NotCurrent ); assert_eq!( - receipt.events[1].visibility(), + receipt.events()[1].visibility(), RadrootsRelayFetchEventVisibility::Visible ); } @@ -2950,15 +2978,15 @@ async fn fetch_maps_one_final_visibility_snapshot_back_to_duplicate_receipts() { .await .expect("fetch ingest"); - assert_eq!(receipt.inserted_count, 2); - assert_eq!(receipt.duplicate_count, 1); - assert_eq!(receipt.valid_stream_eligible_count, 3); - assert_eq!(receipt.visible_count, 1); - assert_eq!(receipt.not_current_count, 2); - assert_eq!(receipt.events.len(), 3); + assert_eq!(receipt.inserted_count(), 2); + assert_eq!(receipt.duplicate_count(), 1); + assert_eq!(receipt.valid_stream_eligible_count(), 3); + assert_eq!(receipt.visible_count(), 1); + assert_eq!(receipt.not_current_count(), 2); + assert_eq!(receipt.events().len(), 3); assert_eq!( receipt - .events + .events() .iter() .map(|event| event.event_id()) .collect::<Vec<_>>(), @@ -2970,7 +2998,7 @@ async fn fetch_maps_one_final_visibility_snapshot_back_to_duplicate_receipts() { ); assert_eq!( receipt - .events + .events() .iter() .map(|event| event.visibility()) .collect::<Vec<_>>(), @@ -2980,7 +3008,7 @@ async fn fetch_maps_one_final_visibility_snapshot_back_to_duplicate_receipts() { RadrootsRelayFetchEventVisibility::NotCurrent, ] ); - assert!(receipt.events[2].was_duplicate()); + assert!(receipt.events()[2].was_duplicate()); } #[tokio::test] @@ -3020,29 +3048,29 @@ async fn fetch_reports_store_suppression_when_deletion_precedes_target_replay() .await .expect("fetch ingest"); - assert_eq!(receipt.inserted_count, 2); - assert_eq!(receipt.valid_stream_eligible_count, 2); - assert_eq!(receipt.visible_count, 1); - assert_eq!(receipt.suppressed_count, 1); - assert_eq!(receipt.not_current_count, 0); + assert_eq!(receipt.inserted_count(), 2); + assert_eq!(receipt.valid_stream_eligible_count(), 2); + assert_eq!(receipt.visible_count(), 1); + assert_eq!(receipt.suppressed_count(), 1); + assert_eq!(receipt.not_current_count(), 0); assert_eq!( - receipt.events[0].visibility(), + receipt.events()[0].visibility(), RadrootsRelayFetchEventVisibility::Visible ); assert_eq!( - receipt.events[1].event_id(), + receipt.events()[1].event_id(), Some(target.id.to_hex().as_str()) ); assert_eq!( - receipt.events[1].admission(), + receipt.events()[1].admission(), RadrootsRelayFetchEventAdmission::Admitted ); assert_eq!( - receipt.events[1].valid_stream(), + receipt.events()[1].valid_stream(), RadrootsRelayFetchEventValidStream::Eligible ); assert_eq!( - receipt.events[1].visibility(), + receipt.events()[1].visibility(), RadrootsRelayFetchEventVisibility::Suppressed ); } @@ -3068,16 +3096,16 @@ async fn fetch_reports_ephemeral_events_as_not_persisted_without_duplicate_or_st .await .expect("ephemeral fetch ingest"); - assert_eq!(receipt.inserted_count, 0); - assert_eq!(receipt.duplicate_count, 0); - assert_eq!(receipt.not_persisted_count, 2); - assert_eq!(receipt.malformed_count, 0); - assert_eq!(receipt.verification_failed_count, 0); - assert_eq!(receipt.admission_unsupported_count, 0); - assert_eq!(receipt.admission_invalid_count, 0); - assert_eq!(receipt.valid_stream_eligible_count, 0); - assert_eq!(receipt.events.len(), 2); - assert!(receipt.events.iter().all(|event| { + assert_eq!(receipt.inserted_count(), 0); + assert_eq!(receipt.duplicate_count(), 0); + assert_eq!(receipt.not_persisted_count(), 2); + assert_eq!(receipt.malformed_count(), 0); + assert_eq!(receipt.verification_failed_count(), 0); + assert_eq!(receipt.admission_unsupported_count(), 0); + assert_eq!(receipt.admission_invalid_count(), 0); + assert_eq!(receipt.valid_stream_eligible_count(), 0); + assert_eq!(receipt.events().len(), 2); + assert!(receipt.events().iter().all(|event| { !event.was_inserted() && !event.was_duplicate() && event.was_not_persisted() @@ -3143,14 +3171,14 @@ async fn fetch_rejects_out_of_filter_events_before_store_mutation() { .await .expect("fetch ingest"); - assert_eq!(receipt.inserted_count, 1); - assert_eq!(receipt.out_of_filter_count, 2); - assert_eq!(receipt.malformed_count, 0); - assert_eq!(receipt.admission_unsupported_count, 0); - assert_eq!(receipt.events.len(), 3); - assert!(receipt.events[0].is_out_of_filter()); - assert!(!receipt.events[1].is_out_of_filter()); - assert!(receipt.events[2].is_out_of_filter()); + assert_eq!(receipt.inserted_count(), 1); + assert_eq!(receipt.out_of_filter_count(), 2); + assert_eq!(receipt.malformed_count(), 0); + assert_eq!(receipt.admission_unsupported_count(), 0); + assert_eq!(receipt.events().len(), 3); + assert!(receipt.events()[0].is_out_of_filter()); + assert!(!receipt.events()[1].is_out_of_filter()); + assert!(receipt.events()[2].is_out_of_filter()); assert!( store .raw_event(accepted.id_str()) @@ -3207,34 +3235,34 @@ async fn fetch_event_cap_counts_accepted_in_filter_events_and_preserves_later_co .await .expect("fetch ingest"); - assert_eq!(receipt.inserted_count, 1); - assert_eq!(receipt.duplicate_count, 0); - assert_eq!(receipt.admission_unsupported_count, 0); - assert_eq!(receipt.malformed_count, 1); - assert_eq!(receipt.out_of_filter_count, 1); - assert_eq!(receipt.skipped_over_limit_count, 1); - assert_eq!(receipt.events.len(), 4); - assert!(receipt.events[0].is_malformed()); - assert!(receipt.events[1].is_out_of_filter()); - assert!(receipt.events[2].was_inserted()); - assert!(receipt.events[3].was_skipped_over_limit()); - assert_eq!(receipt.eose_count, 1); - assert_eq!(receipt.closed_count, 1); - assert_eq!(receipt.notice_count, 1); - assert_eq!(receipt.relay_outcomes.len(), 3); + assert_eq!(receipt.inserted_count(), 1); + assert_eq!(receipt.duplicate_count(), 0); + assert_eq!(receipt.admission_unsupported_count(), 0); + assert_eq!(receipt.malformed_count(), 1); + assert_eq!(receipt.out_of_filter_count(), 1); + assert_eq!(receipt.skipped_over_limit_count(), 1); + assert_eq!(receipt.events().len(), 4); + assert!(receipt.events()[0].is_malformed()); + assert!(receipt.events()[1].is_out_of_filter()); + assert!(receipt.events()[2].was_inserted()); + assert!(receipt.events()[3].was_skipped_over_limit()); + assert_eq!(receipt.eose_count(), 1); + assert_eq!(receipt.closed_count(), 1); + assert_eq!(receipt.notice_count(), 1); + assert_eq!(receipt.relay_outcomes().len(), 3); assert_eq!( - receipt.relay_outcomes[0].kind(), + receipt.relay_outcomes()[0].kind(), RadrootsRelayFetchOutcomeKind::Eose ); assert_eq!( - receipt.relay_outcomes[1] + receipt.relay_outcomes()[1] .relay_outcome() .expect("closed outcome") .kind(), RadrootsRelayOutcomeKind::AuthRequired ); assert_eq!( - receipt.relay_outcomes[2].kind(), + receipt.relay_outcomes()[2].kind(), RadrootsRelayFetchOutcomeKind::Notice ); assert!( @@ -3305,25 +3333,25 @@ async fn fetch_relay_events_applies_shared_filter_limit_and_outcome_evidence() { .expect("fetch events"); assert_eq!( - receipt.target_relays, + receipt.target_relays(), vec![RELAY_PRIMARY_WSS, RELAY_SECONDARY_WSS, RELAY_TERTIARY_WSS] ); - assert_eq!(receipt.connected_relays, vec![RELAY_PRIMARY_WSS]); - assert_eq!(receipt.failed_relays.len(), 1); - assert_eq!(receipt.failed_relays[0].relay_url(), RELAY_SECONDARY_WSS); - assert_eq!(receipt.events.len(), 1); - assert_eq!(receipt.events[0].event().id.to_hex(), accepted_id); - assert_eq!(receipt.malformed_count, 1); - assert_eq!(receipt.out_of_filter_count, 1); - assert_eq!(receipt.skipped_over_limit_count, 1); - assert_eq!(receipt.eose_count, 1); - assert_eq!(receipt.closed_count, 1); - assert_eq!(receipt.notice_count, 1); - assert_eq!(receipt.event_receipts.len(), 4); - assert!(receipt.event_receipts[0].is_malformed()); - assert!(receipt.event_receipts[1].is_out_of_filter()); - assert!(!receipt.event_receipts[2].is_malformed()); - assert!(receipt.event_receipts[3].was_skipped_over_limit()); + assert_eq!(receipt.connected_relays(), vec![RELAY_PRIMARY_WSS]); + assert_eq!(receipt.failed_relays().len(), 1); + assert_eq!(receipt.failed_relays()[0].relay_url(), RELAY_SECONDARY_WSS); + assert_eq!(receipt.events().len(), 1); + assert_eq!(receipt.events()[0].event().id.to_hex(), accepted_id); + assert_eq!(receipt.malformed_count(), 1); + assert_eq!(receipt.out_of_filter_count(), 1); + assert_eq!(receipt.skipped_over_limit_count(), 1); + assert_eq!(receipt.eose_count(), 1); + assert_eq!(receipt.closed_count(), 1); + assert_eq!(receipt.notice_count(), 1); + assert_eq!(receipt.event_receipts().len(), 4); + assert!(receipt.event_receipts()[0].is_malformed()); + assert!(receipt.event_receipts()[1].is_out_of_filter()); + assert!(!receipt.event_receipts()[2].is_malformed()); + assert!(receipt.event_receipts()[3].was_skipped_over_limit()); } #[tokio::test] @@ -3352,12 +3380,13 @@ async fn fetch_raw_scan_limit_bounds_noisy_adapter_output() { .await .expect("fetch ingest"); - assert_eq!(receipt.inserted_count, 0); - assert_eq!(receipt.malformed_count, 1); - assert_eq!(receipt.out_of_filter_count, 1); - assert_eq!(receipt.skipped_over_limit_count, 1); - assert_eq!(receipt.events.len(), 2); - assert_eq!(receipt.eose_count, 1); + assert_eq!(receipt.inserted_count(), 0); + assert_eq!(receipt.malformed_count(), 1); + assert_eq!(receipt.out_of_filter_count(), 1); + assert_eq!(receipt.skipped_over_limit_count(), 1); + assert_eq!(receipt.events().len(), 3); + assert!(receipt.events()[2].was_skipped_over_limit()); + assert_eq!(receipt.eose_count(), 1); assert!( store .raw_event(accepted_id.as_str()) @@ -3395,9 +3424,9 @@ async fn fetch_raw_json_byte_limit_is_exact_global_and_sticky() { ) .await .expect("exact fetch"); - assert_eq!(exact_receipt.events.len(), 2); - assert_eq!(exact_receipt.skipped_over_limit_count, 1); - assert_eq!(exact_receipt.eose_count, 2); + assert_eq!(exact_receipt.events().len(), 2); + assert_eq!(exact_receipt.skipped_over_limit_count(), 1); + assert_eq!(exact_receipt.eose_count(), 2); let crossing_receipt = fetch_relay_events( &RadrootsMockRelayFetchAdapter::new(items), @@ -3408,9 +3437,9 @@ async fn fetch_raw_json_byte_limit_is_exact_global_and_sticky() { ) .await .expect("crossing fetch"); - assert_eq!(crossing_receipt.events.len(), 1); - assert_eq!(crossing_receipt.skipped_over_limit_count, 2); - assert_eq!(crossing_receipt.eose_count, 2); + assert_eq!(crossing_receipt.events().len(), 1); + assert_eq!(crossing_receipt.skipped_over_limit_count(), 2); + assert_eq!(crossing_receipt.eose_count(), 2); } #[tokio::test] @@ -3464,14 +3493,15 @@ async fn fetch_raw_json_budget_charges_every_preparse_event_class() { .await .expect("fetch ingest"); - assert_eq!(receipt.inserted_count, 1); - assert_eq!(receipt.malformed_count, 1); - assert_eq!(receipt.out_of_filter_count, 1); - assert_eq!(receipt.duplicate_count, 1); - assert_eq!(receipt.skipped_over_limit_count, 2); - assert_eq!(receipt.events.len(), 5); - assert!(receipt.events[4].was_skipped_over_limit()); - assert_eq!(receipt.eose_count, 2); + assert_eq!(receipt.inserted_count(), 1); + assert_eq!(receipt.malformed_count(), 1); + assert_eq!(receipt.out_of_filter_count(), 1); + assert_eq!(receipt.duplicate_count(), 1); + assert_eq!(receipt.skipped_over_limit_count(), 2); + assert_eq!(receipt.events().len(), 6); + assert!(receipt.events()[4].was_skipped_over_limit()); + assert!(receipt.events()[5].was_skipped_over_limit()); + assert_eq!(receipt.eose_count(), 2); } #[test] @@ -3517,7 +3547,7 @@ async fn fetch_subscription_mode_and_store_errors_are_propagated() { .await .expect("fetch ingest"); - assert_eq!(receipt.inserted_count, 1); + assert_eq!(receipt.inserted_count(), 1); let observations = store .observations_for_event(signed.id_str()) .await @@ -3886,6 +3916,82 @@ fn fetched_events_bind_raw_bytes_identity_endpoint_and_observation_time() { } #[tokio::test] +async fn fetch_receipts_enforce_outer_item_and_complete_diagnostic_budgets() { + let exact_items = (0..RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX) + .map(|_| bounded_fetch_notice(RELAY_PRIMARY_WSS, "")) + .collect::<Vec<_>>(); + let exact_receipt = fetch_relay_events( + &RadrootsMockRelayFetchAdapter::new(exact_items), + post_relay_fetch_request(1_600, 1), + ) + .await + .expect("exact raw receipt count"); + assert_eq!( + exact_receipt.relay_outcomes().len(), + RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX + ); + + let one_over_items = (0..=RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX) + .map(|_| bounded_fetch_notice(RELAY_PRIMARY_WSS, "")) + .collect::<Vec<_>>(); + assert!(matches!( + fetch_relay_events( + &RadrootsMockRelayFetchAdapter::new(one_over_items), + post_relay_fetch_request(1_600, 1), + ) + .await, + Err(RadrootsRelayTransportError::FetchLimitTooLarge { + field: "raw_item_count", + max: RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX, + actual, + }) if actual == RADROOTS_RELAY_FETCH_RAW_EVENT_LIMIT_MAX + 1 + )); + + let exact_diagnostics = RadrootsMockRelayFetchAdapter::new(vec![ + bounded_fetch_notice( + RELAY_PRIMARY_WSS, + "a".repeat(RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES / 2), + ), + bounded_fetch_notice( + RELAY_PRIMARY_WSS, + "b".repeat(RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES / 2), + ), + ]); + let store = RadrootsEventStore::open_memory().await.expect("store"); + fetch_and_ingest_relay_events( + &exact_diagnostics, + &store, + post_relay_fetch_request(1_600, 1), + ) + .await + .expect("exact aggregate diagnostic bytes"); + + let oversized_diagnostics = RadrootsMockRelayFetchAdapter::new(vec![ + bounded_fetch_notice( + RELAY_PRIMARY_WSS, + "a".repeat(RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES / 2), + ), + bounded_fetch_notice( + RELAY_PRIMARY_WSS, + "b".repeat(RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES / 2 + 1), + ), + ]); + assert!(matches!( + fetch_and_ingest_relay_events( + &oversized_diagnostics, + &store, + post_relay_fetch_request(1_600, 1), + ) + .await, + Err(RadrootsRelayTransportError::DiagnosticLimitExceeded { + field: "fetch_request_diagnostics", + max: RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES, + actual, + }) if actual == RADROOTS_TRANSPORT_DIAGNOSTIC_MAX_BYTES + 1 + )); +} + +#[tokio::test] async fn outbox_publish_persists_partial_success_and_skips_accepted_retry() { let outbox = RadrootsOutbox::open_memory().await.expect("outbox"); let store = RadrootsEventStore::open_memory().await.expect("store"); @@ -5711,16 +5817,16 @@ async fn smoke_relay_fetch_processes_one_thousand_event_receipts() { .await .expect("fetch"); - assert_eq!(receipt.inserted_count, 1_000); - assert_eq!(receipt.duplicate_count, 0); - assert_eq!(receipt.malformed_count, 0); - assert_eq!(receipt.verification_failed_count, 0); - assert_eq!(receipt.admission_unsupported_count, 0); - assert_eq!(receipt.admission_invalid_count, 0); - assert_eq!(receipt.valid_stream_eligible_count, 1_000); - assert_eq!(receipt.visible_count, 1_000); - assert_eq!(receipt.events.len(), 1_000); - assert!(receipt.events.iter().all(|event| event.valid_stream() + assert_eq!(receipt.inserted_count(), 1_000); + assert_eq!(receipt.duplicate_count(), 0); + assert_eq!(receipt.malformed_count(), 0); + assert_eq!(receipt.verification_failed_count(), 0); + assert_eq!(receipt.admission_unsupported_count(), 0); + assert_eq!(receipt.admission_invalid_count(), 0); + assert_eq!(receipt.valid_stream_eligible_count(), 1_000); + assert_eq!(receipt.visible_count(), 1_000); + assert_eq!(receipt.events().len(), 1_000); + assert!(receipt.events().iter().all(|event| event.valid_stream() == RadrootsRelayFetchEventValidStream::Eligible && event.visibility() == RadrootsRelayFetchEventVisibility::Visible)); let replay = store.valid_stream_after(0, 1_000).await.expect("replay");