use crate::error::{Error, Result};
use crate::network::Message;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::time::{Duration, Instant};
#[cfg(feature = "ble")]
use btleplug::api::{
Central, Characteristic, Manager as BtManager, Peripheral as BtPeripheral, ScanFilter,
WriteType,
};
#[cfg(feature = "ble")]
use btleplug::platform::{Adapter, Manager, Peripheral};
#[cfg(feature = "ble")]
use smol::channel::{Receiver, Sender};
#[cfg(feature = "ble")]
use uuid::Uuid;
#[cfg(feature = "ble-esp32")]
use esp32_nimble::{uuid128, BLEAdvertisedDevice, BLEClient, BLEDevice, BLEScan};
#[cfg(feature = "ble-esp32")]
use std::sync::Arc;
pub const AINGLE_SERVICE_UUID: &str = "6e400001-b5a3-f393-e0a9-e50e24dcca9e";
pub const TX_CHAR_UUID: &str = "6e400002-b5a3-f393-e0a9-e50e24dcca9e";
pub const RX_CHAR_UUID: &str = "6e400003-b5a3-f393-e0a9-e50e24dcca9e";
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BleConfig {
pub device_name: String,
pub mesh_relay: bool,
pub tx_power: i8,
pub scan_interval_ms: u32,
pub scan_window_ms: u32,
pub advertising_interval_ms: u32,
pub connection_timeout: Duration,
pub max_connections: usize,
pub passive_scan: bool,
}
impl Default for BleConfig {
fn default() -> Self {
Self {
device_name: "AIngle-Node".to_string(),
mesh_relay: true,
tx_power: 0, scan_interval_ms: 100,
scan_window_ms: 50,
advertising_interval_ms: 100,
connection_timeout: Duration::from_secs(10),
max_connections: 4,
passive_scan: false,
}
}
}
impl BleConfig {
pub fn low_power() -> Self {
Self {
device_name: "AIngle-LP".to_string(),
mesh_relay: false, tx_power: -12, scan_interval_ms: 1000,
scan_window_ms: 100,
advertising_interval_ms: 500,
connection_timeout: Duration::from_secs(30),
max_connections: 2,
passive_scan: true,
}
}
pub fn mesh_relay() -> Self {
Self {
device_name: "AIngle-Relay".to_string(),
mesh_relay: true,
tx_power: 4, scan_interval_ms: 50,
scan_window_ms: 30,
advertising_interval_ms: 50,
connection_timeout: Duration::from_secs(5),
max_connections: 8,
passive_scan: false,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum BleState {
Uninitialized,
Idle,
Scanning,
Advertising,
Connected,
Error,
}
#[derive(Debug, Clone)]
pub struct BlePeer {
pub address: String,
pub name: Option<String>,
pub rssi: i16,
pub discovered_at: Instant,
pub last_seen: Instant,
pub supports_aingle: bool,
pub connected: bool,
}
impl BlePeer {
pub fn new(address: &str, rssi: i16) -> Self {
let now = Instant::now();
Self {
address: address.to_string(),
name: None,
rssi,
discovered_at: now,
last_seen: now,
supports_aingle: false,
connected: false,
}
}
pub fn update_rssi(&mut self, rssi: i16) {
self.rssi = rssi;
self.last_seen = Instant::now();
}
pub fn is_stale(&self, timeout: Duration) -> bool {
self.last_seen.elapsed() > timeout
}
}
#[derive(Debug, Clone, Default)]
pub struct BleStats {
pub messages_sent: u64,
pub messages_received: u64,
pub bytes_sent: u64,
pub bytes_received: u64,
pub messages_relayed: u64,
pub connection_attempts: u64,
pub connections_established: u64,
pub scans_completed: u64,
pub avg_rssi: i16,
}
pub struct BleManager {
config: BleConfig,
state: BleState,
peers: HashMap<String, BlePeer>,
stats: BleStats,
local_address: Option<String>,
running: bool,
#[cfg(feature = "ble")]
adapter: Option<Adapter>,
#[cfg(feature = "ble")]
peripherals: HashMap<String, Peripheral>,
#[cfg(feature = "ble")]
tx_chars: HashMap<String, Characteristic>,
#[cfg(feature = "ble")]
notification_rx: Option<Receiver<(String, Vec<u8>)>>,
#[cfg(feature = "ble")]
notification_tx: Option<Sender<(String, Vec<u8>)>>,
#[cfg(feature = "ble-esp32")]
ble_device: Option<&'static BLEDevice>,
#[cfg(feature = "ble-esp32")]
esp_clients: HashMap<String, Arc<BLEClient>>,
#[cfg(feature = "ble-esp32")]
discovered_devices: Vec<BLEAdvertisedDevice>,
}
impl BleManager {
pub fn new(config: BleConfig) -> Self {
Self {
config,
state: BleState::Uninitialized,
peers: HashMap::new(),
stats: BleStats::default(),
local_address: None,
running: false,
#[cfg(feature = "ble")]
adapter: None,
#[cfg(feature = "ble")]
peripherals: HashMap::new(),
#[cfg(feature = "ble")]
tx_chars: HashMap::new(),
#[cfg(feature = "ble")]
notification_rx: None,
#[cfg(feature = "ble")]
notification_tx: None,
#[cfg(feature = "ble-esp32")]
ble_device: None,
#[cfg(feature = "ble-esp32")]
esp_clients: HashMap::new(),
#[cfg(feature = "ble-esp32")]
discovered_devices: Vec::new(),
}
}
pub async fn init(&mut self) -> Result<()> {
log::info!("Initializing BLE transport: {}", self.config.device_name);
#[cfg(feature = "ble")]
{
let manager = Manager::new()
.await
.map_err(|e| Error::network(format!("Failed to initialize BLE manager: {}", e)))?;
let adapters = manager
.adapters()
.await
.map_err(|e| Error::network(format!("Failed to get BLE adapters: {}", e)))?;
let adapter = adapters
.into_iter()
.next()
.ok_or_else(|| Error::network("No BLE adapter found".to_string()))?;
let info = adapter
.adapter_info()
.await
.map_err(|e| Error::network(format!("Failed to get adapter info: {}", e)))?;
self.local_address = Some(info.to_string());
log::info!("BLE adapter initialized: {}", info);
let (tx, rx) = smol::channel::unbounded();
self.notification_tx = Some(tx);
self.notification_rx = Some(rx);
self.adapter = Some(adapter);
}
#[cfg(feature = "ble-esp32")]
{
let device = BLEDevice::take();
self.local_address = Some(format!("{:?}", device.get_address()));
log::info!("ESP32 BLE device initialized: {:?}", device.get_address());
self.ble_device = Some(device);
}
#[cfg(not(any(feature = "ble", feature = "ble-esp32")))]
{
log::warn!("No BLE feature enabled, using simulated mode");
self.local_address = Some("00:00:00:00:00:00".to_string());
}
self.state = BleState::Idle;
Ok(())
}
pub async fn start(&mut self) -> Result<()> {
if self.running {
return Ok(());
}
log::info!("Starting BLE transport");
self.start_advertising().await?;
self.start_scanning().await?;
self.running = true;
Ok(())
}
pub async fn stop(&mut self) -> Result<()> {
if !self.running {
return Ok(());
}
log::info!("Stopping BLE transport");
for (addr, peer) in self.peers.iter_mut() {
if peer.connected {
log::debug!("Disconnecting from peer: {}", addr);
peer.connected = false;
}
}
self.running = false;
self.state = BleState::Idle;
Ok(())
}
async fn start_advertising(&mut self) -> Result<()> {
log::debug!(
"Starting BLE advertising: {} (interval: {}ms, tx_power: {}dBm)",
self.config.device_name,
self.config.advertising_interval_ms,
self.config.tx_power
);
#[cfg(feature = "ble-esp32")]
{
let device = self
.ble_device
.ok_or_else(|| Error::network("ESP32 BLE device not initialized".to_string()))?;
let advertising = device.get_advertising();
advertising
.lock()
.set_scan_response(false)
.set_min_interval(self.config.advertising_interval_ms as u16)
.set_max_interval((self.config.advertising_interval_ms + 50) as u16);
advertising.lock().start().map_err(|e| {
Error::network(format!("Failed to start ESP32 advertising: {:?}", e))
})?;
log::info!("ESP32 BLE advertising started: {}", self.config.device_name);
}
Ok(())
}
async fn start_scanning(&mut self) -> Result<()> {
log::debug!(
"Starting BLE scan (interval: {}ms, window: {}ms)",
self.config.scan_interval_ms,
self.config.scan_window_ms
);
self.state = BleState::Scanning;
#[cfg(feature = "ble")]
{
let adapter = self
.adapter
.as_ref()
.ok_or_else(|| Error::network("BLE adapter not initialized".to_string()))?;
let service_uuid = Uuid::parse_str(AINGLE_SERVICE_UUID)
.map_err(|e| Error::network(format!("Invalid service UUID: {}", e)))?;
let filter = ScanFilter {
services: vec![service_uuid],
};
adapter
.start_scan(filter)
.await
.map_err(|e| Error::network(format!("Failed to start BLE scan: {}", e)))?;
log::info!("BLE scan started with AIngle service filter");
self.stats.scans_completed += 1;
}
#[cfg(feature = "ble-esp32")]
{
let device = self
.ble_device
.ok_or_else(|| Error::network("ESP32 BLE device not initialized".to_string()))?;
let scan = device.get_scan();
scan.active_scan(true)
.interval(self.config.scan_interval_ms as u16)
.window(self.config.scan_window_ms as u16);
log::info!("ESP32 BLE scan started");
self.stats.scans_completed += 1;
}
Ok(())
}
pub fn on_peer_discovered(&mut self, address: &str, rssi: i16, name: Option<&str>) {
if let Some(peer) = self.peers.get_mut(address) {
peer.update_rssi(rssi);
if name.is_some() {
peer.name = name.map(String::from);
}
} else {
let mut peer = BlePeer::new(address, rssi);
peer.name = name.map(String::from);
log::debug!(
"Discovered BLE peer: {} ({:?}) RSSI: {}",
address,
name,
rssi
);
self.peers.insert(address.to_string(), peer);
}
}
pub async fn connect(&mut self, address: &str) -> Result<()> {
let peer = self
.peers
.get_mut(address)
.ok_or_else(|| Error::network(format!("Unknown peer: {}", address)))?;
if peer.connected {
return Ok(());
}
log::info!("Connecting to BLE peer: {}", address);
self.stats.connection_attempts += 1;
#[cfg(feature = "ble")]
{
let adapter = self
.adapter
.as_ref()
.ok_or_else(|| Error::network("BLE adapter not initialized".to_string()))?;
let peripherals = adapter
.peripherals()
.await
.map_err(|e| Error::network(format!("Failed to get peripherals: {}", e)))?;
let peripheral = peripherals
.into_iter()
.find(|p| {
if let Ok(Some(props)) = smol::block_on(p.properties()) {
props.address.to_string() == address
} else {
false
}
})
.ok_or_else(|| Error::network(format!("Peripheral not found: {}", address)))?;
peripheral
.connect()
.await
.map_err(|e| Error::network(format!("Failed to connect: {}", e)))?;
peripheral
.discover_services()
.await
.map_err(|e| Error::network(format!("Failed to discover services: {}", e)))?;
let tx_uuid = Uuid::parse_str(TX_CHAR_UUID)
.map_err(|e| Error::network(format!("Invalid TX UUID: {}", e)))?;
let rx_uuid = Uuid::parse_str(RX_CHAR_UUID)
.map_err(|e| Error::network(format!("Invalid RX UUID: {}", e)))?;
for service in peripheral.services() {
for char in &service.characteristics {
if char.uuid == tx_uuid {
self.tx_chars.insert(address.to_string(), char.clone());
log::debug!("Found TX characteristic for {}", address);
}
if char.uuid == rx_uuid {
peripheral.subscribe(&char).await.map_err(|e| {
Error::network(format!("Failed to subscribe to RX: {}", e))
})?;
log::debug!("Subscribed to RX characteristic for {}", address);
}
}
}
self.peripherals.insert(address.to_string(), peripheral);
}
#[cfg(feature = "ble-esp32")]
{
let client = Arc::new(BLEClient::new());
client
.connect(address)
.await
.map_err(|e| Error::network(format!("ESP32 BLE connect failed: {:?}", e)))?;
log::info!("ESP32 BLE connected to: {}", address);
let service_uuid = uuid128!("6e400001-b5a3-f393-e0a9-e50e24dcca9e");
if let Some(service) = client.get_service(service_uuid).await {
log::debug!("Found AIngle service on ESP32 peer: {}", address);
let _tx_uuid = uuid128!("6e400002-b5a3-f393-e0a9-e50e24dcca9e");
let rx_uuid = uuid128!("6e400003-b5a3-f393-e0a9-e50e24dcca9e");
if let Some(rx_char) = service.get_characteristic(rx_uuid).await {
rx_char
.subscribe_notify(false, |_| {
log::debug!("ESP32 BLE notification received");
})
.await
.map_err(|e| {
Error::network(format!("Failed to subscribe to RX: {:?}", e))
})?;
}
}
self.esp_clients.insert(address.to_string(), client);
}
let peer = self.peers.get_mut(address).unwrap();
peer.connected = true;
peer.supports_aingle = true;
self.stats.connections_established += 1;
self.state = BleState::Connected;
Ok(())
}
pub async fn disconnect(&mut self, address: &str) -> Result<()> {
let peer = self
.peers
.get_mut(address)
.ok_or_else(|| Error::network(format!("Unknown peer: {}", address)))?;
if !peer.connected {
return Ok(());
}
log::info!("Disconnecting from BLE peer: {}", address);
#[cfg(feature = "ble")]
{
if let Some(peripheral) = self.peripherals.remove(address) {
peripheral
.disconnect()
.await
.map_err(|e| Error::network(format!("Failed to disconnect: {}", e)))?;
}
self.tx_chars.remove(address);
}
#[cfg(feature = "ble-esp32")]
{
if let Some(client) = self.esp_clients.remove(address) {
client
.disconnect()
.map_err(|e| Error::network(format!("ESP32 BLE disconnect failed: {:?}", e)))?;
}
}
peer.connected = false;
if !self.peers.values().any(|p| p.connected) {
self.state = BleState::Scanning;
}
Ok(())
}
pub async fn send(&mut self, address: &str, message: &Message) -> Result<()> {
let peer = self
.peers
.get(address)
.ok_or_else(|| Error::network(format!("Unknown peer: {}", address)))?;
if !peer.connected {
return Err(Error::network(format!("Peer not connected: {}", address)));
}
let payload =
serde_json::to_vec(message).map_err(|e| Error::Serialization(e.to_string()))?;
#[cfg(feature = "ble")]
{
let peripheral = self
.peripherals
.get(address)
.ok_or_else(|| Error::network(format!("Peripheral not found: {}", address)))?;
let tx_char = self.tx_chars.get(address).ok_or_else(|| {
Error::network(format!("TX characteristic not found: {}", address))
})?;
peripheral
.write(tx_char, &payload, WriteType::WithoutResponse)
.await
.map_err(|e| Error::network(format!("Failed to send: {}", e)))?;
}
#[cfg(feature = "ble-esp32")]
{
let client = self
.esp_clients
.get(address)
.ok_or_else(|| Error::network(format!("ESP32 client not found: {}", address)))?;
let service_uuid = uuid128!("6e400001-b5a3-f393-e0a9-e50e24dcca9e");
let tx_uuid = uuid128!("6e400002-b5a3-f393-e0a9-e50e24dcca9e");
if let Some(service) = client.get_service(service_uuid).await {
if let Some(tx_char) = service.get_characteristic(tx_uuid).await {
tx_char
.write_value(&payload, false)
.await
.map_err(|e| Error::network(format!("ESP32 BLE write failed: {:?}", e)))?;
} else {
return Err(Error::network("TX characteristic not found".to_string()));
}
} else {
return Err(Error::network("AIngle service not found".to_string()));
}
}
self.stats.messages_sent += 1;
self.stats.bytes_sent += payload.len() as u64;
log::debug!("Sent BLE message to {} ({} bytes)", address, payload.len());
Ok(())
}
pub async fn broadcast(&mut self, message: &Message) -> Result<usize> {
let connected_peers: Vec<String> = self
.peers
.iter()
.filter(|(_, p)| p.connected)
.map(|(addr, _)| addr.clone())
.collect();
let mut sent_count = 0;
for address in connected_peers {
match self.send(&address, message).await {
Ok(()) => {
log::debug!("Broadcast to {}", address);
sent_count += 1;
}
Err(e) => {
log::warn!("Failed to broadcast to {}: {}", address, e);
}
}
}
Ok(sent_count)
}
pub async fn recv(&mut self) -> Result<Option<(String, Message)>> {
#[cfg(feature = "ble")]
{
if let Some(ref rx) = self.notification_rx {
match rx.try_recv() {
Ok((address, data)) => {
let message: Message = serde_json::from_slice(&data)
.map_err(|e| Error::Serialization(e.to_string()))?;
self.stats.messages_received += 1;
self.stats.bytes_received += data.len() as u64;
log::debug!(
"Received BLE message from {} ({} bytes)",
address,
data.len()
);
return Ok(Some((address, message)));
}
Err(smol::channel::TryRecvError::Empty) => {
return Ok(None);
}
Err(smol::channel::TryRecvError::Closed) => {
return Err(Error::network("Notification channel closed".to_string()));
}
}
}
}
Ok(None)
}
pub async fn relay(&mut self, source: &str, message: &Message) -> Result<usize> {
if !self.config.mesh_relay {
return Ok(0);
}
let connected_peers: Vec<String> = self
.peers
.iter()
.filter(|(addr, p)| p.connected && *addr != source)
.map(|(addr, _)| addr.clone())
.collect();
let mut relayed_count = 0;
for address in connected_peers {
match self.send(&address, message).await {
Ok(()) => {
log::debug!("Relayed from {} to {}", source, address);
relayed_count += 1;
}
Err(e) => {
log::warn!("Failed to relay to {}: {}", address, e);
}
}
}
self.stats.messages_relayed += relayed_count as u64;
Ok(relayed_count)
}
pub fn state(&self) -> BleState {
self.state
}
pub fn is_running(&self) -> bool {
self.running
}
pub fn connected_count(&self) -> usize {
self.peers.values().filter(|p| p.connected).count()
}
pub fn peers(&self) -> impl Iterator<Item = &BlePeer> {
self.peers.values()
}
pub fn connected_peers(&self) -> impl Iterator<Item = &BlePeer> {
self.peers.values().filter(|p| p.connected)
}
pub fn get_peer(&self, address: &str) -> Option<&BlePeer> {
self.peers.get(address)
}
pub fn stats(&self) -> &BleStats {
&self.stats
}
pub fn cleanup_stale_peers(&mut self, timeout: Duration) {
let stale: Vec<String> = self
.peers
.iter()
.filter(|(_, p)| !p.connected && p.is_stale(timeout))
.map(|(addr, _)| addr.clone())
.collect();
for addr in stale {
log::debug!("Removing stale peer: {}", addr);
self.peers.remove(&addr);
}
}
pub fn local_address(&self) -> Option<&str> {
self.local_address.as_deref()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_ble_config_default() {
let config = BleConfig::default();
assert_eq!(config.device_name, "AIngle-Node");
assert!(config.mesh_relay);
assert_eq!(config.tx_power, 0);
assert_eq!(config.max_connections, 4);
}
#[test]
fn test_ble_config_low_power() {
let config = BleConfig::low_power();
assert_eq!(config.device_name, "AIngle-LP");
assert!(!config.mesh_relay);
assert_eq!(config.tx_power, -12);
assert!(config.passive_scan);
}
#[test]
fn test_ble_config_mesh_relay() {
let config = BleConfig::mesh_relay();
assert!(config.mesh_relay);
assert_eq!(config.tx_power, 4);
assert_eq!(config.max_connections, 8);
}
#[test]
fn test_ble_peer_creation() {
let peer = BlePeer::new("AA:BB:CC:DD:EE:FF", -50);
assert_eq!(peer.address, "AA:BB:CC:DD:EE:FF");
assert_eq!(peer.rssi, -50);
assert!(!peer.connected);
assert!(peer.name.is_none());
}
#[test]
fn test_ble_peer_update_rssi() {
let mut peer = BlePeer::new("AA:BB:CC:DD:EE:FF", -50);
std::thread::sleep(Duration::from_millis(10));
peer.update_rssi(-40);
assert_eq!(peer.rssi, -40);
}
#[test]
fn test_ble_peer_stale() {
let peer = BlePeer::new("AA:BB:CC:DD:EE:FF", -50);
assert!(!peer.is_stale(Duration::from_secs(60)));
}
#[test]
fn test_ble_manager_creation() {
let config = BleConfig::default();
let manager = BleManager::new(config);
assert_eq!(manager.state(), BleState::Uninitialized);
assert!(!manager.is_running());
assert_eq!(manager.connected_count(), 0);
}
#[test]
fn test_ble_manager_peer_discovery() {
let mut manager = BleManager::new(BleConfig::default());
manager.on_peer_discovered("AA:BB:CC:DD:EE:FF", -50, Some("Test Device"));
assert_eq!(manager.peers().count(), 1);
let peer = manager.get_peer("AA:BB:CC:DD:EE:FF").unwrap();
assert_eq!(peer.rssi, -50);
assert_eq!(peer.name, Some("Test Device".to_string()));
}
#[test]
fn test_ble_manager_peer_update() {
let mut manager = BleManager::new(BleConfig::default());
manager.on_peer_discovered("AA:BB:CC:DD:EE:FF", -50, None);
manager.on_peer_discovered("AA:BB:CC:DD:EE:FF", -40, Some("Updated Name"));
let peer = manager.get_peer("AA:BB:CC:DD:EE:FF").unwrap();
assert_eq!(peer.rssi, -40);
assert_eq!(peer.name, Some("Updated Name".to_string()));
}
#[test]
fn test_ble_stats_default() {
let stats = BleStats::default();
assert_eq!(stats.messages_sent, 0);
assert_eq!(stats.messages_received, 0);
assert_eq!(stats.messages_relayed, 0);
assert_eq!(stats.connections_established, 0);
}
#[test]
fn test_ble_state_equality() {
assert_eq!(BleState::Idle, BleState::Idle);
assert_ne!(BleState::Connected, BleState::Scanning);
}
}