lib

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

commit 9119bb31c59cd84a4a05c0a75d65d708b49cdd8e
parent ab084e6943b24e24cb7408f05fbb6241ae8f8c83
Author: triesap <tyson@radroots.org>
Date:   Thu, 10 Sep 2026 00:51:44 +0000

transport: Preserve partial coverage at capped relay windows

- Keep exact-cap EOSE distinct from complete window coverage
- Yield scoped older continuation after received equal-time peers
- Reject malformed boundaries and preserve finite shared fetch budgets
- Verify real relay recovery, unchanged APIs, coverage and workspace gates

Diffstat:
Acontracts/architecture/decisions/nostr_fetch_windows.v1.json | 19+++++++++++++++++++
Mcontracts/architecture/deviations.toml | 27+++++++++++++++++++++++++++
Mcrates/transport_nostr/README.md | 22+++++++++++++++++++---
Mcrates/transport_nostr/src/source.rs | 71++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------
Acrates/transport_nostr/src/source_paging_tests.rs | 265+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/transport_nostr/src/source_window.rs | 93+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/transport_nostr/tests/fetch_windows.rs | 161+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/transport_nostr/tests/package_boundary.rs | 2++
8 files changed, 644 insertions(+), 16 deletions(-)

diff --git a/contracts/architecture/decisions/nostr_fetch_windows.v1.json b/contracts/architecture/decisions/nostr_fetch_windows.v1.json @@ -0,0 +1,19 @@ +{ + "schema": "radroots.nostr-fetch-windows.v1", + "owner": "radroots_transport_nostr", + "status": "implemented", + "scope": "Existing bounded Nostr EVENT/REQ discovery with exact target/selector-bound opaque continuation.", + "upstream_event_limit": 1000, + "capped_eose": "EOSE at the requested 1000-event cap is partial coverage. It may prove current relay availability but cannot prove exhaustion of the requested window. Preserve connection-failure backoff separately from a successful capped query.", + "ordinary_pages": "Retain the existing nostr-v2 timestamp/event-ID cursor while collected deduplicated candidates remain. Every received candidate in the bounded window must be available before yielding an older-window continuation.", + "saturated_window": "When capped EOSE leaves no additional collected candidate, return NextPage::Cancelled with an optional explicit older-window cursor. Shared pull returns control; callers must disclose partial coverage before deliberately continuing older. Additional same-time events may remain inaccessible at the cap and are never declared recovered.", + "older_cursor": "The additive nostr-until-v1:<inclusive_until>:<scope> form contains a canonical decimal u64 and the exact existing target/selector scope digest. Its bound is strictly before the maximum of the oldest matching timestamps in capped EOSE relay windows. Match the effective remote until bound defensively before deriving a boundary; never advance from malformed-only or out-of-bound inventory. Timestamp zero has no older continuation. Existing ordinary v2 cursors remain supported without changing their meaning.", + "bounds": "Preserve all existing relay, connection, deadline, raw-byte, parse, event and page limits. Retain only a scalar oldest matching timestamp per current relay batch and one aggregate capped boundary. No ID exclusion set, history buffer, automatic retry or new transport is introduced.", + "compatibility": "Public Rust types, generated SDK surfaces, wire event schemas and package/dependency identities remain unchanged. The new opaque continuation is only interpreted by the concrete owner; callers must preserve typed partial outcomes and pull termination.", + "non_goals": [ + "complete or lossless global history", + "new Nostr extension, NIP-77 or transport", + "app projection, scheduler or database ownership", + "release, Nix, device or deployment activation" + ] +} diff --git a/contracts/architecture/deviations.toml b/contracts/architecture/deviations.toml @@ -2,6 +2,33 @@ schema_version = 1 architecture_id = "radroots.crates.release.v1" [[deviation]] +id = "RCRV1-DEV-016" +date = "2026-09-10" +status = "closed" +approval = "Standing user authorization covers necessary shared-owner repairs, verified checkpoint commits and non-force producer publication." +affected_steps = ["201"] +spec_anchors = [ + "contracts/crates/release_v1/radroots_crates_release_v1.toml#package.radroots_transport_nostr", +] +source_evidence = [ + "The concrete REQ asks for 1000 events, but exactly 1000 followed by EOSE is currently normalized as Complete.", + "After a 500-event first page, another 500 equal-time events produce Complete with no cursor despite additional capped or older history.", +] +replacement_action = "Implement contracts/architecture/decisions/nostr_fetch_windows.v1.json within the existing adapter: partial capped coverage, lossless paging of received candidates and finite explicit older-window continuation." +verification = [ + "Reproduce capped EOSE false-completion, then verify ordinary and capped ties, duplicate relays, older history, malformed/boundary inputs and scoped continuation through injected and loopback sources.", + "Pass affected package and SDK checks, unchanged public API and coverage thresholds, full workspace and required release-preflight gates before publication.", +] +unresolved_risk = "Current owner, SDK, API, coverage and workspace gates pass. A saturated timestamp can still hide additional events; the partial yield and explicit older continuation never claim lossless or global-history recovery." +normative_architecture_change = false +adr_required = false +closure_evidence = [ + "The false-complete regression now passes with 501, 1000 and 1001 equal-time events, cross-relay deduplication, explicit older continuation and canonical scoped cursor rejection.", + "A real WebSocket relay returns exactly 1000 signed same-time events and EOSE; both received pages and the explicit older query pass without reconnect backoff or automatic boundary skipping.", + "Malformed/out-of-bound and zero-time cases fail closed. Public transport and SDK APIs remain byte-identical. Fresh affected coverage, all 45 reports and aggregate, complete workspace and release-preflight checks pass without threshold changes.", +] + +[[deviation]] id = "RCRV1-DEV-015" date = "2026-09-09" status = "closed" diff --git a/crates/transport_nostr/README.md b/crates/transport_nostr/README.md @@ -136,8 +136,9 @@ transport-neutral kind, author, exact indexed single-letter tag, and event-time selectors into Nostr filters, reapplies those selectors defensively, applies the request page bound, deduplicates events by event ID, preserves per-relay provenance, and emits an opaque versioned cursor bound to the exact target set -and selector when more results remain. Equal timestamps are ordered by event ID -so overlap-safe reconnect pagination cannot skip peers. Malformed relay events +and selector when more collected results remain. Equal timestamps are ordered +by event ID so all received peers can be paged without timestamp-only loss. +Malformed relay events are ignored and reported as a partial target outcome rather than admitted. Fetch reports `Complete` only after the exact subscription receives EOSE strictly before its absolute deadline. Deadline expiry is `Cancelled`, and a @@ -152,7 +153,22 @@ limits both frames and complete messages to 512 KiB before upstream decoding. Shared canonical decoding rechecks the event, aggregate-byte and inventory bounds before candidate collection. A fetch returns at most the caller's validated 1,000-event page limit. EOSE describes only the requested relay -subscription and never proves complete global history. +subscription and never proves complete global history. A response reaching the +1,000-event cap is `Partial` even when EOSE follows: the relay may still conceal +other events at the same timestamp or older history. Capped EOSE proves relay +availability and does not start connection-failure backoff. + +After all collected candidates have been paged, a capped window returns +`NextPage::Cancelled` with an optional opaque older-window continuation. Shared +pull yields control at that boundary. Callers must disclose the partial +coverage before explicitly continuing older discovery; additional same-time +events remain unproven and are not declared recovered. The continuation moves +strictly before the capped timestamp using the same exact target/selector scope. +Malformed-only or out-of-bound results cannot invent a continuation; timestamp +zero cannot underflow. Ordinary `nostr-v2` event cursors remain supported; the +additive `nostr-until-v1` form is interpreted only by this concrete adapter. +No automatic retry, unlimited scan, new protocol or global completeness claim +is implied by either form. Live subscriptions use the same explicit readable targets and selector translation. A caller checkpoint is scoped to one exact target and selector; diff --git a/crates/transport_nostr/src/source.rs b/crates/transport_nostr/src/source.rs @@ -22,6 +22,13 @@ use std::time::{SystemTime, UNIX_EPOCH}; mod budget; use budget::FetchBudget; +#[path = "source_window.rs"] +mod window; + +#[cfg(test)] +#[path = "source_paging_tests.rs"] +mod paging_tests; + const UPSTREAM_FETCH_LIMIT: usize = 1_000; const CURSOR_PREFIX: &str = "nostr-v2"; const CURSOR_SCOPE_DOMAIN: &[u8] = b"radroots.transport-nostr.fetch-cursor.v2\0"; @@ -296,9 +303,10 @@ impl EventSource for NostrTransport { let cursor_scope = request_scope(&request); let cursor = request .cursor() - .map(|cursor| parse_cursor(cursor, cursor_scope.as_str())) + .map(|cursor| window::Position::parse(cursor, cursor_scope.as_str())) .transpose()?; - let selector_until = request.selector().until_unix_seconds(); + let effective_until = + window::effective_until(request.selector().until_unix_seconds(), cursor.as_ref()); let now_ms = unix_time_ms(); let remaining_ms = request.bounds().deadline_unix_ms().saturating_sub(now_ms); let timeout_ms = remaining_ms.min(self.config().request_timeout_ms()); @@ -354,12 +362,7 @@ impl EventSource for NostrTransport { .fetch(SourceQuery { relays: targets.keys().cloned().collect(), selector: request.selector().clone(), - until_unix_seconds: match (selector_until, cursor.as_ref()) { - (Some(until), Some(cursor)) => Some(until.min(cursor.created_at_unix_s())), - (Some(until), None) => Some(until), - (None, Some(cursor)) => Some(cursor.created_at_unix_s()), - (None, None) => None, - }, + until_unix_seconds: effective_until, connect_timeout: Duration::from_millis( timeout_ms.min(self.config().connect_timeout_ms()), ), @@ -368,6 +371,8 @@ impl EventSource for NostrTransport { }) .await; let mut candidates = Vec::new(); + let mut capped_window = false; + let mut older_boundary = None; let parse_budget = FetchBudget::default(); let mut malformed_by_relay = BTreeMap::<RelayUrl, usize>::new(); let mut reported = BTreeSet::new(); @@ -399,6 +404,10 @@ impl EventSource for NostrTransport { continue; } }; + let capped_eose = terminal == FetchTargetState::Complete + && raw_events.len() >= UPSTREAM_FETCH_LIMIT; + capped_window |= capped_eose; + let mut oldest_matching = None; for (index, raw) in raw_events.into_iter().enumerate() { if index >= UPSTREAM_FETCH_LIMIT || !parse_budget.event(raw.len()) { if terminal != FetchTargetState::Cancelled { @@ -407,7 +416,15 @@ impl EventSource for NostrTransport { break; } match radroots_event_codec::decode::signed_event(raw.as_str()) { - Ok(event) if request.selector().matches(&event) => { + Ok(event) + if request.selector().matches(&event) + && effective_until + .is_none_or(|until| event.created_at() <= until) => + { + oldest_matching = + Some(oldest_matching.map_or(event.created_at(), |oldest: u64| { + oldest.min(event.created_at()) + })); candidates.push(Candidate { relay: relay.clone(), created_at: event.created_at(), @@ -421,6 +438,12 @@ impl EventSource for NostrTransport { } } } + if capped_eose && let Some(oldest) = oldest_matching { + older_boundary = + Some(older_boundary.map_or(oldest, |boundary: u64| boundary.max(oldest))); + } + // A capped EOSE proves relay availability, while the page's + // coverage stays partial. It must not start reconnect backoff. self.status.record_read( &relay, terminal == FetchTargetState::Complete, @@ -434,9 +457,9 @@ impl EventSource for NostrTransport { FetchTargetState::Cancelled, ) .with_message("relay fetch deadline elapsed before EOSE") - } else if terminal == FetchTargetState::Partial { + } else if terminal == FetchTargetState::Partial || capped_eose { FetchTargetOutcome::new(target.fingerprint().clone(), FetchTargetState::Partial) - .with_message("relay result exceeded the bounded fetch inventory") + .with_message("relay result reached the bounded fetch inventory") } else if malformed == 0 { FetchTargetOutcome::new( target.fingerprint().clone(), @@ -463,7 +486,7 @@ impl EventSource for NostrTransport { } candidates.sort_by(compare_candidate); if let Some(cursor) = &cursor { - candidates.retain(|candidate| candidate_is_after_cursor(candidate, cursor)); + candidates.retain(|candidate| cursor.includes(candidate)); } let mut seen = BTreeSet::new(); candidates.retain(|candidate| seen.insert(candidate.event_id.clone())); @@ -476,6 +499,11 @@ impl EventSource for NostrTransport { "{CURSOR_PREFIX}:{}:{}:{cursor_scope}", last.created_at, last.event_id, ))?) + } else if capped_window { + NextPage::Cancelled { + resume_from: older_boundary + .and_then(|boundary| window::before_boundary(boundary, &cursor_scope)), + } } else { NextPage::Complete }; @@ -1293,13 +1321,30 @@ mod tests { .iter() .filter(|outcome| outcome.state() == FetchTargetState::Complete) .count(), - 4 + 0 ); assert_eq!( page.target_outcomes() .iter() .filter(|outcome| outcome.state() == FetchTargetState::Partial) .count(), + 5 + ); + let report = transport.relay_status(); + assert_eq!( + report + .relays() + .iter() + .filter(|relay| relay.read().state() == crate::RelayEvidenceState::Available) + .count(), + 4 + ); + assert_eq!( + report + .relays() + .iter() + .filter(|relay| relay.read().state() == crate::RelayEvidenceState::Unavailable) + .count(), 1 ); } diff --git a/crates/transport_nostr/src/source_paging_tests.rs b/crates/transport_nostr/src/source_paging_tests.rs @@ -0,0 +1,265 @@ +use super::*; +use crate::{Config, RelayUrlPolicy}; +use radroots_transport::{Target, TargetSet, source::FetchBounds}; +use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering}; + +struct CappedRelay { + history: Vec<String>, + calls: AtomicUsize, +} + +impl RelaySourceClient for CappedRelay { + fn fetch<'a>(&'a self, query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>> { + Box::pin(async move { + self.calls.fetch_add(1, AtomicOrdering::SeqCst); + let events = self + .history + .iter() + .filter(|raw| { + let event = radroots_event_codec::decode::signed_event(raw).unwrap(); + query.selector.matches(&event) + && query + .until_unix_seconds + .is_none_or(|until| event.created_at() <= until) + }) + .take(UPSTREAM_FETCH_LIMIT) + .cloned() + .collect::<Vec<_>>(); + query + .relays + .into_iter() + .map(|relay| RelayFetchBatch { + relay, + result: RelayFetchResult::Complete(events.clone()), + }) + .collect() + }) + } +} + +// Envelope IDs are canonical; signature verification remains in the shared +// ingest owner. This fixture exercises transport ordering only. +fn raw_event(id: usize, time: u64) -> String { + let pubkey = "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df"; + let content = format!("page fixture {id}"); + let canonical = serde_json::json!([0, pubkey, time, 1, [], content]).to_string(); + let event_id = hex_encode(&Sha256::digest(canonical.as_bytes())); + serde_json::json!({ + "id": event_id, "pubkey": pubkey, + "created_at": time, "kind": 1, "tags": [], "content": content, + "sig": "2".repeat(128) + }) + .to_string() +} + +fn fixture( + ties: usize, + time: u64, + older: bool, + duplicate_relay: bool, +) -> (NostrTransport, FetchRequest, Arc<CappedRelay>) { + let mut history = (1..=ties) + .rev() + .map(|id| raw_event(id, time)) + .collect::<Vec<_>>(); + history.sort_by_cached_key(|raw| { + std::cmp::Reverse( + radroots_event_codec::decode::signed_event(raw) + .unwrap() + .id_str() + .to_owned(), + ) + }); + if older { + history.push(raw_event(ties + 1, time - 1)); + } + let source = Arc::new(CappedRelay { + history, + calls: AtomicUsize::new(0), + }); + let mut urls = vec!["wss://one.example"]; + if duplicate_relay { + urls.push("wss://two.example"); + } + let config = Config::from_profile( + crate::profile::test_profile( + crate::RelayProfileKind::Public, + RelayUrlPolicy::Public, + urls.clone(), + ) + .unwrap(), + ); + let transport = NostrTransport::with_source_client(config, source.clone()); + let request = FetchRequest::new( + "capped-page", + TargetSet::new( + urls.into_iter() + .map(|url| Target::nostr_relay(url).unwrap()) + .collect(), + ) + .unwrap(), + FetchBounds::new(500, u64::MAX).unwrap(), + ) + .unwrap(); + (transport, request, source) +} + +#[tokio::test] +async fn ordinary_equal_time_pages_retain_all_peers_and_deduplicate_relays() { + let (transport, request, source) = fixture(501, 100, true, true); + let first = transport.fetch(request.clone()).await.unwrap(); + assert_eq!(first.events().len(), 500); + let NextPage::Cursor(cursor) = first.next_page() else { + panic!("next equal-time page") + }; + let second = transport + .fetch(request.with_cursor(cursor.clone())) + .await + .unwrap(); + assert_eq!(second.events().len(), 2); + assert!(matches!(second.next_page(), NextPage::Complete)); + let ids = first + .events() + .iter() + .chain(second.events()) + .map(|e| e.event().id_str()) + .collect::<BTreeSet<_>>(); + assert_eq!(ids.len(), 502); + assert_eq!(source.calls.load(AtomicOrdering::SeqCst), 2); +} + +#[tokio::test] +async fn capped_eose_yields_partial_coverage_and_explicit_older_backfill() { + for ties in [1000, 1001] { + let (transport, request, source) = fixture(ties, 100, true, true); + let first = transport.fetch(request.clone()).await.unwrap(); + assert_eq!(first.events().len(), 500); + assert!( + first + .target_outcomes() + .iter() + .all(|outcome| outcome.state() == FetchTargetState::Partial) + ); + let NextPage::Cursor(cursor) = first.next_page() else { + panic!("collected peers remain") + }; + let second = transport + .fetch(request.clone().with_cursor(cursor.clone())) + .await + .unwrap(); + assert_eq!(second.events().len(), 500); + let NextPage::Cancelled { + resume_from: Some(older), + } = second.next_page() + else { + panic!("a capped boundary must yield explicit older continuation") + }; + let received = first + .events() + .iter() + .chain(second.events()) + .map(|e| e.event().id_str()) + .collect::<BTreeSet<_>>(); + assert_eq!(received.len(), 1000); + assert_eq!(source.calls.load(AtomicOrdering::SeqCst), 2); + let third = transport + .fetch(request.with_cursor(older.clone())) + .await + .unwrap(); + assert_eq!(third.events().len(), 1); + assert_eq!(third.events()[0].event().created_at(), 99); + assert!(matches!(third.next_page(), NextPage::Complete)); + assert_eq!(source.calls.load(AtomicOrdering::SeqCst), 3); + // The 1001st same-time peer is outside the capped discovery window; + // partial evidence, not an invented completeness claim, describes it. + } +} + +#[tokio::test] +async fn zero_timestamp_yields_partial_without_fabricating_older_history() { + let (transport, request, _) = fixture(1000, 0, false, false); + let first = transport.fetch(request.clone()).await.unwrap(); + let NextPage::Cursor(cursor) = first.next_page() else { + panic!("received peers") + }; + let second = transport + .fetch(request.with_cursor(cursor.clone())) + .await + .unwrap(); + assert_eq!(second.events().len(), 500); + assert_eq!( + second.target_outcomes()[0].state(), + FetchTargetState::Partial + ); + assert!(matches!( + second.next_page(), + NextPage::Cancelled { resume_from: None } + )); +} + +#[tokio::test] +async fn older_window_rejects_changed_scope_before_access() { + let (transport, request, source) = fixture(1000, 100, false, false); + let older = window::before_boundary(100, &request_scope(&request)).unwrap(); + for selector in [ + radroots_transport::source::FetchSelector::all() + .with_kinds(vec![1]) + .unwrap(), + radroots_transport::source::FetchSelector::all() + .with_since_unix_seconds(1) + .unwrap(), + radroots_transport::source::FetchSelector::all() + .with_until_unix_seconds(100) + .unwrap(), + ] { + assert!( + transport + .fetch( + request + .clone() + .with_selector(selector) + .with_cursor(older.clone()) + ) + .await + .is_err() + ); + } + let (_, different, _) = fixture(1, 100, false, true); + assert!(transport.fetch(different.with_cursor(older)).await.is_err()); + assert_eq!(source.calls.load(AtomicOrdering::SeqCst), 0); +} + +struct UnfilteredSource(Vec<String>); +impl RelaySourceClient for UnfilteredSource { + fn fetch<'a>(&'a self, query: SourceQuery) -> BoxFuture<'a, Vec<RelayFetchBatch>> { + Box::pin(async move { + query + .relays + .into_iter() + .map(|relay| RelayFetchBatch { + relay, + result: RelayFetchResult::Complete(self.0.clone()), + }) + .collect() + }) + } +} + +#[tokio::test] +async fn malformed_or_out_of_bound_caps_cannot_fabricate_a_backward_boundary() { + for raw in ["{".to_owned(), raw_event(1, 101)] { + let (base, request, _) = fixture(1, 100, false, false); + let transport = NostrTransport::with_source_client( + base.config().clone(), + Arc::new(UnfilteredSource(vec![raw; 1000])), + ); + let cursor = window::before_boundary(101, &request_scope(&request)).unwrap(); + let page = transport.fetch(request.with_cursor(cursor)).await.unwrap(); + assert!(page.events().is_empty()); + assert_eq!(page.target_outcomes()[0].state(), FetchTargetState::Partial); + assert!(matches!( + page.next_page(), + NextPage::Cancelled { resume_from: None } + )); + } +} diff --git a/crates/transport_nostr/src/source_window.rs b/crates/transport_nostr/src/source_window.rs @@ -0,0 +1,93 @@ +use super::{Candidate, FetchCursor, RelayCursor, candidate_is_after_cursor, parse_cursor}; + +const UNTIL_PREFIX: &str = "nostr-until-v1"; + +/// Concrete-owner positions remain opaque to transport-neutral callers. +pub(super) enum Position { + After(RelayCursor), + Through(u64), +} + +impl Position { + pub(super) fn parse( + cursor: &FetchCursor, + scope: &str, + ) -> Result<Self, radroots_transport::Error> { + if !cursor.as_str().starts_with(UNTIL_PREFIX) { + return parse_cursor(cursor, scope).map(Self::After); + } + let mut parts = cursor.as_str().split(':'); + let prefix = parts.next(); + let time = parts.next(); + let actual_scope = parts.next(); + let until = time.and_then(|value| value.parse::<u64>().ok()); + match (prefix, time, actual_scope, until, parts.next()) { + (Some(UNTIL_PREFIX), Some(time), Some(actual), Some(until), None) + if actual == scope && until.to_string() == time => + { + Ok(Self::Through(until)) + } + _ => Err(radroots_transport::Error::InvalidFetchCursor), + } + } + + pub(super) fn until(&self) -> u64 { + match self { + Self::After(cursor) => cursor.created_at_unix_s(), + Self::Through(until) => *until, + } + } + + pub(super) fn includes(&self, candidate: &Candidate) -> bool { + match self { + Self::After(cursor) => candidate_is_after_cursor(candidate, cursor), + Self::Through(until) => candidate.created_at <= *until, + } + } +} + +pub(super) fn effective_until(selector: Option<u64>, position: Option<&Position>) -> Option<u64> { + match (selector, position) { + (Some(until), Some(position)) => Some(until.min(position.until())), + (Some(until), None) => Some(until), + (None, Some(position)) => Some(position.until()), + (None, None) => None, + } +} + +pub(super) fn before_boundary(boundary: u64, scope: &str) -> Option<FetchCursor> { + let until = boundary.checked_sub(1)?; + FetchCursor::parse(format!("{UNTIL_PREFIX}:{until}:{scope}")).ok() +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn older_positions_are_canonical_bounded_and_never_underflow() { + let scope = "a".repeat(64); + assert!(before_boundary(0, &scope).is_none()); + for boundary in [1, 100, u64::MAX] { + let cursor = before_boundary(boundary, &scope).unwrap(); + let position = Position::parse(&cursor, &scope).unwrap(); + assert_eq!(position.until(), boundary - 1); + assert_eq!(effective_until(None, Some(&position)), Some(boundary - 1)); + assert_eq!(effective_until(Some(0), Some(&position)), Some(0)); + } + assert_eq!(effective_until(Some(10), None), Some(10)); + assert_eq!(effective_until(None, None), None); + for value in ["", "-1", "+1", "01", "18446744073709551616"] { + let cursor = FetchCursor::parse(format!("{UNTIL_PREFIX}:{value}:{scope}")).unwrap(); + assert!(Position::parse(&cursor, &scope).is_err()); + } + for value in [ + format!("{UNTIL_PREFIX}:1:{}", "b".repeat(64)), + format!("{UNTIL_PREFIX}:1:{scope}:extra"), + format!("{UNTIL_PREFIX}:1"), + format!("{UNTIL_PREFIX}x:1:{scope}"), + ] { + assert!(Position::parse(&FetchCursor::parse(value).unwrap(), &scope).is_err()); + } + } +} diff --git a/crates/transport_nostr/tests/fetch_windows.rs b/crates/transport_nostr/tests/fetch_windows.rs @@ -0,0 +1,161 @@ +use futures::{SinkExt, StreamExt}; +use nostr_sdk::prelude::{EventBuilder, JsonUtil, Keys, Timestamp}; +use radroots_transport::{ + EventSource, FetchRequest, TargetSet, + outcome::FetchTargetState, + source::{FetchBounds, NextPage}, +}; +use radroots_transport_nostr::{ + Config, NostrTransport, RelayAccess, RelayEndpoint, RelayProfile, RelayProfileKind, + RelayUrlPolicy, +}; +use serde_json::Value; +use std::{ + collections::BTreeSet, + sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }, + time::{Duration, SystemTime, UNIX_EPOCH}, +}; +use tokio::net::TcpListener; +use tokio_tungstenite::{accept_async, tungstenite::Message}; + +#[tokio::test(flavor = "multi_thread")] +async fn real_capped_eose_preserves_received_ties_and_yields_explicit_older_history() { + let keys = + Keys::parse("0000000000000000000000000000000000000000000000000000000000000001").unwrap(); + let mut history = (0..1001) + .map(|index| { + let event = EventBuilder::text_note(format!("equal-time {index}")) + .custom_created_at(Timestamp::from_secs(100)) + .sign_with_keys(&keys) + .unwrap(); + serde_json::from_str::<Value>(&event.as_json()).unwrap() + }) + .collect::<Vec<_>>(); + history.sort_by(|left, right| right["id"].as_str().cmp(&left["id"].as_str())); + let expected = history[..1000] + .iter() + .map(|event| event["id"].as_str().unwrap().to_owned()) + .collect::<BTreeSet<_>>(); + let older = EventBuilder::text_note("older") + .custom_created_at(Timestamp::from_secs(99)) + .sign_with_keys(&keys) + .unwrap(); + history.push(serde_json::from_str(&older.as_json()).unwrap()); + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let url = format!("ws://{}", listener.local_addr().unwrap()); + let calls = Arc::new(AtomicUsize::new(0)); + let server_calls = calls.clone(); + let server = tokio::spawn(async move { + let (stream, _) = listener.accept().await.unwrap(); + let mut socket = accept_async(stream).await.unwrap(); + while let Some(Ok(message)) = socket.next().await { + let Message::Text(message) = message else { + continue; + }; + let request: Value = serde_json::from_str(&message).unwrap(); + if request[0] != "REQ" { + continue; + } + server_calls.fetch_add(1, Ordering::SeqCst); + assert_eq!(request[2]["limit"], 1000); + let until = request[2]["until"].as_u64().unwrap_or(u64::MAX); + for event in history + .iter() + .filter(|event| event["created_at"].as_u64().unwrap() <= until) + .take(1000) + { + socket + .send(Message::Text( + serde_json::to_string(&("EVENT", &request[1], event)) + .unwrap() + .into(), + )) + .await + .unwrap(); + } + socket + .send(Message::Text( + serde_json::to_string(&("EOSE", &request[1])) + .unwrap() + .into(), + )) + .await + .unwrap(); + } + }); + let profile = RelayProfile::explicit( + RelayProfileKind::Simulator, + [RelayEndpoint::new(&url, RelayUrlPolicy::Local, RelayAccess::ReadOnly).unwrap()], + ) + .unwrap(); + let config = Config::from_profile(profile) + .with_timeouts(5000, 5000, 500) + .unwrap(); + let targets = TargetSet::new( + config + .read_relays() + .map(|relay| relay.to_target().unwrap()) + .collect(), + ) + .unwrap(); + let transport = NostrTransport::new(config); + let deadline = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap() + .as_millis() as u64 + + 20000; + let request = FetchRequest::new( + "live-capped-window", + targets, + FetchBounds::new(500, deadline).unwrap(), + ) + .unwrap(); + tokio::time::timeout(Duration::from_secs(15), async { + let first = transport.fetch(request.clone()).await.unwrap(); + assert_eq!(first.events().len(), 500); + assert_eq!( + first.target_outcomes()[0].state(), + FetchTargetState::Partial + ); + let NextPage::Cursor(cursor) = first.next_page() else { + panic!("received peers remain") + }; + let second = transport + .fetch(request.clone().with_cursor(cursor.clone())) + .await + .unwrap(); + assert_eq!(second.events().len(), 500); + assert_eq!( + second.target_outcomes()[0].state(), + FetchTargetState::Partial + ); + assert_eq!(calls.load(Ordering::SeqCst), 2); + let received = first + .events() + .iter() + .chain(second.events()) + .map(|e| e.event().id_str().to_owned()) + .collect::<BTreeSet<_>>(); + assert_eq!(received, expected); + let NextPage::Cancelled { + resume_from: Some(older), + } = second.next_page() + else { + panic!("explicit partial window yield") + }; + let third = transport + .fetch(request.with_cursor(older.clone())) + .await + .unwrap(); + assert_eq!(third.events().len(), 1); + assert_eq!(third.events()[0].event().created_at(), 99); + assert!(matches!(third.next_page(), NextPage::Complete)); + assert_eq!(calls.load(Ordering::SeqCst), 3); + }) + .await + .unwrap(); + server.abort(); +} diff --git a/crates/transport_nostr/tests/package_boundary.rs b/crates/transport_nostr/tests/package_boundary.rs @@ -321,6 +321,8 @@ fn adapter_owns_no_storage_outbox_or_orchestration_surface() { "sink.rs".to_owned(), "source.rs".to_owned(), "source_budget.rs".to_owned(), + "source_paging_tests.rs".to_owned(), + "source_window.rs".to_owned(), "status.rs".to_owned(), "subscription.rs".to_owned(), ])