1use super::*;
2
3impl Node {
4 fn new_dataplane_node() -> DataplaneNode {
5 DataplaneLiveNode::new(AdmissionConfig::new(
6 1024,
7 FIPS_ENDPOINT_DIRECT_PACKET_QUEUE_MAX_PACKETS,
8 ))
9 }
10
11 pub fn new(config: Config) -> Result<Self, NodeError> {
13 config.validate()?;
14 let identity = config.create_identity()?;
15 let node_addr = *identity.node_addr();
16 let is_leaf_only = config.is_leaf_only();
17
18 let mut startup_epoch = [0u8; 8];
19 rand::rng().fill_bytes(&mut startup_epoch);
20
21 let mut bloom_state = if is_leaf_only {
22 BloomState::leaf_only(node_addr)
23 } else {
24 BloomState::new(node_addr)
25 };
26 bloom_state.set_update_debounce_ms(config.node.bloom.update_debounce_ms);
27
28 let tun_state = if config.tun.enabled {
29 TunState::Configured
30 } else {
31 TunState::Disabled
32 };
33
34 let mut tree_state = TreeState::new(node_addr);
36 tree_state.set_parent_hysteresis(config.node.tree.parent_hysteresis);
37 tree_state.set_hold_down(config.node.tree.hold_down_secs);
38 tree_state.set_flap_dampening(
39 config.node.tree.flap_threshold,
40 config.node.tree.flap_window_secs,
41 config.node.tree.flap_dampening_secs,
42 );
43 tree_state
44 .sign_declaration(&identity)
45 .expect("signing own declaration should never fail");
46
47 let coord_cache = CoordCache::new(
48 config.node.cache.coord_size,
49 config.node.cache.coord_ttl_secs * 1000,
50 );
51 let rl = &config.node.rate_limit;
52 let msg1_rate_limiter = HandshakeRateLimiter::with_params(
53 rate_limit::TokenBucket::with_params(rl.handshake_burst, rl.handshake_rate),
54 config.node.limits.max_pending_inbound,
55 );
56
57 let max_connections = config.node.limits.max_connections;
58 let max_peers = config.node.limits.max_peers;
59 let max_links = config.node.limits.max_links;
60 let coords_response_interval_ms = config.node.session.coords_response_interval_ms;
61 let backoff_base_secs = config.node.discovery.backoff_base_secs;
62 let backoff_max_secs = config.node.discovery.backoff_max_secs;
63 let forward_min_interval_secs = config.node.discovery.forward_min_interval_secs;
64
65 let (host_map, peer_acl) = Self::host_map_and_peer_acl(&config);
66 let configured_peers = ConfiguredPeerLookup::from_config(&config);
67 let local_rendezvous =
68 local_rendezvous::LocalRendezvous::new(&config, &identity, startup_epoch);
69
70 Ok(Self {
71 identity,
72 startup_epoch,
73 started_at: std::time::Instant::now(),
74 config,
75 state: NodeState::Created,
76 is_leaf_only,
77 tree_state,
78 bloom_state,
79 coord_cache,
80 learned_routes: LearnedRouteTable::default(),
81 session_direct_degradation: SessionDirectDegradation::default(),
82 recent_requests: RecentDiscoveryRequests::default(),
83 transports: HashMap::new(),
84 #[cfg(feature = "host-ble-transport")]
85 host_ble_io: None,
86 transport_drops: TransportDropTracker::default(),
87 transport_socket_drops: TransportDropTracker::default(),
88 transport_namespace_drops: TransportDropTracker::default(),
89 links: LinkRegistry::default(),
90 packet_tx: None,
91 packet_rx: None,
92 dataplane: Self::new_dataplane_node(),
93 deferred_dataplane_control_turns: VecDeque::new(),
94 deferred_session_forwards: Default::default(),
95 dataplane_fast_ingress_rx: None,
96 dataplane_direct_fsp_sources: Arc::new(HashMap::new()),
97 dataplane_direct_fsp_sources_dirty: true,
98 dataplane_transport_send_batch_packets: DATAPLANE_TRANSPORT_SEND_BATCH_PACKETS,
99 peers: PeerLifecycleRegistry::default(),
100 sessions: SessionRegistry::default(),
101 identity_cache: IdentityCache::default(),
102 pending_session_traffic: PendingSessionTrafficQueues::default(),
103 pending_lookups: handlers::discovery::PendingDiscoveryLookups::default(),
104 max_connections,
105 max_peers,
106 max_links,
107 next_link_id: 1,
108 next_transport_id: 1,
109 stats: stats::NodeStats::new(),
110 stats_history: stats_history::StatsHistory::new(),
111 tun_state,
112 tun_name: None,
113 tun_tx: None,
114 tun_outbound_rx: None,
115 external_packet_tx: None,
116 endpoint_control_rx: None,
117 endpoint_data_rx: None,
118 endpoint_events: EndpointEventRuntime::default(),
119 endpoint_services: EndpointServiceRuntime::default(),
120 tun_reader_handle: None,
121 tun_writer_handle: None,
122 #[cfg(target_os = "macos")]
123 tun_shutdown_fd: None,
124 dns_identity_rx: None,
125 dns_task: None,
126 index_allocator: IndexAllocator::new(),
127 pending_outbound: PendingOutboundHandshakes::default(),
128 msg1_rate_limiter,
129 icmp_rate_limiter: IcmpRateLimiter::new(),
130 routing_error_rate_limiter: RoutingErrorRateLimiter::new(),
131 coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval(
132 std::time::Duration::from_millis(coords_response_interval_ms),
133 ),
134 discovery_backoff: DiscoveryBackoff::with_params(backoff_base_secs, backoff_max_secs),
135 discovery_forward_limiter: DiscoveryForwardRateLimiter::with_interval(
136 std::time::Duration::from_secs(forward_min_interval_secs),
137 ),
138 pending_connects: Vec::new(),
139 retry_pending: retry::PendingRouteRetries::default(),
140 nostr_discovery: None,
141 pending_mesh_signals: HashMap::new(),
142 nostr_discovery_started_at_ms: None,
143 lan_discovery: None,
144 local_rendezvous,
145 startup_open_discovery_sweep_done: false,
146 bootstrap_transports: BootstrapTransports::default(),
147 discovery_fallback_transit: DiscoveryFallbackTransit::default(),
148 last_parent_reeval: None,
149 last_congestion_log: None,
150 estimated_mesh_size: None,
151 last_mesh_size_log: None,
152 last_self_warn: None,
153 local_send_failures: LocalSendFailures::default(),
154 last_rx_loop_maintenance_timeout_at: None,
155 peer_aliases: HashMap::new(),
156 configured_peers,
157 peer_acl,
158 host_map,
159 path_mtu_lookup: Arc::new(std::sync::RwLock::new(HashMap::new())),
160 })
161 }
162
163 pub fn with_identity(identity: Identity, config: Config) -> Result<Self, NodeError> {
168 config.validate()?;
169 let node_addr = *identity.node_addr();
170
171 let mut startup_epoch = [0u8; 8];
172 rand::rng().fill_bytes(&mut startup_epoch);
173
174 let tun_state = if config.tun.enabled {
175 TunState::Configured
176 } else {
177 TunState::Disabled
178 };
179
180 let mut tree_state = TreeState::new(node_addr);
182 tree_state.set_parent_hysteresis(config.node.tree.parent_hysteresis);
183 tree_state.set_hold_down(config.node.tree.hold_down_secs);
184 tree_state.set_flap_dampening(
185 config.node.tree.flap_threshold,
186 config.node.tree.flap_window_secs,
187 config.node.tree.flap_dampening_secs,
188 );
189 tree_state
190 .sign_declaration(&identity)
191 .expect("signing own declaration should never fail");
192
193 let mut bloom_state = BloomState::new(node_addr);
194 bloom_state.set_update_debounce_ms(config.node.bloom.update_debounce_ms);
195
196 let coord_cache = CoordCache::new(
197 config.node.cache.coord_size,
198 config.node.cache.coord_ttl_secs * 1000,
199 );
200 let rl = &config.node.rate_limit;
201 let msg1_rate_limiter = HandshakeRateLimiter::with_params(
202 rate_limit::TokenBucket::with_params(rl.handshake_burst, rl.handshake_rate),
203 config.node.limits.max_pending_inbound,
204 );
205
206 let max_connections = config.node.limits.max_connections;
207 let max_peers = config.node.limits.max_peers;
208 let max_links = config.node.limits.max_links;
209 let coords_response_interval_ms = config.node.session.coords_response_interval_ms;
210
211 let (host_map, peer_acl) = Self::host_map_and_peer_acl(&config);
212 let configured_peers = ConfiguredPeerLookup::from_config(&config);
213 let local_rendezvous =
214 local_rendezvous::LocalRendezvous::new(&config, &identity, startup_epoch);
215
216 Ok(Self {
217 identity,
218 startup_epoch,
219 started_at: std::time::Instant::now(),
220 config,
221 state: NodeState::Created,
222 is_leaf_only: false,
223 tree_state,
224 bloom_state,
225 coord_cache,
226 learned_routes: LearnedRouteTable::default(),
227 session_direct_degradation: SessionDirectDegradation::default(),
228 recent_requests: RecentDiscoveryRequests::default(),
229 transports: HashMap::new(),
230 #[cfg(feature = "host-ble-transport")]
231 host_ble_io: None,
232 transport_drops: TransportDropTracker::default(),
233 transport_socket_drops: TransportDropTracker::default(),
234 transport_namespace_drops: TransportDropTracker::default(),
235 links: LinkRegistry::default(),
236 packet_tx: None,
237 packet_rx: None,
238 dataplane: Self::new_dataplane_node(),
239 deferred_dataplane_control_turns: VecDeque::new(),
240 deferred_session_forwards: Default::default(),
241 dataplane_fast_ingress_rx: None,
242 dataplane_direct_fsp_sources: Arc::new(HashMap::new()),
243 dataplane_direct_fsp_sources_dirty: true,
244 dataplane_transport_send_batch_packets: DATAPLANE_TRANSPORT_SEND_BATCH_PACKETS,
245 peers: PeerLifecycleRegistry::default(),
246 sessions: SessionRegistry::default(),
247 identity_cache: IdentityCache::default(),
248 pending_session_traffic: PendingSessionTrafficQueues::default(),
249 pending_lookups: handlers::discovery::PendingDiscoveryLookups::default(),
250 max_connections,
251 max_peers,
252 max_links,
253 next_link_id: 1,
254 next_transport_id: 1,
255 stats: stats::NodeStats::new(),
256 stats_history: stats_history::StatsHistory::new(),
257 tun_state,
258 tun_name: None,
259 tun_tx: None,
260 tun_outbound_rx: None,
261 external_packet_tx: None,
262 endpoint_control_rx: None,
263 endpoint_data_rx: None,
264 endpoint_events: EndpointEventRuntime::default(),
265 endpoint_services: EndpointServiceRuntime::default(),
266 tun_reader_handle: None,
267 tun_writer_handle: None,
268 #[cfg(target_os = "macos")]
269 tun_shutdown_fd: None,
270 dns_identity_rx: None,
271 dns_task: None,
272 index_allocator: IndexAllocator::new(),
273 pending_outbound: PendingOutboundHandshakes::default(),
274 msg1_rate_limiter,
275 icmp_rate_limiter: IcmpRateLimiter::new(),
276 routing_error_rate_limiter: RoutingErrorRateLimiter::new(),
277 coords_response_rate_limiter: RoutingErrorRateLimiter::with_interval(
278 std::time::Duration::from_millis(coords_response_interval_ms),
279 ),
280 discovery_backoff: DiscoveryBackoff::new(),
281 discovery_forward_limiter: DiscoveryForwardRateLimiter::new(),
282 pending_connects: Vec::new(),
283 retry_pending: retry::PendingRouteRetries::default(),
284 nostr_discovery: None,
285 pending_mesh_signals: HashMap::new(),
286 nostr_discovery_started_at_ms: None,
287 lan_discovery: None,
288 local_rendezvous,
289 startup_open_discovery_sweep_done: false,
290 bootstrap_transports: BootstrapTransports::default(),
291 discovery_fallback_transit: DiscoveryFallbackTransit::default(),
292 last_parent_reeval: None,
293 last_congestion_log: None,
294 estimated_mesh_size: None,
295 last_mesh_size_log: None,
296 last_self_warn: None,
297 local_send_failures: LocalSendFailures::default(),
298 last_rx_loop_maintenance_timeout_at: None,
299 peer_aliases: HashMap::new(),
300 configured_peers,
301 peer_acl,
302 host_map,
303 path_mtu_lookup: Arc::new(std::sync::RwLock::new(HashMap::new())),
304 })
305 }
306
307 pub fn leaf_only(config: Config) -> Result<Self, NodeError> {
309 let mut node = Self::new(config)?;
310 node.is_leaf_only = true;
311 node.bloom_state = BloomState::leaf_only(*node.identity.node_addr());
312 Ok(node)
313 }
314
315 pub(super) fn host_map_and_peer_acl(config: &Config) -> (Arc<HostMap>, acl::PeerAclReloader) {
316 let base_host_map = HostMap::from_peer_configs(config.peers());
317 if !config.node.system_files_enabled {
318 return (
319 Arc::new(base_host_map.clone()),
320 acl::PeerAclReloader::memory_only(base_host_map),
321 );
322 }
323
324 let mut host_map = base_host_map.clone();
325 let hosts_path = std::path::PathBuf::from(crate::upper::hosts::DEFAULT_HOSTS_PATH);
326 let hosts_file = HostMap::load_hosts_file(std::path::Path::new(
327 crate::upper::hosts::DEFAULT_HOSTS_PATH,
328 ));
329 host_map.merge(hosts_file);
330 let peer_acl = acl::PeerAclReloader::with_alias_sources(
331 std::path::PathBuf::from(acl::DEFAULT_PEERS_ALLOW_PATH),
332 std::path::PathBuf::from(acl::DEFAULT_PEERS_DENY_PATH),
333 base_host_map,
334 hosts_path,
335 );
336 (Arc::new(host_map), peer_acl)
337 }
338
339 pub(super) async fn create_transports(&mut self, packet_tx: &PacketTx) -> Vec<TransportHandle> {
343 let mut transports = Vec::new();
344
345 let udp_instances: Vec<_> = self
347 .config
348 .transports
349 .udp
350 .iter()
351 .map(|(name, config)| (name.map(|s| s.to_string()), config.clone()))
352 .collect();
353
354 for (name, udp_config) in udp_instances {
356 let transport_id = self.allocate_transport_id();
357 let udp = UdpTransport::new(transport_id, name, udp_config, packet_tx.clone());
358 transports.push(TransportHandle::Udp(udp));
359 }
360
361 #[cfg(feature = "sim-transport")]
362 {
363 let sim_instances: Vec<_> = self
364 .config
365 .transports
366 .sim
367 .iter()
368 .map(|(name, config)| (name.map(|s| s.to_string()), config.clone()))
369 .collect();
370
371 for (name, sim_config) in sim_instances {
372 let transport_id = self.allocate_transport_id();
373 let sim = crate::transport::sim::SimTransport::new(
374 transport_id,
375 name,
376 sim_config,
377 packet_tx.clone(),
378 );
379 transports.push(TransportHandle::Sim(sim));
380 }
381 }
382
383 #[cfg(any(target_os = "linux", target_os = "macos"))]
385 {
386 let eth_instances: Vec<_> = self
387 .config
388 .transports
389 .ethernet
390 .iter()
391 .map(|(name, config)| (name.map(|s| s.to_string()), config.clone()))
392 .collect();
393 let xonly = self.identity.pubkey();
394 for (name, eth_config) in eth_instances {
395 let mut eth_config = eth_config;
396 if eth_config.discovery_scope.is_none() {
397 eth_config.discovery_scope = self.lan_discovery_scope();
398 }
399 let transport_id = self.allocate_transport_id();
400 let mut eth =
401 EthernetTransport::new(transport_id, name, eth_config, packet_tx.clone());
402 eth.set_local_pubkey(xonly);
403 transports.push(TransportHandle::Ethernet(eth));
404 }
405 }
406
407 let tcp_instances: Vec<_> = self
409 .config
410 .transports
411 .tcp
412 .iter()
413 .map(|(name, config)| (name.map(|s| s.to_string()), config.clone()))
414 .collect();
415
416 for (name, tcp_config) in tcp_instances {
417 let transport_id = self.allocate_transport_id();
418 let tcp = TcpTransport::new(transport_id, name, tcp_config, packet_tx.clone());
419 transports.push(TransportHandle::Tcp(tcp));
420 }
421
422 let websocket_instances: Vec<_> = self
423 .config
424 .transports
425 .websocket
426 .iter()
427 .map(|(name, config)| (name.map(str::to_string), config.clone()))
428 .collect();
429 for (name, websocket_config) in websocket_instances {
430 let transport_id = self.allocate_transport_id();
431 let websocket = WebSocketTransport::new(
432 transport_id,
433 name,
434 websocket_config,
435 packet_tx.clone(),
436 &self.identity,
437 );
438 transports.push(TransportHandle::WebSocket(Box::new(websocket)));
439 }
440
441 let tor_instances: Vec<_> = self
443 .config
444 .transports
445 .tor
446 .iter()
447 .map(|(name, config)| (name.map(|s| s.to_string()), config.clone()))
448 .collect();
449
450 for (name, tor_config) in tor_instances {
451 let transport_id = self.allocate_transport_id();
452 let tor = TorTransport::new(transport_id, name, tor_config, packet_tx.clone());
453 transports.push(TransportHandle::Tor(tor));
454 }
455
456 let webrtc_instances: Vec<_> = self
457 .config
458 .transports
459 .webrtc
460 .iter()
461 .map(|(name, config)| (name.map(|s| s.to_string()), config.clone()))
462 .collect();
463
464 #[cfg(feature = "webrtc-transport")]
465 {
466 for (name, webrtc_config) in webrtc_instances {
467 let transport_id = self.allocate_transport_id();
468 match WebRtcTransport::new(
469 transport_id,
470 name,
471 webrtc_config,
472 packet_tx.clone(),
473 &self.identity,
474 &self.config.node.discovery.nostr,
475 ) {
476 Ok(webrtc) => transports.push(TransportHandle::WebRtc(Box::new(webrtc))),
477 Err(err) => {
478 warn!(
479 transport_id = %transport_id,
480 error = %err,
481 "failed to initialize WebRTC transport"
482 );
483 }
484 }
485 }
486 }
487 #[cfg(not(feature = "webrtc-transport"))]
488 if !webrtc_instances.is_empty() {
489 warn!("WebRTC transport configured but this build lacks WebRTC transport support");
490 }
491
492 #[cfg(feature = "host-ble-transport")]
493 {
494 if let Some(io) = self.host_ble_io.take() {
495 let (name, ble_config) = self
496 .config
497 .transports
498 .ble
499 .iter()
500 .next()
501 .map(|(name, config)| (name.map(str::to_string), config.clone()))
502 .unwrap_or((None, crate::config::BleConfig::default()));
503 let transport_id = self.allocate_transport_id();
504 let mut ble = crate::transport::ble::BleTransport::new(
505 transport_id,
506 name,
507 ble_config,
508 io,
509 packet_tx.clone(),
510 );
511 ble.set_local_pubkey(self.identity.pubkey().serialize());
512 transports.push(TransportHandle::Ble(ble));
513 } else if !self.config.transports.ble.is_empty() {
514 warn!("host BLE configured without a platform adapter");
515 }
516 }
517
518 #[cfg(all(bluer_available, not(feature = "host-ble-transport")))]
520 {
521 let ble_instances: Vec<_> = self
522 .config
523 .transports
524 .ble
525 .iter()
526 .map(|(name, config)| (name.map(|s| s.to_string()), config.clone()))
527 .collect();
528
529 #[cfg(all(bluer_available, not(test)))]
530 for (name, ble_config) in ble_instances {
531 let transport_id = self.allocate_transport_id();
532 let adapter = ble_config.adapter().to_string();
533 let mtu = ble_config.mtu();
534 match crate::transport::ble::io::BluerIo::new(&adapter, mtu).await {
535 Ok(io) => {
536 let mut ble = crate::transport::ble::BleTransport::new(
537 transport_id,
538 name,
539 ble_config,
540 io,
541 packet_tx.clone(),
542 );
543 ble.set_local_pubkey(self.identity.pubkey().serialize());
544 transports.push(TransportHandle::Ble(ble));
545 }
546 Err(e) => {
547 tracing::warn!(adapter = %adapter, error = %e, "failed to initialize BLE adapter");
548 }
549 }
550 }
551
552 #[cfg(any(not(bluer_available), test))]
553 if !ble_instances.is_empty() {
554 #[cfg(not(test))]
555 tracing::warn!("BLE transport configured but this build lacks BlueZ support");
556 }
557 }
558
559 transports
560 }
561
562 #[cfg(feature = "host-ble-transport")]
563 pub(crate) fn set_host_ble_io(&mut self, io: crate::transport::ble::host::HostBleIo) {
564 self.host_ble_io = Some(io);
565 }
566
567 pub(super) fn find_transport_for_type(&self, transport_type: &str) -> Option<TransportId> {
577 self.transports
578 .iter()
579 .filter(|(id, handle)| {
580 handle.transport_type().name == transport_type
581 && handle.is_operational()
582 && !self.bootstrap_transports.contains(id)
583 && !self.is_local_rendezvous_transport(id)
584 })
585 .min_by_key(|(id, _)| id.as_u32())
586 .map(|(id, _)| *id)
587 }
588
589 #[cfg(any(target_os = "linux", target_os = "macos"))]
595 pub(super) fn resolve_ethernet_addr(
596 &self,
597 addr_str: &str,
598 ) -> Result<(TransportId, TransportAddr), NodeError> {
599 let (iface, mac_str) = addr_str.split_once('/').ok_or_else(|| {
600 NodeError::NoTransportForType(format!(
601 "invalid Ethernet address format '{}': expected 'interface/mac'",
602 addr_str
603 ))
604 })?;
605
606 let transport_id = self
608 .transports
609 .iter()
610 .find(|(_, handle)| {
611 handle.transport_type().name == "ethernet"
612 && handle.is_operational()
613 && handle.interface_name() == Some(iface)
614 })
615 .map(|(id, _)| *id)
616 .ok_or_else(|| {
617 NodeError::NoTransportForType(format!(
618 "no operational Ethernet transport for interface '{}'",
619 iface
620 ))
621 })?;
622
623 let mac = crate::transport::ethernet::parse_mac_string(mac_str).map_err(|e| {
624 NodeError::NoTransportForType(format!("invalid MAC in '{}': {}", addr_str, e))
625 })?;
626
627 Ok((transport_id, TransportAddr::from_bytes(&mac)))
628 }
629
630 #[cfg(not(any(target_os = "linux", target_os = "macos")))]
631 pub(super) fn resolve_ethernet_addr(
632 &self,
633 _addr_str: &str,
634 ) -> Result<(TransportId, TransportAddr), NodeError> {
635 Err(NodeError::NoTransportForType(
636 "Ethernet transport is not supported on this platform".to_string(),
637 ))
638 }
639
640 #[cfg(bluer_available)]
644 pub(super) fn resolve_ble_addr(
645 &self,
646 addr_str: &str,
647 ) -> Result<(TransportId, TransportAddr), NodeError> {
648 let ta = TransportAddr::from_string(addr_str);
649 let adapter = crate::transport::ble::addr::adapter_from_addr(&ta).ok_or_else(|| {
650 NodeError::NoTransportForType(format!(
651 "invalid BLE address format '{}': expected 'adapter/mac'",
652 addr_str
653 ))
654 })?;
655
656 let transport_id = self
658 .transports
659 .iter()
660 .find(|(_, handle)| handle.transport_type().name == "ble" && handle.is_operational())
661 .map(|(id, _)| *id)
662 .ok_or_else(|| {
663 NodeError::NoTransportForType(format!(
664 "no operational BLE transport for adapter '{}'",
665 adapter
666 ))
667 })?;
668
669 crate::transport::ble::addr::BleAddr::parse(addr_str).map_err(|e| {
671 NodeError::NoTransportForType(format!("invalid BLE address '{}': {}", addr_str, e))
672 })?;
673
674 Ok((transport_id, TransportAddr::from_string(addr_str)))
675 }
676}