Deliver node-local channel directory events to the gRPC subscribe stream
This commit is contained in:
parent
af5896743d
commit
876bc55e2b
3 changed files with 34 additions and 1 deletions
|
|
@ -851,8 +851,12 @@ impl EventLog {
|
|||
}
|
||||
if global {
|
||||
self.versions.insert(entry.origin.clone(), entry.seq);
|
||||
self.notify(&entry);
|
||||
}
|
||||
// Push every committed entry to subscribers: the gossip forwarder filters to
|
||||
// Global before sending to a peer, but the gRPC directory stream wants the
|
||||
// node-local channel events too (a website mirroring the account + channel
|
||||
// directory, which registers/drops/founder-sets never reached before).
|
||||
self.notify(&entry);
|
||||
self.entries.push(entry);
|
||||
Ok(())
|
||||
}
|
||||
|
|
|
|||
|
|
@ -458,6 +458,26 @@
|
|||
assert_eq!(db.certfps("bob"), &[kept][..], "kept cert survives");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn channel_directory_events_reach_the_outbound_subscriber() {
|
||||
// A gRPC directory subscriber (a website mirror) reads the outbound stream.
|
||||
// Account events are Global and always flowed; channel events are Local and
|
||||
// were dropped before reaching outbound, so a mirror never saw registrations.
|
||||
let (tx, mut rx) = tokio::sync::broadcast::channel(64);
|
||||
let mut db = Db::open(tmp("chan-sub"), "N1");
|
||||
db.scram_iterations = 4096;
|
||||
db.set_outbound(tx);
|
||||
db.register("alice", "pw", None).unwrap();
|
||||
db.register_channel("#c", "alice").unwrap();
|
||||
let mut saw_channel = false;
|
||||
while let Ok(entry) = rx.try_recv() {
|
||||
if matches!(entry.event(), Event::ChannelRegistered { name, .. } if name == "#c") {
|
||||
saw_channel = true;
|
||||
}
|
||||
}
|
||||
assert!(saw_channel, "channel registration must reach the outbound directory stream");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn compaction_preserves_channel_last_used_and_levels() {
|
||||
let p = tmp("compact-chan");
|
||||
|
|
|
|||
|
|
@ -13,6 +13,8 @@ use tokio_rustls::rustls::pki_types::{CertificateDer, PrivateKeyDer, ServerName}
|
|||
use tokio_rustls::rustls::{self, ClientConfig, RootCertStore, ServerConfig};
|
||||
use tokio_rustls::{TlsAcceptor, TlsConnector};
|
||||
|
||||
use crate::engine::db::Scope;
|
||||
|
||||
use crate::config::{Gossip, Peer, Tls};
|
||||
use crate::engine::db::LogEntry;
|
||||
use crate::engine::Engine;
|
||||
|
|
@ -205,6 +207,13 @@ where
|
|||
loop {
|
||||
match rx.recv().await {
|
||||
Ok(entry) => {
|
||||
// Only global (account/identity) state replicates to peers;
|
||||
// channel state is node-local, so never forward it — the
|
||||
// outbound stream now also carries local events for the gRPC
|
||||
// directory subscriber.
|
||||
if entry.event().scope() != Scope::Global {
|
||||
continue;
|
||||
}
|
||||
if send(&tx, &Msg::Entry { entry: Box::new(entry) }).await.is_err() {
|
||||
break;
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue