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