use super::*;
use running_process::broker::protocol::read_frame;
use running_process::broker::protocol_v2::ServiceDefinitionBuilder;
use std::sync::atomic::AtomicBool;
use tempfile::tempdir;
fn make_hello(service: &str, wanted: &str) -> Hello {
Hello {
client_min_protocol: ENVELOPE_VERSION as u32,
client_max_protocol: ENVELOPE_VERSION as u32,
service_name: service.to_string(),
wanted_version: wanted.to_string(),
client_version: "test".to_string(),
client_capabilities: 0,
auth_token: Vec::new(),
request_id: "test".to_string(),
connection_id: 42,
peer_pid: 1234,
client_lib_name: "test".to_string(),
client_lib_version: "test".to_string(),
peer_attestation_nonce: Vec::new(),
capability_token: Vec::new(),
client_keepalive_secs: 0,
}
}
#[test]
fn accept_poll_exits_without_accepting_after_shutdown_request() {
static ALREADY_ASKED: AtomicBool = AtomicBool::new(true);
let shutdown = ShutdownRequest::watching(&ALREADY_ASKED);
let result = poll_accept_until_shutdown(&shutdown, || -> std::io::Result<()> {
panic!("accept must not run after shutdown")
})
.unwrap();
assert!(result.is_none());
}
#[test]
fn accept_poll_observes_shutdown_while_listener_would_block() {
static ASKED_MID_POLL: AtomicBool = AtomicBool::new(false);
ASKED_MID_POLL.store(false, Ordering::Relaxed);
let shutdown = ShutdownRequest::watching(&ASKED_MID_POLL);
let signaler = thread::spawn(move || {
thread::sleep(ACCEPT_POLL_INTERVAL);
ASKED_MID_POLL.store(true, Ordering::Relaxed);
});
let start = Instant::now();
let result = poll_accept_until_shutdown(&shutdown, || -> std::io::Result<()> {
Err(std::io::ErrorKind::WouldBlock.into())
})
.unwrap();
signaler.join().unwrap();
assert!(result.is_none());
assert!(start.elapsed() < Duration::from_secs(1));
}
#[test]
fn parse_cli_defaults() {
let args = vec!["bin".to_owned()];
let opts = parse_cli(&args).unwrap();
assert!(!opts.no_bind);
assert!(!opts.once);
assert_eq!(opts.program, DEFAULT_PROGRAM);
}
#[test]
fn parse_cli_program_arg() {
let args = vec![
"bin".to_owned(),
"--program".to_owned(),
"zccache".to_owned(),
];
let opts = parse_cli(&args).unwrap();
assert_eq!(opts.program, "zccache");
}
#[test]
fn parse_cli_once_flag() {
let args = vec!["bin".to_owned(), "--once".to_owned()];
let opts = parse_cli(&args).unwrap();
assert!(opts.once);
}
#[test]
fn parse_cli_program_missing_value_errs() {
let args = vec!["bin".to_owned(), "--program".to_owned()];
assert!(parse_cli(&args).is_err());
}
#[test]
fn parse_cli_unknown_arg_errs() {
let args = vec!["bin".to_owned(), "--bogus".to_owned()];
assert!(parse_cli(&args).is_err());
}
#[test]
fn the_http_surface_is_off_unless_asked_for() {
let opts = parse_cli(&["bin".to_owned()]).unwrap();
assert!(opts.http_port.is_none());
}
#[test]
fn http_port_accepts_dynamic_and_a_number() {
assert_eq!(parse_http_port("dynamic").unwrap(), BrokerHttpPort::Dynamic);
assert_eq!(parse_http_port("DYNAMIC").unwrap(), BrokerHttpPort::Dynamic);
assert_eq!(
parse_http_port("8080").unwrap(),
BrokerHttpPort::StaticOrFallback { preferred: 8080 }
);
}
#[test]
fn http_port_zero_is_dynamic() {
assert_eq!(parse_http_port("0").unwrap(), BrokerHttpPort::Dynamic);
}
#[test]
fn http_port_rejects_a_non_port() {
for bad in ["", "http", "-1", "65536", "80x"] {
assert!(
parse_http_port(bad).is_err(),
"{bad:?} should not parse as a port"
);
}
}
#[test]
fn http_port_requires_a_value() {
let args = vec!["bin".to_owned(), "--http-port".to_owned()];
assert!(parse_cli(&args).is_err());
}
#[test]
fn parse_cli_threads_the_http_port_through() {
let args = vec![
"bin".to_owned(),
"--http-port".to_owned(),
"dynamic".to_owned(),
];
let opts = parse_cli(&args).unwrap();
assert_eq!(opts.http_port, Some(BrokerHttpPort::Dynamic));
}
#[test]
fn the_page_lists_a_backend_once_it_is_tracked() {
use std::io::{Read as _, Write as _};
let dir = tempdir().unwrap();
let started = start_http_surface(BrokerHttpPort::Dynamic, "track-program", dir.path()).unwrap();
started.registry.track("zccache".to_string());
let (ip, port) = broker_http_discovery::read_http_port(dir.path(), "track-program")
.unwrap()
.expect("endpoint published");
let mut stream = std::net::TcpStream::connect(std::net::SocketAddr::new(ip, port)).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(10)))
.unwrap();
stream
.write_all(b"GET / HTTP/1.0\r\nHost: localhost\r\n\r\n")
.unwrap();
let mut body = Vec::new();
stream.read_to_end(&mut body).unwrap();
let page = String::from_utf8_lossy(&body);
assert!(
page.contains("zccache"),
"backend missing from page: {page}"
);
assert!(
!page.contains("no backends registered yet"),
"page still claims nothing is registered: {page}"
);
}
#[test]
fn starting_the_surface_publishes_an_endpoint_that_answers() {
use std::io::{Read as _, Write as _};
let dir = tempdir().unwrap();
let started = start_http_surface(BrokerHttpPort::Dynamic, "test-program", dir.path()).unwrap();
assert_eq!(started.program, "test-program");
let (ip, port) = broker_http_discovery::read_http_port(dir.path(), "test-program")
.unwrap()
.expect("the surface publishes its endpoint");
assert_ne!(port, 0, "a published port of 0 is not reachable");
let mut stream = std::net::TcpStream::connect(std::net::SocketAddr::new(ip, port))
.expect("the published endpoint accepts connections");
stream
.set_read_timeout(Some(Duration::from_secs(10)))
.unwrap();
stream
.write_all(b"GET / HTTP/1.0\r\nHost: localhost\r\n\r\n")
.unwrap();
let mut response = Vec::new();
stream.read_to_end(&mut response).unwrap();
let text = String::from_utf8_lossy(&response);
assert!(
text.starts_with("HTTP/"),
"expected an HTTP response, got {text:?}"
);
}
#[test]
fn unpublishing_leaves_no_endpoint_behind() {
let dir = tempdir().unwrap();
broker_http_discovery::publish_http_port(
dir.path(),
"test-program",
std::net::IpAddr::from([127, 0, 0, 1]),
1234,
)
.unwrap();
broker_http_discovery::unpublish_http_port(dir.path(), "test-program").unwrap();
assert_eq!(
broker_http_discovery::read_http_port(dir.path(), "test-program").unwrap(),
None
);
}
#[test]
fn a_published_daemon_identity_becomes_the_backend_pipe() {
use running_process::broker::backend_lifecycle::identity::DaemonProcess;
use running_process::broker::backend_sdk::{
remove_daemon_identity_file, write_daemon_identity_file,
};
use running_process::broker::lifecycle::names_v2::daemon_identity_path;
use running_process::broker::protocol::Endpoint;
use running_process::broker::secure_dir::ensure_private_dir;
let service = format!("resolve-test-{}", std::process::id());
assert_eq!(
resolve_backend_pipe(&service),
"",
"nothing published yet, so there is no pipe to report"
);
let path = daemon_identity_path(&service);
ensure_private_dir(path.parent().expect("parent")).expect("private dir");
let endpoint = Endpoint {
namespace_id: String::new(),
path: "daemon-endpoint-under-test".to_string(),
};
let daemon = DaemonProcess::current_process(endpoint, None).expect("identity");
write_daemon_identity_file(&path, &daemon).expect("publish");
assert_eq!(resolve_backend_pipe(&service), "daemon-endpoint-under-test");
remove_daemon_identity_file(&path);
assert_eq!(
resolve_backend_pipe(&service),
"",
"a retracted daemon must stop being advertised"
);
}
#[test]
fn build_hello_reply_refuses_unknown_service() {
let dir = tempdir().unwrap();
let loader = ServiceDefinitionLoader::new(dir.path());
let hello = make_hello("nosuch", "1.0.0");
let reply = build_hello_reply(&hello, &loader, "");
match reply.result {
Some(hello_reply::Result::Refused(r)) => {
assert_eq!(r.code, ErrorCode::ErrorServiceUnknown as i32);
}
other => panic!("expected Refused, got {other:?}"),
}
}
#[test]
fn build_hello_reply_negotiates_registered_service() {
let dir = tempdir().unwrap();
ServiceDefinitionBuilder::shared_broker("zccache", "/usr/bin/zccache-daemon")
.install_in(dir.path())
.unwrap();
let loader = ServiceDefinitionLoader::new(dir.path());
let hello = make_hello("zccache", "1.0.0");
let reply = build_hello_reply(&hello, &loader, "");
match reply.result {
Some(hello_reply::Result::Negotiated(n)) => {
assert_eq!(n.connection_id, 42);
assert!(n.backend_pipe.is_empty());
}
other => panic!("expected Negotiated, got {other:?}"),
}
}
#[test]
fn hello_handler_round_trips_the_client_v2_frame_contract() {
let dir = tempdir().unwrap();
ServiceDefinitionBuilder::shared_broker("zccache", "/usr/bin/zccache-daemon")
.install_in(dir.path())
.unwrap();
let loader = ServiceDefinitionLoader::new(dir.path());
let hello = make_hello("zccache", "1.0.0");
let mut request =
Frame::request(CONTROL_PAYLOAD_PROTOCOL, hello.encode_to_vec()).with_request_id(0x930);
request.traceparent = "00-test-trace".to_string();
request.tracestate = "vendor=test".to_string();
let mut wire = Vec::new();
let service = handle_hello_bytes(&mut wire, &loader, request.encode_to_vec())
.expect("Frame-wrapped Hello negotiates");
assert_eq!(service, "zccache");
let response_bytes = read_frame(&mut std::io::Cursor::new(wire)).expect("framed response");
let response = Frame::decode(response_bytes.as_slice()).expect("response Frame");
validate_frame_envelope(&response, FrameKind::Response, CONTROL_PAYLOAD_PROTOCOL)
.expect("response envelope");
assert_eq!(response.request_id, 0x930);
assert_eq!(response.traceparent, request.traceparent);
assert_eq!(response.tracestate, request.tracestate);
let reply = HelloReply::decode(response.payload.as_slice()).expect("HelloReply payload");
assert!(matches!(
reply.result,
Some(hello_reply::Result::Negotiated(_))
));
}
#[test]
fn hello_handler_rejects_the_obsolete_raw_hello_contract() {
let dir = tempdir().unwrap();
let loader = ServiceDefinitionLoader::new(dir.path());
let hello = make_hello("zccache", "1.0.0");
let mut wire = Vec::new();
let error = handle_hello_bytes(&mut wire, &loader, hello.encode_to_vec())
.expect_err("raw Hello must not masquerade as a Frame");
assert!(
error.contains("Hello Frame"),
"error should identify the envelope boundary: {error}"
);
let response_bytes = read_frame(&mut std::io::Cursor::new(wire)).expect("framed refusal");
let response = Frame::decode(response_bytes.as_slice()).expect("response Frame");
validate_frame_envelope(&response, FrameKind::Response, CONTROL_PAYLOAD_PROTOCOL)
.expect("response envelope");
let reply = HelloReply::decode(response.payload.as_slice()).expect("HelloReply payload");
assert!(matches!(
reply.result,
Some(hello_reply::Result::Refused(_))
));
}
#[test]
fn hello_handler_returns_a_correlated_refusal_for_an_invalid_envelope() {
let dir = tempdir().unwrap();
let loader = ServiceDefinitionLoader::new(dir.path());
let hello = make_hello("zccache", "1.0.0");
let mut request =
Frame::request(CONTROL_PAYLOAD_PROTOCOL, hello.encode_to_vec()).with_request_id(0x931);
request.kind = FrameKind::Event as i32;
request.traceparent = "00-invalid-envelope".to_string();
request.tracestate = "vendor=invalid".to_string();
let mut wire = Vec::new();
let error = handle_hello_bytes(&mut wire, &loader, request.encode_to_vec())
.expect_err("invalid envelope must be refused");
assert!(error.contains("validate Hello Frame"), "got: {error}");
let response_bytes = read_frame(&mut std::io::Cursor::new(wire)).expect("framed refusal");
let response = Frame::decode(response_bytes.as_slice()).expect("response Frame");
validate_frame_envelope(&response, FrameKind::Response, CONTROL_PAYLOAD_PROTOCOL)
.expect("response envelope");
assert_eq!(response.request_id, request.request_id);
assert_eq!(response.traceparent, request.traceparent);
assert_eq!(response.tracestate, request.tracestate);
let reply = HelloReply::decode(response.payload.as_slice()).expect("HelloReply payload");
match reply.result {
Some(hello_reply::Result::Refused(refused)) => {
assert_eq!(refused.code, ErrorCode::ErrorPeerRejected as i32);
assert_eq!(refused.reason, "Hello frame kind must be REQUEST");
}
other => panic!("expected Refused, got {other:?}"),
}
}
#[test]
fn build_hello_reply_blocks_below_min_version() {
let dir = tempdir().unwrap();
ServiceDefinitionBuilder::shared_broker("zccache", "/usr/bin/zccache-daemon")
.min_version("2.0.0")
.install_in(dir.path())
.unwrap();
let loader = ServiceDefinitionLoader::new(dir.path());
let hello = make_hello("zccache", "1.0.0");
let reply = build_hello_reply(&hello, &loader, "");
match reply.result {
Some(hello_reply::Result::Refused(r)) => {
assert_eq!(r.code, ErrorCode::ErrorVersionBlocked as i32);
assert!(r.reason.contains("min_version"), "got: {}", r.reason);
}
other => panic!("expected Refused, got {other:?}"),
}
}
#[test]
fn build_hello_reply_blocks_outside_version_allow_list() {
let dir = tempdir().unwrap();
ServiceDefinitionBuilder::shared_broker("zccache", "/usr/bin/zccache-daemon")
.version_allow_list(["1.0.0", "1.1.0"])
.install_in(dir.path())
.unwrap();
let loader = ServiceDefinitionLoader::new(dir.path());
let hello = make_hello("zccache", "1.2.0");
let reply = build_hello_reply(&hello, &loader, "");
match reply.result {
Some(hello_reply::Result::Refused(r)) => {
assert_eq!(r.code, ErrorCode::ErrorVersionBlocked as i32);
assert!(r.reason.contains("allow_list"), "got: {}", r.reason);
}
other => panic!("expected Refused, got {other:?}"),
}
}
#[test]
fn build_hello_reply_allows_version_in_allow_list() {
let dir = tempdir().unwrap();
ServiceDefinitionBuilder::shared_broker("zccache", "/usr/bin/zccache-daemon")
.version_allow_list(["1.0.0", "1.1.0"])
.install_in(dir.path())
.unwrap();
let loader = ServiceDefinitionLoader::new(dir.path());
let hello = make_hello("zccache", "1.1.0");
let reply = build_hello_reply(&hello, &loader, "");
assert!(matches!(
reply.result,
Some(hello_reply::Result::Negotiated(_))
));
}