1use crate::config::{Config, NostrEventTransport};
7#[cfg(feature = "experimental-decentralized-pubsub")]
8use crate::nostr_relay::NostrRelay;
9use crate::storage::{HashtreeStore, StorageRouter};
10use anyhow::{Context, Result};
11use hashtree_core::{BlobRoute, StoreBlobRoute};
12use hashtree_fips_transport::{
13 bind_fips_endpoint, bind_fips_endpoint_at_local_rendezvous, set_fips_peer_configs,
14 BoundFipsEndpoint, FipsBlobRoute, FipsEndpoint, FipsEndpointOptions, FipsPeerConfig,
15 PeerIdentity, TcpBlobTransport, TcpBlobTransportConfig, WebSocketConfig,
16 DEFAULT_FIPS_DISCOVERY_SCOPE,
17};
18use hashtree_network::{BlobRouteEntry, BlobRouter, BlobRouterConfig, MeshForwardingRoute};
19use hashtree_nostr_pubsub::HashtreeNostrBoundedEventCache;
20use nostr::nips::nip19::ToBech32;
21use nostr::PublicKey;
22#[cfg(feature = "experimental-decentralized-pubsub")]
23use nostr_pubsub::{EventBus, Filter, VerifiedEvent};
24use nostr_pubsub::{
25 EventPolicyContext, EventRetentionPolicy, EventSource, NostrPubsubRouter, PolicyDecision,
26 PubsubPolicy, RouterLiveSource, RouterPublishSource, RouterQuerySource, SourcePolicyContext,
27 SourceRoute,
28};
29use nostr_pubsub_fips::{FipsPubsubClient, FipsPubsubClientOptions};
30use std::collections::HashSet;
31use std::sync::{Arc, Mutex};
32use std::time::Duration;
33#[cfg(feature = "experimental-decentralized-pubsub")]
34use tokio::sync::mpsc;
35#[cfg(feature = "experimental-decentralized-pubsub")]
36use tokio::task::JoinHandle;
37
38pub type DaemonBlobResolver = BlobRouter;
39type DaemonBlobTransport = TcpBlobTransport<StorageRouter>;
40pub type DaemonNostrCache = HashtreeNostrBoundedEventCache<StorageRouter>;
41
42const DAEMON_SAME_HOST_PROVIDER_PRIORITY: i16 = 100;
43const DAEMON_NOSTR_CACHE_EVENTS: usize = 4_096;
44const BLOB_RESOLVER_REPLY_MARGIN: Duration = Duration::from_secs(1);
45
46pub struct DaemonFipsHandle {
47 pub endpoint: Arc<FipsEndpoint>,
48 pub endpoint_npub: String,
49 pub discovery_scope: String,
50 pub pubsub_client: Option<Arc<FipsPubsubClient>>,
51 pub blob_resolver: Arc<DaemonBlobResolver>,
52 blob_transport: Mutex<Option<Arc<DaemonBlobTransport>>>,
53}
54
55impl DaemonFipsHandle {
56 pub async fn shutdown(&self) {
57 if let Ok(mut transport) = self.blob_transport.lock() {
58 transport.take();
59 }
60 if let Err(err) = self.endpoint.shutdown().await {
61 tracing::warn!("failed to stop embedded FIPS endpoint: {err}");
62 }
63 }
64}
65
66#[cfg(feature = "experimental-decentralized-pubsub")]
67pub struct DaemonNostrPubsubHandle {
68 _client: Arc<FipsPubsubClient>,
69 relay: Arc<NostrRelay>,
70 ingest_task: JoinHandle<()>,
71 outbound_task: JoinHandle<()>,
72}
73
74#[cfg(feature = "experimental-decentralized-pubsub")]
75impl DaemonNostrPubsubHandle {
76 pub fn shutdown(&self) {
77 self.relay.set_decentralized_pubsub_sender(None);
78 self.ingest_task.abort();
79 self.outbound_task.abort();
80 }
81}
82
83pub async fn start_daemon_fips_transport(
84 config: &Config,
85 keys: &nostr::Keys,
86 store: Arc<HashtreeStore>,
87 peer_ids: Vec<String>,
88) -> Result<Option<DaemonFipsHandle>> {
89 if !config.server.enable_fips || !config.server.mode.hash_get_enabled() {
90 return Ok(None);
91 }
92
93 let active_relays = config.nostr.active_relays();
94 let relays = config.server.resolved_fips_relays(&active_relays);
95 let discovery_scope = normalized_discovery_scope(&config.server.fips_discovery_scope);
96 let peer_configs = daemon_fips_peer_configs(config, peer_ids);
97 let identity_nsec = keys
98 .secret_key()
99 .to_bech32()
100 .context("Failed to encode daemon identity for FIPS endpoint")?;
101 let mut options = FipsEndpointOptions::new(identity_nsec);
102 options.discovery_scope = discovery_scope.clone();
103 options.relays = relays;
104 options.enable_udp = config.server.enable_fips_udp;
105 options.enable_webrtc = config.server.enable_fips_webrtc;
106 let websocket_seed_urls = config.server.resolved_fips_websocket_seed_urls();
107 if !websocket_seed_urls.is_empty() {
108 options.websocket = Some(WebSocketConfig {
109 seed_urls: websocket_seed_urls,
110 ..Default::default()
111 });
112 }
113 options.enable_local_rendezvous = true;
114 options.enable_lan_discovery = config.server.enable_fips_lan_discovery;
115 options.ethernet_interfaces = config.server.fips_ethernet_interfaces.clone();
116 options.udp_bind_addr = config.server.fips_udp_bind_addr.clone();
117 options.udp_public = config.server.fips_udp_public;
118 options.udp_external_addr = config.server.fips_udp_external_addr.clone();
119 options.webrtc_auto_connect = options.enable_webrtc && !peer_configs.is_empty();
123 options.webrtc_max_connections = hashtree_fips_transport::DEFAULT_FIPS_WEBRTC_MAX_CONNECTIONS;
124 options.open_discovery_max_pending = config.server.fips_open_discovery_max_pending;
125 options.packet_channel_capacity = 1024;
126 let endpoint = if let Some(raw) = config.server.fips_local_rendezvous_addr.as_deref() {
127 let rendezvous_addr = raw
128 .parse()
129 .context("server.fips_local_rendezvous_addr must be an IPv4 socket address")?;
130 bind_fips_endpoint_at_local_rendezvous(options, rendezvous_addr).await
131 } else {
132 bind_fips_endpoint(options).await
133 }
134 .context("Failed to start FIPS endpoint")?;
135 let request_timeout = Duration::from_millis(config.server.fips_request_timeout_ms.max(1));
136 let pubsub_client = if daemon_fips_pubsub_required(config) {
137 let options = FipsPubsubClientOptions {
138 query_timeout: request_timeout,
139 max_frame_bytes: config
140 .nostr
141 .decentralized_pubsub_max_event_bytes
142 .min(nostr_pubsub_fips::FIPS_NOSTR_PUBSUB_MAX_FRAME_BYTES),
143 ..Default::default()
144 };
145 Some(Arc::new(
146 FipsPubsubClient::start(endpoint.native_endpoint.clone(), options)
147 .await
148 .context("Failed to start FIPS Nostr pubsub provider")?,
149 ))
150 } else {
151 None
152 };
153 if !peer_configs.is_empty() {
154 set_fips_peer_configs(endpoint.native_endpoint.as_ref(), peer_configs.clone())
155 .await
156 .context("Failed to configure FIPS peers")?;
157 }
158 let (blob_resolver, blob_transport) =
159 bind_daemon_blob_resolver(&endpoint, store.store_arc(), &peer_configs, request_timeout)
160 .await?;
161
162 Ok(Some(DaemonFipsHandle {
163 endpoint: endpoint.native_endpoint,
164 endpoint_npub: endpoint.local_peer_id,
165 discovery_scope: endpoint.discovery_scope,
166 pubsub_client,
167 blob_resolver,
168 blob_transport: Mutex::new(Some(blob_transport)),
169 }))
170}
171
172fn daemon_fips_pubsub_required(config: &Config) -> bool {
173 config.nostr.event_transport == NostrEventTransport::FipsLocalOnly
174 || config.nostr.decentralized_pubsub_enabled()
175}
176
177async fn bind_daemon_blob_resolver(
178 endpoint: &BoundFipsEndpoint,
179 store: Arc<StorageRouter>,
180 peers: &[FipsPeerConfig],
181 request_timeout: Duration,
182) -> Result<(Arc<DaemonBlobResolver>, Arc<DaemonBlobTransport>)> {
183 let mut routes = vec![BlobRouteEntry::new(
184 "configured-store",
185 Arc::new(StoreBlobRoute::new(store.clone())),
186 )];
187 let resolver = Arc::new(
188 BlobRouter::new(
189 routes.clone(),
190 Some(store.clone()),
191 BlobRouterConfig {
192 request_timeout,
193 ..Default::default()
194 },
195 )
196 .map_err(anyhow::Error::msg)
197 .context("Failed to configure the daemon blob router")?,
198 );
199 let inbound_route: Arc<dyn BlobRoute> = resolver.clone();
203 let transport = Arc::new(
204 TcpBlobTransport::bind_advertised_route_with_config(
205 endpoint.native_endpoint.clone(),
206 store.clone(),
207 inbound_route,
208 blob_transport_config(request_timeout),
209 DAEMON_SAME_HOST_PROVIDER_PRIORITY,
210 )
211 .await
212 .context("Failed to advertise htree's blob resolver")?,
213 );
214 let peer_identities: Vec<_> = peers
215 .iter()
216 .filter_map(|peer| PeerIdentity::from_npub(&peer.npub).ok())
217 .collect();
218 if !peer_identities.is_empty() {
219 let max_provider_attempts = peer_identities.len().min(4);
220 let fips_route =
221 FipsBlobRoute::explicit(transport.clone(), peer_identities, max_provider_attempts)
222 .map_err(anyhow::Error::msg)
223 .context("Failed to configure the daemon FIPS blob route")?;
224 routes.push(BlobRouteEntry::new(
225 "configured-fips-peers",
226 Arc::new(MeshForwardingRoute::new(Arc::new(fips_route))),
227 ));
228 }
229 resolver
230 .set_routes(routes)
231 .await
232 .map_err(anyhow::Error::msg)
233 .context("Failed to install daemon blob routes")?;
234 Ok((resolver, transport))
235}
236
237fn blob_transport_config(resolver_timeout: Duration) -> TcpBlobTransportConfig {
238 TcpBlobTransportConfig {
239 idle_timeout: resolver_timeout.saturating_add(BLOB_RESOLVER_REPLY_MARGIN),
240 }
241}
242
243pub fn new_daemon_nostr_cache(store: Arc<StorageRouter>) -> Arc<DaemonNostrCache> {
244 Arc::new(
245 HashtreeNostrBoundedEventCache::new(
246 store,
247 None,
248 EventSource::local_index("hashtree-nostr-cache"),
249 EventRetentionPolicy::new(DAEMON_NOSTR_CACHE_EVENTS, Vec::new()),
250 )
251 .with_priority(nostr_pubsub::SOURCE_PRIORITY_LOCAL_INDEX),
252 )
253}
254
255struct AllowDaemonPubsubRoutes;
256
257#[async_trait::async_trait]
258impl PubsubPolicy for AllowDaemonPubsubRoutes {
259 async fn check_event(
260 &self,
261 _context: EventPolicyContext<'_>,
262 ) -> nostr_pubsub::Result<PolicyDecision> {
263 Ok(PolicyDecision::allow_with_priority(0))
264 }
265
266 async fn check_source(
267 &self,
268 _context: SourcePolicyContext<'_>,
269 ) -> nostr_pubsub::Result<PolicyDecision> {
270 Ok(PolicyDecision::allow_with_priority(0))
271 }
272}
273
274pub async fn start_daemon_nostr_provider(
275 config: &Config,
276 fips_handle: Option<&DaemonFipsHandle>,
277 cache: Option<Arc<DaemonNostrCache>>,
278) -> Result<Option<Arc<dyn nostr_pubsub::PubsubProvider>>> {
279 let mut router = NostrPubsubRouter::new(Arc::new(AllowDaemonPubsubRoutes));
280 if let Some(cache) = cache {
281 let route = cache
282 .source_route("hashtree-local-cache")
283 .context("Failed to configure Hashtree Nostr cache route")?;
284 router = router.with_query_source(RouterQuerySource::new(route, cache));
285 }
286
287 let router = match config.nostr.event_transport {
288 NostrEventTransport::Relay => {
289 let relays = config.nostr.active_relays();
290 if relays.is_empty() {
291 return Ok(None);
292 }
293 let route = SourceRoute::relay(relays.join(","))
294 .with_dataset("configured-relays")
295 .context("Failed to configure Nostr relay route")?;
296 let event_bus = Arc::new(
297 nostr_pubsub_relay::RelayEventBus::new(
298 relays.clone(),
299 Duration::from_millis(config.server.fips_request_timeout_ms.max(1)),
300 )
301 .await
302 .context("Failed to start Nostr relay event provider")?,
303 );
304 router
305 .with_query_source(RouterQuerySource::new(
306 route.clone(),
307 Arc::clone(&event_bus),
308 ))
309 .with_publish_source(RouterPublishSource::new(
310 route.clone(),
311 Arc::clone(&event_bus),
312 ))
313 .with_live_source(RouterLiveSource::new(route, event_bus))
314 }
315 NostrEventTransport::FipsLocalOnly => {
316 let client = fips_handle
317 .and_then(|handle| handle.pubsub_client.clone())
318 .ok_or_else(|| {
319 anyhow::anyhow!(
320 "nostr.event_transport=fips-local-only requires the local FIPS Nostr pubsub provider"
321 )
322 })?;
323 let route = SourceRoute::fips_peer_default("connected-fips-mesh")
324 .with_dataset("fips-mesh")
325 .context("Failed to configure FIPS Nostr route")?;
326 router
327 .with_query_source(RouterQuerySource::new(route.clone(), Arc::clone(&client)))
328 .with_publish_source(RouterPublishSource::new(route.clone(), Arc::clone(&client)))
329 .with_live_source(RouterLiveSource::new(route, client))
330 }
331 };
332 Ok(Some(Arc::new(router)))
333}
334
335#[cfg(feature = "experimental-decentralized-pubsub")]
336pub async fn start_daemon_nostr_pubsub(
337 config: &Config,
338 fips_handle: Option<&DaemonFipsHandle>,
339 relay: Option<Arc<NostrRelay>>,
340 cache: Arc<DaemonNostrCache>,
341) -> Result<Option<Arc<DaemonNostrPubsubHandle>>> {
342 if !daemon_decentralized_pubsub_ready(config, fips_handle.is_some(), relay.is_some()) {
343 return Ok(None);
344 }
345
346 let fips_handle = fips_handle.expect("checked by daemon_decentralized_pubsub_ready");
347 let relay = relay.expect("checked by daemon_decentralized_pubsub_ready");
348 let client = fips_handle.pubsub_client.clone().ok_or_else(|| {
349 anyhow::anyhow!("decentralized Nostr pubsub requires the shared FIPS pubsub client")
350 })?;
351
352 let (outbound_tx, outbound_rx) = mpsc::unbounded_channel();
353 relay.set_decentralized_pubsub_sender(Some(outbound_tx));
354 let ingest_task = spawn_daemon_nostr_pubsub_ingest(
355 Arc::clone(&client),
356 Arc::clone(&fips_handle.endpoint),
357 Arc::clone(&relay),
358 Arc::clone(&cache),
359 );
360 let outbound_task = spawn_daemon_nostr_pubsub_outbound(Arc::clone(&client), outbound_rx, cache);
361
362 Ok(Some(Arc::new(DaemonNostrPubsubHandle {
363 _client: client,
364 relay,
365 ingest_task,
366 outbound_task,
367 })))
368}
369
370#[cfg(feature = "experimental-decentralized-pubsub")]
371fn daemon_decentralized_pubsub_ready(config: &Config, has_fips: bool, has_relay: bool) -> bool {
372 config.nostr.decentralized_pubsub_enabled() && has_fips && has_relay
373}
374
375#[cfg(feature = "experimental-decentralized-pubsub")]
376fn spawn_daemon_nostr_pubsub_ingest(
377 client: Arc<FipsPubsubClient>,
378 endpoint: Arc<FipsEndpoint>,
379 relay: Arc<NostrRelay>,
380 cache: Arc<DaemonNostrCache>,
381) -> JoinHandle<()> {
382 tokio::spawn(async move {
383 loop {
384 let subscribed_peers = connected_fips_peer_ids(&endpoint).await;
385 if subscribed_peers.is_empty() {
386 tokio::time::sleep(Duration::from_millis(250)).await;
387 continue;
388 }
389 let mut subscription = match client.subscribe(vec![Filter::new()]).await {
390 Ok(subscription) => subscription,
391 Err(err) => {
392 tracing::debug!("waiting to subscribe to FIPS Nostr peers: {err}");
393 tokio::time::sleep(Duration::from_millis(250)).await;
394 continue;
395 }
396 };
397 loop {
398 tokio::select! {
399 delivery = subscription.recv() => {
400 let Some(delivery) = delivery else {
401 break;
402 };
403 let source = delivery.source.clone();
404 let source_id = source.id.as_str().to_string();
405 let verified = delivery.event;
406 let event = verified.as_event().clone();
407 let event_id = event.id.to_hex();
408 match relay.ingest_peer_event_silent(event).await {
409 Ok(true) => {
410 if let Err(error) = cache.publish(verified, source).await {
411 tracing::warn!(event_id, %error, "failed to cache decentralized Nostr event");
412 }
413 tracing::debug!(
414 source = source_id,
415 event_id,
416 "ingested decentralized Nostr pubsub event"
417 );
418 }
419 Ok(false) => {}
420 Err(err) => tracing::warn!(
421 source = source_id,
422 event_id,
423 "nostr decentralized pubsub ingest failed: {err:#}"
424 ),
425 }
426 }
427 () = tokio::time::sleep(Duration::from_millis(500)) => {
428 if connected_fips_peer_ids(&endpoint).await != subscribed_peers {
429 break;
430 }
431 }
432 }
433 }
434 }
435 })
436}
437
438#[cfg(feature = "experimental-decentralized-pubsub")]
439async fn connected_fips_peer_ids(endpoint: &FipsEndpoint) -> Vec<String> {
440 let mut peers = endpoint
441 .peers()
442 .await
443 .unwrap_or_default()
444 .into_iter()
445 .filter(|peer| peer.connected)
446 .map(|peer| peer.npub)
447 .collect::<Vec<_>>();
448 peers.sort_unstable();
449 peers.dedup();
450 peers
451}
452
453#[cfg(feature = "experimental-decentralized-pubsub")]
454fn spawn_daemon_nostr_pubsub_outbound(
455 client: Arc<FipsPubsubClient>,
456 mut outbound_rx: mpsc::UnboundedReceiver<nostr::Event>,
457 cache: Arc<DaemonNostrCache>,
458) -> JoinHandle<()> {
459 tokio::spawn(async move {
460 while let Some(event) = outbound_rx.recv().await {
461 let event_id = event.id.to_hex();
462 let verified = match VerifiedEvent::try_from(event) {
463 Ok(event) => event,
464 Err(err) => {
465 tracing::warn!(event_id, "invalid outbound Nostr event: {err}");
466 continue;
467 }
468 };
469 let source = EventSource::local_index("htree-relay");
470 if let Err(error) = cache.publish(verified.clone(), source.clone()).await {
471 tracing::warn!(event_id, %error, "failed to cache outbound Nostr event");
472 }
473 match client.publish(verified, source).await {
474 Ok(report) => tracing::debug!(
475 event_id,
476 accepted = report.accepted,
477 "published Nostr event over decentralized pubsub"
478 ),
479 Err(err) => {
480 tracing::warn!(event_id, "nostr decentralized pubsub publish failed: {err}")
481 }
482 }
483 }
484 })
485}
486
487pub fn fips_peer_ids_from_pubkeys(pubkeys: Vec<[u8; 32]>) -> Vec<String> {
488 pubkeys
489 .into_iter()
490 .filter_map(|pubkey| PublicKey::from_slice(&pubkey).ok())
491 .filter_map(|pubkey| pubkey.to_bech32().ok())
492 .collect()
493}
494
495pub fn daemon_fips_peer_configs(
496 config: &Config,
497 discovered_peer_ids: Vec<String>,
498) -> Vec<FipsPeerConfig> {
499 let mut seen = HashSet::new();
500 let mut peers = Vec::new();
501
502 for peer in &config.server.fips_peers {
503 let npub = peer.npub.trim().to_string();
504 if npub.is_empty() || !seen.insert(npub.clone()) {
505 continue;
506 }
507 let udp_addresses = peer
508 .udp_addresses
509 .iter()
510 .map(|addr| addr.trim().to_string())
511 .filter(|addr| !addr.is_empty())
512 .collect();
513 peers.push(FipsPeerConfig {
514 npub,
515 udp_addresses,
516 });
517 }
518
519 for peer_id in discovered_peer_ids {
520 let npub = peer_id.trim().to_string();
521 if npub.is_empty() || !seen.insert(npub.clone()) {
522 continue;
523 }
524 peers.push(FipsPeerConfig::new(npub));
525 }
526
527 peers
528}
529
530fn normalized_discovery_scope(scope: &str) -> String {
531 let scope = scope.trim();
532 if scope.is_empty() {
533 DEFAULT_FIPS_DISCOVERY_SCOPE.to_string()
534 } else {
535 scope.to_string()
536 }
537}
538
539#[cfg(test)]
540mod tests {
541 use super::*;
542 use hashtree_core::Store;
543 #[cfg(feature = "experimental-decentralized-pubsub")]
544 use nostr::{EventBuilder, Kind, Tag};
545 #[cfg(feature = "experimental-decentralized-pubsub")]
546 use nostr_pubsub::{PubsubProviderMode, QueryOptions};
547 use sha2::{Digest, Sha256};
548 use tokio::time::timeout;
549
550 struct RecordingStoreRoute {
551 store: Arc<hashtree_core::MemoryStore>,
552 requests: Arc<std::sync::Mutex<Vec<hashtree_core::BlobRequest>>>,
553 }
554
555 #[async_trait::async_trait]
556 impl BlobRoute for RecordingStoreRoute {
557 async fn route(
558 &self,
559 request: hashtree_core::BlobRequest,
560 ) -> Result<hashtree_core::BlobReply, hashtree_core::StoreError> {
561 self.requests.lock().unwrap().push(request);
562 StoreBlobRoute::new(self.store.clone()).route(request).await
563 }
564 }
565
566 #[test]
567 fn fips_peer_ids_from_pubkeys_encodes_npbus() {
568 let keys = nostr::Keys::generate();
569 let expected = keys.public_key().to_bech32().unwrap();
570
571 assert_eq!(
572 fips_peer_ids_from_pubkeys(vec![keys.public_key().to_bytes()]),
573 vec![expected]
574 );
575 }
576
577 #[test]
578 fn empty_discovery_scope_uses_hashtree_default() {
579 assert_eq!(
580 normalized_discovery_scope(" "),
581 DEFAULT_FIPS_DISCOVERY_SCOPE.to_string()
582 );
583 }
584
585 #[tokio::test]
586 async fn fips_local_only_event_transport_requires_local_provider() {
587 let mut config = Config::default();
588 config.nostr.event_transport = NostrEventTransport::FipsLocalOnly;
589
590 let error = start_daemon_nostr_provider(&config, None, None)
591 .await
592 .err()
593 .expect("missing local FIPS provider must fail");
594
595 assert!(error.to_string().contains("fips-local-only"));
596 assert!(error.to_string().contains("requires"));
597 }
598
599 #[cfg(feature = "experimental-decentralized-pubsub")]
600 #[tokio::test]
601 async fn one_native_endpoint_serves_roots_and_trusted_decentralized_events() {
602 let scope = format!("htree-native-pubsub-{}", uuid::Uuid::new_v4());
603 let daemon_addr = reserve_udp_addr();
604 let rendezvous_addr = reserve_udp_addr();
605 let (remote_endpoint, remote_addr) = udp_endpoint(&scope).await;
606 let daemon_keys = nostr::Keys::generate();
607 let trusted_keys = nostr::Keys::generate();
608 let untrusted_keys = nostr::Keys::generate();
609 let outbound_keys = nostr::Keys::generate();
610 let temp = tempfile::tempdir().unwrap();
611 let store = Arc::new(HashtreeStore::new(temp.path().join("blobs")).unwrap());
612 let relay = trusted_test_relay(
613 temp.path(),
614 &[trusted_keys.public_key(), outbound_keys.public_key()],
615 )
616 .await;
617
618 let mut config = Config::default();
619 config.server.fips_discovery_scope = scope;
620 config.server.fips_relays = Some(Vec::new());
621 config.server.fips_udp_bind_addr = Some(daemon_addr.clone());
622 config.server.fips_local_rendezvous_addr = Some(rendezvous_addr);
623 config.server.enable_fips_udp = true;
624 config.server.enable_fips_webrtc = false;
625 config.server.fips_request_timeout_ms = 2_000;
626 config.server.fips_peers = vec![crate::config::ConfiguredFipsPeer {
627 npub: remote_endpoint.local_peer_id.clone(),
628 udp_addresses: vec![remote_addr],
629 }];
630 config.nostr.enabled = true;
631 config.nostr.event_transport = NostrEventTransport::FipsLocalOnly;
632 config.nostr.decentralized_pubsub = true;
633
634 let cache = new_daemon_nostr_cache(store.store_arc());
635 let daemon = start_daemon_fips_transport(&config, &daemon_keys, store, Vec::new())
636 .await
637 .unwrap()
638 .expect("daemon FIPS endpoint");
639 let provider =
640 start_daemon_nostr_provider(&config, Some(&daemon), Some(Arc::clone(&cache)))
641 .await
642 .unwrap()
643 .expect("FIPS root provider");
644 assert_eq!(provider.mode(), PubsubProviderMode::Router);
645 let decentralized =
646 start_daemon_nostr_pubsub(&config, Some(&daemon), Some(relay.clone()), cache)
647 .await
648 .unwrap()
649 .expect("decentralized pubsub");
650
651 set_fips_peer_configs(
652 remote_endpoint.native_endpoint.as_ref(),
653 vec![FipsPeerConfig {
654 npub: daemon.endpoint_npub.clone(),
655 udp_addresses: vec![daemon_addr],
656 }],
657 )
658 .await
659 .unwrap();
660 let remote_client = Arc::new(
661 FipsPubsubClient::start(
662 remote_endpoint.native_endpoint.clone(),
663 FipsPubsubClientOptions {
664 query_timeout: Duration::from_secs(2),
665 ..Default::default()
666 },
667 )
668 .await
669 .unwrap(),
670 );
671
672 wait_for_native_peer(&daemon.endpoint, &remote_endpoint.local_peer_id).await;
673 wait_for_peer(&remote_endpoint, &daemon.endpoint_npub).await;
674 timeout(Duration::from_secs(5), async {
675 while remote_client.peer_subscription_count().unwrap() == 0 {
676 tokio::time::sleep(Duration::from_millis(20)).await;
677 }
678 })
679 .await
680 .expect("daemon did not subscribe to the authenticated peer");
681
682 let published_root = root_event(&outbound_keys, "published-root", "aa");
685 let mut published_root_rx = remote_client
686 .subscribe(vec![Filter::new().id(published_root.id)])
687 .await
688 .unwrap();
689 provider
690 .publish(
691 VerifiedEvent::try_from(published_root.clone()).unwrap(),
692 EventSource::local_index("root-publisher"),
693 )
694 .await
695 .unwrap();
696 let delivered_root = timeout(Duration::from_secs(5), published_root_rx.recv())
697 .await
698 .unwrap()
699 .expect("published root delivery");
700 assert_eq!(delivered_root.event.as_event().id, published_root.id);
701
702 let resolved_root = root_event(&trusted_keys, "resolved-root", "bb");
703 remote_client
704 .publish(
705 VerifiedEvent::try_from(resolved_root.clone()).unwrap(),
706 EventSource::local_index("remote-root"),
707 )
708 .await
709 .unwrap();
710 let resolution = provider
711 .query(
712 vec![Filter::new().id(resolved_root.id)],
713 QueryOptions { limit: Some(1) },
714 )
715 .await
716 .unwrap();
717 assert_eq!(resolution.events.len(), 1);
718 assert_eq!(resolution.events[0].event.as_event().id, resolved_root.id);
719
720 let rejected = EventBuilder::new(Kind::TextNote, "untrusted")
724 .sign_with_keys(&untrusted_keys)
725 .unwrap();
726 remote_client
727 .publish(
728 VerifiedEvent::try_from(rejected.clone()).unwrap(),
729 EventSource::local_index("remote-untrusted"),
730 )
731 .await
732 .unwrap();
733 let accepted = EventBuilder::new(Kind::TextNote, "trusted")
734 .sign_with_keys(&trusted_keys)
735 .unwrap();
736 remote_client
737 .publish(
738 VerifiedEvent::try_from(accepted.clone()).unwrap(),
739 EventSource::local_index("remote-trusted"),
740 )
741 .await
742 .unwrap();
743 wait_for_relay_event(&relay, accepted.id).await;
744 assert!(relay
745 .query_events(&Filter::new().id(rejected.id), 1)
746 .await
747 .is_empty());
748
749 let outbound = EventBuilder::new(Kind::TextNote, "relay outbound")
752 .sign_with_keys(&outbound_keys)
753 .unwrap();
754 let mut outbound_rx = remote_client
755 .subscribe(vec![Filter::new().id(outbound.id)])
756 .await
757 .unwrap();
758 let (client_tx, _client_rx) = mpsc::unbounded_channel();
759 relay.register_client(1, client_tx, None).await;
760 relay
761 .handle_client_message(1, nostr::ClientMessage::event(outbound.clone()))
762 .await;
763 let delivered = timeout(Duration::from_secs(5), outbound_rx.recv())
764 .await
765 .unwrap()
766 .expect("outbound relay publication");
767 assert_eq!(delivered.event.as_event().id, outbound.id);
768
769 decentralized.shutdown();
770 daemon.shutdown().await;
771 remote_endpoint.native_endpoint.shutdown().await.unwrap();
772 }
773
774 #[tokio::test]
775 async fn daemon_blob_resolver_advertises_and_serves_storage_router() {
776 let rendezvous_addr = reserve_local_rendezvous_addr();
777 let test_id = uuid::Uuid::new_v4();
778 let provider_scope = format!("htree-provider-test-{test_id}");
779 let observer_scope = format!("iris-drive-observer-test-{test_id}");
780 let provider_endpoint = local_only_endpoint(rendezvous_addr, &provider_scope).await;
781 let observer_endpoint = local_only_endpoint(rendezvous_addr, &observer_scope).await;
782 let temp = tempfile::tempdir().unwrap();
783 let store = Arc::new(HashtreeStore::new(temp.path()).unwrap());
784 let data = b"served directly from htree StorageRouter".to_vec();
785 let hash = Sha256::digest(&data).into();
786 store.store_arc().put(hash, data.clone()).await.unwrap();
787 let (blob_resolver, provider) = bind_daemon_blob_resolver(
788 &provider_endpoint,
789 store.store_arc(),
790 &[],
791 Duration::from_secs(1),
792 )
793 .await
794 .expect("bind htree provider");
795 let handle = DaemonFipsHandle {
796 endpoint: provider_endpoint.native_endpoint.clone(),
797 endpoint_npub: provider_endpoint.local_peer_id.clone(),
798 discovery_scope: provider_endpoint.discovery_scope.clone(),
799 pubsub_client: None,
800 blob_resolver,
801 blob_transport: Mutex::new(Some(provider)),
802 };
803
804 timeout(Duration::from_secs(10), async {
805 loop {
806 if provider_advertised(&observer_endpoint, &provider_endpoint.local_peer_id) {
807 break;
808 }
809 tokio::time::sleep(Duration::from_millis(20)).await;
810 }
811 })
812 .await
813 .expect("htree provider capability did not converge");
814
815 let observer_store = Arc::new(hashtree_core::MemoryStore::new());
816 let observer_transport = Arc::new(
817 TcpBlobTransport::bind_route(
818 observer_endpoint.native_endpoint.clone(),
819 observer_store.clone(),
820 Arc::new(StoreBlobRoute::new(observer_store.clone())),
821 )
822 .await
823 .unwrap(),
824 );
825 let observer_fips_route = FipsBlobRoute::discovered(
826 observer_endpoint.native_endpoint.clone(),
827 observer_transport.clone(),
828 4,
829 )
830 .unwrap();
831 let observer_router = BlobRouter::new(
832 vec![BlobRouteEntry::new(
833 "same-host-fips",
834 Arc::new(observer_fips_route),
835 )],
836 Some(observer_store),
837 BlobRouterConfig::default(),
838 )
839 .unwrap();
840 assert_eq!(observer_router.get(&hash, None).await.unwrap(), Some(data));
841
842 drop(observer_router);
843 drop(observer_transport);
844 handle.shutdown().await;
845 timeout(Duration::from_secs(10), async {
846 while provider_advertised(&observer_endpoint, &provider_endpoint.local_peer_id) {
847 tokio::time::sleep(Duration::from_millis(20)).await;
848 }
849 })
850 .await
851 .expect("htree provider capability survived daemon shutdown");
852 drop(handle);
853 observer_endpoint.native_endpoint.shutdown().await.unwrap();
854 }
855
856 #[tokio::test]
857 async fn daemon_mesh_forwarding_observes_two_one_zero_and_exhaustion() {
858 let scope = format!("htree-inbound-forwarding-{}", uuid::Uuid::new_v4());
859 let (upstream_endpoint, upstream_addr) = udp_endpoint(&scope).await;
860 let (provider_endpoint, provider_addr) = udp_endpoint(&scope).await;
861 let (observer_endpoint, observer_addr) = udp_endpoint(&scope).await;
862 let temp = tempfile::tempdir().unwrap();
863 let provider_store = Arc::new(HashtreeStore::new(temp.path().join("provider")).unwrap());
864 let upstream_store = Arc::new(hashtree_core::MemoryStore::new());
865 let data = b"served through the daemon's configured FIPS peer".to_vec();
866 let hash = Sha256::digest(&data).into();
867 upstream_store.put(hash, data.clone()).await.unwrap();
868 let upstream_requests = Arc::new(std::sync::Mutex::new(Vec::new()));
869 let upstream_transport = Arc::new(
870 TcpBlobTransport::bind_route(
871 upstream_endpoint.native_endpoint.clone(),
872 upstream_store.clone(),
873 Arc::new(RecordingStoreRoute {
874 store: upstream_store,
875 requests: upstream_requests.clone(),
876 }),
877 )
878 .await
879 .unwrap(),
880 );
881 let upstream_peer = FipsPeerConfig {
882 npub: upstream_endpoint.local_peer_id.clone(),
883 udp_addresses: vec![upstream_addr],
884 };
885 set_fips_peer_configs(
886 provider_endpoint.native_endpoint.as_ref(),
887 vec![
888 upstream_peer.clone(),
889 FipsPeerConfig {
890 npub: observer_endpoint.local_peer_id.clone(),
891 udp_addresses: vec![observer_addr],
892 },
893 ],
894 )
895 .await
896 .unwrap();
897 set_fips_peer_configs(
898 upstream_endpoint.native_endpoint.as_ref(),
899 vec![FipsPeerConfig {
900 npub: provider_endpoint.local_peer_id.clone(),
901 udp_addresses: vec![provider_addr.clone()],
902 }],
903 )
904 .await
905 .unwrap();
906 set_fips_peer_configs(
907 observer_endpoint.native_endpoint.as_ref(),
908 vec![FipsPeerConfig {
909 npub: provider_endpoint.local_peer_id.clone(),
910 udp_addresses: vec![provider_addr],
911 }],
912 )
913 .await
914 .unwrap();
915 let (provider_resolver, provider_transport) = bind_daemon_blob_resolver(
916 &provider_endpoint,
917 provider_store.store_arc(),
918 &[upstream_peer],
919 Duration::from_secs(2),
920 )
921 .await
922 .unwrap();
923 let observer_store = Arc::new(hashtree_core::MemoryStore::new());
924 let observer_transport = Arc::new(
925 TcpBlobTransport::bind_route(
926 observer_endpoint.native_endpoint.clone(),
927 observer_store.clone(),
928 Arc::new(StoreBlobRoute::new(observer_store)),
929 )
930 .await
931 .unwrap(),
932 );
933 let observer_fips_route = FipsBlobRoute::explicit(
934 observer_transport.clone(),
935 vec![PeerIdentity::from_npub(&provider_endpoint.local_peer_id).unwrap()],
936 1,
937 )
938 .unwrap();
939 let observer_route = MeshForwardingRoute::new(Arc::new(observer_fips_route));
940
941 wait_for_native_peer(
942 provider_endpoint.native_endpoint.as_ref(),
943 &upstream_endpoint.local_peer_id,
944 )
945 .await;
946 wait_for_native_peer(
947 observer_endpoint.native_endpoint.as_ref(),
948 &provider_endpoint.local_peer_id,
949 )
950 .await;
951 assert_eq!(
952 observer_route
953 .route(hashtree_core::BlobRequest { hash, htl: 0 })
954 .await
955 .unwrap(),
956 hashtree_core::BlobReply::NoResult,
957 );
958 assert!(upstream_requests.lock().unwrap().is_empty());
959 assert_eq!(
960 observer_route
961 .route(hashtree_core::BlobRequest { hash, htl: 1 })
962 .await
963 .unwrap(),
964 hashtree_core::BlobReply::NoResult,
965 );
966 assert!(
967 upstream_requests.lock().unwrap().is_empty(),
968 "HTL 1 reaches the intermediate daemon as 0 and cannot be forwarded"
969 );
970 assert_eq!(
971 observer_route
972 .route(hashtree_core::BlobRequest { hash, htl: 2 })
973 .await
974 .unwrap(),
975 hashtree_core::BlobReply::Data(data),
976 );
977 assert_eq!(
978 upstream_requests.lock().unwrap().as_slice(),
979 &[hashtree_core::BlobRequest { hash, htl: 0 }],
980 "three Hashtree nodes must observe the two-to-one-to-zero forwarding budget while each FIPS carrier preserves its request",
981 );
982
983 drop(observer_route);
984 drop(observer_transport);
985 drop(provider_resolver);
986 drop(provider_transport);
987 drop(upstream_transport);
988 observer_endpoint.native_endpoint.shutdown().await.unwrap();
989 provider_endpoint.native_endpoint.shutdown().await.unwrap();
990 upstream_endpoint.native_endpoint.shutdown().await.unwrap();
991 }
992
993 fn provider_advertised(endpoint: &BoundFipsEndpoint, npub: &str) -> bool {
994 endpoint
995 .native_endpoint
996 .local_instance_advertisements()
997 .unwrap()
998 .iter()
999 .any(|advert| {
1000 advert.npub == npub
1001 && advert
1002 .capability(hashtree_fips_transport::TCP_BLOB_CAPABILITY)
1003 .and_then(|capability| capability.fsp_port)
1004 == Some(hashtree_fips_transport::TCP_BLOB_SERVICE_PORT)
1005 })
1006 }
1007
1008 async fn local_only_endpoint(
1009 rendezvous_addr: std::net::SocketAddrV4,
1010 scope: &str,
1011 ) -> hashtree_fips_transport::BoundFipsEndpoint {
1012 let keys = nostr::Keys::generate();
1013 let mut options = FipsEndpointOptions::new(keys.secret_key().to_bech32().unwrap());
1014 options.discovery_scope = scope.to_string();
1015 options.enable_udp = false;
1016 options.enable_webrtc = false;
1017 options.enable_local_rendezvous = true;
1018 options.enable_lan_discovery = false;
1019 options.share_local_candidates = false;
1020 bind_fips_endpoint_at_local_rendezvous(options, rendezvous_addr)
1021 .await
1022 .unwrap()
1023 }
1024
1025 async fn udp_endpoint(scope: &str) -> (BoundFipsEndpoint, String) {
1026 let keys = nostr::Keys::generate();
1027 let mut options = FipsEndpointOptions::new(keys.secret_key().to_bech32().unwrap());
1028 options.discovery_scope = scope.to_string();
1029 options.enable_udp = true;
1030 options.udp_bind_addr = Some("127.0.0.1:0".to_string());
1031 options.enable_webrtc = false;
1032 options.enable_local_rendezvous = false;
1033 options.enable_lan_discovery = false;
1034 options.share_local_candidates = false;
1035 let endpoint = bind_fips_endpoint(options).await.unwrap();
1036 let addrs = endpoint
1037 .native_endpoint
1038 .bound_udp_listen_addrs()
1039 .await
1040 .unwrap();
1041 let [addr] = addrs.as_slice() else {
1042 panic!("expected exactly one bound UDP listener, got {addrs:?}");
1043 };
1044 assert!(addr.ip().is_loopback());
1045 assert_ne!(addr.port(), 0);
1046 (endpoint, addr.to_string())
1047 }
1048
1049 #[cfg(feature = "experimental-decentralized-pubsub")]
1050 fn reserve_udp_addr() -> String {
1051 let socket = std::net::UdpSocket::bind("127.0.0.1:0").unwrap();
1052 socket.local_addr().unwrap().to_string()
1053 }
1054
1055 fn reserve_local_rendezvous_addr() -> std::net::SocketAddrV4 {
1056 let socket = std::net::UdpSocket::bind("127.0.0.1:0").unwrap();
1057 match socket.local_addr().unwrap() {
1058 std::net::SocketAddr::V4(addr) => addr,
1059 std::net::SocketAddr::V6(_) => unreachable!("loopback IPv4 bind returned IPv6"),
1060 }
1061 }
1062
1063 #[cfg(feature = "experimental-decentralized-pubsub")]
1064 async fn wait_for_peer(endpoint: &BoundFipsEndpoint, npub: &str) {
1065 wait_for_native_peer(endpoint.native_endpoint.as_ref(), npub).await;
1066 }
1067
1068 async fn wait_for_native_peer(endpoint: &FipsEndpoint, npub: &str) {
1069 let result = timeout(Duration::from_secs(10), async {
1070 loop {
1071 if endpoint
1072 .peers()
1073 .await
1074 .unwrap()
1075 .iter()
1076 .any(|peer| peer.npub == npub && peer.connected)
1077 {
1078 break;
1079 }
1080 tokio::time::sleep(Duration::from_millis(20)).await;
1081 }
1082 })
1083 .await;
1084 if result.is_err() {
1085 panic!(
1086 "FIPS peer {npub} did not connect; peers={:?}",
1087 endpoint.peers().await
1088 );
1089 }
1090 }
1091
1092 #[cfg(feature = "experimental-decentralized-pubsub")]
1093 async fn trusted_test_relay(
1094 root: &std::path::Path,
1095 allowed: &[nostr::PublicKey],
1096 ) -> Arc<NostrRelay> {
1097 let graph_dir = root.join("graph");
1098 let data_dir = root.join("nostr");
1099 std::fs::create_dir_all(&data_dir).unwrap();
1100 let graph = {
1101 let _guard = crate::socialgraph::test_lock().await;
1102 crate::socialgraph::open_test_social_graph_store_with_mapsize(
1103 &graph_dir,
1104 Some(128 * 1024 * 1024),
1105 )
1106 .unwrap()
1107 };
1108 let backend: Arc<dyn crate::socialgraph::SocialGraphBackend> = graph;
1109 let allowed = allowed
1110 .iter()
1111 .map(|pubkey| pubkey.to_hex())
1112 .collect::<HashSet<_>>();
1113 let access = Arc::new(crate::socialgraph::SocialGraphAccessControl::new(
1114 Arc::clone(&backend),
1115 0,
1116 allowed.clone(),
1117 ));
1118 Arc::new(
1119 NostrRelay::new(
1120 backend,
1121 data_dir,
1122 allowed,
1123 Some(access),
1124 crate::nostr_relay::NostrRelayConfig {
1125 spambox_db_max_bytes: 0,
1126 ..Default::default()
1127 },
1128 )
1129 .unwrap(),
1130 )
1131 }
1132
1133 #[cfg(feature = "experimental-decentralized-pubsub")]
1134 fn root_event(keys: &nostr::Keys, name: &str, hash: &str) -> nostr::Event {
1135 EventBuilder::new(Kind::Custom(30_078), hash)
1136 .tags([
1137 Tag::parse(["d", name]).unwrap(),
1138 Tag::parse(["l", "hashtree"]).unwrap(),
1139 Tag::parse(["hash", hash]).unwrap(),
1140 ])
1141 .sign_with_keys(keys)
1142 .unwrap()
1143 }
1144
1145 #[cfg(feature = "experimental-decentralized-pubsub")]
1146 async fn wait_for_relay_event(relay: &NostrRelay, id: nostr::EventId) {
1147 timeout(Duration::from_secs(5), async {
1148 loop {
1149 if relay
1150 .query_events(&Filter::new().id(id), 1)
1151 .await
1152 .iter()
1153 .any(|event| event.id == id)
1154 {
1155 break;
1156 }
1157 tokio::time::sleep(Duration::from_millis(20)).await;
1158 }
1159 })
1160 .await
1161 .expect("trusted decentralized event was not indexed");
1162 }
1163
1164 #[test]
1165 fn daemon_fips_peer_configs_prefer_configured_peers() {
1166 let mut config = Config::default();
1167 config.server.fips_peers = vec![
1168 crate::config::ConfiguredFipsPeer {
1169 npub: " origin ".to_string(),
1170 udp_addresses: vec![" udp:192.0.2.10:2121 ".to_string(), " ".to_string()],
1171 },
1172 crate::config::ConfiguredFipsPeer {
1173 npub: "origin".to_string(),
1174 udp_addresses: vec!["udp:ignored:2121".to_string()],
1175 },
1176 ];
1177
1178 let peers = daemon_fips_peer_configs(
1179 &config,
1180 vec![
1181 "origin".to_string(),
1182 " followed ".to_string(),
1183 " ".to_string(),
1184 ],
1185 );
1186
1187 assert_eq!(
1188 peers,
1189 vec![
1190 FipsPeerConfig {
1191 npub: "origin".to_string(),
1192 udp_addresses: vec!["udp:192.0.2.10:2121".to_string()],
1193 },
1194 FipsPeerConfig::new("followed"),
1195 ]
1196 );
1197 }
1198
1199 #[test]
1200 fn blob_carrier_deadline_outlives_the_complete_resolver_budget() {
1201 for resolver_timeout in [Duration::from_millis(1), Duration::from_secs(20)] {
1202 assert!(blob_transport_config(resolver_timeout).idle_timeout > resolver_timeout);
1203 }
1204 }
1205
1206 #[cfg(feature = "experimental-decentralized-pubsub")]
1207 #[test]
1208 fn daemon_decentralized_pubsub_requires_config_fips_and_relay() {
1209 let mut config = Config::default();
1210 assert!(!daemon_decentralized_pubsub_ready(&config, true, true));
1211
1212 config.nostr.decentralized_pubsub = true;
1213 assert!(daemon_decentralized_pubsub_ready(&config, true, true));
1214 assert!(!daemon_decentralized_pubsub_ready(&config, false, true));
1215 assert!(!daemon_decentralized_pubsub_ready(&config, true, false));
1216
1217 config.nostr.enabled = false;
1218 assert!(!daemon_decentralized_pubsub_ready(&config, true, true));
1219 }
1220}