use std::path::Path;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use serde_json::{Value, json};
use tokio::sync::mpsc;
use crate::memory_rpc::MAX_FRAME_BYTES;
use crate::uds::send_framed_stream_request_capped;
use super::parsers::parse_memory_event;
use super::types::MemoryEvent;
const METHOD: &str = "memory.activity_stream";
const FRAME_TIMEOUT: Duration = Duration::from_secs(300);
const FEED_BUFFER: usize = 256;
#[derive(Debug)]
pub struct ActivityFeed {
rx: mpsc::Receiver<MemoryEvent>,
last_error: Arc<Mutex<Option<String>>>,
live: Arc<std::sync::atomic::AtomicBool>,
}
impl ActivityFeed {
pub async fn open(socket: &Path) -> anyhow::Result<Self> {
let request = json!({
"jsonrpc": "2.0",
"id": 1,
"method": METHOD,
"stream": true,
});
let mut stream = send_framed_stream_request_capped::<_, Value>(
socket,
&request,
FRAME_TIMEOUT,
MAX_FRAME_BYTES,
)
.await?;
let (tx, rx) = mpsc::channel(FEED_BUFFER);
let last_error = Arc::new(Mutex::new(None));
let live = Arc::new(std::sync::atomic::AtomicBool::new(true));
let task_error = Arc::clone(&last_error);
let task_live = Arc::clone(&live);
tokio::spawn(async move {
loop {
match stream.next_frame().await {
Some(Ok(frame)) => {
if frame.get("type").and_then(Value::as_str) == Some("lagged") {
let n = frame.get("lagged").and_then(Value::as_u64).unwrap_or(0);
*task_error.lock().unwrap_or_else(|e| e.into_inner()) =
Some(format!("activity feed lagged {n} events"));
continue;
}
let Some(event) = parse_memory_event(&frame) else {
continue;
};
if tx.send(event).await.is_err() {
break;
}
}
Some(Err(e)) => {
*task_error.lock().unwrap_or_else(|e| e.into_inner()) =
Some(format!("{e}"));
break;
}
None => {
*task_error.lock().unwrap_or_else(|e| e.into_inner()) =
Some("the daemon closed the activity stream".to_string());
break;
}
}
}
task_live.store(false, std::sync::atomic::Ordering::Relaxed);
});
Ok(Self {
rx,
last_error,
live,
})
}
pub fn is_live(&self) -> bool {
self.live.load(std::sync::atomic::Ordering::Relaxed)
}
pub fn take_last_error(&self) -> Option<String> {
self.last_error
.lock()
.unwrap_or_else(|e| e.into_inner())
.take()
}
pub fn drain(&mut self) -> Vec<MemoryEvent> {
let mut out = Vec::new();
while let Ok(event) = self.rx.try_recv() {
out.push(event);
}
out
}
}
#[cfg(test)]
#[path = "feed_tests.rs"]
mod tests;