use crate::config::Config;
use crate::crypto::Keypair;
use crate::error::Result;
use crate::gossip::GossipManager;
use crate::network::{Message, Network};
use crate::storage_factory::DynamicStorage;
use crate::storage_trait::StorageBackend;
use crate::sync::SyncManager;
use crate::types::*;
use serde::{Deserialize, Serialize};
use std::net::SocketAddr;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
const PEERS_METADATA_KEY: &str = "known_peers";
const PEER_SAVE_INTERVAL_SECS: u64 = 300;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PeerRecord {
pub addr: String,
pub latest_seq: u32,
pub quality: u8,
pub last_seen_secs: u64,
}
pub struct MinimalNode {
config: Config,
keypair: Keypair,
storage: DynamicStorage,
network: Network,
gossip: GossipManager,
sync: SyncManager,
running: Arc<AtomicBool>,
start_time: Instant,
last_peer_save: Instant,
}
impl MinimalNode {
pub fn new(config: Config) -> Result<Self> {
config.validate()?;
let keypair = Keypair::generate();
let storage = DynamicStorage::from_config(config.storage.clone())?;
log::info!("Using {} storage backend", storage.backend_name());
let node_id = keypair.public_key().to_hex();
let network = Network::new(config.transport.clone(), config.gossip.clone(), node_id);
let gossip = GossipManager::new(config.gossip.clone());
let sync = SyncManager::new(config.gossip.loop_delay * 2);
let mut node = Self {
config,
keypair,
storage,
network,
gossip,
sync,
running: Arc::new(AtomicBool::new(false)),
start_time: Instant::now(),
last_peer_save: Instant::now(),
};
if let Err(e) = node.load_peers() {
log::warn!("Failed to load persisted peers: {}", e);
}
Ok(node)
}
pub fn public_key(&self) -> AgentPubKey {
self.keypair.public_key()
}
pub fn stats(&self) -> Result<NodeStats> {
let storage_stats = self.storage.stats()?;
Ok(NodeStats {
entries_count: storage_stats.entry_count,
actions_count: storage_stats.action_count,
memory_used: 0, storage_used: storage_stats.db_size,
peer_count: self.network.peer_count(),
uptime_secs: self.start_time.elapsed().as_secs(),
})
}
pub fn create_entry<T: serde::Serialize>(&mut self, content: T) -> Result<Hash> {
let entry = Entry::app(content)?;
let entry_hash = entry.hash();
let seq = self.storage.get_latest_seq()? + 1;
let prev_action = if seq > 1 {
None
} else {
None
};
let action = Action {
action_type: ActionType::Create,
author: self.keypair.public_key(),
timestamp: Timestamp::now(),
seq,
prev_action,
entry_hash: Some(entry_hash.clone()),
signature: self.sign_action_data(seq, &entry_hash),
};
let record = Record {
action,
entry: Some(entry),
};
let action_hash = self.storage.put_record(&record)?;
self.gossip.announce(action_hash.clone());
self.sync.add_local_hash(action_hash.clone());
if self.config.publish_interval == Duration::ZERO {
self.publish_pending()?;
}
Ok(action_hash)
}
pub fn create_entries_batch<T: serde::Serialize>(
&mut self,
contents: &[T],
) -> Result<Vec<Hash>> {
if contents.is_empty() {
return Ok(Vec::new());
}
let base_seq = self.storage.get_latest_seq()? + 1;
let author = self.keypair.public_key();
let timestamp = Timestamp::now();
let mut records = Vec::with_capacity(contents.len());
for (i, content) in contents.iter().enumerate() {
let entry = Entry::app(content)?;
let entry_hash = entry.hash();
let seq = base_seq + i as u32;
let action = Action {
action_type: ActionType::Create,
author: author.clone(),
timestamp,
seq,
prev_action: None,
entry_hash: Some(entry_hash.clone()),
signature: self.sign_action_data(seq, &entry_hash),
};
records.push(Record {
action,
entry: Some(entry),
});
}
let hashes = self.storage.put_records_batch(&records)?;
for hash in &hashes {
self.gossip.announce(hash.clone());
self.sync.add_local_hash(hash.clone());
}
if self.config.publish_interval == Duration::ZERO {
self.publish_pending()?;
}
Ok(hashes)
}
pub fn get_entry(&self, hash: &Hash) -> Result<Option<Entry>> {
self.storage.get_entry(hash)
}
fn sign_action_data(&self, seq: u32, entry_hash: &Hash) -> Signature {
let mut data = Vec::new();
data.extend_from_slice(&seq.to_be_bytes());
data.extend_from_slice(entry_hash.as_bytes());
self.keypair.sign(&data)
}
fn publish_pending(&mut self) -> Result<()> {
let announcements = self.gossip.take_announcements(50);
for hash in announcements {
let message = Message::NewRecord { hash };
smol::block_on(self.network.broadcast(&message))?;
}
Ok(())
}
pub async fn run(&mut self) -> Result<()> {
self.running.store(true, Ordering::SeqCst);
log::info!(
"Starting minimal node: {}",
self.keypair.public_key().to_hex()
);
self.network.start().await?;
if self.config.enable_mdns {
let port = match &self.config.transport {
crate::config::TransportConfig::Coap { port, .. } => *port,
crate::config::TransportConfig::Quic { port, .. } => *port,
#[cfg(feature = "webrtc")]
crate::config::TransportConfig::WebRtc { signaling_port, .. } => *signaling_port,
_ => 5683,
};
if let Err(e) = self.network.start_discovery(port) {
log::warn!("Failed to start mDNS discovery: {}", e);
}
}
let mut discovery_sync_counter = 0u32;
while self.running.load(Ordering::SeqCst) {
discovery_sync_counter = discovery_sync_counter.wrapping_add(1);
if discovery_sync_counter % 100 == 0 {
self.network.sync_discovered_peers();
}
if self.gossip.should_gossip() {
self.run_gossip_round().await;
}
if self.config.publish_interval > Duration::ZERO {
self.publish_pending()?;
}
if self.last_peer_save.elapsed().as_secs() >= PEER_SAVE_INTERVAL_SECS {
if let Err(e) = self.save_peers() {
log::warn!("Failed to save peers: {}", e);
}
}
smol::Timer::after(Duration::from_millis(10)).await;
}
if let Err(e) = self.save_peers() {
log::warn!("Failed to save peers on shutdown: {}", e);
}
self.network.stop().await?;
log::info!("Node stopped");
Ok(())
}
async fn run_gossip_round(&mut self) {
let peers = self.network.gossip_peers();
let mut success_count = 0;
for addr in peers {
let latest_seq = self.storage.get_latest_seq().unwrap_or(0);
match self
.sync
.sync_with_peer(&addr, &mut self.network, &self.storage, &mut self.gossip)
.await
{
Ok(result) => {
self.network.update_peer(addr, latest_seq);
success_count += 1;
log::debug!(
"Sync with {} complete: sent_filter={}, records_sent={}, records_received={}",
addr,
result.sent_filter,
result.records_sent,
result.records_received
);
}
Err(e) => {
log::warn!("Sync with {} failed: {}", addr, e);
self.network.mark_peer_failed(&addr);
}
}
}
self.gossip.gossip_complete(success_count > 0);
}
pub fn stop(&self) {
self.running.store(false, Ordering::SeqCst);
}
pub fn is_running(&self) -> bool {
self.running.load(Ordering::SeqCst)
}
pub fn add_peer(&mut self, addr: std::net::SocketAddr) {
self.network.add_peer(addr);
}
pub fn save_peers(&mut self) -> Result<()> {
let peers = self.network.active_peers();
if peers.is_empty() {
log::debug!("No active peers to save");
return Ok(());
}
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let peer_records: Vec<PeerRecord> = peers
.iter()
.map(|p| {
let elapsed_secs = p.last_seen.elapsed().as_secs();
PeerRecord {
addr: p.addr.to_string(),
latest_seq: p.latest_seq,
quality: p.quality,
last_seen_secs: now.saturating_sub(elapsed_secs),
}
})
.collect();
let json = serde_json::to_string(&peer_records)
.map_err(|e| crate::error::Error::Serialization(e.to_string()))?;
self.storage.set_metadata(PEERS_METADATA_KEY, &json)?;
self.last_peer_save = Instant::now();
log::info!("Saved {} peers to storage", peer_records.len());
Ok(())
}
pub fn load_peers(&mut self) -> Result<()> {
let json = match self.storage.get_metadata(PEERS_METADATA_KEY)? {
Some(data) => data,
None => {
log::debug!("No persisted peers found");
return Ok(());
}
};
let peer_records: Vec<PeerRecord> = serde_json::from_str(&json)
.map_err(|e| crate::error::Error::Serialization(e.to_string()))?;
let mut loaded_count = 0;
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
for record in peer_records {
if record.quality < 10 {
log::debug!(
"Skipping low-quality peer: {} (quality={})",
record.addr,
record.quality
);
continue;
}
let age_secs = now.saturating_sub(record.last_seen_secs);
if age_secs > 24 * 60 * 60 {
log::debug!(
"Skipping stale peer: {} (last seen {} hours ago)",
record.addr,
age_secs / 3600
);
continue;
}
if let Ok(addr) = record.addr.parse::<SocketAddr>() {
self.network.add_peer(addr);
let quality_boosts = record.quality.saturating_sub(50) / 5;
for _ in 0..quality_boosts {
self.network.update_peer(addr, record.latest_seq);
}
loaded_count += 1;
} else {
log::warn!("Invalid peer address in storage: {}", record.addr);
}
}
log::info!("Loaded {} peers from storage", loaded_count);
Ok(())
}
pub fn get_known_peers(&self) -> Vec<PeerRecord> {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
self.network
.active_peers()
.iter()
.map(|p| {
let elapsed_secs = p.last_seen.elapsed().as_secs();
PeerRecord {
addr: p.addr.to_string(),
latest_seq: p.latest_seq,
quality: p.quality,
last_seen_secs: now.saturating_sub(elapsed_secs),
}
})
.collect()
}
pub fn sync_stats(&self) -> crate::sync::SyncStats {
self.sync.stats()
}
pub fn gossip_stats(&self) -> crate::gossip::GossipStats {
self.gossip.stats()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_node_creation() {
let config = Config::test_mode();
let node = MinimalNode::new(config);
assert!(node.is_ok());
}
#[test]
fn test_create_entry() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let content = serde_json::json!({
"sensor_id": "temp_001",
"value": 23.5,
"unit": "celsius"
});
let hash = node.create_entry(content);
assert!(hash.is_ok());
}
#[test]
fn test_node_stats() {
let config = Config::test_mode();
let node = MinimalNode::new(config).unwrap();
let stats = node.stats();
assert!(stats.is_ok());
assert_eq!(stats.unwrap().entries_count, 0);
}
#[test]
fn test_public_key() {
let config = Config::test_mode();
let node = MinimalNode::new(config).unwrap();
let pubkey = node.public_key();
assert_eq!(pubkey.0.len(), 32);
assert_eq!(pubkey.to_hex().len(), 64);
}
#[test]
fn test_get_entry() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let content = "test content";
let hash = node.create_entry(content).unwrap();
let entry = node.get_entry(&hash);
assert!(entry.is_ok());
}
#[test]
fn test_get_entry_not_found() {
let config = Config::test_mode();
let node = MinimalNode::new(config).unwrap();
let hash = Hash([0u8; 32]);
let entry = node.get_entry(&hash);
assert!(entry.is_ok());
assert!(entry.unwrap().is_none());
}
#[test]
fn test_stop_and_is_running() {
let config = Config::test_mode();
let node = MinimalNode::new(config).unwrap();
assert!(!node.is_running());
node.stop();
assert!(!node.is_running());
}
#[test]
fn test_add_peer() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let addr: std::net::SocketAddr = "192.168.1.100:5683".parse().unwrap();
node.add_peer(addr);
let stats = node.stats().unwrap();
assert!(stats.peer_count >= 0); }
#[test]
fn test_add_multiple_peers() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
node.add_peer("192.168.1.100:5683".parse().unwrap());
node.add_peer("192.168.1.101:5683".parse().unwrap());
node.add_peer("192.168.1.102:5683".parse().unwrap());
let stats = node.stats().unwrap();
assert!(stats.peer_count >= 0);
}
#[test]
fn test_sync_stats() {
let config = Config::test_mode();
let node = MinimalNode::new(config).unwrap();
let sync_stats = node.sync_stats();
assert_eq!(sync_stats.total_successful_syncs, 0);
assert_eq!(sync_stats.total_failed_syncs, 0);
assert_eq!(sync_stats.peer_count, 0);
}
#[test]
fn test_gossip_stats() {
let config = Config::test_mode();
let node = MinimalNode::new(config).unwrap();
let gossip_stats = node.gossip_stats();
assert_eq!(gossip_stats.round, 0);
assert_eq!(gossip_stats.pending_announcements, 0);
assert_eq!(gossip_stats.queue_length, 0);
}
#[test]
fn test_create_multiple_entries() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let hash1 = node.create_entry("entry 1").unwrap();
let hash2 = node.create_entry("entry 2").unwrap();
let hash3 = node.create_entry("entry 3").unwrap();
assert_ne!(hash1, hash2);
assert_ne!(hash2, hash3);
assert_ne!(hash1, hash3);
let stats = node.stats().unwrap();
assert!(stats.actions_count >= 3);
}
#[test]
fn test_node_stats_after_entries() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let initial_stats = node.stats().unwrap();
let initial_actions = initial_stats.actions_count;
node.create_entry("test").unwrap();
let updated_stats = node.stats().unwrap();
assert!(updated_stats.actions_count > initial_actions);
}
#[test]
fn test_node_uptime() {
let config = Config::test_mode();
let node = MinimalNode::new(config).unwrap();
std::thread::sleep(std::time::Duration::from_millis(10));
let stats = node.stats().unwrap();
assert!(stats.uptime_secs >= 0);
}
#[test]
fn test_node_with_iot_config() {
let config = Config::iot_mode();
let node = MinimalNode::new(config);
assert!(node.is_ok());
let node = node.unwrap();
assert!(!node.is_running());
}
#[test]
fn test_create_entry_with_struct() {
use serde::{Deserialize, Serialize};
#[derive(Serialize, Deserialize)]
struct SensorReading {
temperature: f64,
humidity: f64,
timestamp: u64,
}
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let reading = SensorReading {
temperature: 23.5,
humidity: 65.2,
timestamp: 1234567890,
};
let hash = node.create_entry(reading);
assert!(hash.is_ok());
}
#[test]
fn test_create_entry_with_nested_struct() {
use serde::{Deserialize, Serialize};
#[derive(Serialize, Deserialize)]
struct Location {
latitude: f64,
longitude: f64,
}
#[derive(Serialize, Deserialize)]
struct DeviceData {
device_id: String,
location: Location,
readings: Vec<f64>,
}
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let data = DeviceData {
device_id: "device_001".to_string(),
location: Location {
latitude: 40.7128,
longitude: -74.0060,
},
readings: vec![1.0, 2.0, 3.0, 4.0, 5.0],
};
let hash = node.create_entry(data);
assert!(hash.is_ok());
}
#[test]
fn test_peer_record_serialization() {
let record = PeerRecord {
addr: "192.168.1.100:5683".to_string(),
latest_seq: 42,
quality: 75,
last_seen_secs: 1234567890,
};
let json = serde_json::to_string(&record).unwrap();
assert!(json.contains("192.168.1.100:5683"));
assert!(json.contains("42"));
assert!(json.contains("75"));
let parsed: PeerRecord = serde_json::from_str(&json).unwrap();
assert_eq!(parsed.addr, record.addr);
assert_eq!(parsed.latest_seq, record.latest_seq);
assert_eq!(parsed.quality, record.quality);
assert_eq!(parsed.last_seen_secs, record.last_seen_secs);
}
#[test]
fn test_peer_record_vec_serialization() {
let records = vec![
PeerRecord {
addr: "192.168.1.100:5683".to_string(),
latest_seq: 10,
quality: 80,
last_seen_secs: 1000,
},
PeerRecord {
addr: "10.0.0.50:5683".to_string(),
latest_seq: 20,
quality: 60,
last_seen_secs: 2000,
},
];
let json = serde_json::to_string(&records).unwrap();
let parsed: Vec<PeerRecord> = serde_json::from_str(&json).unwrap();
assert_eq!(parsed.len(), 2);
assert_eq!(parsed[0].addr, "192.168.1.100:5683");
assert_eq!(parsed[1].addr, "10.0.0.50:5683");
}
#[test]
fn test_save_peers_empty() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let result = node.save_peers();
assert!(result.is_ok());
}
#[test]
fn test_save_peers_with_peers() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
node.add_peer("192.168.1.100:5683".parse().unwrap());
node.add_peer("192.168.1.101:5683".parse().unwrap());
let result = node.save_peers();
assert!(result.is_ok());
}
#[test]
fn test_load_peers_empty_storage() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let result = node.load_peers();
assert!(result.is_ok());
}
#[test]
fn test_save_and_load_peers_roundtrip() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
node.add_peer("192.168.1.100:5683".parse().unwrap());
node.add_peer("10.0.0.50:5683".parse().unwrap());
node.save_peers().unwrap();
let peer_count_before = node.stats().unwrap().peer_count;
node.load_peers().unwrap();
let peer_count_after = node.stats().unwrap().peer_count;
assert!(peer_count_after >= peer_count_before);
}
#[test]
fn test_get_known_peers_empty() {
let config = Config::test_mode();
let node = MinimalNode::new(config).unwrap();
let peers = node.get_known_peers();
assert!(peers.is_empty());
}
#[test]
fn test_get_known_peers_with_peers() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
node.add_peer("192.168.1.100:5683".parse().unwrap());
node.add_peer("192.168.1.101:5683".parse().unwrap());
let peers = node.get_known_peers();
assert_eq!(peers.len(), 2);
let addrs: Vec<&str> = peers.iter().map(|p| p.addr.as_str()).collect();
assert!(addrs.contains(&"192.168.1.100:5683"));
assert!(addrs.contains(&"192.168.1.101:5683"));
for peer in &peers {
assert!(peer.quality > 0);
}
}
#[test]
fn test_peer_record_debug() {
let record = PeerRecord {
addr: "127.0.0.1:5683".to_string(),
latest_seq: 0,
quality: 50,
last_seen_secs: 0,
};
let debug_str = format!("{:?}", record);
assert!(debug_str.contains("PeerRecord"));
assert!(debug_str.contains("127.0.0.1:5683"));
}
#[test]
fn test_peer_record_clone() {
let record1 = PeerRecord {
addr: "192.168.1.1:5683".to_string(),
latest_seq: 100,
quality: 90,
last_seen_secs: 999999,
};
let record2 = record1.clone();
assert_eq!(record1.addr, record2.addr);
assert_eq!(record1.latest_seq, record2.latest_seq);
assert_eq!(record1.quality, record2.quality);
assert_eq!(record1.last_seen_secs, record2.last_seen_secs);
}
#[test]
fn test_peer_persistence_skips_low_quality() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let records = vec![PeerRecord {
addr: "192.168.1.100:5683".to_string(),
latest_seq: 0,
quality: 5, last_seen_secs: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs(),
}];
let json = serde_json::to_string(&records).unwrap();
node.storage
.set_metadata(PEERS_METADATA_KEY, &json)
.unwrap();
node.load_peers().unwrap();
let peers = node.get_known_peers();
assert!(peers.is_empty());
}
#[test]
fn test_peer_persistence_skips_stale() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let old_time = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs()
.saturating_sub(25 * 60 * 60);
let records = vec![PeerRecord {
addr: "192.168.1.100:5683".to_string(),
latest_seq: 100,
quality: 80, last_seen_secs: old_time,
}];
let json = serde_json::to_string(&records).unwrap();
node.storage
.set_metadata(PEERS_METADATA_KEY, &json)
.unwrap();
node.load_peers().unwrap();
let peers = node.get_known_peers();
assert!(peers.is_empty());
}
#[test]
fn test_peer_persistence_accepts_recent_quality_peer() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let recent_time = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs();
let records = vec![PeerRecord {
addr: "192.168.1.100:5683".to_string(),
latest_seq: 50,
quality: 60, last_seen_secs: recent_time,
}];
let json = serde_json::to_string(&records).unwrap();
node.storage
.set_metadata(PEERS_METADATA_KEY, &json)
.unwrap();
node.load_peers().unwrap();
let peers = node.get_known_peers();
assert_eq!(peers.len(), 1);
assert_eq!(peers[0].addr, "192.168.1.100:5683");
}
#[test]
fn test_peer_persistence_invalid_address() {
let config = Config::test_mode();
let mut node = MinimalNode::new(config).unwrap();
let recent_time = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_secs();
let records = vec![PeerRecord {
addr: "not-a-valid-address".to_string(), latest_seq: 50,
quality: 80,
last_seen_secs: recent_time,
}];
let json = serde_json::to_string(&records).unwrap();
node.storage
.set_metadata(PEERS_METADATA_KEY, &json)
.unwrap();
let result = node.load_peers();
assert!(result.is_ok());
let peers = node.get_known_peers();
assert!(peers.is_empty());
}
}