diff --git a/docs/roadmap.md b/docs/roadmap.md index 561a8a2..3f6e25b 100644 --- a/docs/roadmap.md +++ b/docs/roadmap.md @@ -57,6 +57,30 @@ --- +## Milestone 1.9: Async Refactor + +**Goal:** Prepare codebase for concurrent CLI + network operation. + +### Deliverables + +**Phase 1: Store Actor (sync)** ✓ +- [x] Store actor pattern: dedicated thread owns Store, receives commands via `std::sync::mpsc` +- [x] StoreHandle wraps channel sender, keeps current API +- [x] Validate: CLI works as before with actor + +**Phase 2: Async Runtime** +- [ ] Add tokio runtime (`#[tokio::main]`) +- [ ] Migrate `std::sync::mpsc` → `tokio::sync::mpsc` +- [ ] Async CLI using `tokio::io::stdin()` or `rustyline` async + +### Success Criteria + +- [x] CLI still works as before +- [x] Store operations serialized (no data races) +- [ ] Ready for concurrent network tasks + +--- + ## Milestone 2: Two-Node Sync **Goal:** Two nodes can sync their logs over the network. @@ -68,6 +92,7 @@ - [ ] Iroh integration (peer discovery, connection) - [ ] Multi-author log merging - [ ] CLI: `peers`, `connect`/`join` commands +- [ ] Background sync task (tokio::spawn) ### Success Criteria diff --git a/lattice-cli/src/main.rs b/lattice-cli/src/main.rs index df76338..15123da 100644 --- a/lattice-cli/src/main.rs +++ b/lattice-cli/src/main.rs @@ -2,6 +2,7 @@ mod node; mod commands; +mod store_actor; use commands::CommandResult; use node::{LatticeNodeBuilder, StoreHandle}; diff --git a/lattice-cli/src/node.rs b/lattice-cli/src/node.rs index 99871d2..6c58402 100644 --- a/lattice-cli/src/node.rs +++ b/lattice-cli/src/node.rs @@ -1,8 +1,7 @@ //! Local Lattice node API with multi-store support use lattice_core::{ - DataDir, EntryBuilder, MetaStore, Node, SigChain, Store, Uuid, - hlc::HLC, + DataDir, MetaStore, Node, SigChain, Store, Uuid, log::LogError, meta_store::MetaStoreError, sigchain::SigChainError, @@ -10,7 +9,6 @@ use lattice_core::{ }; use std::path::Path; use std::rc::Rc; -use std::cell::RefCell; use thiserror::Error; #[derive(Error, Debug)] @@ -35,6 +33,12 @@ pub enum NodeError { #[error("Already initialized")] AlreadyInitialized, + + #[error("Channel closed")] + ChannelClosed, + + #[error("Actor error: {0}")] + Actor(String), } pub struct NodeInfo { @@ -167,85 +171,101 @@ impl LatticeNode { }; let info = StoreInfo { store_id, entries_replayed }; + + // Spawn actor thread - actor owns store, sigchain, and node copy + let (tx, actor_handle) = crate::store_actor::spawn_store_actor( + store_id, + store, + sigchain, + (*self.node).clone(), + ); + let handle = StoreHandle { store_id, - node: Rc::clone(&self.node), - sigchain: RefCell::new(sigchain), - store, + tx, + actor_handle, }; Ok((handle, info)) } } -/// A handle to a specific store with KV operations +/// A handle to a specific store - wraps channel to actor thread pub struct StoreHandle { store_id: Uuid, - node: Rc, - sigchain: RefCell, - store: Store, + tx: std::sync::mpsc::Sender, + #[allow(dead_code)] + actor_handle: std::thread::JoinHandle<()>, } impl StoreHandle { pub fn id(&self) -> Uuid { self.store_id } pub fn get(&self, key: &[u8]) -> Result>, NodeError> { - Ok(self.store.get(key)?) + use crate::store_actor::StoreCmd; + let (resp_tx, resp_rx) = std::sync::mpsc::channel(); + self.tx.send(StoreCmd::Get { key: key.to_vec(), resp: resp_tx }) + .map_err(|_| NodeError::ChannelClosed)?; + resp_rx.recv() + .map_err(|_| NodeError::ChannelClosed)? + .map_err(NodeError::Store) } pub fn get_heads(&self, key: &[u8]) -> Result, NodeError> { - Ok(self.store.get_heads(key)?) + use crate::store_actor::StoreCmd; + let (resp_tx, resp_rx) = std::sync::mpsc::channel(); + self.tx.send(StoreCmd::GetHeads { key: key.to_vec(), resp: resp_tx }) + .map_err(|_| NodeError::ChannelClosed)?; + resp_rx.recv() + .map_err(|_| NodeError::ChannelClosed)? + .map_err(NodeError::Store) } pub fn list(&self) -> Result, Vec)>, NodeError> { - Ok(self.store.list_all()?) + use crate::store_actor::StoreCmd; + let (resp_tx, resp_rx) = std::sync::mpsc::channel(); + self.tx.send(StoreCmd::List { resp: resp_tx }) + .map_err(|_| NodeError::ChannelClosed)?; + resp_rx.recv() + .map_err(|_| NodeError::ChannelClosed)? + .map_err(NodeError::Store) } pub fn log_seq(&self) -> u64 { - self.sigchain.borrow().len() + use crate::store_actor::StoreCmd; + let (resp_tx, resp_rx) = std::sync::mpsc::channel(); + let _ = self.tx.send(StoreCmd::LogSeq { resp: resp_tx }); + resp_rx.recv().unwrap_or(0) } pub fn applied_seq(&self) -> Result { - let author = self.node.public_key_bytes(); - Ok(self.store.author_state(&author)? - .map(|s| s.seq) - .unwrap_or(0)) + use crate::store_actor::StoreCmd; + let (resp_tx, resp_rx) = std::sync::mpsc::channel(); + self.tx.send(StoreCmd::AppliedSeq { resp: resp_tx }) + .map_err(|_| NodeError::ChannelClosed)?; + resp_rx.recv() + .map_err(|_| NodeError::ChannelClosed)? + .map_err(NodeError::Store) } pub fn put(&self, key: &[u8], value: &[u8]) -> Result { - // Get current heads for this key to cite as parents - let heads = self.store.get_heads(key)?; - let parent_hashes: Vec> = heads.iter().map(|h| h.hash.clone()).collect(); - - self.commit_entry(parent_hashes, |b| b.put(key.to_vec(), value.to_vec())) + use crate::store_actor::StoreCmd; + let (resp_tx, resp_rx) = std::sync::mpsc::channel(); + self.tx.send(StoreCmd::Put { key: key.to_vec(), value: value.to_vec(), resp: resp_tx }) + .map_err(|_| NodeError::ChannelClosed)?; + resp_rx.recv() + .map_err(|_| NodeError::ChannelClosed)? + .map_err(|e| NodeError::Actor(e.to_string())) } pub fn delete(&self, key: &[u8]) -> Result { - // Get current heads for this key to cite as parents - let heads = self.store.get_heads(key)?; - let parent_hashes: Vec> = heads.iter().map(|h| h.hash.clone()).collect(); - - self.commit_entry(parent_hashes, |b| b.delete(key.to_vec())) - } - - fn commit_entry(&self, parent_hashes: Vec>, build: F) -> Result - where - F: FnOnce(EntryBuilder) -> EntryBuilder, - { - let mut sigchain = self.sigchain.borrow_mut(); - let seq = sigchain.len() + 1; - let prev_hash = sigchain.last_hash(); - - let builder = EntryBuilder::new(seq, HLC::now()) - .store_id(self.store_id.as_bytes().to_vec()) - .prev_hash(prev_hash.to_vec()) - .parent_hashes(parent_hashes); - let entry = build(builder).sign(&self.node); - - sigchain.append(&entry)?; - self.store.apply_entry(&entry)?; - - Ok(seq) + use crate::store_actor::StoreCmd; + let (resp_tx, resp_rx) = std::sync::mpsc::channel(); + self.tx.send(StoreCmd::Delete { key: key.to_vec(), resp: resp_tx }) + .map_err(|_| NodeError::ChannelClosed)?; + resp_rx.recv() + .map_err(|_| NodeError::ChannelClosed)? + .map_err(|e| NodeError::Actor(e.to_string())) } } diff --git a/lattice-cli/src/store_actor.rs b/lattice-cli/src/store_actor.rs new file mode 100644 index 0000000..af5d49e --- /dev/null +++ b/lattice-cli/src/store_actor.rs @@ -0,0 +1,189 @@ +//! Store Actor - dedicated thread that owns Store and processes commands via channel + +use lattice_core::{ + EntryBuilder, HeadInfo, Node, SigChain, Store, Uuid, + hlc::HLC, + proto::AuthorState, + sigchain::SigChainError, + store::StoreError, +}; +use std::sync::mpsc::{self, Receiver, Sender}; +use std::thread::{self, JoinHandle}; + +/// Commands sent to the store actor +pub enum StoreCmd { + Get { + key: Vec, + resp: std::sync::mpsc::Sender>, StoreError>>, + }, + GetHeads { + key: Vec, + resp: std::sync::mpsc::Sender, StoreError>>, + }, + List { + resp: std::sync::mpsc::Sender, Vec)>, StoreError>>, + }, + Put { + key: Vec, + value: Vec, + resp: std::sync::mpsc::Sender>, + }, + Delete { + key: Vec, + resp: std::sync::mpsc::Sender>, + }, + LogSeq { + resp: std::sync::mpsc::Sender, + }, + AppliedSeq { + resp: std::sync::mpsc::Sender>, + }, + AuthorState { + author: [u8; 32], + resp: std::sync::mpsc::Sender, StoreError>>, + }, + Shutdown, +} + +#[derive(Debug)] +pub enum StoreActorError { + Store(StoreError), + SigChain(SigChainError), + ChannelClosed, +} + +impl From for StoreActorError { + fn from(e: StoreError) -> Self { + StoreActorError::Store(e) + } +} + +impl From for StoreActorError { + fn from(e: SigChainError) -> Self { + StoreActorError::SigChain(e) + } +} + +impl std::fmt::Display for StoreActorError { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + StoreActorError::Store(e) => write!(f, "Store error: {}", e), + StoreActorError::SigChain(e) => write!(f, "SigChain error: {}", e), + StoreActorError::ChannelClosed => write!(f, "Channel closed"), + } + } +} + +impl std::error::Error for StoreActorError {} + +/// The store actor - runs in its own thread, owns Store and SigChain +pub struct StoreActor { + store_id: Uuid, + store: Store, + sigchain: SigChain, + node: Node, + rx: Receiver, +} + +impl StoreActor { + /// Create a new store actor (but don't start the thread yet) + pub fn new( + store_id: Uuid, + store: Store, + sigchain: SigChain, + node: Node, + rx: Receiver, + ) -> Self { + Self { + store_id, + store, + sigchain, + node, + rx, + } + } + + /// Run the actor loop - processes commands until Shutdown received + pub fn run(mut self) { + while let Ok(cmd) = self.rx.recv() { + match cmd { + StoreCmd::Get { key, resp } => { + let _ = resp.send(self.store.get(&key)); + } + StoreCmd::GetHeads { key, resp } => { + let _ = resp.send(self.store.get_heads(&key)); + } + StoreCmd::List { resp } => { + let _ = resp.send(self.store.list_all()); + } + StoreCmd::Put { key, value, resp } => { + let result = self.do_put(&key, &value); + let _ = resp.send(result); + } + StoreCmd::Delete { key, resp } => { + let result = self.do_delete(&key); + let _ = resp.send(result); + } + StoreCmd::LogSeq { resp } => { + let _ = resp.send(self.sigchain.len()); + } + StoreCmd::AppliedSeq { resp } => { + let author = self.node.public_key_bytes(); + let result = self.store.author_state(&author) + .map(|s| s.map(|a| a.seq).unwrap_or(0)); + let _ = resp.send(result); + } + StoreCmd::AuthorState { author, resp } => { + let _ = resp.send(self.store.author_state(&author)); + } + StoreCmd::Shutdown => { + break; + } + } + } + } + + fn do_put(&mut self, key: &[u8], value: &[u8]) -> Result { + let heads = self.store.get_heads(key)?; + let parent_hashes: Vec> = heads.iter().map(|h| h.hash.clone()).collect(); + self.commit_entry(parent_hashes, |b| b.put(key.to_vec(), value.to_vec())) + } + + fn do_delete(&mut self, key: &[u8]) -> Result { + let heads = self.store.get_heads(key)?; + let parent_hashes: Vec> = heads.iter().map(|h| h.hash.clone()).collect(); + self.commit_entry(parent_hashes, |b| b.delete(key.to_vec())) + } + + fn commit_entry(&mut self, parent_hashes: Vec>, build: F) -> Result + where + F: FnOnce(EntryBuilder) -> EntryBuilder, + { + let seq = self.sigchain.len() + 1; + let prev_hash = self.sigchain.last_hash(); + + let builder = EntryBuilder::new(seq, HLC::now()) + .store_id(self.store_id.as_bytes().to_vec()) + .prev_hash(prev_hash.to_vec()) + .parent_hashes(parent_hashes); + let entry = build(builder).sign(&self.node); + + self.sigchain.append(&entry)?; + self.store.apply_entry(&entry)?; + + Ok(seq) + } +} + +/// Spawn a store actor in a new thread, returns (sender, join_handle) +pub fn spawn_store_actor( + store_id: Uuid, + store: Store, + sigchain: SigChain, + node: Node, +) -> (Sender, JoinHandle<()>) { + let (tx, rx) = mpsc::channel(); + let actor = StoreActor::new(store_id, store, sigchain, node, rx); + let handle = thread::spawn(move || actor.run()); + (tx, handle) +} diff --git a/lattice-core/src/node.rs b/lattice-core/src/node.rs index 8c22bfb..96958cc 100644 --- a/lattice-core/src/node.rs +++ b/lattice-core/src/node.rs @@ -28,6 +28,7 @@ pub enum NodeError { /// /// Each node has an Ed25519 keypair used for signing sigchain entries /// and establishing trust within the network. +#[derive(Clone)] pub struct Node { signing_key: SigningKey, }