Ambient system companions over one privacy-preserving signal daemon (aggregate-only, no keystroke content): a git-driven terminal garden and IOKit hardware collectors. ambient daemon macos privacy terminal

crates/signald/tests/streaming.rs

79 lines · 2526 bytes

 1//! Live-streaming acceptance test: a subscriber gets the current
 2//! value of each metric immediately on connect (last-value cache), then receives
 3//! updates as values change — over the real Unix socket, not just the hub.
 4
 5use std::io::BufReader;
 6use std::os::unix::net::UnixStream;
 7use std::path::{Path, PathBuf};
 8use std::thread;
 9use std::time::{Duration, Instant};
10
11use signal_schema::{wire, Signal, SignalName, Source, Value, SCHEMA_VERSION};
12use signald::hub::Hub;
13use signald::publish;
14
15fn sig(name: SignalName, value: f64) -> Signal {
16    Signal {
17        schema_version: SCHEMA_VERSION,
18        ts: 1,
19        source: Source::Terminal,
20        name,
21        value: Value(value),
22        tag: None,
23    }
24}
25
26#[test]
27fn subscriber_sees_cached_value_then_a_live_update() {
28    let socket = unique_socket();
29    let hub = Hub::new();
30
31    // A value published before anyone connects must still reach a subscriber
32    // (the last-value cache).
33    hub.publish(sig(SignalName::KeysPerMin, 12.0));
34
35    let serve_hub = hub.clone();
36    let serve_socket = socket.clone();
37    thread::spawn(move || {
38        let _ = publish::serve(&serve_socket, serve_hub);
39    });
40
41    let stream = connect(&socket);
42    let mut reader = BufReader::new(stream);
43
44    // 1. On connect: the cached snapshot arrives immediately.
45    let wire::Frame::Signal(first) = wire::read_frame(&mut reader).unwrap() else {
46        panic!("expected a snapshot frame");
47    };
48    assert_eq!(first.name, SignalName::KeysPerMin);
49    assert_eq!(first.value, Value(12.0));
50
51    // 2. A change after subscribing is streamed live.
52    hub.publish(sig(SignalName::KeysPerMin, 99.0));
53    let wire::Frame::Signal(update) = wire::read_frame(&mut reader).unwrap() else {
54        panic!("expected a live update frame");
55    };
56    assert_eq!(update.name, SignalName::KeysPerMin);
57    assert_eq!(update.value, Value(99.0));
58
59    let _ = std::fs::remove_file(&socket);
60}
61
62fn connect(socket: &Path) -> UnixStream {
63    let deadline = Instant::now() + Duration::from_secs(5);
64    loop {
65        if let Ok(s) = UnixStream::connect(socket) {
66            return s;
67        }
68        assert!(Instant::now() < deadline, "signald socket never came up");
69        thread::sleep(Duration::from_millis(20));
70    }
71}
72
73fn unique_socket() -> PathBuf {
74    let nanos = std::time::SystemTime::now()
75        .duration_since(std::time::UNIX_EPOCH)
76        .unwrap()
77        .as_nanos();
78    std::env::temp_dir().join(format!("signald-stream-{}-{nanos}.sock", std::process::id()))
79}