gossip: replicate account state only, keep channels node-local
Each event now has a scope: account identity is global and gossips to every peer; channel state is local to the node that authored it. The version vector, push, digest and compaction only carry global events, and ingest refuses any non-global entry a peer sends. A node can no longer be handed ownership of a channel registered on a network it was never part of.
This commit is contained in:
parent
e8b55bdb27
commit
d4b9172455
3 changed files with 132 additions and 17 deletions
107
src/engine/db.rs
107
src/engine/db.rs
|
|
@ -50,6 +50,34 @@ pub enum Event {
|
|||
ChannelEntryMsgSet { channel: String, msg: String },
|
||||
}
|
||||
|
||||
// Whether an event replicates across the federation. Account identity is Global
|
||||
// (one owner, gossiped everywhere); channel state is Local (scoped to the one
|
||||
// network that authored it, so a node can't be handed ownership of a channel it
|
||||
// never saw registered). Exhaustive on purpose: a new event must pick a side.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
enum Scope {
|
||||
Global,
|
||||
Local,
|
||||
}
|
||||
|
||||
impl Event {
|
||||
fn scope(&self) -> Scope {
|
||||
match self {
|
||||
Event::AccountRegistered(_) | Event::CertAdded { .. } | Event::CertRemoved { .. } => Scope::Global,
|
||||
Event::ChannelRegistered { .. }
|
||||
| Event::ChannelDropped { .. }
|
||||
| Event::ChannelMlock { .. }
|
||||
| Event::ChannelAccessAdd { .. }
|
||||
| Event::ChannelAccessDel { .. }
|
||||
| Event::ChannelAkickAdd { .. }
|
||||
| Event::ChannelAkickDel { .. }
|
||||
| Event::ChannelFounderSet { .. }
|
||||
| Event::ChannelDescSet { .. }
|
||||
| Event::ChannelEntryMsgSet { .. } => Scope::Local,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// An access-list entry: an account and its level ("op" or "voice").
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct ChanAccess {
|
||||
|
|
@ -203,8 +231,12 @@ impl EventLog {
|
|||
(log, events)
|
||||
}
|
||||
|
||||
// Roll the clock and version vector forward over an entry.
|
||||
// Roll the clock and version vector forward over an entry. Local (channel)
|
||||
// entries carry no gossip identity, so they never touch the vector or clock.
|
||||
fn absorb(&mut self, entry: &LogEntry) {
|
||||
if entry.event.scope() != Scope::Global {
|
||||
return;
|
||||
}
|
||||
self.lamport = self.lamport.max(entry.lamport);
|
||||
self.versions.entry(entry.origin.clone()).and_modify(|s| *s = (*s).max(entry.seq)).or_insert(entry.seq);
|
||||
}
|
||||
|
|
@ -214,14 +246,22 @@ impl EventLog {
|
|||
self.versions.get(&self.origin).map_or(0, |s| s + 1)
|
||||
}
|
||||
|
||||
// Stamp a locally-authored event with the next seq + a ticked Lamport clock,
|
||||
// then persist it.
|
||||
// Persist a locally-authored event. Global events get the next seq + a ticked
|
||||
// Lamport clock and are pushed to peers; local (channel) events are written
|
||||
// for restart but never gossiped and carry no version-vector identity.
|
||||
fn append(&mut self, event: Event) -> std::io::Result<()> {
|
||||
self.lamport += 1;
|
||||
let entry = LogEntry { origin: self.origin.clone(), seq: self.next_seq(), lamport: self.lamport, event };
|
||||
let global = event.scope() == Scope::Global;
|
||||
let entry = if global {
|
||||
self.lamport += 1;
|
||||
LogEntry { origin: self.origin.clone(), seq: self.next_seq(), lamport: self.lamport, event }
|
||||
} else {
|
||||
LogEntry { origin: self.origin.clone(), seq: 0, lamport: 0, event }
|
||||
};
|
||||
self.persist(&entry)?;
|
||||
self.versions.insert(entry.origin.clone(), entry.seq);
|
||||
self.notify(&entry);
|
||||
if global {
|
||||
self.versions.insert(entry.origin.clone(), entry.seq);
|
||||
self.notify(&entry);
|
||||
}
|
||||
self.entries.push(entry);
|
||||
Ok(())
|
||||
}
|
||||
|
|
@ -230,6 +270,12 @@ impl EventLog {
|
|||
// event to fold into state, or None if already applied (idempotent, so
|
||||
// re-delivery converges). Assumes per-origin in-order delivery.
|
||||
fn ingest(&mut self, entry: LogEntry) -> std::io::Result<Option<Event>> {
|
||||
// A node never accepts another node's channel state — only global (account)
|
||||
// identity replicates. This is the guarantee: you can't be handed ownership
|
||||
// of a channel that was registered on a network you're not part of.
|
||||
if entry.event.scope() != Scope::Global {
|
||||
return Ok(None);
|
||||
}
|
||||
if self.versions.get(&entry.origin).is_some_and(|&s| entry.seq <= s) {
|
||||
return Ok(None); // already have it
|
||||
}
|
||||
|
|
@ -259,10 +305,12 @@ impl EventLog {
|
|||
self.versions.clone()
|
||||
}
|
||||
|
||||
// Entries a peer is missing, given the version vector it advertised.
|
||||
// Global entries a peer is missing, given the version vector it advertised.
|
||||
// Local (channel) entries are never offered — they don't leave the node.
|
||||
fn missing_for(&self, peer: &HashMap<String, u64>) -> Vec<LogEntry> {
|
||||
self.entries
|
||||
.iter()
|
||||
.filter(|e| e.event.scope() == Scope::Global)
|
||||
.filter(|e| peer.get(&e.origin).map_or(true, |&s| e.seq > s))
|
||||
.cloned()
|
||||
.collect()
|
||||
|
|
@ -286,16 +334,19 @@ impl EventLog {
|
|||
// leaves the old log.
|
||||
fn compact(&mut self, events: Vec<Event>) -> std::io::Result<()> {
|
||||
let mut seq = self.next_seq();
|
||||
let mut last_global = None;
|
||||
let mut snapshot = Vec::with_capacity(events.len());
|
||||
for event in events {
|
||||
self.lamport += 1;
|
||||
snapshot.push(LogEntry {
|
||||
origin: self.origin.clone(),
|
||||
seq,
|
||||
lamport: self.lamport,
|
||||
event,
|
||||
});
|
||||
seq += 1;
|
||||
let entry = if event.scope() == Scope::Global {
|
||||
self.lamport += 1;
|
||||
let e = LogEntry { origin: self.origin.clone(), seq, lamport: self.lamport, event };
|
||||
last_global = Some(seq);
|
||||
seq += 1;
|
||||
e
|
||||
} else {
|
||||
LogEntry { origin: self.origin.clone(), seq: 0, lamport: 0, event }
|
||||
};
|
||||
snapshot.push(entry);
|
||||
}
|
||||
let tmp = self.path.with_extension("compact");
|
||||
{
|
||||
|
|
@ -307,7 +358,7 @@ impl EventLog {
|
|||
}
|
||||
std::fs::rename(&tmp, &self.path)?;
|
||||
self.versions = HashMap::new();
|
||||
if let Some(last) = seq.checked_sub(1) {
|
||||
if let Some(last) = last_global {
|
||||
self.versions.insert(self.origin.clone(), last);
|
||||
}
|
||||
self.entries = snapshot;
|
||||
|
|
@ -829,6 +880,28 @@ mod tests {
|
|||
assert!(!glob_match("alice!*@*", "bob!~b@h"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn channel_state_is_node_local_but_persists() {
|
||||
let path = tmp("scope");
|
||||
{
|
||||
let mut db = Db::open(&path, "N1");
|
||||
db.register("alice", "pw", None).unwrap(); // global
|
||||
db.register_channel("#c", "alice").unwrap(); // local
|
||||
db.set_mlock("#c", "nt", "").unwrap(); // local
|
||||
// Only the account advances the version vector.
|
||||
assert_eq!(db.version_vector().get("N1"), Some(&0), "channels don't advance the vector");
|
||||
// A fresh peer is offered the account, never the channel.
|
||||
let missing = db.missing_for(&HashMap::new());
|
||||
assert_eq!(missing.len(), 1, "only the global account entry is offered: {missing:?}");
|
||||
assert!(matches!(missing[0].event, Event::AccountRegistered(_)));
|
||||
}
|
||||
// Reopen: node-local channel state replays from disk unchanged.
|
||||
let db = Db::open(&path, "N1");
|
||||
assert!(db.exists("alice"));
|
||||
assert_eq!(db.channel("#c").map(|c| c.lock_on.clone()), Some("nt".to_string()), "channel state persists locally");
|
||||
assert_eq!(db.version_vector().get("N1"), Some(&0), "vector unchanged after replay");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn akick_add_del_and_match() {
|
||||
let mut db = Db::open(&tmp("akick"), "N1");
|
||||
|
|
|
|||
|
|
@ -127,6 +127,16 @@ impl Engine {
|
|||
self.db.exists(name)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn test_register_channel(&mut self, name: &str, founder: &str) {
|
||||
self.db.register_channel(name, founder).unwrap();
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn test_has_channel(&self, name: &str) -> bool {
|
||||
self.db.channel(name).is_some()
|
||||
}
|
||||
|
||||
// Insert or refresh a client's in-progress SASL session, stamped now.
|
||||
fn stash_sasl(&mut self, client: String, session: SaslSession) {
|
||||
self.sasl_sessions.insert(client, TimedSession { touched: Instant::now(), session });
|
||||
|
|
|
|||
|
|
@ -327,6 +327,38 @@ mod tests {
|
|||
assert!(got, "a post-connect write should push to B well under the 10s digest");
|
||||
}
|
||||
|
||||
// Accounts federate, but channel registrations stay on the node that made
|
||||
// them: A registers both; B receives the account and never the channel.
|
||||
#[tokio::test]
|
||||
async fn channels_stay_node_local() {
|
||||
let (a, atx) = engine("A", "local-a");
|
||||
let (b, btx) = engine("B", "local-b");
|
||||
a.lock().await.test_register("alice");
|
||||
a.lock().await.test_register_channel("#secret", "alice");
|
||||
|
||||
let (ca, cb) = tokio::io::duplex(64 * 1024);
|
||||
let sa = tokio::spawn(session(ca, a.clone(), "s3cret".into(), "A".into(), atx));
|
||||
let sb = tokio::spawn(session(cb, b.clone(), "s3cret".into(), "B".into(), btx));
|
||||
|
||||
// Wait for the account to converge (proves the link works), then assert
|
||||
// the channel never crossed it.
|
||||
let mut got_account = false;
|
||||
for _ in 0..100 {
|
||||
if b.lock().await.test_has_account("alice") {
|
||||
got_account = true;
|
||||
break;
|
||||
}
|
||||
tokio::time::sleep(Duration::from_millis(20)).await;
|
||||
}
|
||||
// Give any (erroneous) channel replication ample time to arrive too.
|
||||
tokio::time::sleep(Duration::from_millis(200)).await;
|
||||
let leaked_channel = b.lock().await.test_has_channel("#secret");
|
||||
sa.abort();
|
||||
sb.abort();
|
||||
assert!(got_account, "the account should federate");
|
||||
assert!(!leaked_channel, "the channel must NOT federate to a node that never saw it registered");
|
||||
}
|
||||
|
||||
// A wrong secret is rejected before any state is exchanged.
|
||||
#[tokio::test]
|
||||
async fn bad_secret_is_rejected() {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue