Archived
d1ec384b27
- crates/ → engrams/ (Rust engrams live here) - bindings/ → receptors/ (cross-language access points into the graph) - Cargo.toml workspace paths updated
266 lines
9.7 KiB
Rust
266 lines
9.7 KiB
Rust
/// Engram Sync — swarm memory protocol for distributed Engram instances.
|
|
///
|
|
/// This crate turns Engram from a local-first database into a distributed
|
|
/// swarm memory protocol. Multiple independent Engram instances (peers) can
|
|
/// share memory across each other using delta sync and can fan spreading
|
|
/// activation out across the swarm, merging results by strength.
|
|
///
|
|
/// # Architecture
|
|
///
|
|
/// ```text
|
|
/// [Neuron-A Engram] <-> sync <-> [Neuron-B Engram]
|
|
/// |
|
|
/// [Neuron-C Engram]
|
|
///
|
|
/// Swarm activation: seed on A → propagate locally → fan-out to B and C
|
|
/// → merge all results → unified ranked response
|
|
/// ```
|
|
///
|
|
/// # Protocol
|
|
///
|
|
/// - Each peer is local and authoritative. There is no central server.
|
|
/// - Peers sync via delta exchange: "give me everything since timestamp T".
|
|
/// - Only configured memory tiers flow between peers (Semantic by default;
|
|
/// Episodic and Working are private unless explicitly enabled).
|
|
/// - Trusted peers get the full configured tier set; untrusted peers get
|
|
/// Semantic only.
|
|
/// - Swarm activation fans out to all trusted peers in parallel, deduplicates
|
|
/// by UUID (keeping strongest activation), and re-ranks.
|
|
|
|
pub mod client;
|
|
pub mod engine;
|
|
pub mod types;
|
|
|
|
// Public surface
|
|
pub use engine::{merge_activation_results, SyncEngine};
|
|
pub use types::{
|
|
MergedActivatedNode, Peer, PeerActivationResult, PeerStatus, PeerSyncResult,
|
|
SerializableActivatedNode, SyncConfig, SyncDelta, SyncReport, SwarmActivateRequest,
|
|
SwarmActivateResponse,
|
|
};
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
use engram_core::types::{MemoryTier, Node, NodeType};
|
|
use uuid::Uuid;
|
|
|
|
fn make_node(tier: MemoryTier) -> Node {
|
|
Node::new(
|
|
NodeType::Concept,
|
|
vec![0.1, 0.2, 0.3],
|
|
b"test content".to_vec(),
|
|
tier,
|
|
0.5,
|
|
)
|
|
}
|
|
|
|
fn make_activated(node: Node, strength: f32, hops: u8) -> SerializableActivatedNode {
|
|
SerializableActivatedNode {
|
|
node,
|
|
activation_strength: strength,
|
|
hops,
|
|
}
|
|
}
|
|
|
|
// ── merge_activation_results tests ───────────────────────────────────────
|
|
|
|
#[test]
|
|
fn merge_empty() {
|
|
let merged = merge_activation_results(&[], &[], 10);
|
|
assert!(merged.is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn merge_local_only() {
|
|
let node = make_node(MemoryTier::Semantic);
|
|
let node_id = node.id;
|
|
let local = vec![make_activated(node, 0.8, 1)];
|
|
let merged = merge_activation_results(&local, &[], 10);
|
|
assert_eq!(merged.len(), 1);
|
|
assert_eq!(merged[0].node.id, node_id);
|
|
assert!(merged[0].source_peer.is_none(), "local nodes have no source_peer");
|
|
assert!((merged[0].activation_strength - 0.8).abs() < f32::EPSILON);
|
|
}
|
|
|
|
#[test]
|
|
fn merge_peer_result_included() {
|
|
let local: Vec<SerializableActivatedNode> = vec![];
|
|
let node = make_node(MemoryTier::Semantic);
|
|
let peer_id = Uuid::new_v4();
|
|
let peer_results = vec![PeerActivationResult {
|
|
peer_id,
|
|
peer_name: "peer-a".into(),
|
|
results: vec![make_activated(node, 0.6, 2)],
|
|
error: None,
|
|
}];
|
|
let merged = merge_activation_results(&local, &peer_results, 10);
|
|
assert_eq!(merged.len(), 1);
|
|
assert_eq!(merged[0].source_peer, Some(peer_id));
|
|
}
|
|
|
|
#[test]
|
|
fn merge_deduplicates_by_uuid_keeps_strongest() {
|
|
let node = make_node(MemoryTier::Semantic);
|
|
let id = node.id;
|
|
|
|
// Same node appears locally at 0.4 and from a peer at 0.9
|
|
let local = vec![make_activated(node.clone(), 0.4, 1)];
|
|
let peer_id = Uuid::new_v4();
|
|
let peer_results = vec![PeerActivationResult {
|
|
peer_id,
|
|
peer_name: "peer-a".into(),
|
|
results: vec![make_activated(node, 0.9, 1)],
|
|
error: None,
|
|
}];
|
|
|
|
let merged = merge_activation_results(&local, &peer_results, 10);
|
|
// Should be deduplicated to 1 result
|
|
assert_eq!(merged.len(), 1);
|
|
// The stronger version (0.9, from peer) should win
|
|
assert!((merged[0].activation_strength - 0.9).abs() < f32::EPSILON);
|
|
assert_eq!(merged[0].source_peer, Some(peer_id));
|
|
}
|
|
|
|
#[test]
|
|
fn merge_respects_limit() {
|
|
let local: Vec<SerializableActivatedNode> = (0..20)
|
|
.map(|i| make_activated(make_node(MemoryTier::Semantic), i as f32 / 20.0, 1))
|
|
.collect();
|
|
let merged = merge_activation_results(&local, &[], 5);
|
|
assert_eq!(merged.len(), 5, "limit must be respected");
|
|
}
|
|
|
|
#[test]
|
|
fn merge_sorted_by_strength_descending() {
|
|
let strengths = vec![0.3f32, 0.9, 0.1, 0.7, 0.5];
|
|
let local: Vec<SerializableActivatedNode> = strengths
|
|
.iter()
|
|
.map(|&s| make_activated(make_node(MemoryTier::Semantic), s, 1))
|
|
.collect();
|
|
let merged = merge_activation_results(&local, &[], 10);
|
|
let result_strengths: Vec<f32> = merged.iter().map(|m| m.activation_strength).collect();
|
|
// Should be in descending order
|
|
for window in result_strengths.windows(2) {
|
|
assert!(window[0] >= window[1], "results must be sorted descending");
|
|
}
|
|
}
|
|
|
|
#[test]
|
|
fn merge_skips_errored_peers() {
|
|
let peer_results = vec![PeerActivationResult {
|
|
peer_id: Uuid::new_v4(),
|
|
peer_name: "failed-peer".into(),
|
|
results: vec![make_activated(make_node(MemoryTier::Semantic), 0.9, 1)],
|
|
error: Some("connection refused".into()),
|
|
}];
|
|
let merged = merge_activation_results(&[], &peer_results, 10);
|
|
// Errored peers should be excluded
|
|
assert!(merged.is_empty());
|
|
}
|
|
|
|
// ── SyncDelta serialization ───────────────────────────────────────────────
|
|
|
|
#[test]
|
|
fn sync_delta_roundtrips_json() {
|
|
let delta = SyncDelta {
|
|
peer_id: Uuid::new_v4(),
|
|
since: 1000,
|
|
nodes: vec![make_node(MemoryTier::Semantic)],
|
|
edges: vec![],
|
|
tombstones: vec![Uuid::new_v4()],
|
|
generated_at: 2000,
|
|
};
|
|
let json = serde_json::to_string(&delta).expect("serialize delta");
|
|
let decoded: SyncDelta = serde_json::from_str(&json).expect("deserialize delta");
|
|
assert_eq!(delta.peer_id, decoded.peer_id);
|
|
assert_eq!(delta.since, decoded.since);
|
|
assert_eq!(delta.nodes.len(), decoded.nodes.len());
|
|
assert_eq!(delta.tombstones.len(), decoded.tombstones.len());
|
|
}
|
|
|
|
// ── SyncEngine peer management ────────────────────────────────────────────
|
|
|
|
#[test]
|
|
fn engine_add_remove_peer() {
|
|
use engram_core::EngramDb;
|
|
use std::sync::{Arc, Mutex};
|
|
|
|
let dir = tempfile::tempdir().unwrap();
|
|
let db = Arc::new(Mutex::new(EngramDb::open(dir.path()).unwrap()));
|
|
let config = SyncConfig::default();
|
|
let mut engine = SyncEngine::new(db, config);
|
|
|
|
assert_eq!(engine.list_peers().len(), 0);
|
|
|
|
let peer = Peer {
|
|
id: Uuid::new_v4(),
|
|
name: "test-peer".into(),
|
|
address: "http://localhost:9999".into(),
|
|
api_key: "secret".into(),
|
|
sync_tiers: vec![MemoryTier::Semantic],
|
|
last_sync_at: 0,
|
|
trusted: true,
|
|
};
|
|
let peer_id = peer.id;
|
|
engine.add_peer(peer);
|
|
assert_eq!(engine.list_peers().len(), 1);
|
|
|
|
engine.remove_peer(peer_id);
|
|
assert_eq!(engine.list_peers().len(), 0);
|
|
}
|
|
|
|
#[test]
|
|
fn engine_add_peer_replaces_existing() {
|
|
use engram_core::EngramDb;
|
|
use std::sync::{Arc, Mutex};
|
|
|
|
let dir = tempfile::tempdir().unwrap();
|
|
let db = Arc::new(Mutex::new(EngramDb::open(dir.path()).unwrap()));
|
|
let mut engine = SyncEngine::new(db, SyncConfig::default());
|
|
|
|
let id = Uuid::new_v4();
|
|
for i in 0..3 {
|
|
engine.add_peer(Peer {
|
|
id,
|
|
name: format!("peer-v{}", i),
|
|
address: "http://localhost:1234".into(),
|
|
api_key: "k".into(),
|
|
sync_tiers: vec![],
|
|
last_sync_at: 0,
|
|
trusted: false,
|
|
});
|
|
}
|
|
// Should still be just one peer (latest version)
|
|
assert_eq!(engine.list_peers().len(), 1);
|
|
assert_eq!(engine.list_peers()[0].name, "peer-v2");
|
|
}
|
|
|
|
// ── generate_delta ────────────────────────────────────────────────────────
|
|
|
|
#[test]
|
|
fn generate_delta_filters_by_tier() {
|
|
use engram_core::types::MemoryTier;
|
|
use engram_core::EngramDb;
|
|
use std::sync::{Arc, Mutex};
|
|
|
|
let dir = tempfile::tempdir().unwrap();
|
|
let db = Arc::new(Mutex::new(EngramDb::open(dir.path()).unwrap()));
|
|
|
|
// Insert nodes in different tiers
|
|
{
|
|
let db_locked = db.lock().unwrap();
|
|
db_locked.put_node(make_node(MemoryTier::Semantic)).unwrap();
|
|
db_locked.put_node(make_node(MemoryTier::Episodic)).unwrap();
|
|
db_locked.put_node(make_node(MemoryTier::Working)).unwrap();
|
|
}
|
|
|
|
let engine = SyncEngine::new(db, SyncConfig::default());
|
|
|
|
// Only request Semantic tier
|
|
let delta = engine.generate_delta(0, &[MemoryTier::Semantic]).unwrap();
|
|
assert_eq!(delta.nodes.len(), 1, "only Semantic node should be in delta");
|
|
assert!(delta.nodes[0].tier == MemoryTier::Semantic);
|
|
}
|
|
}
|