use std::sync::Arc;
use std::time::Duration;
use crate::store::{BusMessage, SharedStore};
use crate::wave::runtime::WaveRuntime;
pub const POLL_CADENCE: Duration = Duration::from_millis(250);
#[derive(Debug)]
pub struct BusListener {
runtime: Arc<WaveRuntime>,
store: SharedStore,
subscriber: String,
cursor: tokio::sync::Mutex<i64>,
}
impl BusListener {
pub fn new(runtime: Arc<WaveRuntime>, store: SharedStore) -> Self {
let subscriber = runtime.channel_name().to_string();
Self {
runtime,
store,
subscriber,
cursor: tokio::sync::Mutex::new(0),
}
}
pub async fn attach(&self) -> anyhow::Result<()> {
let stored = self.store.bus_cursor(self.subscriber.clone()).await?;
let cursor = match stored {
Some(cursor) => self.report_gap(cursor).await?,
None => self.store.bus_head().await?,
};
*self.cursor.lock().await = cursor;
self.store
.set_bus_cursor(self.subscriber.clone(), cursor)
.await?;
Ok(())
}
pub async fn run(self: Arc<Self>, cadence: Duration) {
if let Err(err) = self.attach().await {
tracing::warn!(error = %err, "bus subscription failed to attach; the mind is deaf");
return;
}
loop {
if let Err(err) = self.poll_once().await {
tracing::warn!(error = %err, "bus poll failed; retrying");
}
tokio::time::sleep(cadence).await;
}
}
pub async fn poll_once(&self) -> anyhow::Result<()> {
let mut cursor = self.cursor.lock().await;
let messages = self.store.read_bus_after(*cursor).await?;
for message in &messages {
if self.hears(message) {
self.runtime
.deliver_say(message.text.clone(), message.byline.clone());
}
*cursor = message.id;
self.store
.set_bus_cursor(self.subscriber.clone(), *cursor)
.await?;
}
Ok(())
}
fn hears(&self, message: &BusMessage) -> bool {
self.runtime.in_family(&message.channel) && message.byline != self.subscriber
}
async fn report_gap(&self, cursor: i64) -> anyhow::Result<i64> {
let Some(floor) = self.store.bus_floor().await? else {
return Ok(cursor);
};
if floor <= cursor + 1 {
return Ok(cursor);
}
let missed = floor - cursor - 1;
self.runtime.deliver_say(
format!(
"bus cursor jumped {cursor} → {floor}: {missed} broadcast(s) aged past the \
sweep window while this mind was asleep. The PRs and `lf runs` hold what \
the bus dropped."
),
"bus".to_string(),
);
Ok(floor - 1)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::store::sqlite::BUS_WINDOW_SECS;
use crate::store::{open_store, StorageConfig};
async fn temp_store(dir: &std::path::Path) -> SharedStore {
Arc::new(
open_store(&StorageConfig::sqlite(dir.join("loopflow.db")))
.await
.expect("open sqlite store"),
)
}
#[tokio::test]
async fn a_sleeping_mind_catches_up_exactly_once() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let origin = tmp.path().join("repo");
std::fs::create_dir_all(&origin).unwrap();
let runtime = WaveRuntime::open("ship".into(), origin.clone()).expect("runtime");
let listener = BusListener::new(runtime.clone(), store.clone());
listener.attach().await.expect("attach");
listener.poll_once().await.expect("poll");
assert!(runtime.thread_snapshot().is_empty());
drop(runtime);
store
.publish_bus("ship.148e".into(), "ship.148e".into(), "landed PR".into())
.await
.expect("publish");
let runtime = WaveRuntime::open("ship".into(), origin.clone()).expect("runtime");
let listener = BusListener::new(runtime.clone(), store.clone());
listener.attach().await.expect("attach");
listener.poll_once().await.expect("poll");
let thread = runtime.thread_snapshot();
assert_eq!(thread.len(), 1);
assert_eq!(thread[0].text, "landed PR");
assert_eq!(thread[0].from.as_deref(), Some("ship.148e"));
drop(runtime);
let runtime = WaveRuntime::open("ship".into(), origin).expect("runtime");
let listener = BusListener::new(runtime.clone(), store.clone());
listener.attach().await.expect("attach");
listener.poll_once().await.expect("poll");
assert_eq!(runtime.thread_snapshot().len(), 1, "exactly once");
}
#[tokio::test]
async fn the_mind_folds_family_reports_and_ignores_its_own_voice() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let origin = tmp.path().join("repo");
std::fs::create_dir_all(&origin).unwrap();
let runtime = WaveRuntime::open("ship".into(), origin).expect("runtime");
let listener = BusListener::new(runtime.clone(), store.clone());
listener.attach().await.expect("attach");
for (channel, byline, text) in [
("goals", "goals", "another family"),
("ship.148e", "ship", "the mind steering its hand"),
("ship.148e", "ship.148e", "the hand reporting"),
("ship", "concerto", "a child wave escalating"),
] {
store
.publish_bus(channel.into(), byline.into(), text.into())
.await
.expect("publish");
}
listener.poll_once().await.expect("poll");
let thread = runtime.thread_snapshot();
let texts: Vec<&str> = thread.iter().map(|turn| turn.text.as_str()).collect();
assert_eq!(texts, vec!["the hand reporting", "a child wave escalating"]);
}
#[tokio::test]
async fn a_swept_report_leaves_a_visible_cursor_jump() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let origin = tmp.path().join("repo");
std::fs::create_dir_all(&origin).unwrap();
let runtime = WaveRuntime::open("ship".into(), origin).expect("runtime");
store
.set_bus_cursor("ship".into(), 0)
.await
.expect("cursor");
for text in ["one", "two"] {
store
.publish_bus("ship.a".into(), "ship.a".into(), text.into())
.await
.expect("publish");
}
let cutoff = time::OffsetDateTime::now_utc().unix_timestamp() + 1;
assert_eq!(store.sweep_bus(cutoff).await.expect("sweep"), 2);
store
.publish_bus("ship.a".into(), "ship.a".into(), "three".into())
.await
.expect("publish");
let listener = Arc::new(BusListener::new(runtime.clone(), store.clone()));
listener.attach().await.expect("attach");
listener.poll_once().await.expect("poll");
let thread = runtime.thread_snapshot();
assert_eq!(thread.len(), 2, "the jump note, then the surviving report");
assert_eq!(thread[0].from.as_deref(), Some("bus"));
assert!(
thread[0].text.contains("cursor jumped 0 → 3")
&& thread[0].text.contains("2 broadcast(s)"),
"the miss names what it lost: {}",
thread[0].text
);
assert_eq!(thread[1].text, "three");
}
#[tokio::test]
async fn a_lone_report_expires_on_a_quiet_bus() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = temp_store(tmp.path()).await;
let origin = tmp.path().join("repo");
std::fs::create_dir_all(&origin).unwrap();
let runtime = WaveRuntime::open("ship".into(), origin.clone()).expect("runtime");
store
.set_bus_cursor("ship".into(), 0)
.await
.expect("cursor");
store
.publish_bus("ship.a".into(), "ship.a".into(), "stale report".into())
.await
.expect("publish");
age_whole_bus(&tmp.path().join("loopflow.db"), BUS_WINDOW_SECS + 60);
let listener = BusListener::new(runtime.clone(), store.clone());
listener.attach().await.expect("attach");
listener.poll_once().await.expect("poll");
let thread = runtime.thread_snapshot();
assert!(
!thread.iter().any(|turn| turn.text == "stale report"),
"an expired report is never delivered"
);
assert_eq!(thread.len(), 1, "just the jump note");
assert_eq!(thread[0].from.as_deref(), Some("bus"));
assert!(
thread[0].text.contains("cursor jumped 0 → 2")
&& thread[0].text.contains("1 broadcast(s)"),
"the miss is visible on an emptied bus: {}",
thread[0].text
);
assert_eq!(store.bus_cursor("ship".into()).await.unwrap(), Some(1));
drop(runtime);
let runtime = WaveRuntime::open("ship".into(), origin).expect("runtime");
let listener = BusListener::new(runtime.clone(), store);
listener.attach().await.expect("reattach");
listener.poll_once().await.expect("poll");
assert_eq!(runtime.thread_snapshot().len(), 1, "gap announced once");
}
fn age_whole_bus(db: &std::path::Path, seconds: i64) {
let conn = rusqlite::Connection::open(db).expect("open store file");
conn.execute("UPDATE bus_messages SET at = at - ?1", [seconds])
.expect("age the bus");
}
}