use tokio::sync::broadcast::Receiver;
use tokio::sync::broadcast::error::RecvError;
use crate::MemoryLayer;
pub(crate) const DEFAULT_EVENT_CAPACITY: usize = 1024;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum MemoryEventKind {
Remembered,
Invalidated,
Forgotten,
Consolidated,
}
#[derive(Debug, Clone)]
pub struct MemoryEvent {
pub agent_id: String,
pub kind: MemoryEventKind,
pub layer: MemoryLayer,
pub id: String,
}
#[derive(Debug)]
pub struct MemorySubscription {
rx: Receiver<MemoryEvent>,
agent_id: String,
layer: Option<MemoryLayer>,
}
impl MemorySubscription {
pub(crate) fn new(rx: Receiver<MemoryEvent>, agent_id: String, layer: Option<MemoryLayer>) -> Self {
Self { rx, agent_id, layer }
}
pub async fn recv(&mut self) -> Option<MemoryEvent> {
loop {
match self.rx.recv().await {
Ok(ev) if self.delivers(&ev) => return Some(ev),
Ok(_) => continue,
Err(RecvError::Lagged(_)) => continue,
Err(RecvError::Closed) => return None,
}
}
}
fn delivers(&self, ev: &MemoryEvent) -> bool {
ev.agent_id == self.agent_id && self.layer.is_none_or(|l| l == ev.layer)
}
}