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