1pub mod addr;
22pub mod bootstrap;
23pub mod discovery;
24pub mod framing;
25#[cfg(feature = "host-ble-transport")]
26pub mod host;
27pub mod io;
28pub mod pool;
29pub mod stats;
30mod tasks;
31
32use tasks::{
33 AcceptLoopContext, ScanProbeContext, accept_loop, attach_receive_loop, pubkey_exchange,
34 scan_probe_loop,
35};
36
37use super::{
38 ConnectionState, DiscoveredPeer, PacketTx, Transport, TransportAddr, TransportError,
39 TransportId, TransportState, TransportType,
40};
41use crate::config::BleConfig;
42use crate::identity::NodeAddr;
43use addr::BleAddr;
44use discovery::DiscoveryBuffer;
45use framing::FramedBleStream;
46use io::{BleAcceptor, BleIo, BleStream};
47use pool::{BleConnection, ConnectionPool};
48use stats::BleStats;
49
50use secp256k1::XOnlyPublicKey;
51use std::collections::HashMap;
52use std::sync::Arc;
53use tokio::sync::Mutex;
54use tokio::task::JoinHandle;
55use tracing::{debug, info, warn};
56
57pub(super) type SharedBlePool<S> = Arc<Mutex<ConnectionPool<Arc<FramedBleStream<S>>>>>;
58
59pub const DEFAULT_PSM: u16 = 0x0085;
63
64#[cfg(feature = "host-ble-transport")]
69pub type DefaultBleTransport = BleTransport<host::HostBleIo>;
70
71#[cfg(all(not(feature = "host-ble-transport"), bluer_available, not(test)))]
72pub type DefaultBleTransport = BleTransport<io::BluerIo>;
73
74#[cfg(all(not(feature = "host-ble-transport"), any(not(bluer_available), test)))]
75pub type DefaultBleTransport = BleTransport<io::MockBleIo>;
76
77pub struct BleTransport<I: BleIo> {
87 transport_id: TransportId,
89 name: Option<String>,
91 config: BleConfig,
93 state: TransportState,
95 io: Arc<I>,
97 pool: SharedBlePool<I::Stream>,
99 connecting: Arc<Mutex<HashMap<TransportAddr, ConnectingEntry>>>,
101 packet_tx: PacketTx,
103 accept_task: Option<JoinHandle<()>>,
105 scan_probe_task: Option<JoinHandle<()>>,
107 discovery_buffer: Arc<DiscoveryBuffer>,
109 stats: Arc<BleStats>,
111 local_pubkey: Option<[u8; 32]>,
118}
119
120struct ConnectingEntry {
122 task: JoinHandle<()>,
123}
124
125impl<I: BleIo> BleTransport<I> {
126 pub fn new(
128 transport_id: TransportId,
129 name: Option<String>,
130 config: BleConfig,
131 io: I,
132 packet_tx: PacketTx,
133 ) -> Self {
134 let max_conns = config.max_connections();
135 Self {
136 transport_id,
137 name,
138 config,
139 state: TransportState::Configured,
140 io: Arc::new(io),
141 pool: Arc::new(Mutex::new(ConnectionPool::new(max_conns))),
142 connecting: Arc::new(Mutex::new(HashMap::new())),
143 packet_tx,
144 accept_task: None,
145 scan_probe_task: None,
146 discovery_buffer: Arc::new(DiscoveryBuffer::new(transport_id)),
147 stats: Arc::new(BleStats::new()),
148 local_pubkey: None,
149 }
150 }
151
152 pub fn name(&self) -> Option<&str> {
154 self.name.as_deref()
155 }
156
157 pub fn stats(&self) -> &Arc<BleStats> {
159 &self.stats
160 }
161
162 pub fn io(&self) -> &Arc<I> {
164 &self.io
165 }
166
167 pub fn set_local_pubkey(&mut self, pubkey: [u8; 32]) {
173 self.local_pubkey = Some(pubkey);
174 }
175
176 pub async fn start_async(&mut self) -> Result<(), TransportError> {
178 if !self.state.can_start() {
179 return Err(TransportError::AlreadyStarted);
180 }
181 self.state = TransportState::Starting;
182
183 let preferred_psm = self.config.psm();
184 let mut listener_psm = preferred_psm;
185 let adapter = self.io.adapter_name().to_string();
186
187 let local_node_addr = self.local_pubkey.and_then(|pk| {
189 XOnlyPublicKey::from_slice(&pk)
190 .ok()
191 .map(|xonly| NodeAddr::from_pubkey(&xonly))
192 });
193
194 if self.config.accept_connections() {
196 match self.io.listen(preferred_psm).await {
197 Ok(acceptor) => {
198 listener_psm = acceptor.psm();
199 self.accept_task = Some(tokio::spawn(accept_loop(
200 acceptor,
201 AcceptLoopContext {
202 pool: Arc::clone(&self.pool),
203 packet_tx: self.packet_tx.clone(),
204 transport_id: self.transport_id,
205 stats: Arc::clone(&self.stats),
206 local_pubkey: self.local_pubkey,
207 discovery_buffer: Arc::clone(&self.discovery_buffer),
208 local_node_addr,
209 max_packet: self.config.mtu(),
210 },
211 )));
212 debug!(adapter = %adapter, psm = listener_psm, "BLE accept loop started");
213 }
214 Err(e) => {
215 warn!(adapter = %adapter, error = %e, "failed to start BLE listener");
216 self.state = TransportState::Failed;
217 return Err(e);
218 }
219 }
220 }
221
222 if self.config.advertise() {
224 let bootstrap = crate::transport::ble::bootstrap::BleBootstrap::new(
225 listener_psm,
226 self.config.mtu(),
227 )
228 .map_err(|error| TransportError::StartFailed(error.to_string()))?;
229 if let Err(e) = self.io.start_advertising(bootstrap).await {
230 warn!(adapter = %adapter, error = %e, "failed to start BLE advertising");
231 } else {
232 self.stats.record_advertisement();
233 debug!(adapter = %adapter, "BLE advertising started (continuous)");
234 }
235 }
236
237 if self.config.scan() {
239 match self.io.start_scanning().await {
240 Ok(scanner) => {
241 self.scan_probe_task = Some(tokio::spawn(scan_probe_loop::<I>(
242 scanner,
243 ScanProbeContext {
244 io: Arc::clone(&self.io),
245 pool: Arc::clone(&self.pool),
246 buffer: Arc::clone(&self.discovery_buffer),
247 stats: Arc::clone(&self.stats),
248 local_pubkey: self.local_pubkey,
249 connect_timeout_ms: self.config.connect_timeout_ms(),
250 cooldown_secs: self.config.probe_cooldown_secs(),
251 local_node_addr,
252 packet_tx: self.packet_tx.clone(),
253 transport_id: self.transport_id,
254 max_packet: self.config.mtu(),
255 },
256 )));
257 debug!(adapter = %adapter, "BLE scan+probe loop started");
258 }
259 Err(e) => {
260 warn!(adapter = %adapter, error = %e, "failed to start BLE scanning");
261 }
262 }
263 }
264
265 self.state = TransportState::Up;
266 info!(adapter = %adapter, psm = listener_psm, "BLE transport started");
267 Ok(())
268 }
269
270 pub async fn stop_async(&mut self) -> Result<(), TransportError> {
272 let _ = self.io.stop_advertising().await;
274
275 if let Some(task) = self.accept_task.take() {
277 task.abort();
278 }
279
280 if let Some(task) = self.scan_probe_task.take() {
282 task.abort();
283 }
284
285 {
287 let mut connecting = self.connecting.lock().await;
288 for (_, entry) in connecting.drain() {
289 entry.task.abort();
290 }
291 }
292
293 {
295 let mut pool = self.pool.lock().await;
296 for addr in pool.addrs() {
297 pool.remove(&addr);
298 }
299 }
300
301 self.state = TransportState::Down;
302 info!("BLE transport stopped");
303 Ok(())
304 }
305
306 pub async fn send_async(
313 &self,
314 addr: &TransportAddr,
315 data: &[u8],
316 ) -> Result<usize, TransportError> {
317 let pool = self.pool.lock().await;
318 let conn = match pool.get(addr) {
319 Some(c) => c,
320 None => {
321 drop(pool);
323 let _ = self.connect_async(addr).await;
325 return Err(TransportError::SendFailed("not connected".into()));
326 }
327 };
328
329 let mtu = conn.effective_mtu() as usize;
331 if data.len() > mtu {
332 self.stats.record_mtu_exceeded();
333 return Err(TransportError::MtuExceeded {
334 packet_size: data.len(),
335 mtu: mtu as u16,
336 });
337 }
338
339 match conn.stream.send(data).await {
340 Ok(()) => {
341 self.stats.record_send(data.len());
342 Ok(data.len())
343 }
344 Err(e) => {
345 self.stats.record_send_error();
346 drop(pool);
348 let mut pool = self.pool.lock().await;
349 pool.remove(addr);
350 warn!(addr = %addr, error = %e, "BLE send failed, connection removed");
351 Err(e)
352 }
353 }
354 }
355
356 pub async fn connect_async(&self, addr: &TransportAddr) -> Result<(), TransportError> {
361 let pool_guard = self.pool.lock().await;
362 if pool_guard.contains(addr) {
363 return Ok(());
364 }
365 let active_connections = pool_guard.len();
366 let max_connections = pool_guard.max_connections();
367 let mut connecting_guard = self.connecting.lock().await;
368 if connecting_guard.contains_key(addr) {
369 return Ok(());
370 }
371 if active_connections.saturating_add(connecting_guard.len()) >= max_connections {
372 return Err(TransportError::ConnectionRefused);
373 }
374 drop(pool_guard);
375
376 let ble_addr = BleAddr::parse(
377 addr.as_str()
378 .ok_or_else(|| TransportError::InvalidAddress("not valid UTF-8".into()))?,
379 )?;
380
381 let io = Arc::clone(&self.io);
382 let pool = Arc::clone(&self.pool);
383 let connecting = Arc::clone(&self.connecting);
384 let packet_tx = self.packet_tx.clone();
385 let transport_id = self.transport_id;
386 let stats = Arc::clone(&self.stats);
387 let advertised_bootstrap = self.discovery_buffer.bootstrap_for(&ble_addr);
388 let psm = advertised_bootstrap
389 .map(|bootstrap| bootstrap.psm)
390 .unwrap_or_else(|| self.config.psm());
391 let timeout_ms = self.config.connect_timeout_ms();
392 let addr_clone = addr.clone();
393 let local_pubkey = self.local_pubkey;
394 let discovery_buffer = Arc::clone(&self.discovery_buffer);
395 let max_packet = advertised_bootstrap
396 .map(|bootstrap| self.config.mtu().min(bootstrap.max_packet))
397 .unwrap_or_else(|| self.config.mtu());
398 let (start_tx, start_rx) = tokio::sync::oneshot::channel();
399
400 let task = tokio::spawn(async move {
401 if start_rx.await.is_err() {
402 return;
403 }
404 async {
405 let result = tokio::time::timeout(
406 std::time::Duration::from_millis(timeout_ms),
407 io.connect(&ble_addr, psm),
408 )
409 .await;
410
411 match result {
412 Ok(Ok(stream)) => {
413 let stream = FramedBleStream::new(stream, max_packet);
414 if let Some(ref our_pubkey) = local_pubkey {
416 match pubkey_exchange(&stream, our_pubkey).await {
417 Ok(peer_pubkey) => {
418 debug!(addr = %addr_clone, "BLE outbound pubkey exchange complete");
419 discovery_buffer.add_peer_with_pubkey(&ble_addr, peer_pubkey);
420 }
421 Err(e) => {
422 warn!(
423 addr = %addr_clone, error = %e,
424 "BLE outbound pubkey exchange failed"
425 );
426 return;
427 }
428 }
429 }
430
431 let send_mtu = stream.send_mtu();
432 let recv_mtu = stream.recv_mtu();
433 let stream = Arc::new(stream);
434 let conn = BleConnection {
435 stream: Arc::clone(&stream),
436 recv_task: None,
437 send_mtu,
438 recv_mtu,
439 established_at: tokio::time::Instant::now(),
440 is_static: false,
441 addr: ble_addr,
442 };
443
444 match pool.lock().await.insert(addr_clone.clone(), conn) {
445 Ok(Some(evicted)) => {
446 stats.record_pool_eviction();
447 debug!(addr = %addr_clone, evicted = %evicted, "BLE connection established (evicted peer)");
448 }
449 Ok(None) => {
450 debug!(addr = %addr_clone, "BLE connection established");
451 }
452 Err(e) => {
453 warn!(addr = %addr_clone, error = %e, "BLE pool full, connection dropped");
454 stats.record_connection_rejected();
455 return;
456 }
457 }
458 if !attach_receive_loop(
459 stream,
460 addr_clone.clone(),
461 pool,
462 packet_tx,
463 transport_id,
464 Arc::clone(&stats),
465 recv_mtu,
466 )
467 .await
468 {
469 return;
470 }
471 stats.record_connection_established();
472 }
473 Ok(Err(e)) => {
474 debug!(addr = %addr_clone, error = %e, "BLE connect failed");
475 }
476 Err(_) => {
477 stats.record_connect_timeout();
478 debug!(addr = %addr_clone, "BLE connect timeout");
479 }
480 }
481 }
482 .await;
483 connecting.lock().await.remove(&addr_clone);
484 });
485
486 connecting_guard.insert(addr.clone(), ConnectingEntry { task });
487 drop(connecting_guard);
488 let _ = start_tx.send(());
489
490 Ok(())
491 }
492
493 pub fn connection_state_sync(&self, addr: &TransportAddr) -> ConnectionState {
495 if let Ok(pool) = self.pool.try_lock()
497 && pool.contains(addr)
498 {
499 return ConnectionState::Connected;
500 }
501
502 if let Ok(connecting) = self.connecting.try_lock()
504 && connecting.contains_key(addr)
505 {
506 return ConnectionState::Connecting;
507 }
508
509 ConnectionState::None
510 }
511
512 pub async fn close_connection_async(&self, addr: &TransportAddr) {
514 let mut pool = self.pool.lock().await;
515 if let Some(conn) = pool.remove(addr) {
516 debug!(addr = %addr, "BLE connection closed");
517 drop(conn); }
519 }
520
521 pub fn link_mtu(&self, addr: &TransportAddr) -> u16 {
523 if let Ok(pool) = self.pool.try_lock()
524 && let Some(conn) = pool.get(addr)
525 {
526 return conn.effective_mtu();
527 }
528 self.config.mtu()
529 }
530}
531
532impl<I: BleIo> Transport for BleTransport<I> {
533 fn transport_id(&self) -> TransportId {
534 self.transport_id
535 }
536
537 fn transport_type(&self) -> &TransportType {
538 &TransportType::BLE
539 }
540
541 fn state(&self) -> TransportState {
542 self.state
543 }
544
545 fn mtu(&self) -> u16 {
546 self.config.mtu()
547 }
548
549 fn link_mtu(&self, addr: &TransportAddr) -> u16 {
550 self.link_mtu(addr)
551 }
552
553 fn start(&mut self) -> Result<(), TransportError> {
554 Err(TransportError::NotSupported(
555 "use start_async() for BLE transport".into(),
556 ))
557 }
558
559 fn stop(&mut self) -> Result<(), TransportError> {
560 Err(TransportError::NotSupported(
561 "use stop_async() for BLE transport".into(),
562 ))
563 }
564
565 fn send(&self, _addr: &TransportAddr, _data: &[u8]) -> Result<(), TransportError> {
566 Err(TransportError::NotSupported(
567 "use send_async() for BLE transport".into(),
568 ))
569 }
570
571 fn discover(&self) -> Result<Vec<DiscoveredPeer>, TransportError> {
572 Ok(self.discovery_buffer.take())
573 }
574
575 fn auto_connect(&self) -> bool {
576 self.config.auto_connect()
577 }
578
579 fn accept_connections(&self) -> bool {
580 self.config.accept_connections()
581 }
582
583 fn close_connection(&self, _addr: &TransportAddr) {
584 }
586}
587
588#[cfg(test)]
593mod tests;