commit 773bc8afb53eae5efb44506e892f62a89b275fc4
parent a7d0e471bd32d891248762b8f6d8d6428fa16f9d
Author: triesap <tyson@radroots.org>
Date: Wed, 9 Sep 2026 22:22:10 +0000
transport_nostr: Bound shared fetch work and queued deadlines
- Apply one absolute deadline to all relay scheduling and collection
- Cap wire messages and aggregate event bytes and parsing inventory
- Preserve scoped partial outcomes and already collected evidence
- Verify loopback regressions, public API, coverage and workspace preflight
Diffstat:
9 files changed, 707 insertions(+), 38 deletions(-)
diff --git a/contracts/architecture/decisions/nostr_fetch_bounds.v1.json b/contracts/architecture/decisions/nostr_fetch_bounds.v1.json
@@ -0,0 +1,24 @@
+{
+ "schema": "radroots.nostr-fetch-bounds.v1",
+ "owner": "radroots_transport_nostr",
+ "status": "implemented",
+ "request_scope": "One EventSource::fetch call across every selected relay, including queued batches.",
+ "max_relays": 64,
+ "default_max_connections": 8,
+ "max_filters_per_relay": 1,
+ "max_wire_message_bytes": 524288,
+ "max_wire_frame_bytes": 524288,
+ "max_event_json_bytes": 262144,
+ "max_events_per_relay": 1000,
+ "max_events_per_fetch": 4096,
+ "max_event_json_bytes_per_fetch": 8388608,
+ "max_notifications_per_fetch": 8192,
+ "max_returned_events": 1000,
+ "deadline": "The minimum of the caller absolute deadline and configured request timeout is frozen once before relay scheduling. Connect, REQ, collection and queued starts consume its remaining duration; no later relay receives a fresh request timeout.",
+ "collection": "Charge event inventory and JSON bytes across all relays before retaining raw events; duplicates and malformed events consume work. Defensively enforce the same limits before shared event decoding and candidate collection.",
+ "parse_boundary": "The WebSocket message/frame bound precedes upstream message decoding. The shared event-JSON bound and aggregate inventory precede canonical decoding. Notification processing is finite; upstream socket tasks remain subject to bounded messages and request auto-close.",
+ "completion": "EOSE confirms only the selected relay subscription within the request budget. Exhausted resources remain Partial, elapsed deadlines remain Cancelled, and neither implies global-history completeness.",
+ "partial_finalization": "After bounded network collection, bounded local normalization retains previously collected admissible events and each relay outcome. One slow relay must not discard another relay's earlier evidence. No further network request is started during finalization.",
+ "cancellation": "Unpolled futures perform no I/O. Dropping polled work cannot await cleanup; published subscriptions retain the original bounded auto-close deadline. No durable admission or publication rollback is claimed.",
+ "compatibility": "No new public Rust type, feature, dependency, transport or product policy is required. Callers retain bounded FetchRequest and existing typed outcomes; large or continuous results can now terminate earlier as partial evidence."
+}
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-013"
+date = "2026-09-09"
+status = "closed"
+approval = "Standing user approval of the complete reviewed refactor and its necessary same-owner producer repairs."
+affected_steps = ["201"]
+spec_anchors = [
+ "contracts/crates/release_v1/radroots_crates_release_v1.toml#package.radroots_transport_nostr",
+]
+source_evidence = [
+ "Current source.rs starts each buffered relay batch with a fresh timeout although the public fetch contract promises one absolute request deadline.",
+ "Per-relay event counts are bounded, but aggregate raw JSON bytes and parse inventory are not explicitly capped before shared candidate collection; relay.rs inherits dependency-default wire limits.",
+]
+replacement_action = "Repair the current concrete adapter within its existing bounded-fetch charter, using contracts/architecture/decisions/nostr_fetch_bounds.v1.json. This is a current source repair, not a reopening or replacement of historical release qualification."
+verification = [
+ "Deterministic limit, aggregate competition, queued-deadline and retained partial-outcome regressions plus bounded loopback source tests.",
+ "Canonical package/workspace checks, contracts, architecture, public API review and the unchanged affected-package coverage requirement.",
+]
+unresolved_risk = "Current package and workspace verification, byte-identical public API, all 49 required coverage reports and release preflight pass on the qualified source. Per-target partial/cancelled evidence is bounded; no global-history, power-loss, external release, device or new-platform qualification is claimed."
+normative_architecture_change = false
+adr_required = false
+closure_evidence = [
+ "Fresh transport coverage passes all four unchanged thresholds; the public API is byte-identical to its baseline.",
+ "All 49 required coverage reports and the aggregate pass, together with workspace check, tests, Clippy, Rustdoc, contracts, freshness, dependency graph, portable checks and release preflight.",
+ "Deterministic maxima and aggregate budget tests plus real loopback queued-deadline, mixed-relay, oversized-wire and cancellation regressions pass.",
+]
+
+[[deviation]]
id = "RCRV1-DEV-014"
date = "2026-09-09"
status = "closed"
diff --git a/crates/transport_nostr/README.md b/crates/transport_nostr/README.md
@@ -144,6 +144,16 @@ strictly before its absolute deadline. Deadline expiry is `Cancelled`, and a
relay result that exceeds the bounded inventory is `Partial`; neither state is
rewritten as completion even when it carries admissible events.
+One fetch shares an 8 MiB raw event-JSON budget, a 4,096-event inventory and an
+8,192-notification work limit across all relay batches. Each relay retains at
+most 1,000 events; duplicate and malformed observations still consume the
+budget. An event is at most 256 KiB, and the WebSocket connector explicitly
+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.
+
Live subscriptions use the same explicit readable targets and selector
translation. A caller checkpoint is scoped to one exact target and selector;
the adapter reconnects with Nostr's inclusive `since` timestamp, suppresses
@@ -175,6 +185,10 @@ or delivery work.
The absolute deadline in each generic request bounds the complete operation.
The configured connection and request timeouts are upper bounds within that
remaining budget. An already-expired request performs no relay work.
+Queued relay batches consume that same frozen deadline; they never receive a
+new timeout after an earlier relay stalls. Bounded local normalization retains
+the events and distinct outcomes already collected when network work ends, so
+one timed-out relay cannot erase another relay's earlier successful evidence.
Dropping an unpolled fetch, subscription-start, or delivery future performs no
I/O. Once polled, cancellation is best effort at the socket boundary. For
diff --git a/crates/transport_nostr/src/relay.rs b/crates/transport_nostr/src/relay.rs
@@ -19,6 +19,7 @@ use tokio::net::TcpStream;
use url::Url;
const MAX_RESOLVED_ADDRESSES: usize = 32;
+const MAX_WIRE_MESSAGE_BYTES: usize = 512 * 1024;
/// Validated canonical Nostr relay URL.
#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
@@ -153,9 +154,17 @@ impl WebSocketTransport for HardenedWebsocketTransport {
.validate_resolved_addresses(policy, addresses.iter().map(SocketAddr::ip))
.map_err(|_| policy_error("relay DNS result is denied by network policy"))?;
let tcp = connect_pinned(addresses.as_slice()).await?;
- let (stream, _) = tokio_tungstenite::client_async_tls(relay.as_str(), tcp)
- .await
- .map_err(TransportError::backend)?;
+ let config = tokio_tungstenite::tungstenite::protocol::WebSocketConfig::default()
+ .max_message_size(Some(MAX_WIRE_MESSAGE_BYTES))
+ .max_frame_size(Some(MAX_WIRE_MESSAGE_BYTES));
+ let (stream, _) = tokio_tungstenite::client_async_tls_with_config(
+ relay.as_str(),
+ tcp,
+ Some(config),
+ None,
+ )
+ .await
+ .map_err(TransportError::backend)?;
let socket = WebSocket::Tokio(stream);
let (tx, rx) = socket.split();
let sink: WebSocketSink = Box::new(HardenedTransportSink(tx));
diff --git a/crates/transport_nostr/src/source.rs b/crates/transport_nostr/src/source.rs
@@ -15,8 +15,13 @@ use radroots_transport::{
};
use sha2::{Digest, Sha256};
use std::collections::{BTreeMap, BTreeSet};
+use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
+#[path = "source_budget.rs"]
+mod budget;
+use budget::FetchBudget;
+
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";
@@ -27,7 +32,7 @@ pub(crate) struct SourceQuery {
selector: radroots_transport::source::FetchSelector,
until_unix_seconds: Option<u64>,
connect_timeout: Duration,
- timeout: Duration,
+ deadline: tokio::time::Instant,
max_connections: usize,
}
@@ -78,15 +83,19 @@ impl RelaySourceClient for LiveRelaySourceClient {
selector,
until_unix_seconds,
connect_timeout,
- timeout,
+ deadline,
max_connections,
} = query;
+ let budget = Arc::new(FetchBudget::default());
stream::iter(relays.into_iter().map(|relay| {
let selector = selector.clone();
+ let budget = Arc::clone(&budget);
async move {
let url = relay.as_str().to_owned();
- let started_at = tokio::time::Instant::now();
let result = async {
+ if tokio::time::Instant::now() >= deadline {
+ return Ok(RelayFetchResult::Timeout(Vec::new()));
+ }
let kinds = selector
.kinds()
.iter()
@@ -122,17 +131,26 @@ impl RelaySourceClient for LiveRelaySourceClient {
if let Some(until) = until_unix_seconds {
filter = filter.until(Timestamp::from_secs(until));
}
- self.client
- .add_relay(url.as_str())
- .await
- .map_err(|error| error.to_string())?;
- self.client
- .try_connect_relay(url.as_str(), connect_timeout)
- .await
- .map_err(|error| error.to_string())?;
- let remaining = timeout.saturating_sub(started_at.elapsed());
- if remaining.is_zero() {
- return Ok(RelayFetchResult::Timeout(Vec::new()));
+ let connected = tokio::time::timeout_at(deadline, async {
+ self.client
+ .add_relay(url.as_str())
+ .await
+ .map_err(|error| error.to_string())?;
+ self.client
+ .try_connect_relay(
+ url.as_str(),
+ connect_timeout.min(
+ deadline
+ .saturating_duration_since(tokio::time::Instant::now()),
+ ),
+ )
+ .await
+ .map_err(|error| error.to_string())
+ })
+ .await;
+ match connected {
+ Ok(result) => result?,
+ Err(_) => return Ok(RelayFetchResult::Timeout(Vec::new())),
}
let relay = self
.client
@@ -141,20 +159,34 @@ impl RelaySourceClient for LiveRelaySourceClient {
.map_err(|error| error.to_string())?;
let mut notifications = relay.notifications();
let subscription_id = SubscriptionId::generate();
- let eose_deadline = tokio::time::Instant::now() + remaining;
+ let remaining =
+ deadline.saturating_duration_since(tokio::time::Instant::now());
+ if remaining.is_zero() {
+ return Ok(RelayFetchResult::Timeout(Vec::new()));
+ }
let options = SubscribeOptions::default().close_on(Some(
SubscribeAutoCloseOptions::default()
.exit_policy(ReqExitPolicy::ExitOnEOSE)
.timeout(Some(remaining)),
));
- relay
- .subscribe_with_id(subscription_id.clone(), filter, options)
- .await
- .map_err(|error| error.to_string())?;
- Ok(
- collect_until_eose(&mut notifications, &subscription_id, eose_deadline)
- .await,
+ match tokio::time::timeout_at(
+ deadline,
+ relay.subscribe_with_id(subscription_id.clone(), filter, options),
)
+ .await
+ {
+ Ok(result) => {
+ result.map_err(|error| error.to_string())?;
+ }
+ Err(_) => return Ok(RelayFetchResult::Timeout(Vec::new())),
+ }
+ Ok(collect_until_eose(
+ &mut notifications,
+ &subscription_id,
+ deadline,
+ &budget,
+ )
+ .await)
}
.await;
RelayFetchBatch {
@@ -174,6 +206,7 @@ async fn collect_until_eose(
notifications: &mut tokio::sync::broadcast::Receiver<RelayNotification>,
subscription_id: &SubscriptionId,
deadline: tokio::time::Instant,
+ budget: &FetchBudget,
) -> RelayFetchResult {
let mut events = Vec::new();
loop {
@@ -185,6 +218,9 @@ async fn collect_until_eose(
Ok(Err(error)) => return RelayFetchResult::Failed(error.to_string()),
Err(_) => return RelayFetchResult::Timeout(events),
};
+ if !budget.notification() {
+ return RelayFetchResult::ResourceLimit(events);
+ }
match notification {
RelayNotification::Message {
message:
@@ -196,7 +232,11 @@ async fn collect_until_eose(
if events.len() >= UPSTREAM_FETCH_LIMIT {
return RelayFetchResult::ResourceLimit(events);
}
- events.push(event.as_ref().as_json());
+ let raw = event.as_ref().as_json();
+ if !budget.event(raw.len()) {
+ return RelayFetchResult::ResourceLimit(events);
+ }
+ events.push(raw);
}
RelayNotification::Message {
message: RelayMessage::EndOfStoredEvents(observed_subscription),
@@ -262,6 +302,7 @@ impl EventSource for NostrTransport {
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());
+ let deadline = tokio::time::Instant::now() + Duration::from_millis(timeout_ms);
let mut targets = BTreeMap::new();
let mut outcomes = Vec::new();
@@ -322,11 +363,12 @@ impl EventSource for NostrTransport {
connect_timeout: Duration::from_millis(
timeout_ms.min(self.config().connect_timeout_ms()),
),
- timeout: Duration::from_millis(timeout_ms),
+ deadline,
max_connections: self.config().max_connections(),
})
.await;
let mut candidates = Vec::new();
+ let parse_budget = FetchBudget::default();
let mut malformed_by_relay = BTreeMap::<RelayUrl, usize>::new();
let mut reported = BTreeSet::new();
let observed_at_unix_ms = unix_time_ms().max(now_ms);
@@ -338,7 +380,7 @@ impl EventSource for NostrTransport {
if !reported.insert(relay.clone()) {
return Err(radroots_transport::Error::DuplicateFetchTargetOutcome);
}
- let (raw_events, terminal) = match result {
+ let (raw_events, mut terminal) = match result {
RelayFetchResult::Complete(events) => (events, FetchTargetState::Complete),
RelayFetchResult::Timeout(events) => (events, FetchTargetState::Cancelled),
RelayFetchResult::ResourceLimit(events) => (events, FetchTargetState::Partial),
@@ -357,13 +399,13 @@ impl EventSource for NostrTransport {
continue;
}
};
- self.status.record_read(
- &relay,
- terminal == FetchTargetState::Complete,
- terminal != FetchTargetState::Complete,
- observed_at_unix_ms,
- );
- for raw in raw_events {
+ for (index, raw) in raw_events.into_iter().enumerate() {
+ if index >= UPSTREAM_FETCH_LIMIT || !parse_budget.event(raw.len()) {
+ if terminal != FetchTargetState::Cancelled {
+ terminal = FetchTargetState::Partial;
+ }
+ break;
+ }
match radroots_event_codec::decode::signed_event(raw.as_str()) {
Ok(event) if request.selector().matches(&event) => {
candidates.push(Candidate {
@@ -379,6 +421,12 @@ impl EventSource for NostrTransport {
}
}
}
+ self.status.record_read(
+ &relay,
+ terminal == FetchTargetState::Complete,
+ terminal != FetchTargetState::Complete,
+ observed_at_unix_ms,
+ );
let malformed = malformed_by_relay.get(&relay).copied().unwrap_or_default();
let outcome = if terminal == FetchTargetState::Cancelled {
FetchTargetOutcome::new(
@@ -883,6 +931,7 @@ mod tests {
&mut receiver,
&subscription_id,
tokio::time::Instant::now() + Duration::from_secs(1),
+ &FetchBudget::default(),
)
.await,
RelayFetchResult::Complete(Vec::new())
@@ -890,7 +939,13 @@ mod tests {
let (_sender, mut receiver) = tokio::sync::broadcast::channel(2);
assert_eq!(
- collect_until_eose(&mut receiver, &subscription_id, tokio::time::Instant::now(),).await,
+ collect_until_eose(
+ &mut receiver,
+ &subscription_id,
+ tokio::time::Instant::now(),
+ &FetchBudget::default()
+ )
+ .await,
RelayFetchResult::Timeout(Vec::new())
);
}
@@ -1048,10 +1103,204 @@ mod tests {
selector,
until_unix_seconds: None,
connect_timeout: Duration::from_millis(1),
- timeout: Duration::from_millis(1),
+ deadline: tokio::time::Instant::now() + Duration::from_millis(1),
max_connections: 1,
}));
assert_eq!(batches.len(), 1);
assert_eq!(batches[0].result, RelayFetchResult::Complete(vec![]));
}
+
+ #[tokio::test]
+ async fn continuous_notifications_terminate_at_the_shared_work_limit() {
+ let id = SubscriptionId::generate();
+ let (sender, mut receiver) =
+ tokio::sync::broadcast::channel(budget::MAX_FETCH_NOTIFICATIONS + 1);
+ for _ in 0..=budget::MAX_FETCH_NOTIFICATIONS {
+ sender
+ .send(RelayNotification::RelayStatus {
+ status: nostr_sdk::prelude::RelayStatus::Connected,
+ })
+ .unwrap();
+ }
+ assert_eq!(
+ collect_until_eose(
+ &mut receiver,
+ &id,
+ tokio::time::Instant::now() + Duration::from_secs(10),
+ &FetchBudget::default()
+ )
+ .await,
+ RelayFetchResult::ResourceLimit(Vec::new())
+ );
+ }
+
+ #[tokio::test]
+ async fn repeated_events_consume_inventory_before_deduplication() {
+ let id = SubscriptionId::generate();
+ let event = nostr_sdk::prelude::Event::from_json(FIRST).unwrap();
+ let (sender, mut receiver) = tokio::sync::broadcast::channel(UPSTREAM_FETCH_LIMIT + 1);
+ for _ in 0..=UPSTREAM_FETCH_LIMIT {
+ sender
+ .send(RelayNotification::Message {
+ message: RelayMessage::Event {
+ subscription_id: Cow::Owned(id.clone()),
+ event: Cow::Owned(event.clone()),
+ },
+ })
+ .unwrap();
+ }
+ let result = collect_until_eose(
+ &mut receiver,
+ &id,
+ tokio::time::Instant::now() + Duration::from_secs(10),
+ &FetchBudget::default(),
+ )
+ .await;
+ let RelayFetchResult::ResourceLimit(events) = result else {
+ panic!("bounded inventory");
+ };
+ assert_eq!(events.len(), UPSTREAM_FETCH_LIMIT);
+ }
+
+ #[test]
+ fn defensive_parse_budget_preserves_earlier_events_and_refuses_excess_work() {
+ let relay = RelayUrl::parse("wss://one.example", RelayUrlPolicy::Public).unwrap();
+ let mut records = vec![FIRST.to_owned()];
+ records.extend((1..UPSTREAM_FETCH_LIMIT).map(|_| "{".to_owned()));
+ records.push(SECOND.to_owned());
+ for records in [
+ records,
+ vec![
+ FIRST.to_owned(),
+ "x".repeat(budget::MAX_EVENT_BYTES + 1),
+ SECOND.to_owned(),
+ ],
+ ] {
+ let page = futures::executor::block_on(
+ scripted(vec![RelayFetchBatch {
+ relay: relay.clone(),
+ result: RelayFetchResult::Complete(records),
+ }])
+ .fetch(single_request(10)),
+ )
+ .unwrap();
+ assert_eq!(page.events().len(), 1);
+ assert_eq!(
+ page.events()[0].event().id_str(),
+ radroots_event_codec::decode::signed_event(FIRST)
+ .unwrap()
+ .id_str()
+ );
+ assert_eq!(page.target_outcomes()[0].state(), FetchTargetState::Partial);
+ }
+ }
+
+ #[tokio::test]
+ async fn completed_relay_batches_do_not_refund_the_shared_byte_budget() {
+ let mut wire: serde_json::Value = serde_json::from_str(FIRST).unwrap();
+ wire["content"] = serde_json::json!("");
+ let empty = nostr_sdk::prelude::Event::from_json(wire.to_string()).unwrap();
+ wire["content"] =
+ serde_json::json!("x".repeat(budget::MAX_EVENT_BYTES - empty.as_json().len()));
+ // Collection bounds precede canonical event admission; the fixture only
+ // needs the upstream event structure and an exact serialized size.
+ let event = nostr_sdk::prelude::Event::from_json(wire.to_string()).unwrap();
+ assert_eq!(event.as_json().len(), budget::MAX_EVENT_BYTES);
+ let shared = FetchBudget::default();
+ let mut retained_bytes = 0;
+ for batch in 0..2 {
+ let id = SubscriptionId::generate();
+ let (sender, mut receiver) = tokio::sync::broadcast::channel(18);
+ for _ in 0..17 {
+ sender
+ .send(RelayNotification::Message {
+ message: RelayMessage::Event {
+ subscription_id: Cow::Owned(id.clone()),
+ event: Cow::Owned(event.clone()),
+ },
+ })
+ .unwrap();
+ }
+ sender
+ .send(RelayNotification::Message {
+ message: RelayMessage::EndOfStoredEvents(Cow::Owned(id.clone())),
+ })
+ .unwrap();
+ let result = collect_until_eose(
+ &mut receiver,
+ &id,
+ tokio::time::Instant::now() + Duration::from_secs(10),
+ &shared,
+ )
+ .await;
+ let events = match (batch, result) {
+ (0, RelayFetchResult::Complete(events)) => {
+ assert_eq!(events.len(), 17);
+ events
+ }
+ (1, RelayFetchResult::ResourceLimit(events)) => {
+ assert_eq!(events.len(), 15);
+ events
+ }
+ _ => panic!("the second relay must exhaust the shared byte budget"),
+ };
+ retained_bytes += events.iter().map(String::len).sum::<usize>();
+ }
+ assert_eq!(retained_bytes, budget::MAX_FETCH_BYTES);
+ assert!(!shared.event(1));
+ }
+
+ #[test]
+ fn duplicate_inventory_is_bounded_across_all_relay_batches_before_deduplication() {
+ let urls = (0..5)
+ .map(|index| format!("wss://relay{index}.example"))
+ .collect::<Vec<_>>();
+ let config = Config::from_profile(
+ crate::profile::test_profile(
+ crate::RelayProfileKind::Public,
+ RelayUrlPolicy::Public,
+ urls.iter().map(String::as_str),
+ )
+ .unwrap(),
+ );
+ let targets = TargetSet::new(
+ config
+ .read_relays()
+ .map(|relay| relay.to_target().unwrap())
+ .collect(),
+ )
+ .unwrap();
+ let request = FetchRequest::new(
+ "aggregate-duplicates",
+ targets,
+ FetchBounds::new(10, unix_time_ms() + 10_000).unwrap(),
+ )
+ .unwrap();
+ let batches = urls
+ .iter()
+ .map(|url| RelayFetchBatch {
+ relay: RelayUrl::parse(url, RelayUrlPolicy::Public).unwrap(),
+ result: RelayFetchResult::Complete(vec![FIRST.to_owned(); UPSTREAM_FETCH_LIMIT]),
+ })
+ .collect();
+ let transport =
+ NostrTransport::with_source_client(config, Arc::new(ScriptedSourceClient(batches)));
+ let page = futures::executor::block_on(transport.fetch(request)).unwrap();
+ assert_eq!(page.events().len(), 1);
+ assert_eq!(page.target_outcomes().len(), 5);
+ assert_eq!(
+ page.target_outcomes()
+ .iter()
+ .filter(|outcome| outcome.state() == FetchTargetState::Complete)
+ .count(),
+ 4
+ );
+ assert_eq!(
+ page.target_outcomes()
+ .iter()
+ .filter(|outcome| outcome.state() == FetchTargetState::Partial)
+ .count(),
+ 1
+ );
+ }
}
diff --git a/crates/transport_nostr/src/source_budget.rs b/crates/transport_nostr/src/source_budget.rs
@@ -0,0 +1,87 @@
+use std::sync::Mutex;
+
+pub(super) const MAX_EVENT_BYTES: usize = radroots_event_codec::decode::MAX_EVENT_JSON_BYTES;
+pub(super) const MAX_FETCH_BYTES: usize = 8 * 1024 * 1024;
+pub(super) const MAX_FETCH_EVENTS: usize = 4096;
+pub(super) const MAX_FETCH_NOTIFICATIONS: usize = 8192;
+
+#[derive(Debug, Default)]
+struct Usage {
+ bytes: usize,
+ events: usize,
+ notifications: usize,
+}
+
+/// One monotonic inventory shared by all relay batches. Reservations are not
+/// refunded when a duplicate, malformed event or completed batch is discarded.
+#[derive(Debug, Default)]
+pub(super) struct FetchBudget(Mutex<Usage>);
+
+impl FetchBudget {
+ pub(super) fn notification(&self) -> bool {
+ let Ok(mut usage) = self.0.lock() else {
+ return false;
+ };
+ if usage.notifications == MAX_FETCH_NOTIFICATIONS {
+ return false;
+ }
+ usage.notifications += 1;
+ true
+ }
+
+ pub(super) fn event(&self, bytes: usize) -> bool {
+ if bytes > MAX_EVENT_BYTES {
+ return false;
+ }
+ let Ok(mut usage) = self.0.lock() else {
+ return false;
+ };
+ if usage.events == MAX_FETCH_EVENTS || bytes > MAX_FETCH_BYTES - usage.bytes {
+ return false;
+ }
+ usage.events += 1;
+ usage.bytes += bytes;
+ true
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ #[test]
+ fn each_maximum_is_accepted_and_the_next_unit_is_rejected() {
+ let bytes = FetchBudget::default();
+ assert!(!bytes.event(MAX_EVENT_BYTES + 1));
+ for _ in 0..MAX_FETCH_BYTES / MAX_EVENT_BYTES {
+ assert!(bytes.event(MAX_EVENT_BYTES));
+ }
+ assert!(!bytes.event(1));
+ let events = FetchBudget::default();
+ for _ in 0..MAX_FETCH_EVENTS {
+ assert!(events.event(1));
+ }
+ assert!(!events.event(1));
+ let notifications = FetchBudget::default();
+ for _ in 0..MAX_FETCH_NOTIFICATIONS {
+ assert!(notifications.notification());
+ }
+ assert!(!notifications.notification());
+ }
+
+ #[test]
+ fn competing_relays_cannot_overreserve_the_aggregate_budget() {
+ let budget = FetchBudget::default();
+ let accepted = std::thread::scope(|scope| {
+ let tasks = (0..8)
+ .map(|_| scope.spawn(|| (0..32).filter(|_| budget.event(MAX_EVENT_BYTES)).count()))
+ .collect::<Vec<_>>();
+ tasks
+ .into_iter()
+ .map(|task| task.join().unwrap())
+ .sum::<usize>()
+ });
+ assert_eq!(accepted * MAX_EVENT_BYTES, MAX_FETCH_BYTES);
+ assert!(!budget.event(1));
+ }
+}
diff --git a/crates/transport_nostr/tests/fetch_bounds.rs b/crates/transport_nostr/tests/fetch_bounds.rs
@@ -0,0 +1,256 @@
+use futures::{SinkExt, StreamExt};
+use nostr_sdk::prelude::{EventBuilder, JsonUtil, Keys};
+use radroots_transport::{
+ EventSource, FetchRequest, TargetSet, outcome::FetchTargetState, source::FetchBounds,
+};
+use radroots_transport_nostr::{
+ Config, NostrTransport, RelayAccess, RelayEndpoint, RelayProfile, RelayProfileKind,
+ RelayUrlPolicy,
+};
+use serde_json::Value;
+use std::{
+ sync::{
+ Arc,
+ atomic::{AtomicUsize, Ordering},
+ },
+ time::{Duration, SystemTime, UNIX_EPOCH},
+};
+use tokio::{net::TcpListener, sync::oneshot};
+use tokio_tungstenite::{accept_async, tungstenite::Message};
+
+fn transport(urls: &[String], connections: usize) -> (NostrTransport, FetchRequest) {
+ let endpoints = urls
+ .iter()
+ .map(|url| RelayEndpoint::new(url, RelayUrlPolicy::Local, RelayAccess::ReadOnly).unwrap());
+ let profile = RelayProfile::explicit(RelayProfileKind::Simulator, endpoints).unwrap();
+ let config = Config::from_profile(profile)
+ .with_timeouts(1000, 1000, 500)
+ .unwrap()
+ .with_max_connections(connections)
+ .unwrap();
+ let targets = TargetSet::new(
+ config
+ .read_relays()
+ .map(|relay| relay.to_target().unwrap())
+ .collect(),
+ )
+ .unwrap();
+ let now = SystemTime::now()
+ .duration_since(UNIX_EPOCH)
+ .unwrap()
+ .as_millis() as u64;
+ let request = FetchRequest::new(
+ "bounded-loopback",
+ targets,
+ FetchBounds::new(10, now + 5000).unwrap(),
+ )
+ .unwrap();
+ (NostrTransport::new(config), request)
+}
+
+async fn listener() -> (TcpListener, String) {
+ let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
+ let url = format!("ws://{}", listener.local_addr().unwrap());
+ (listener, url)
+}
+
+fn event(content: &str) -> Value {
+ let keys =
+ Keys::parse("0000000000000000000000000000000000000000000000000000000000000001").unwrap();
+ let event = EventBuilder::text_note(content)
+ .sign_with_keys(&keys)
+ .unwrap();
+ serde_json::from_str(&event.as_json()).unwrap()
+}
+
+async fn serve(
+ listener: TcpListener,
+ payload: Option<Value>,
+ eose: bool,
+ requests: Arc<AtomicUsize>,
+) {
+ let (stream, _) = listener.accept().await.unwrap();
+ let mut socket = accept_async(stream).await.unwrap();
+ while let Some(message) = socket.next().await {
+ let Ok(Message::Text(message)) = message else {
+ continue;
+ };
+ let values: Value = serde_json::from_str(&message).unwrap();
+ if values[0] != "REQ" {
+ continue;
+ }
+ requests.fetch_add(1, Ordering::SeqCst);
+ if let Some(payload) = &payload {
+ socket
+ .send(Message::Text(
+ serde_json::to_string(&("EVENT", &values[1], payload))
+ .unwrap()
+ .into(),
+ ))
+ .await
+ .unwrap();
+ }
+ if eose {
+ socket
+ .send(Message::Text(
+ serde_json::to_string(&("EOSE", &values[1])).unwrap().into(),
+ ))
+ .await
+ .unwrap();
+ }
+ std::future::pending::<()>().await;
+ }
+}
+
+#[tokio::test(flavor = "multi_thread")]
+async fn queued_relays_do_not_receive_a_fresh_timeout_after_a_stalled_relay() {
+ let (first, first_url) = listener().await;
+ let (second, second_url) = listener().await;
+ // Adapter scheduling is canonical URL order, independent of profile order.
+ let mut relays = [(first_url, first), (second_url, second)];
+ relays.sort_by(|left, right| left.0.cmp(&right.0));
+ let [(first_url, first), (second_url, second)] = relays;
+ let first_requests = Arc::new(AtomicUsize::new(0));
+ let second_requests = Arc::new(AtomicUsize::new(0));
+ let first_task = tokio::spawn(serve(first, None, false, Arc::clone(&first_requests)));
+ let second_task = tokio::spawn(serve(second, None, true, Arc::clone(&second_requests)));
+ let (transport, request) = transport(&[first_url, second_url], 1);
+ let page = tokio::time::timeout(Duration::from_secs(5), transport.fetch(request))
+ .await
+ .unwrap()
+ .unwrap();
+ first_task.abort();
+ second_task.abort();
+ assert_eq!(first_requests.load(Ordering::SeqCst), 1);
+ assert_eq!(second_requests.load(Ordering::SeqCst), 0);
+ assert_eq!(page.target_outcomes().len(), 2);
+ assert!(
+ page.target_outcomes()
+ .iter()
+ .all(|outcome| outcome.state() == FetchTargetState::Cancelled)
+ );
+}
+
+#[tokio::test(flavor = "multi_thread")]
+async fn a_slow_relay_preserves_the_other_relays_completed_evidence_and_collected_events() {
+ let (good, good_url) = listener().await;
+ let (slow, slow_url) = listener().await;
+ let good_event = event("complete relay");
+ let slow_event = event("partial relay");
+ let expected = [
+ good_event["id"].as_str().unwrap().to_owned(),
+ slow_event["id"].as_str().unwrap().to_owned(),
+ ];
+ let calls = Arc::new(AtomicUsize::new(0));
+ let good_task = tokio::spawn(serve(good, Some(good_event), true, Arc::clone(&calls)));
+ let slow_task = tokio::spawn(serve(slow, Some(slow_event), false, Arc::clone(&calls)));
+ let (transport, request) = transport(&[good_url, slow_url], 2);
+ let page = tokio::time::timeout(Duration::from_secs(5), transport.fetch(request))
+ .await
+ .unwrap()
+ .unwrap();
+ good_task.abort();
+ slow_task.abort();
+ assert_eq!(calls.load(Ordering::SeqCst), 2);
+ assert_eq!(page.events().len(), 2);
+ for id in expected {
+ assert!(
+ page.events()
+ .iter()
+ .any(|event| event.event().id_str() == id)
+ );
+ }
+ assert_eq!(
+ page.target_outcomes()
+ .iter()
+ .filter(|outcome| outcome.state() == FetchTargetState::Complete)
+ .count(),
+ 1
+ );
+ assert_eq!(
+ page.target_outcomes()
+ .iter()
+ .filter(|outcome| outcome.state() == FetchTargetState::Cancelled)
+ .count(),
+ 1
+ );
+}
+
+#[tokio::test(flavor = "multi_thread")]
+async fn oversized_wire_messages_never_become_completed_fetch_evidence() {
+ let (listener, url) = listener().await;
+ let server = tokio::spawn(async move {
+ let (stream, _) = listener.accept().await.unwrap();
+ let mut socket = accept_async(stream).await.unwrap();
+ while let Some(message) = socket.next().await {
+ let message = match message {
+ Ok(Message::Text(message)) => message,
+ Ok(Message::Close(_)) | Err(_) => return,
+ _ => continue,
+ };
+ let values: Value = serde_json::from_str(&message).unwrap();
+ if values[0] == "REQ" {
+ let _ = socket
+ .send(Message::Text(
+ serde_json::to_string(&("NOTICE", "x".repeat(512 * 1024)))
+ .unwrap()
+ .into(),
+ ))
+ .await;
+ }
+ }
+ });
+ let (transport, request) = transport(&[url], 1);
+ let page = tokio::time::timeout(Duration::from_secs(5), transport.fetch(request))
+ .await
+ .unwrap()
+ .unwrap();
+ tokio::time::timeout(Duration::from_secs(5), server)
+ .await
+ .unwrap()
+ .unwrap();
+ assert!(page.events().is_empty());
+ assert_eq!(page.target_outcomes().len(), 1);
+ assert_ne!(
+ page.target_outcomes()[0].state(),
+ FetchTargetState::Complete
+ );
+}
+
+#[tokio::test(flavor = "multi_thread")]
+async fn dropping_a_polled_fetch_retains_the_original_remote_auto_close_bound() {
+ let (listener, url) = listener().await;
+ let (started, observed) = oneshot::channel();
+ let server = tokio::spawn(async move {
+ let (stream, _) = listener.accept().await.unwrap();
+ let mut socket = accept_async(stream).await.unwrap();
+ let mut started = Some(started);
+ while let Some(message) = socket.next().await {
+ let Ok(Message::Text(message)) = message else {
+ continue;
+ };
+ let values: Value = serde_json::from_str(&message).unwrap();
+ if values[0] == "REQ" {
+ started.take().unwrap().send(()).unwrap();
+ } else if values[0] == "CLOSE" {
+ return;
+ }
+ }
+ panic!("the published subscription must receive CLOSE");
+ });
+ let (transport, request) = transport(&[url], 1);
+ let mut fetch = Box::pin(transport.fetch(request));
+ tokio::time::timeout(Duration::from_secs(5), async {
+ tokio::select! {
+ _ = &mut fetch => panic!("fetch completed before its relay response"),
+ result = observed => { result.unwrap(); }
+ }
+ })
+ .await
+ .unwrap();
+ drop(fetch);
+ tokio::time::timeout(Duration::from_secs(5), server)
+ .await
+ .unwrap()
+ .unwrap();
+}
diff --git a/crates/transport_nostr/tests/network_hardening.rs b/crates/transport_nostr/tests/network_hardening.rs
@@ -11,7 +11,9 @@ const CLIENT_SOURCE: &str = include_str!("../src/client.rs");
fn tls_verification_and_pinned_dns_are_non_configurable_live_defaults() {
assert!(WORKSPACE_MANIFEST.contains("rustls-tls-webpki-roots"));
for required in [
- "client_async_tls(relay.as_str(), tcp)",
+ "client_async_tls_with_config(",
+ ".max_message_size(Some(MAX_WIRE_MESSAGE_BYTES))",
+ ".max_frame_size(Some(MAX_WIRE_MESSAGE_BYTES))",
"validate_resolved_addresses(",
"connect_pinned(addresses.as_slice())",
"ConnectionMode::Direct",
diff --git a/crates/transport_nostr/tests/package_boundary.rs b/crates/transport_nostr/tests/package_boundary.rs
@@ -320,6 +320,7 @@ fn adapter_owns_no_storage_outbox_or_orchestration_surface() {
"relay.rs".to_owned(),
"sink.rs".to_owned(),
"source.rs".to_owned(),
+ "source_budget.rs".to_owned(),
"status.rs".to_owned(),
"subscription.rs".to_owned(),
])