use crate::daemon::frame;
use crate::daemon::protocol::{
BackendRecoveryStatus, Request, Response, WarmBackendIdentity, WarmBackendStatus,
};
use keyhog_scanner::hw_probe::ScanBackend;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
fn ready_warm_backend() -> WarmBackendStatus {
WarmBackendStatus {
ready: true,
daemon_generation: "frame-test".into(),
identity: WarmBackendIdentity {
engine: "test-engine".into(),
gpu_artifact: Some("test-gpu-artifact".into()),
binary_sha256: "test-binary".into(),
detector_rules_digest: "test-detectors".into(),
config_digest: "test-config".into(),
},
required_backends: vec!["gpu-cuda".into()],
initialized_backends: vec!["gpu-cuda".into()],
reason: None,
repair_command: None,
}
}
async fn wire_bytes_for(request: &Request) -> Vec<u8> {
let (mut w, mut r) = tokio::io::duplex(1 << 20);
frame::write_request(&mut w, request)
.await
.expect("write request into capture duplex");
drop(w);
let mut wire = Vec::new();
r.read_to_end(&mut wire).await.expect("drain framed bytes");
wire
}
#[tokio::test]
async fn clean_eof_between_frames_returns_none_not_an_error() {
let (client, mut server) = tokio::io::duplex(64);
drop(client); let result = frame::read_request(&mut server).await;
assert!(
matches!(result, Ok(None)),
"an empty clean EOF must be Ok(None), got {result:?}"
);
}
#[tokio::test]
async fn scan_path_request_roundtrips() {
let (mut client, mut server) = tokio::io::duplex(64 * 1024);
let sent = Request::ScanPath {
path: "src/main.rs".into(),
working_dir: Some("/tmp/project".into()),
dogfood: true,
profile: false,
};
frame::write_request(&mut client, &sent)
.await
.expect("write ScanPath");
let got = frame::read_request(&mut server)
.await
.expect("read ScanPath")
.expect("a frame");
match got {
Request::ScanPath {
path,
working_dir,
dogfood,
..
} => {
assert_eq!(path, "src/main.rs");
assert_eq!(working_dir.as_deref(), Some("/tmp/project"));
assert!(dogfood);
}
other => panic!("expected ScanPath, got {other:?}"),
}
}
#[tokio::test]
async fn health_and_shutdown_unit_requests_roundtrip() {
for sent in [Request::Health, Request::Shutdown] {
let (mut client, mut server) = tokio::io::duplex(1024);
frame::write_request(&mut client, &sent)
.await
.expect("write unit request");
let got = frame::read_request(&mut server)
.await
.expect("read unit request")
.expect("a frame");
match (&sent, &got) {
(Request::Health, Request::Health) | (Request::Shutdown, Request::Shutdown) => {}
_ => panic!("unit request did not roundtrip: sent {sent:?}, got {got:?}"),
}
}
}
#[tokio::test]
async fn health_response_roundtrips() {
let (mut server, mut client) = tokio::io::duplex(64 * 1024);
frame::write_response(
&mut server,
&Response::Health {
uptime_secs: 42,
scans_served: 7,
active_scans: 3,
detector_count: 900,
backend_recoveries: 1,
last_backend_fault: Some(BackendRecoveryStatus {
failed_backend: ScanBackend::GpuCuda.label().to_string(),
recovery_backend: ScanBackend::CpuFallback.label().to_string(),
recovered_ranges: vec![keyhog::daemon::protocol::RecoveredInputRangeStatus {
chunk_index: 0,
byte_start: 0,
byte_end: 128,
}],
recovered_chunks: 1,
recovered_bytes: 128,
reason: "test recovery".to_string(),
}),
guard_roots_registered: 0,
guard_roots_current: 0,
guard_roots_blocked: 0,
guard_roots_degraded: 0,
guard_active_transactions: 0,
warm_backend: ready_warm_backend(),
},
)
.await
.expect("write Health response");
let got = frame::read_response(&mut client)
.await
.expect("read Health response")
.expect("a frame");
match got {
Response::Health {
uptime_secs,
scans_served,
active_scans,
detector_count,
backend_recoveries,
last_backend_fault,
warm_backend: _,
..
} => {
assert_eq!(uptime_secs, 42);
assert_eq!(scans_served, 7);
assert_eq!(active_scans, 3);
assert_eq!(detector_count, 900);
assert_eq!(backend_recoveries, 1);
assert_eq!(
last_backend_fault,
Some(BackendRecoveryStatus {
failed_backend: ScanBackend::GpuCuda.label().to_string(),
recovery_backend: ScanBackend::CpuFallback.label().to_string(),
recovered_ranges: vec![keyhog::daemon::protocol::RecoveredInputRangeStatus {
chunk_index: 0,
byte_start: 0,
byte_end: 128,
}],
recovered_chunks: 1,
recovered_bytes: 128,
reason: "test recovery".to_string(),
})
);
}
other => panic!("expected Health response, got {other:?}"),
}
}
#[tokio::test]
async fn error_response_roundtrips_carrying_its_message() {
let (mut server, mut client) = tokio::io::duplex(64 * 1024);
frame::write_response(
&mut server,
&Response::Error {
message: "scanner refused: path outside working_dir".into(),
},
)
.await
.expect("write Error response");
let got = frame::read_response(&mut client)
.await
.expect("read Error response")
.expect("a frame");
match got {
Response::Error { message } => {
assert_eq!(message, "scanner refused: path outside working_dir");
}
other => panic!("expected Error response, got {other:?}"),
}
}
#[tokio::test]
async fn a_valid_prefix_with_a_non_json_body_fails_closed_as_a_parse_error() {
let (mut client, mut server) = tokio::io::duplex(64);
let body = b"hello";
client
.write_all(&(body.len() as u32).to_be_bytes())
.await
.expect("write length prefix");
client.write_all(body).await.expect("write non-JSON body");
drop(client);
let err = frame::read_request(&mut server)
.await
.expect_err("non-JSON body must be a parse error, not a frame");
let message = err.to_string();
assert!(
message.contains("parse request") && message.contains("5 bytes"),
"parse error must name the operation and byte count; got {message}"
);
}
#[tokio::test]
async fn a_zero_length_frame_is_a_parse_error_not_a_phantom_success() {
let (mut client, mut server) = tokio::io::duplex(64);
client
.write_all(&0u32.to_be_bytes())
.await
.expect("write zero length prefix");
drop(client);
let err = frame::read_request(&mut server)
.await
.expect_err("a zero-length body must not parse to a Request");
assert!(
err.to_string().contains("parse request"),
"empty body must surface a parse error; got {err}"
);
}
#[tokio::test]
async fn a_body_delivered_one_byte_at_a_time_reassembles_into_the_exact_frame() {
let sent = Request::ScanText {
path: Some("trickle.txt".into()),
text: "a slowly delivered scan body".into(),
dogfood: true,
profile: false,
};
let wire = wire_bytes_for(&sent).await;
assert!(
wire.len() > 4,
"frame must have a length prefix plus a body"
);
let (mut writer, mut reader) = tokio::io::duplex(64 * 1024);
let feed = tokio::spawn(async move {
writer
.write_all(&wire[..4])
.await
.expect("write length prefix");
for byte in &wire[4..] {
tokio::task::yield_now().await;
writer
.write_all(std::slice::from_ref(byte))
.await
.expect("write body byte");
}
});
let got = frame::read_request(&mut reader)
.await
.expect("read trickled request")
.expect("a fully reassembled frame");
feed.await.expect("feeder task");
match got {
Request::ScanText {
path,
text,
dogfood,
..
} => {
assert_eq!(path.as_deref(), Some("trickle.txt"));
assert_eq!(text, "a slowly delivered scan body");
assert!(dogfood);
}
other => panic!("trickled frame did not reassemble to the original: {other:?}"),
}
}