pub use super::super::adopt::AdoptError;
#[cfg(feature = "client-async")]
pub use super::super::adopt::{
AsyncBrokerSession, IntoBackendIoError, OwnedBackendIo, OwnedConnectRequest,
};
pub use super::super::client::{BackendConnectionRoute, BrokerClientError, RefusalKind};
pub fn refusal_kind(error: &super::super::client_v2::BrokerV2Error) -> Option<RefusalKind> {
match error {
super::super::client_v2::BrokerV2Error::Refused { details, .. } => {
Some(RefusalKind::from_code(details.code()))
}
_ => None,
}
}
#[cfg(test)]
mod tests {
use super::*;
#[cfg(feature = "client-async")]
use prost::Message as _;
fn refused_v2(
code: crate::broker::protocol::ErrorCode,
) -> super::super::super::client_v2::BrokerV2Error {
let mut refused = crate::broker::protocol::Refused {
reason: "nope".into(),
..Default::default()
};
refused.set_code(code);
super::super::super::client_v2::BrokerV2Error::Refused {
reason: "nope".into(),
retry_after_ms: 0,
details: Box::new(refused),
}
}
#[test]
fn a_v2_refusal_classifies_the_same_as_a_v1_one() {
use crate::broker::protocol::ErrorCode;
for code in [
ErrorCode::ErrorVersionUnsupported,
ErrorCode::ErrorVersionBlocked,
ErrorCode::ErrorServiceUnknown,
ErrorCode::ErrorRateLimited,
ErrorCode::ErrorShuttingDown,
] {
assert_eq!(
refusal_kind(&refused_v2(code)),
Some(RefusalKind::from_code(code)),
"v2 refusal for {code:?} classified differently from v1"
);
}
}
#[test]
fn an_unknown_code_stays_unknown_rather_than_becoming_a_named_refusal() {
use crate::broker::protocol::ErrorCode;
let kind = refusal_kind(&refused_v2(ErrorCode::Unspecified));
assert_eq!(kind, Some(RefusalKind::from_code(ErrorCode::Unspecified)));
assert!(matches!(kind, Some(RefusalKind::Other(_))));
}
#[test]
fn a_transport_failure_is_not_a_refusal() {
let io = super::super::super::client_v2::BrokerV2Error::Io(std::io::Error::other("boom"));
assert_eq!(refusal_kind(&io), None);
}
#[test]
fn v1_client_adopt_types_are_aliased_under_v2_namespace() {
use std::any::TypeId;
assert_eq!(
TypeId::of::<super::super::super::adopt::AdoptError>(),
TypeId::of::<AdoptError>(),
"AdoptError aliased"
);
#[cfg(feature = "client-async")]
{
assert_eq!(
TypeId::of::<super::super::super::adopt::OwnedConnectRequest>(),
TypeId::of::<OwnedConnectRequest>(),
"OwnedConnectRequest aliased"
);
}
assert_eq!(
TypeId::of::<super::super::super::client::BackendConnectionRoute>(),
TypeId::of::<BackendConnectionRoute>(),
"BackendConnectionRoute aliased"
);
assert_eq!(
TypeId::of::<super::super::super::client::BrokerClientError>(),
TypeId::of::<BrokerClientError>(),
"BrokerClientError aliased"
);
assert_eq!(
TypeId::of::<super::super::super::client::RefusalKind>(),
TypeId::of::<RefusalKind>(),
"RefusalKind aliased"
);
}
#[cfg(feature = "client-async")]
#[test]
fn async_session_keeps_canonical_type_identity() {
use std::any::TypeId;
assert_eq!(
TypeId::of::<super::super::super::adopt::AsyncBrokerSession>(),
TypeId::of::<AsyncBrokerSession>(),
"the v2 wire swap must not change public type identity"
);
}
#[cfg(feature = "client-async")]
fn test_endpoint(label: &str) -> String {
let nonce = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.expect("clock after epoch")
.as_nanos();
crate::broker::server::singleton_bind::resolve_path_scoped_socket_path(&format!(
"rp-v2-compat-{label}-{}-{nonce}",
std::process::id()
))
.expect("resolve test endpoint")
}
#[cfg(feature = "client-async")]
fn bind_test_listener(endpoint: &str) -> crate::platform::ipc::Listener {
crate::broker::server::singleton_bind::bind_singleton(endpoint).expect("bind test listener")
}
#[cfg(feature = "client-async")]
#[tokio::test]
async fn compat_adopt_speaks_v2_and_reaches_the_negotiated_backend() {
use crate::broker::protocol::{
hello_reply, read_frame, write_frame, Frame, FrameKind, Hello, HelloReply, Negotiated,
PayloadEncoding, CONTROL_PAYLOAD_PROTOCOL, PROTOCOL_VERSION,
};
let broker_endpoint = test_endpoint("broker");
let backend_endpoint = test_endpoint("backend");
let broker_listener = bind_test_listener(&broker_endpoint);
let backend_listener = bind_test_listener(&backend_endpoint);
let (hello_tx, hello_rx) = std::sync::mpsc::channel();
let backend = std::thread::spawn(move || {
let mut stream = backend_listener.accept().expect("accept backend client");
let bytes = read_frame(&mut stream).expect("read backend request");
let request = Frame::decode(bytes.as_slice()).expect("decode backend request");
let response = Frame {
envelope_version: PROTOCOL_VERSION,
kind: FrameKind::Response as i32,
payload_protocol: request.payload_protocol,
payload: b"pong".to_vec(),
request_id: request.request_id,
payload_encoding: PayloadEncoding::None as i32,
deadline_unix_ms: 0,
traceparent: String::new(),
tracestate: String::new(),
};
write_frame(&mut stream, &response.encode_to_vec()).expect("write backend response");
});
let backend_for_broker = backend_endpoint.clone();
let broker = std::thread::spawn(move || {
let mut stream = broker_listener.accept().expect("accept broker client");
let bytes = read_frame(&mut stream).expect("read Hello frame");
let request_frame = Frame::decode(bytes.as_slice()).expect("decode request Frame");
let hello = Hello::decode(request_frame.payload.as_slice()).expect("decode Hello");
hello_tx
.send(hello)
.expect("report observed Hello contract");
let reply = HelloReply {
result: Some(hello_reply::Result::Negotiated(Negotiated {
negotiated_protocol: PROTOCOL_VERSION,
daemon_version: "test-daemon".into(),
backend_pipe: backend_for_broker,
..Default::default()
})),
};
let response = Frame {
envelope_version: PROTOCOL_VERSION,
kind: FrameKind::Response as i32,
payload_protocol: CONTROL_PAYLOAD_PROTOCOL,
payload: reply.encode_to_vec(),
request_id: request_frame.request_id,
payload_encoding: PayloadEncoding::None as i32,
deadline_unix_ms: 0,
traceparent: String::new(),
tracestate: String::new(),
};
write_frame(&mut stream, &response.encode_to_vec()).expect("write HelloReply frame");
});
let mut request =
OwnedConnectRequest::new(&broker_endpoint, "compat-service", "1.2.3", "1.2.3");
request.client_version = "consumer-9.8.7".into();
request.client_lib_name = "consumer-broker-adapter".into();
request.client_lib_version = "6.5.4".into();
request.client_keepalive_secs = 17;
let mut session = AsyncBrokerSession::adopt(request)
.await
.expect("v2-compatible adoption");
assert_eq!(session.route(), BackendConnectionRoute::BrokerNegotiated);
assert_eq!(session.endpoint(), backend_endpoint);
let response = session
.request(0xCAFE, b"ping".to_vec())
.await
.expect("backend round trip");
assert_eq!(response.payload, b"pong");
let hello = hello_rx.recv().expect("observed Hello contract");
assert!(
hello.request_id.starts_with("client_v2-compat-service-"),
"compat adapter sent a v1 Hello request id: {:?}",
hello.request_id
);
assert_eq!(hello.service_name, "compat-service");
assert_eq!(hello.wanted_version, "1.2.3");
assert_eq!(hello.client_version, "consumer-9.8.7");
assert_eq!(hello.client_lib_name, "consumer-broker-adapter");
assert_eq!(hello.client_lib_version, "6.5.4");
assert_eq!(hello.client_keepalive_secs, 17);
broker.join().expect("broker stub exits cleanly");
backend.join().expect("backend stub exits cleanly");
let _ = std::fs::remove_file(broker_endpoint);
let _ = std::fs::remove_file(backend_endpoint);
}
}