use std::sync::Arc;
use std::time::Duration;
use chrono::Utc;
use tokio::io::AsyncWriteExt;
use tokio::net::UnixStream;
use trusty_common::control_bus::{EventId, HarnessEvent, HarnessPayload, HarnessSource};
use trusty_common::uds::{SocketVerdict, connect_hardened, probe_socket_verdict};
use super::bus::{EventBus, EventBusConfig, IngestOutcome};
use super::ingest::{bind_ingest, serve_ingest};
fn make_event(id: Option<EventId>) -> HarnessEvent {
HarnessEvent {
source: HarnessSource::Mpm,
session: None,
seq: 0,
at: Utc::now(),
payload: HarnessPayload::Ping,
id: id.unwrap_or_default(),
parent_id: None,
}
}
#[test]
fn zero_capacity_is_raised_to_one() {
let bus = EventBus::new(EventBusConfig { capacity: 0 });
assert_eq!(bus.ingest(make_event(None)), IngestOutcome::Ingested);
assert_eq!(bus.len(), 1);
}
#[test]
fn duplicate_id_is_deduped() {
let bus = EventBus::new(EventBusConfig { capacity: 8 });
let id = EventId::new();
assert_eq!(bus.ingest(make_event(Some(id))), IngestOutcome::Ingested);
assert_eq!(bus.ingest(make_event(Some(id))), IngestOutcome::Deduped);
assert_eq!(bus.len(), 1);
let metrics = bus.metrics();
assert_eq!(metrics.ingested, 1);
assert_eq!(metrics.deduped, 1);
assert_eq!(metrics.evicted, 0);
}
#[test]
fn eviction_at_capacity_drops_the_oldest() {
let bus = EventBus::new(EventBusConfig { capacity: 2 });
let first = make_event(None);
let first_id = first.id;
let second = make_event(None);
let third = make_event(None);
assert_eq!(bus.ingest(first), IngestOutcome::Ingested);
assert_eq!(bus.ingest(second), IngestOutcome::Ingested);
assert!(bus.contains(first_id));
assert_eq!(bus.ingest(third), IngestOutcome::Ingested);
assert_eq!(bus.len(), 2, "ring never exceeds its configured capacity");
assert!(
!bus.contains(first_id),
"the oldest event is evicted to make room for the newest"
);
let metrics = bus.metrics();
assert_eq!(metrics.ingested, 3);
assert_eq!(metrics.evicted, 1);
}
#[tokio::test]
async fn a_subscriber_receives_ingested_events() {
let bus = EventBus::new(EventBusConfig { capacity: 8 });
let mut rx = bus.subscribe();
let event = make_event(None);
let sent_id = event.id;
bus.ingest(event);
let received = tokio::time::timeout(Duration::from_secs(1), rx.recv())
.await
.expect("recv did not time out")
.expect("channel still open");
assert_eq!(received.event.id, sent_id);
}
#[tokio::test]
async fn ingest_assigns_sequential_console_seqs_starting_at_one() {
let bus = EventBus::new(EventBusConfig { capacity: 8 });
let mut rx = bus.subscribe();
bus.ingest(make_event(None));
bus.ingest(make_event(None));
let first = tokio::time::timeout(Duration::from_secs(1), rx.recv())
.await
.expect("recv 1")
.expect("open");
let second = tokio::time::timeout(Duration::from_secs(1), rx.recv())
.await
.expect("recv 2")
.expect("open");
assert_eq!(
first.event.seq, 1,
"a fresh bus with no recovered log starts at seq 1"
);
assert_eq!(
second.event.seq, 2,
"seq is console-assigned and monotonic, overwriting the producer's own \
(always 0 from `make_event`)"
);
}
#[tokio::test]
async fn live_fanout_frames_are_never_marked_persisted() {
let bus = EventBus::new(EventBusConfig { capacity: 8 });
let mut rx = bus.subscribe();
bus.ingest(make_event(None));
let received = tokio::time::timeout(Duration::from_secs(1), rx.recv())
.await
.expect("recv did not time out")
.expect("channel still open");
assert!(
!received.persisted,
"a bus with no durable log configured must never claim persistence"
);
}
#[tokio::test]
async fn concurrent_producers_write_the_durable_log_in_seq_order() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let (log, recovered) = super::log::DurableLog::open(super::log::LogConfig {
dir: tmp.path().to_path_buf(),
retain_days: 7,
})
.await
.expect("open a fresh durable log");
let bus = Arc::new(EventBus::with_log(
EventBusConfig { capacity: 256 },
Some(log),
recovered.next_seq,
));
const PRODUCERS: usize = 8;
let handles: Vec<_> = (0..PRODUCERS)
.map(|_| {
let bus = Arc::clone(&bus);
std::thread::spawn(move || {
bus.ingest(make_event(None));
})
})
.collect();
for h in handles {
h.join().expect("producer thread panicked");
}
let deadline = tokio::time::Instant::now() + Duration::from_secs(2);
let seqs = loop {
let items = bus
.replay_since(0)
.await
.expect("bus has a durable log")
.expect("replay succeeds");
let seqs: Vec<u64> = items
.into_iter()
.map(|i| match i {
super::log::ReplayItem::Event(e) => e.seq,
super::log::ReplayItem::Gap { .. } => {
panic!("no backpressure expected at this capacity")
}
})
.collect();
if seqs.len() == PRODUCERS {
break seqs;
}
assert!(
tokio::time::Instant::now() < deadline,
"writer did not durably write all {PRODUCERS} events within budget"
);
tokio::time::sleep(Duration::from_millis(10)).await;
};
let mut sorted = seqs.clone();
sorted.sort_unstable();
assert_eq!(
seqs, sorted,
"on-disk seq order must match assignment order under concurrent producers"
);
assert_eq!(
seqs,
(1..=PRODUCERS as u64).collect::<Vec<_>>(),
"no gaps and no duplicates across the concurrent producers"
);
}
async fn spawn_test_listener(tmp: &std::path::Path) -> (std::path::PathBuf, Arc<EventBus>) {
let socket = tmp.join("sockets").join("trusty-console.sock");
let listener = bind_ingest(&socket).await.expect("bind ingest socket");
let bus = Arc::new(EventBus::new(EventBusConfig::default()));
let served_bus = Arc::clone(&bus);
tokio::spawn(async move {
serve_ingest(listener, served_bus, std::future::pending()).await;
});
(socket, bus)
}
async fn dial(socket: &std::path::Path) -> UnixStream {
connect_hardened(socket)
.await
.expect("connect to ingest socket")
}
#[tokio::test]
async fn ingest_over_uds_socket_reaches_the_bus() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let (socket, bus) = spawn_test_listener(tmp.path()).await;
let event = make_event(None);
let sent_id = event.id;
let mut line = serde_json::to_vec(&event).expect("serialize event");
line.push(b'\n');
let mut stream = dial(&socket).await;
stream.write_all(&line).await.expect("write frame");
stream.shutdown().await.expect("half-close");
wait_until(Duration::from_secs(1), || bus.metrics().ingested >= 1).await;
assert!(bus.contains(sent_id));
assert_eq!(bus.metrics().deduped, 0);
}
#[tokio::test]
async fn malformed_line_does_not_kill_the_listener() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let (socket, bus) = spawn_test_listener(tmp.path()).await;
let mut stream = dial(&socket).await;
stream
.write_all(b"not valid json at all\n")
.await
.expect("write malformed line");
let good = make_event(None);
let good_id = good.id;
let mut line = serde_json::to_vec(&good).expect("serialize event");
line.push(b'\n');
stream.write_all(&line).await.expect("write valid frame");
stream.shutdown().await.expect("half-close");
wait_until(Duration::from_secs(1), || bus.metrics().ingested >= 1).await;
assert!(
bus.contains(good_id),
"the valid frame after a malformed one is still ingested"
);
let other = make_event(None);
let other_id = other.id;
let mut other_line = serde_json::to_vec(&other).expect("serialize event");
other_line.push(b'\n');
let mut second_stream = dial(&socket).await;
second_stream
.write_all(&other_line)
.await
.expect("write second connection's frame");
second_stream.shutdown().await.expect("half-close");
wait_until(Duration::from_secs(1), || bus.contains(other_id)).await;
}
#[tokio::test]
async fn oversized_line_ends_only_its_own_connection() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let (socket, bus) = spawn_test_listener(tmp.path()).await;
let oversized = vec![b'a'; super::ingest::MAX_LINE_BYTES + 1];
let mut stream = dial(&socket).await;
let _ = stream.write_all(&oversized).await;
let _ = stream.write_all(b"\n").await;
let _ = stream.shutdown().await;
let event = make_event(None);
let event_id = event.id;
let mut line = serde_json::to_vec(&event).expect("serialize event");
line.push(b'\n');
let mut second_stream = dial(&socket).await;
second_stream.write_all(&line).await.expect("write frame");
second_stream.shutdown().await.expect("half-close");
wait_until(Duration::from_secs(1), || bus.contains(event_id)).await;
assert_eq!(
bus.metrics().ingested,
1,
"only the second connection's well-formed frame was ever ingested"
);
}
#[tokio::test]
async fn unterminated_line_never_grows_past_the_line_cap() {
let (mut writer, reader) = UnixStream::pair().expect("create a connected pair");
let total = super::ingest::MAX_LINE_BYTES * 2;
let write_task = tokio::spawn(async move {
let chunk = vec![b'a'; 64 * 1024];
let mut sent = 0usize;
while sent < total {
if writer.write_all(&chunk).await.is_err() {
break;
}
sent += chunk.len();
}
});
let mut buffered = tokio::io::BufReader::new(reader);
let mut line = Vec::new();
let read = tokio::time::timeout(
Duration::from_secs(5),
super::ingest::read_capped_line(&mut buffered, &mut line),
)
.await
.expect("read_capped_line must return once its budget is exhausted, not hang")
.expect("a budget-exhausted read is not an I/O error");
assert!(
line.len() <= super::ingest::MAX_LINE_BYTES,
"the line buffer must never grow past the per-line cap, got {} bytes",
line.len()
);
assert_eq!(read, line.len());
assert!(
!line.ends_with(b"\n"),
"no newline was ever sent, so the read must have stopped on the budget alone"
);
write_task.abort();
}
#[tokio::test]
async fn stale_socket_file_is_reclaimed_on_bind() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = tmp.path().join("sockets").join("trusty-console.sock");
let dead = bind_ingest(&socket).await.expect("first bind");
drop(dead); assert!(socket.exists(), "the corpse must still be on disk");
let deadline = tokio::time::Instant::now() + Duration::from_secs(5);
while probe_socket_verdict(&socket, Duration::from_millis(50)).await
!= SocketVerdict::NotServing
{
assert!(
tokio::time::Instant::now() < deadline,
"the dropped listener never stopped answering connects"
);
tokio::time::sleep(Duration::from_millis(10)).await;
}
let listener = bind_ingest(&socket)
.await
.expect("a stale socket file must be reclaimed, not refused forever");
drop(listener);
}
#[tokio::test]
async fn idle_connection_is_dropped_after_the_read_timeout() {
let (writer, reader) = UnixStream::pair().expect("create a connected pair");
let bus = Arc::new(EventBus::new(EventBusConfig::default()));
let handle = tokio::spawn(super::ingest::handle_connection_with_timeout(
reader,
Arc::clone(&bus),
Duration::from_millis(50),
));
tokio::time::timeout(Duration::from_secs(5), handle)
.await
.expect("handle_connection_with_timeout must return once idle")
.expect("the connection task must not panic");
drop(writer);
assert_eq!(
bus.metrics().ingested,
0,
"an idle connection that never sent anything ingests nothing"
);
}
#[tokio::test]
async fn connections_beyond_the_limit_wait_for_a_free_slot() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = tmp.path().join("sockets").join("trusty-console.sock");
let listener = bind_ingest(&socket).await.expect("bind ingest socket");
let bus = Arc::new(EventBus::new(EventBusConfig::default()));
let served_bus = Arc::clone(&bus);
tokio::spawn(async move {
super::ingest::serve_ingest_with_limit(listener, served_bus, 1, std::future::pending())
.await;
});
let occupant = dial(&socket).await;
let event = make_event(None);
let event_id = event.id;
let mut line = serde_json::to_vec(&event).expect("serialize event");
line.push(b'\n');
let mut second = dial(&socket).await;
second.write_all(&line).await.expect("write frame");
second.shutdown().await.expect("half-close");
assert_stays_false(Duration::from_millis(300), || bus.contains(event_id)).await;
drop(occupant);
wait_until(Duration::from_secs(2), || bus.contains(event_id)).await;
}
#[tokio::test]
async fn shutdown_is_observed_while_the_connection_pool_is_saturated() {
let tmp = tempfile::TempDir::new().expect("tempdir");
let socket = tmp.path().join("sockets").join("trusty-console.sock");
let listener = bind_ingest(&socket).await.expect("bind ingest socket");
let bus = Arc::new(EventBus::new(EventBusConfig::default()));
let served_bus = Arc::clone(&bus);
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
let serve_handle = tokio::spawn(async move {
super::ingest::serve_ingest_with_limit(listener, served_bus, 1, async move {
let _ = shutdown_rx.await;
})
.await;
});
let occupant = dial(&socket).await;
let event = make_event(None);
let event_id = event.id;
let mut line = serde_json::to_vec(&event).expect("serialize event");
line.push(b'\n');
let mut second = dial(&socket).await;
second.write_all(&line).await.expect("write frame");
second.shutdown().await.expect("half-close");
assert_stays_false(Duration::from_millis(300), || bus.contains(event_id)).await;
shutdown_tx
.send(())
.expect("serve loop is still awaiting shutdown");
tokio::time::timeout(Duration::from_secs(2), serve_handle)
.await
.expect(
"serve_ingest_with_limit must return once shutdown resolves, even with \
every permit held and a connection stuck on acquire",
)
.expect("the serve task must not panic");
drop(occupant);
}
async fn wait_until(budget: Duration, mut predicate: impl FnMut() -> bool) {
let deadline = tokio::time::Instant::now() + budget;
loop {
if predicate() {
return;
}
assert!(
tokio::time::Instant::now() < deadline,
"condition did not become true within {budget:?}"
);
tokio::time::sleep(Duration::from_millis(10)).await;
}
}
async fn assert_stays_false(budget: Duration, mut predicate: impl FnMut() -> bool) {
let deadline = tokio::time::Instant::now() + budget;
while tokio::time::Instant::now() < deadline {
assert!(!predicate(), "condition became true within {budget:?}");
tokio::time::sleep(Duration::from_millis(10)).await;
}
}