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 provider_endpoint = local_only_endpoint("htree-provider-test").await;
777 let observer_endpoint = local_only_endpoint("iris-drive-observer-test").await;
778 let temp = tempfile::tempdir().unwrap();
779 let store = Arc::new(HashtreeStore::new(temp.path()).unwrap());
780 let data = b"served directly from htree StorageRouter".to_vec();
781 let hash = Sha256::digest(&data).into();
782 store.store_arc().put(hash, data.clone()).await.unwrap();
783 let (blob_resolver, provider) = bind_daemon_blob_resolver(
784 &provider_endpoint,
785 store.store_arc(),
786 &[],
787 Duration::from_secs(1),
788 )
789 .await
790 .expect("bind htree provider");
791 let handle = DaemonFipsHandle {
792 endpoint: provider_endpoint.native_endpoint.clone(),
793 endpoint_npub: provider_endpoint.local_peer_id.clone(),
794 discovery_scope: provider_endpoint.discovery_scope.clone(),
795 pubsub_client: None,
796 blob_resolver,
797 blob_transport: Mutex::new(Some(provider)),
798 };
799
800 timeout(Duration::from_secs(10), async {
801 loop {
802 if provider_advertised(&observer_endpoint, &provider_endpoint.local_peer_id) {
803 break;
804 }
805 tokio::time::sleep(Duration::from_millis(20)).await;
806 }
807 })
808 .await
809 .expect("htree provider capability did not converge");
810
811 let observer_store = Arc::new(hashtree_core::MemoryStore::new());
812 let observer_transport = Arc::new(
813 TcpBlobTransport::bind_route(
814 observer_endpoint.native_endpoint.clone(),
815 observer_store.clone(),
816 Arc::new(StoreBlobRoute::new(observer_store.clone())),
817 )
818 .await
819 .unwrap(),
820 );
821 let observer_fips_route = FipsBlobRoute::discovered(
822 observer_endpoint.native_endpoint.clone(),
823 observer_transport.clone(),
824 4,
825 )
826 .unwrap();
827 let observer_router = BlobRouter::new(
828 vec![BlobRouteEntry::new(
829 "same-host-fips",
830 Arc::new(observer_fips_route),
831 )],
832 Some(observer_store),
833 BlobRouterConfig::default(),
834 )
835 .unwrap();
836 assert_eq!(observer_router.get(&hash, None).await.unwrap(), Some(data));
837
838 drop(observer_router);
839 drop(observer_transport);
840 handle.shutdown().await;
841 timeout(Duration::from_secs(10), async {
842 while provider_advertised(&observer_endpoint, &provider_endpoint.local_peer_id) {
843 tokio::time::sleep(Duration::from_millis(20)).await;
844 }
845 })
846 .await
847 .expect("htree provider capability survived daemon shutdown");
848 drop(handle);
849 observer_endpoint.native_endpoint.shutdown().await.unwrap();
850 }
851
852 #[tokio::test]
853 async fn daemon_mesh_forwarding_observes_two_one_zero_and_exhaustion() {
854 let scope = format!("htree-inbound-forwarding-{}", uuid::Uuid::new_v4());
855 let (upstream_endpoint, upstream_addr) = udp_endpoint(&scope).await;
856 let (provider_endpoint, provider_addr) = udp_endpoint(&scope).await;
857 let (observer_endpoint, observer_addr) = udp_endpoint(&scope).await;
858 let temp = tempfile::tempdir().unwrap();
859 let provider_store = Arc::new(HashtreeStore::new(temp.path().join("provider")).unwrap());
860 let upstream_store = Arc::new(hashtree_core::MemoryStore::new());
861 let data = b"served through the daemon's configured FIPS peer".to_vec();
862 let hash = Sha256::digest(&data).into();
863 upstream_store.put(hash, data.clone()).await.unwrap();
864 let upstream_requests = Arc::new(std::sync::Mutex::new(Vec::new()));
865 let upstream_transport = Arc::new(
866 TcpBlobTransport::bind_route(
867 upstream_endpoint.native_endpoint.clone(),
868 upstream_store.clone(),
869 Arc::new(RecordingStoreRoute {
870 store: upstream_store,
871 requests: upstream_requests.clone(),
872 }),
873 )
874 .await
875 .unwrap(),
876 );
877 let upstream_peer = FipsPeerConfig {
878 npub: upstream_endpoint.local_peer_id.clone(),
879 udp_addresses: vec![upstream_addr],
880 };
881 set_fips_peer_configs(
882 provider_endpoint.native_endpoint.as_ref(),
883 vec![
884 upstream_peer.clone(),
885 FipsPeerConfig {
886 npub: observer_endpoint.local_peer_id.clone(),
887 udp_addresses: vec![observer_addr],
888 },
889 ],
890 )
891 .await
892 .unwrap();
893 set_fips_peer_configs(
894 upstream_endpoint.native_endpoint.as_ref(),
895 vec![FipsPeerConfig {
896 npub: provider_endpoint.local_peer_id.clone(),
897 udp_addresses: vec![provider_addr.clone()],
898 }],
899 )
900 .await
901 .unwrap();
902 set_fips_peer_configs(
903 observer_endpoint.native_endpoint.as_ref(),
904 vec![FipsPeerConfig {
905 npub: provider_endpoint.local_peer_id.clone(),
906 udp_addresses: vec![provider_addr],
907 }],
908 )
909 .await
910 .unwrap();
911 let (provider_resolver, provider_transport) = bind_daemon_blob_resolver(
912 &provider_endpoint,
913 provider_store.store_arc(),
914 &[upstream_peer],
915 Duration::from_secs(2),
916 )
917 .await
918 .unwrap();
919 let observer_store = Arc::new(hashtree_core::MemoryStore::new());
920 let observer_transport = Arc::new(
921 TcpBlobTransport::bind_route(
922 observer_endpoint.native_endpoint.clone(),
923 observer_store.clone(),
924 Arc::new(StoreBlobRoute::new(observer_store)),
925 )
926 .await
927 .unwrap(),
928 );
929 let observer_fips_route = FipsBlobRoute::explicit(
930 observer_transport.clone(),
931 vec![PeerIdentity::from_npub(&provider_endpoint.local_peer_id).unwrap()],
932 1,
933 )
934 .unwrap();
935 let observer_route = MeshForwardingRoute::new(Arc::new(observer_fips_route));
936
937 wait_for_native_peer(
938 provider_endpoint.native_endpoint.as_ref(),
939 &upstream_endpoint.local_peer_id,
940 )
941 .await;
942 wait_for_native_peer(
943 observer_endpoint.native_endpoint.as_ref(),
944 &provider_endpoint.local_peer_id,
945 )
946 .await;
947 assert_eq!(
948 observer_route
949 .route(hashtree_core::BlobRequest { hash, htl: 0 })
950 .await
951 .unwrap(),
952 hashtree_core::BlobReply::NoResult,
953 );
954 assert!(upstream_requests.lock().unwrap().is_empty());
955 assert_eq!(
956 observer_route
957 .route(hashtree_core::BlobRequest { hash, htl: 1 })
958 .await
959 .unwrap(),
960 hashtree_core::BlobReply::NoResult,
961 );
962 assert!(
963 upstream_requests.lock().unwrap().is_empty(),
964 "HTL 1 reaches the intermediate daemon as 0 and cannot be forwarded"
965 );
966 assert_eq!(
967 observer_route
968 .route(hashtree_core::BlobRequest { hash, htl: 2 })
969 .await
970 .unwrap(),
971 hashtree_core::BlobReply::Data(data),
972 );
973 assert_eq!(
974 upstream_requests.lock().unwrap().as_slice(),
975 &[hashtree_core::BlobRequest { hash, htl: 0 }],
976 "three Hashtree nodes must observe the two-to-one-to-zero forwarding budget while each FIPS carrier preserves its request",
977 );
978
979 drop(observer_route);
980 drop(observer_transport);
981 drop(provider_resolver);
982 drop(provider_transport);
983 drop(upstream_transport);
984 observer_endpoint.native_endpoint.shutdown().await.unwrap();
985 provider_endpoint.native_endpoint.shutdown().await.unwrap();
986 upstream_endpoint.native_endpoint.shutdown().await.unwrap();
987 }
988
989 fn provider_advertised(endpoint: &BoundFipsEndpoint, npub: &str) -> bool {
990 endpoint
991 .native_endpoint
992 .local_instance_advertisements()
993 .unwrap()
994 .iter()
995 .any(|advert| {
996 advert.npub == npub
997 && advert
998 .capability(hashtree_fips_transport::TCP_BLOB_CAPABILITY)
999 .and_then(|capability| capability.fsp_port)
1000 == Some(hashtree_fips_transport::TCP_BLOB_SERVICE_PORT)
1001 })
1002 }
1003
1004 async fn local_only_endpoint(scope: &str) -> hashtree_fips_transport::BoundFipsEndpoint {
1005 let keys = nostr::Keys::generate();
1006 let mut options = FipsEndpointOptions::new(keys.secret_key().to_bech32().unwrap());
1007 options.discovery_scope = scope.to_string();
1008 options.enable_udp = false;
1009 options.enable_webrtc = false;
1010 options.enable_local_rendezvous = true;
1011 options.enable_lan_discovery = false;
1012 options.share_local_candidates = false;
1013 bind_fips_endpoint(options).await.unwrap()
1014 }
1015
1016 async fn udp_endpoint(scope: &str) -> (BoundFipsEndpoint, String) {
1017 let addr = reserve_udp_addr();
1018 let keys = nostr::Keys::generate();
1019 let mut options = FipsEndpointOptions::new(keys.secret_key().to_bech32().unwrap());
1020 options.discovery_scope = scope.to_string();
1021 options.enable_udp = true;
1022 options.udp_bind_addr = Some(addr.clone());
1023 options.enable_webrtc = false;
1024 options.enable_local_rendezvous = false;
1025 options.enable_lan_discovery = false;
1026 options.share_local_candidates = false;
1027 let endpoint = bind_fips_endpoint(options).await.unwrap();
1028 (endpoint, addr)
1029 }
1030
1031 fn reserve_udp_addr() -> String {
1032 let socket = std::net::UdpSocket::bind("127.0.0.1:0").unwrap();
1033 socket.local_addr().unwrap().to_string()
1034 }
1035
1036 #[cfg(feature = "experimental-decentralized-pubsub")]
1037 async fn wait_for_peer(endpoint: &BoundFipsEndpoint, npub: &str) {
1038 wait_for_native_peer(endpoint.native_endpoint.as_ref(), npub).await;
1039 }
1040
1041 async fn wait_for_native_peer(endpoint: &FipsEndpoint, npub: &str) {
1042 let result = timeout(Duration::from_secs(10), async {
1043 loop {
1044 if endpoint
1045 .peers()
1046 .await
1047 .unwrap()
1048 .iter()
1049 .any(|peer| peer.npub == npub && peer.connected)
1050 {
1051 break;
1052 }
1053 tokio::time::sleep(Duration::from_millis(20)).await;
1054 }
1055 })
1056 .await;
1057 if result.is_err() {
1058 panic!(
1059 "FIPS peer {npub} did not connect; peers={:?}",
1060 endpoint.peers().await
1061 );
1062 }
1063 }
1064
1065 #[cfg(feature = "experimental-decentralized-pubsub")]
1066 async fn trusted_test_relay(
1067 root: &std::path::Path,
1068 allowed: &[nostr::PublicKey],
1069 ) -> Arc<NostrRelay> {
1070 let graph_dir = root.join("graph");
1071 let data_dir = root.join("nostr");
1072 std::fs::create_dir_all(&data_dir).unwrap();
1073 let graph = {
1074 let _guard = crate::socialgraph::test_lock().await;
1075 crate::socialgraph::open_test_social_graph_store_with_mapsize(
1076 &graph_dir,
1077 Some(128 * 1024 * 1024),
1078 )
1079 .unwrap()
1080 };
1081 let backend: Arc<dyn crate::socialgraph::SocialGraphBackend> = graph;
1082 let allowed = allowed
1083 .iter()
1084 .map(|pubkey| pubkey.to_hex())
1085 .collect::<HashSet<_>>();
1086 let access = Arc::new(crate::socialgraph::SocialGraphAccessControl::new(
1087 Arc::clone(&backend),
1088 0,
1089 allowed.clone(),
1090 ));
1091 Arc::new(
1092 NostrRelay::new(
1093 backend,
1094 data_dir,
1095 allowed,
1096 Some(access),
1097 crate::nostr_relay::NostrRelayConfig {
1098 spambox_db_max_bytes: 0,
1099 ..Default::default()
1100 },
1101 )
1102 .unwrap(),
1103 )
1104 }
1105
1106 #[cfg(feature = "experimental-decentralized-pubsub")]
1107 fn root_event(keys: &nostr::Keys, name: &str, hash: &str) -> nostr::Event {
1108 EventBuilder::new(Kind::Custom(30_078), hash)
1109 .tags([
1110 Tag::parse(["d", name]).unwrap(),
1111 Tag::parse(["l", "hashtree"]).unwrap(),
1112 Tag::parse(["hash", hash]).unwrap(),
1113 ])
1114 .sign_with_keys(keys)
1115 .unwrap()
1116 }
1117
1118 #[cfg(feature = "experimental-decentralized-pubsub")]
1119 async fn wait_for_relay_event(relay: &NostrRelay, id: nostr::EventId) {
1120 timeout(Duration::from_secs(5), async {
1121 loop {
1122 if relay
1123 .query_events(&Filter::new().id(id), 1)
1124 .await
1125 .iter()
1126 .any(|event| event.id == id)
1127 {
1128 break;
1129 }
1130 tokio::time::sleep(Duration::from_millis(20)).await;
1131 }
1132 })
1133 .await
1134 .expect("trusted decentralized event was not indexed");
1135 }
1136
1137 #[test]
1138 fn daemon_fips_peer_configs_prefer_configured_peers() {
1139 let mut config = Config::default();
1140 config.server.fips_peers = vec![
1141 crate::config::ConfiguredFipsPeer {
1142 npub: " origin ".to_string(),
1143 udp_addresses: vec![" udp:192.0.2.10:2121 ".to_string(), " ".to_string()],
1144 },
1145 crate::config::ConfiguredFipsPeer {
1146 npub: "origin".to_string(),
1147 udp_addresses: vec!["udp:ignored:2121".to_string()],
1148 },
1149 ];
1150
1151 let peers = daemon_fips_peer_configs(
1152 &config,
1153 vec![
1154 "origin".to_string(),
1155 " followed ".to_string(),
1156 " ".to_string(),
1157 ],
1158 );
1159
1160 assert_eq!(
1161 peers,
1162 vec![
1163 FipsPeerConfig {
1164 npub: "origin".to_string(),
1165 udp_addresses: vec!["udp:192.0.2.10:2121".to_string()],
1166 },
1167 FipsPeerConfig::new("followed"),
1168 ]
1169 );
1170 }
1171
1172 #[test]
1173 fn blob_carrier_deadline_outlives_the_complete_resolver_budget() {
1174 for resolver_timeout in [Duration::from_millis(1), Duration::from_secs(20)] {
1175 assert!(blob_transport_config(resolver_timeout).idle_timeout > resolver_timeout);
1176 }
1177 }
1178
1179 #[cfg(feature = "experimental-decentralized-pubsub")]
1180 #[test]
1181 fn daemon_decentralized_pubsub_requires_config_fips_and_relay() {
1182 let mut config = Config::default();
1183 assert!(!daemon_decentralized_pubsub_ready(&config, true, true));
1184
1185 config.nostr.decentralized_pubsub = true;
1186 assert!(daemon_decentralized_pubsub_ready(&config, true, true));
1187 assert!(!daemon_decentralized_pubsub_ready(&config, false, true));
1188 assert!(!daemon_decentralized_pubsub_ready(&config, true, false));
1189
1190 config.nostr.enabled = false;
1191 assert!(!daemon_decentralized_pubsub_ready(&config, true, true));
1192 }
1193}