refactor: Consolidate network sync operations into LatticeServer methods

This commit is contained in:
2025-12-23 01:41:34 +01:00
parent 2de54d0033
commit 4ccdbc97f5
10 changed files with 336 additions and 254 deletions
+1 -1
View File
@@ -15,7 +15,7 @@ pub mod mesh;
pub use endpoint::{LatticeEndpoint, PublicKey};
pub use framing::{MessageSink, MessageStream};
pub use lattice_core::proto::{SyncRequest, SyncResponse, SyncEntry, SyncDone, SyncState, Frontier};
pub use mesh::{spawn_accept_loop, join_mesh, sync_with_peer, sync_all, SyncResult};
pub use mesh::{LatticeServer, SyncResult};
/// Parse a PublicKey (NodeId) from hex or base32 string
pub fn parse_node_id(s: &str) -> Result<PublicKey, String> {
+2 -5
View File
@@ -1,13 +1,10 @@
//! Mesh networking - peer-to-peer join and sync operations
//!
//! - **server**: Accept incoming connections and handle join/sync requests
//! - **sync**: Outgoing join and sync operations
//! - **server**: LatticeServer for mesh networking (join, sync, accept loop)
//! - **protocol**: Shared send/receive entry logic
mod server;
mod sync;
mod protocol;
pub use server::spawn_accept_loop;
pub use sync::{join_mesh, sync_with_peer, sync_all, SyncResult};
pub use server::{LatticeServer, SyncResult};
pub use protocol::{send_missing_entries, receive_entries};
+184 -22
View File
@@ -1,37 +1,199 @@
//! Server - handle incoming peer connections for join and sync
//! Server - LatticeServer for mesh networking
use crate::{MessageSink, MessageStream};
use lattice_core::{Node, PeerStatus, Uuid};
use iroh::Endpoint;
use crate::{MessageSink, MessageStream, LatticeEndpoint, parse_node_id};
use lattice_core::{Node, NodeError, PeerStatus, Uuid, StoreHandle};
use iroh::endpoint::Connection;
use std::sync::Arc;
use lattice_core::proto::{PeerMessage, peer_message, JoinResponse};
use lattice_core::proto::{PeerMessage, peer_message, JoinRequest, JoinResponse};
use super::protocol;
/// Spawn the accept loop for incoming connections.
pub fn spawn_accept_loop(
/// Result of a sync operation with a peer
pub struct SyncResult {
pub entries_applied: u64,
pub entries_sent_by_peer: u64,
}
/// LatticeServer wraps Node + Endpoint and provides mesh networking methods.
/// Spawns accept loop on creation to handle incoming connections.
pub struct LatticeServer {
node: Arc<Node>,
endpoint: Endpoint,
) {
tokio::spawn(async move {
loop {
if let Some(incoming) = endpoint.accept().await {
match incoming.await {
Ok(conn) => {
let node = node.clone();
tokio::spawn(async move {
if let Err(e) = handle_connection(node, conn).await {
eprintln!("[Accept] Error: {}", e);
}
});
endpoint: LatticeEndpoint,
}
impl LatticeServer {
/// Create a new LatticeServer from just a Node (creates endpoint internally).
pub async fn new_from_node(node: Arc<Node>) -> Result<Self, String> {
let endpoint = LatticeEndpoint::new(node.secret_key_bytes()).await
.map_err(|e| format!("Failed to create endpoint: {}", e))?;
Ok(Self::new(node, endpoint))
}
/// Create a new LatticeServer with existing endpoint and spawn the accept loop.
pub fn new(node: Arc<Node>, endpoint: LatticeEndpoint) -> Self {
let server = Self { node, endpoint };
server.spawn_accept_loop();
server
}
/// Access the underlying node
pub fn node(&self) -> &Node {
&self.node
}
/// Access the underlying endpoint
pub fn endpoint(&self) -> &LatticeEndpoint {
&self.endpoint
}
/// Spawn the accept loop for incoming connections.
fn spawn_accept_loop(&self) {
let node = self.node.clone();
let endpoint = self.endpoint.endpoint().clone();
tokio::spawn(async move {
loop {
if let Some(incoming) = endpoint.accept().await {
match incoming.await {
Ok(conn) => {
let node = node.clone();
tokio::spawn(async move {
if let Err(e) = handle_connection(node, conn).await {
eprintln!("[Accept] Error: {}", e);
}
});
}
Err(e) => eprintln!("[Accept] Handshake error: {:?}", e),
}
Err(e) => eprintln!("[Accept] Handshake error: {:?}", e),
}
}
});
}
/// Join an existing mesh by connecting to a peer.
pub async fn join_mesh(&self, peer_id: iroh::PublicKey) -> Result<StoreHandle, NodeError> {
let conn = self.endpoint.connect(peer_id).await
.map_err(|e| NodeError::Actor(format!("Connection failed: {}", e)))?;
let (send, recv) = conn.open_bi().await
.map_err(|e| NodeError::Actor(format!("Failed to open stream: {}", e)))?;
let mut sink = MessageSink::new(send);
let mut stream = MessageStream::new(recv);
// Send JoinRequest
let req = PeerMessage {
message: Some(peer_message::Message::JoinRequest(JoinRequest {
node_pubkey: self.node.node_id().to_vec(),
})),
};
sink.send(&req).await.map_err(|e| NodeError::Actor(e))?;
sink.finish().await.map_err(|e| NodeError::Actor(e))?;
// Receive JoinResponse
let msg = stream.recv().await
.map_err(|e| NodeError::Actor(e))?
.ok_or_else(|| NodeError::Actor("Peer closed stream".to_string()))?;
match msg.message {
Some(peer_message::Message::JoinResponse(resp)) => {
let store_uuid = lattice_core::Uuid::from_slice(&resp.store_uuid)
.map_err(|_| NodeError::Actor("Invalid UUID from peer".to_string()))?;
let handle = self.node.complete_join(store_uuid).await?;
// Sync with peer to get initial data
println!("[Join] Syncing with peer to get initial data...");
if let Ok(result) = self.sync_with_peer(&handle, peer_id).await {
println!("[Join] Initial sync complete: {} entries", result.entries_applied);
}
Ok(handle)
}
_ => Err(NodeError::Actor("Unexpected response".to_string())),
}
});
}
/// Sync with a specific peer.
pub async fn sync_with_peer(&self, store: &StoreHandle, peer_id: iroh::PublicKey) -> Result<SyncResult, NodeError> {
let conn = self.endpoint.connect(peer_id).await
.map_err(|e| NodeError::Actor(format!("Connection failed: {}", e)))?;
let (send, recv) = conn.open_bi().await
.map_err(|e| NodeError::Actor(format!("Failed to open stream: {}", e)))?;
let mut sink = MessageSink::new(send);
let mut stream = MessageStream::new(recv);
let my_state = store.sync_state().await?;
// Send SyncRequest
let req = PeerMessage {
message: Some(peer_message::Message::SyncRequest(lattice_core::proto::SyncRequest {
store_id: store.id().as_bytes().to_vec(),
state: Some(my_state.to_proto()),
full_sync: false,
})),
};
sink.send(&req).await.map_err(|e| NodeError::Actor(e))?;
// Receive SyncResponse
let resp_msg = stream.recv().await.map_err(|e| NodeError::Actor(e))?
.ok_or_else(|| NodeError::Actor("Peer closed stream".to_string()))?;
let peer_state = match resp_msg.message {
Some(peer_message::Message::SyncResponse(resp)) => {
resp.state.map(|s| lattice_core::sync_state::SyncState::from_proto(&s))
.unwrap_or_default()
}
_ => return Err(NodeError::Actor("Expected SyncResponse".to_string())),
};
// Exchange entries
let _entries_sent = protocol::send_missing_entries(&mut sink, store, &my_state, &peer_state).await
.map_err(|e| NodeError::Actor(e))?;
let (entries_applied, entries_sent_by_peer) = protocol::receive_entries(&mut stream, store).await
.map_err(|e| NodeError::Actor(e))?;
sink.finish().await.map_err(|e| NodeError::Actor(e))?;
Ok(SyncResult { entries_applied, entries_sent_by_peer })
}
/// Sync with all active peers.
pub async fn sync_all(&self, store: &StoreHandle) -> Result<Vec<SyncResult>, NodeError> {
let peers = self.node.list_peers().await?;
let mut results = Vec::new();
for peer in peers {
if peer.status != PeerStatus::Active {
continue;
}
let peer_id = match parse_node_id(&peer.pubkey) {
Ok(id) => id,
Err(e) => {
eprintln!("[Sync] Failed to parse peer {}: {}", peer.pubkey, e);
continue;
}
};
println!("[Sync] Syncing with {}...", peer_id.fmt_short());
match self.sync_with_peer(store, peer_id).await {
Ok(result) => {
println!("[Sync] Applied {} entries", result.entries_applied);
results.push(result);
}
Err(e) => eprintln!("[Sync] Failed: {}", e),
}
}
Ok(results)
}
}
// --- Connection handling ---
// --- Connection handling ---
/// Handle a single incoming connection
async fn handle_connection(
node: Arc<Node>,
-181
View File
@@ -1,181 +0,0 @@
//! Sync - outgoing mesh join and sync operations
use crate::{MessageSink, MessageStream, LatticeEndpoint, parse_node_id};
use lattice_core::{Node, NodeError, StoreHandle, PeerStatus};
use lattice_core::proto::{peer_message, PeerMessage, JoinRequest, SignedEntry};
use prost::Message;
use super::protocol;
/// Result of a sync operation with a peer
pub struct SyncResult {
pub entries_applied: u64,
pub entries_sent_by_peer: u64,
}
/// Join an existing mesh by connecting to a peer.
/// Returns the new StoreHandle on success.
/// After joining, automatically syncs with the peer to get initial data.
pub async fn join_mesh(
node: &Node,
endpoint: &LatticeEndpoint,
peer_id: iroh::PublicKey,
) -> Result<StoreHandle, NodeError> {
// Connect to peer
let conn = endpoint.connect(peer_id).await
.map_err(|e| NodeError::Actor(format!("Connection failed: {}", e)))?;
// Open stream
let (send, recv) = conn.open_bi().await
.map_err(|e| NodeError::Actor(format!("Failed to open stream: {}", e)))?;
let mut sink = MessageSink::new(send);
let mut stream = MessageStream::new(recv);
// Send JoinRequest
let req = PeerMessage {
message: Some(peer_message::Message::JoinRequest(JoinRequest {
node_pubkey: node.node_id().to_vec(),
})),
};
sink.send(&req).await
.map_err(|e| NodeError::Actor(format!("Failed to send: {}", e)))?;
sink.finish().await
.map_err(|e| NodeError::Actor(format!("Failed to finish: {}", e)))?;
// Receive JoinResponse
let msg = stream.recv().await
.map_err(|e| NodeError::Actor(format!("Recv error: {}", e)))?
.ok_or_else(|| NodeError::Actor("Peer closed stream".to_string()))?;
match msg.message {
Some(peer_message::Message::JoinResponse(resp)) => {
let store_uuid = lattice_core::Uuid::from_slice(&resp.store_uuid)
.map_err(|_| NodeError::Actor("Invalid UUID from peer".to_string()))?;
// Complete join - creates store, sets as root, caches handle
let handle = node.complete_join(store_uuid).await?;
// Immediately sync with the peer to get initial data
println!("[Join] Syncing with peer to get initial data...");
match sync_with_peer(endpoint, &handle, peer_id).await {
Ok(result) => {
println!("[Join] Initial sync complete: applied {} entries", result.entries_applied);
}
Err(e) => {
eprintln!("[Join] Warning: Initial sync failed: {}", e);
}
}
Ok(handle)
}
_ => Err(NodeError::Actor("Unexpected response message type".to_string())),
}
}
/// Sync with a specific peer (bidirectional).
/// Both sides exchange states and send missing entries to each other.
pub async fn sync_with_peer(
endpoint: &LatticeEndpoint,
store: &StoreHandle,
peer_id: iroh::PublicKey,
) -> Result<SyncResult, NodeError> {
// Connect
let conn = endpoint.connect(peer_id).await
.map_err(|e| NodeError::Actor(format!("Connection failed: {}", e)))?;
// Open stream
let (send, recv) = conn.open_bi().await
.map_err(|e| NodeError::Actor(format!("Failed to open stream: {}", e)))?;
let mut sink = MessageSink::new(send);
let mut stream = MessageStream::new(recv);
// Get our sync state
let my_state = store.sync_state().await?;
// 1. Send SyncRequest with our state (don't finish yet - we'll send entries later)
let req = PeerMessage {
message: Some(peer_message::Message::SyncRequest(lattice_core::proto::SyncRequest {
store_id: store.id().as_bytes().to_vec(),
state: Some(my_state.to_proto()),
full_sync: false,
})),
};
sink.send(&req).await
.map_err(|e| NodeError::Actor(format!("Failed to send: {}", e)))?;
// 2. Receive SyncResponse (peer's state) and entries until SyncDone
let mut entries_applied = 0u64;
let mut entries_sent_by_peer = 0u64;
let mut peer_state = lattice_core::sync_state::SyncState::default();
loop {
match stream.recv().await {
Ok(Some(msg)) => match msg.message {
Some(peer_message::Message::SyncResponse(resp)) => {
// Peer's sync state - we'll use this to compute what to send
if let Some(s) = resp.state {
peer_state = lattice_core::sync_state::SyncState::from_proto(&s);
}
}
Some(peer_message::Message::SyncEntry(entry)) => {
if let Ok(signed) = SignedEntry::decode(&entry.signed_entry[..]) {
if store.apply_entry(signed).await.is_ok() {
entries_applied += 1;
}
}
}
Some(peer_message::Message::SyncDone(done)) => {
entries_sent_by_peer = done.entries_sent;
break;
}
_ => {}
}
Ok(None) => break,
Err(_) => break,
}
}
// 3. Send entries peer is missing (using shared protocol)
let entries_sent = protocol::send_missing_entries(&mut sink, store, &my_state, &peer_state).await
.map_err(|e| NodeError::Actor(e))?;
sink.finish().await
.map_err(|e| NodeError::Actor(format!("Failed to finish: {}", e)))?;
println!("[Sync] Applied {} entries, sent {} entries", entries_applied, entries_sent);
Ok(SyncResult {
entries_applied,
entries_sent_by_peer,
})
}
/// Sync with all active peers from the node.
pub async fn sync_all(
node: &Node,
endpoint: &LatticeEndpoint,
store: &StoreHandle,
) -> Result<Vec<SyncResult>, NodeError> {
let my_pubkey = hex::encode(node.node_id());
// Get all active peers using node.list_peers()
let peers = node.list_peers().await?;
let mut results = Vec::new();
for peer in peers {
if peer.status == PeerStatus::Active && peer.pubkey != my_pubkey {
if let Ok(peer_id) = parse_node_id(&peer.pubkey) {
match sync_with_peer(endpoint, store, peer_id).await {
Ok(result) => results.push(result),
Err(e) => {
eprintln!("Sync with {} failed: {}", peer_id.fmt_short(), e);
}
}
}
}
}
Ok(results)
}