//! Live-streaming acceptance test (spec §1.3): a subscriber gets the current //! value of each metric immediately on connect (last-value cache), then receives //! updates as values change — over the real Unix socket, not just the hub. use std::io::BufReader; use std::os::unix::net::UnixStream; use std::path::{Path, PathBuf}; use std::thread; use std::time::{Duration, Instant}; use signal_schema::{wire, Signal, SignalName, Source, Value, SCHEMA_VERSION}; use signald::hub::Hub; use signald::publish; fn sig(name: SignalName, value: f64) -> Signal { Signal { schema_version: SCHEMA_VERSION, ts: 1, source: Source::Terminal, name, value: Value(value), tag: None, } } #[test] fn subscriber_sees_cached_value_then_a_live_update() { let socket = unique_socket(); let hub = Hub::new(); // A value published before anyone connects must still reach a subscriber // (the last-value cache). hub.publish(sig(SignalName::KeysPerMin, 12.0)); let serve_hub = hub.clone(); let serve_socket = socket.clone(); thread::spawn(move || { let _ = publish::serve(&serve_socket, serve_hub); }); let stream = connect(&socket); let mut reader = BufReader::new(stream); // 1. On connect: the cached snapshot arrives immediately. let first = wire::read_frame(&mut reader).unwrap().expect("snapshot frame"); assert_eq!(first.name, SignalName::KeysPerMin); assert_eq!(first.value, Value(12.0)); // 2. A change after subscribing is streamed live. hub.publish(sig(SignalName::KeysPerMin, 99.0)); let update = wire::read_frame(&mut reader).unwrap().expect("live update frame"); assert_eq!(update.name, SignalName::KeysPerMin); assert_eq!(update.value, Value(99.0)); let _ = std::fs::remove_file(&socket); } fn connect(socket: &Path) -> UnixStream { let deadline = Instant::now() + Duration::from_secs(5); loop { if let Ok(s) = UnixStream::connect(socket) { return s; } assert!(Instant::now() < deadline, "signald socket never came up"); thread::sleep(Duration::from_millis(20)); } } fn unique_socket() -> PathBuf { let nanos = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .unwrap() .as_nanos(); std::env::temp_dir().join(format!("signald-stream-{}-{nanos}.sock", std::process::id())) }