use super::*;
use crate::broadcast::fan_out_evicting;
pub(super) fn violates_broadcast_origin_contract(
session_id: Option<u64>,
msg: &DaemonMessage,
) -> bool {
match msg {
DaemonMessage::Session {
session_id: envelope_id,
..
} => match (session_id, envelope_id) {
(Some(origin), Some(inner)) => origin != *inner,
(Some(_), None) | (None, Some(_)) => true,
(None, None) => false,
},
_ => session_id.is_some(),
}
}
impl DaemonState {
pub(super) fn broadcast(&mut self, msg: DaemonMessage) {
let (evict_activity, evict_activity_largest) = fan_out_evicting(
&mut self.activity_subscribers,
&msg,
&self.lag_limits,
&self.global_lag,
|_| false, );
self.finish_evictions(evict_activity, evict_activity_largest);
let (evict_clients, evict_largest) = fan_out_evicting(
&mut self.summary_subscribers,
&msg,
&self.lag_limits,
&self.global_lag,
|client_id| {
self.activity_subscribers.contains_key(&client_id)
},
);
self.finish_evictions(evict_clients, evict_largest);
}
pub(super) fn finish_evictions(&mut self, evict_clients: Vec<u64>, evict_largest: bool) {
for client_id in evict_clients {
self.handle_evict_client(client_id);
}
if evict_largest {
self.handle_evict_largest_lagging();
}
}
pub(super) fn handle_register_summary_subscriber(
&mut self,
client_id: u64,
writer: SubscriberSink,
) {
self.summary_subscribers.insert(client_id, writer);
}
pub(super) fn handle_unregister_summary_subscriber(&mut self, client_id: u64) {
self.summary_subscribers.remove(&client_id);
}
pub(super) fn handle_broadcast_session_status(
&mut self,
session_id: u64,
status: SessionStatus,
) {
let last_modified = match self.session_metadata.get_mut(&session_id) {
Some(meta) => {
meta.status = status.clone();
meta.last_modified
}
None => 0,
};
let msg = DaemonMessage::Session {
session_id: Some(session_id),
event: SessionEvent::SessionStatusChanged {
status,
last_modified,
},
};
if self.session_metadata.contains_key(&session_id) {
let (evict_clients, evict_largest) = fan_out_evicting(
&mut self.summary_subscribers,
&msg,
&self.lag_limits,
&self.global_lag,
|client_id| {
if self
.client_subscribed_sessions
.get(&client_id)
.is_some_and(|sessions| sessions.contains(&session_id))
{
return true;
}
self.activity_subscribers.contains_key(&client_id)
},
);
self.finish_evictions(evict_clients, evict_largest);
}
}
pub(super) fn handle_register_activity_subscriber(
&mut self,
client_id: u64,
writer: SubscriberSink,
) {
info!("registering activity subscriber: client_id={}", client_id);
self.activity_subscribers.insert(client_id, writer.clone());
let providers = catalog_provider_pairs();
let _ = writer.enqueue(
&DaemonMessage::CatalogUpdated { providers },
&self.lag_limits,
&self.global_lag,
);
let lock_msg = self.current_lock_message();
let _ = writer.enqueue(&lock_msg, &self.lag_limits, &self.global_lag);
}
pub(super) fn handle_unregister_activity_subscriber(&mut self, client_id: u64) {
debug!("unregistering activity subscriber: client_id={}", client_id);
self.activity_subscribers.remove(&client_id);
}
pub(super) fn current_lock_message(&self) -> DaemonMessage {
if self.locked {
DaemonMessage::Locked
} else {
DaemonMessage::Unlocked
}
}
pub(super) fn broadcast_lock_state(&mut self) {
let msg = self.current_lock_message();
self.handle_broadcast_activity(None, msg);
}
pub(super) fn handle_register_client_writer(&mut self, client_id: u64, writer: SubscriberSink) {
debug!("registering client writer: client_id={}", client_id);
self.client_writers.insert(client_id, writer);
}
pub(super) fn handle_evict_client(&mut self, client_id: u64) {
let Some(sink) = self.client_writers.get(&client_id) else {
return;
};
warn!(
"evicting lagging client: client_id={}, backlog_bytes={}",
client_id,
sink.bytes_in_flight.load(Ordering::Relaxed)
);
let _ = sink.send_unchecked(&DaemonMessage::Evicted, &self.global_lag);
self.summary_subscribers.remove(&client_id);
self.activity_subscribers.remove(&client_id);
self.remove_client_from_sessions(client_id);
self.client_writers.remove(&client_id);
crate::metrics::record_eviction();
}
pub(super) fn handle_evict_largest_lagging(&mut self) {
let best = self
.client_writers
.iter()
.filter(|(_, sink)| sink.bytes_in_flight.load(Ordering::Relaxed) > 0)
.max_by_key(|(_, sink)| sink.bytes_in_flight.load(Ordering::Relaxed));
if let Some((client_id, _)) = best {
self.handle_evict_client(*client_id);
}
}
pub(super) fn handle_broadcast_shutting_down(&mut self) {
let clients = self.client_writers.len();
info!("broadcasting ShuttingDown to {clients} client(s)");
self.client_writers.retain(|client_id, sink| {
if sink.send_unchecked(&DaemonMessage::ShuttingDown, &self.global_lag) {
true
} else {
warn!("removing disconnected client {client_id} during shutdown");
false
}
});
}
pub(super) fn handle_client_disconnected(&mut self, client_id: u64) {
info!("client disconnected cleanup: client_id={}", client_id);
self.summary_subscribers.remove(&client_id);
self.activity_subscribers.remove(&client_id);
self.remove_client_from_sessions(client_id);
self.client_writers.remove(&client_id);
}
pub(super) fn remove_client_from_sessions(&mut self, client_id: u64) {
if let Some(sessions) = self.client_subscribed_sessions.remove(&client_id) {
for session_id in &sessions {
if let Some(entry) = self.active_sessions.get(session_id) {
let _ = entry
.cmd_tx
.send(SessionCommand::RemoveSubscriber { client_id });
}
}
}
}
pub(super) fn handle_track_session_subscription(&mut self, client_id: u64, session_id: u64) {
debug!(
"track session subscription: client_id={}, session_id={}",
client_id, session_id
);
self.client_subscribed_sessions
.entry(client_id)
.or_default()
.insert(session_id);
}
pub(super) fn handle_untrack_session_subscription(&mut self, client_id: u64, session_id: u64) {
debug!(
"untrack session subscription: client_id={}, session_id={}",
client_id, session_id
);
if let std::collections::hash_map::Entry::Occupied(mut entry) =
self.client_subscribed_sessions.entry(client_id)
{
entry.get_mut().remove(&session_id);
if entry.get().is_empty() {
entry.remove();
}
}
}
pub(super) fn handle_broadcast_activity(
&mut self,
session_id: Option<u64>,
msg: DaemonMessage,
) {
if violates_broadcast_origin_contract(session_id, &msg) {
warn!(
session_id,
"BroadcastActivity violates the origin contract: command provenance and \
message origin disagree, so some subscriber class will miss this message \
or receive it twice — no current producer does this, inspect the caller"
);
}
let (evict_clients, evict_largest) = fan_out_evicting(
&mut self.activity_subscribers,
&msg,
&self.lag_limits,
&self.global_lag,
|client_id| {
if let Some(sid) = session_id
&& let Some(sessions) = self.client_subscribed_sessions.get(&client_id)
&& sessions.contains(&sid)
{
return true;
}
false
},
);
self.finish_evictions(evict_clients, evict_largest);
}
}