use crate::frame::Frame;
pub const LENGTH_PREFIX_BYTES: usize = 4;
pub(crate) const INCONSISTENT_SCOPE_ERROR_PREFIX: &str = "error frame violates the id/scope rule: ";
pub const DEFAULT_MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
#[derive(Debug, thiserror::Error, PartialEq, Eq)]
pub enum CodecError {
#[error("truncated length prefix: got {available} of {LENGTH_PREFIX_BYTES} bytes")]
TruncatedLengthPrefix { available: usize },
#[error("truncated payload: declared {declared} bytes, got {available}")]
TruncatedPayload { declared: usize, available: usize },
#[error("frame of {declared} bytes exceeds the {max} byte maximum")]
FrameTooLarge { declared: usize, max: usize },
#[error("frame of {declared} bytes exceeds the u32 length prefix's {max} byte capacity")]
U32PrefixLimitExceeded { declared: usize, max: usize },
#[error("payload is not valid JSON: {0}")]
InvalidJson(String),
#[error("payload has no string \"kind\" discriminant field")]
MissingKind,
#[error("unknown frame kind: {0:?}")]
UnknownFrameKind(String),
#[error("frame kind {kind:?}: {detail}")]
InvalidFields { kind: String, detail: String },
#[error("error frame violates the id/scope rule: {detail}")]
InconsistentErrorScope { detail: String },
#[error(
"fallback error frame (unrecognized code {code:?}) is not re-encodable: \
re-encoding would emit \"internal\" and discard the newer wire code"
)]
FallbackFrameNotEncodable { code: String },
}
impl CodecError {
pub const fn wire_code(&self) -> crate::error::WireErrorCode {
match self {
CodecError::FrameTooLarge { .. } | CodecError::U32PrefixLimitExceeded { .. } => {
crate::error::WireErrorCode::FrameTooLarge
}
CodecError::TruncatedLengthPrefix { .. }
| CodecError::TruncatedPayload { .. }
| CodecError::InvalidJson(_)
| CodecError::MissingKind
| CodecError::UnknownFrameKind(_)
| CodecError::InvalidFields { .. }
| CodecError::InconsistentErrorScope { .. } => {
crate::error::WireErrorCode::MalformedFrame
}
CodecError::FallbackFrameNotEncodable { .. } => crate::error::WireErrorCode::Internal,
}
}
}
impl From<&CodecError> for crate::error::WireErrorCode {
fn from(err: &CodecError) -> Self {
err.wire_code()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct FrameCodec {
max_frame_bytes: usize,
}
impl FrameCodec {
pub const fn new(max_frame_bytes: usize) -> Self {
Self { max_frame_bytes }
}
pub const fn max_frame_bytes(&self) -> usize {
self.max_frame_bytes
}
pub fn encode(&self, frame: &Frame) -> Result<Vec<u8>, CodecError> {
encode_frame_with_max(frame, self.max_frame_bytes)
}
pub fn decode(&self, buf: &[u8]) -> Result<Frame, CodecError> {
decode_frame(buf, self.max_frame_bytes)
}
pub fn decode_with_consumed(&self, buf: &[u8]) -> Result<(Frame, usize), CodecError> {
decode_frame_with_consumed(buf, self.max_frame_bytes)
}
}
impl Default for FrameCodec {
fn default() -> Self {
Self::new(DEFAULT_MAX_FRAME_BYTES)
}
}
pub fn encode_frame(frame: &Frame) -> Result<Vec<u8>, CodecError> {
encode_frame_with_max(frame, DEFAULT_MAX_FRAME_BYTES)
}
pub fn encode_frame_with_max(frame: &Frame, max_frame_bytes: usize) -> Result<Vec<u8>, CodecError> {
validate_frame_for_wire(frame)?;
let payload = serde_json::to_vec(frame).map_err(|e| CodecError::InvalidJson(e.to_string()))?;
check_encode_payload_len(payload.len(), max_frame_bytes)?;
let mut buf = Vec::with_capacity(LENGTH_PREFIX_BYTES + payload.len());
buf.extend_from_slice(&(payload.len() as u32).to_be_bytes());
buf.extend_from_slice(&payload);
Ok(buf)
}
fn validate_frame_for_wire(frame: &Frame) -> Result<(), CodecError> {
use crate::error::TerminalScope;
fn check_id(kind: &str, id: &crate::frame::OperationId) -> Result<(), CodecError> {
if id.0.is_empty() {
return Err(CodecError::InvalidFields {
kind: kind.to_string(),
detail: "operation id must be a non-empty string".to_string(),
});
}
Ok(())
}
match frame {
Frame::Handshake { version } => {
if version.get() == 0 {
return Err(CodecError::InvalidFields {
kind: "handshake".to_string(),
detail: "protocol version 0 does not exist".to_string(),
});
}
}
Frame::HandshakeAck { version } => {
if version.get() == 0 {
return Err(CodecError::InvalidFields {
kind: "handshake_ack".to_string(),
detail: "protocol version 0 does not exist".to_string(),
});
}
}
Frame::Event { .. } => {}
Frame::Request { id, .. } => check_id("request", id)?,
Frame::Response { id, .. } => check_id("response", id)?,
Frame::Cancel { id } => check_id("cancel", id)?,
Frame::Subscribe { id, .. } => check_id("subscribe", id)?,
Frame::SubscribeAck { id, .. } => check_id("subscribe_ack", id)?,
Frame::Unsubscribe { id, .. } => check_id("unsubscribe", id)?,
Frame::UnsubscribeAck { id, .. } => check_id("unsubscribe_ack", id)?,
Frame::Error {
id,
code,
unrecognized_code,
..
} => {
if let Some(raw_code) = unrecognized_code {
return Err(CodecError::FallbackFrameNotEncodable {
code: raw_code.clone(),
});
}
if let Some(id) = id {
check_id("error", id)?;
}
match (code.terminal_scope(), id) {
(TerminalScope::Connection, Some(id)) => {
return Err(CodecError::InconsistentErrorScope {
detail: format!(
"connection-terminal code {code} must not carry an operation id, got {id}"
),
});
}
(TerminalScope::Request, None) => {
return Err(CodecError::InconsistentErrorScope {
detail: format!(
"request-terminal code {code} must echo the operation id it terminates"
),
});
}
(TerminalScope::Connection, None) | (TerminalScope::Request, Some(_)) => {}
}
}
}
Ok(())
}
fn check_encode_payload_len(payload_len: usize, max_frame_bytes: usize) -> Result<(), CodecError> {
if payload_len > max_frame_bytes {
return Err(CodecError::FrameTooLarge {
declared: payload_len,
max: max_frame_bytes,
});
}
if payload_len > u32::MAX as usize {
return Err(CodecError::U32PrefixLimitExceeded {
declared: payload_len,
max: u32::MAX as usize,
});
}
Ok(())
}
pub fn decode_frame(buf: &[u8], max_frame_bytes: usize) -> Result<Frame, CodecError> {
decode_frame_with_consumed(buf, max_frame_bytes).map(|(frame, _)| frame)
}
pub fn decode_frame_with_consumed(
buf: &[u8],
max_frame_bytes: usize,
) -> Result<(Frame, usize), CodecError> {
if buf.len() < LENGTH_PREFIX_BYTES {
return Err(CodecError::TruncatedLengthPrefix {
available: buf.len(),
});
}
let mut len_bytes = [0u8; LENGTH_PREFIX_BYTES];
len_bytes.copy_from_slice(&buf[..LENGTH_PREFIX_BYTES]);
let declared = u32::from_be_bytes(len_bytes) as usize;
if declared > max_frame_bytes {
return Err(CodecError::FrameTooLarge {
declared,
max: max_frame_bytes,
});
}
let available = buf.len() - LENGTH_PREFIX_BYTES;
if available < declared {
return Err(CodecError::TruncatedPayload {
declared,
available,
});
}
let payload = &buf[LENGTH_PREFIX_BYTES..LENGTH_PREFIX_BYTES + declared];
let frame = decode_payload(payload)?;
Ok((frame, LENGTH_PREFIX_BYTES + declared))
}
pub(crate) fn decode_payload(payload: &[u8]) -> Result<Frame, CodecError> {
let value: serde_json::Value =
serde_json::from_slice(payload).map_err(|e| CodecError::InvalidJson(e.to_string()))?;
let kind = value
.as_object()
.and_then(|obj| obj.get("kind"))
.and_then(|k| k.as_str())
.ok_or(CodecError::MissingKind)?
.to_string();
if !crate::frame::FRAME_KINDS.contains(&kind.as_str()) {
return Err(CodecError::UnknownFrameKind(kind));
}
serde_json::from_value(value).map_err(|e| {
let detail = e.to_string();
match detail.strip_prefix(INCONSISTENT_SCOPE_ERROR_PREFIX) {
Some(scope_detail) => CodecError::InconsistentErrorScope {
detail: scope_detail.to_string(),
},
None => CodecError::InvalidFields { kind, detail },
}
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::frame::OperationId;
fn sample_frames() -> Vec<Frame> {
vec![
Frame::Handshake {
version: crate::version::CURRENT_VERSION,
},
Frame::HandshakeAck {
version: crate::version::CURRENT_VERSION,
},
Frame::Request {
id: OperationId::from("op-1"),
ops: "stats()".to_string(),
deadline_ms: Some(5000),
namespace: Some("research".to_string()),
actor_id: Some("lambda".to_string()),
visible_namespaces: Some(vec!["research".to_string(), "ops".to_string()]),
},
Frame::Response {
id: OperationId::from("op-1"),
result: serde_json::json!({"ok": true, "tool": "stats", "result": {"entities": 3}}),
},
Frame::Error {
id: Some(OperationId::from("op-1")),
code: crate::error::WireErrorCode::PeerClassDenied,
message: "denied".to_string(),
unrecognized_code: None,
},
Frame::Cancel {
id: OperationId::from("op-1"),
},
Frame::Subscribe {
id: OperationId::from("op-2"),
topic: "comm.message_created".to_string(),
resume_cursor: Some(42),
},
Frame::SubscribeAck {
id: OperationId::from("op-2"),
topic: "comm.message_created".to_string(),
start_cursor: 42,
},
Frame::Unsubscribe {
id: OperationId::from("op-3"),
topic: "comm.message_created".to_string(),
},
Frame::UnsubscribeAck {
id: OperationId::from("op-3"),
topic: "comm.message_created".to_string(),
},
Frame::Event {
topic: "comm.message_created".to_string(),
cursor: 43,
occurred_at: "2026-08-04T11:00:00Z".to_string(),
payload: serde_json::json!({"message_id": "m-1"}),
},
]
}
#[test]
fn round_trips_a_cancel_frame() {
let frame = Frame::Cancel {
id: OperationId::from("op-1"),
};
let codec = FrameCodec::default();
let wire = codec.encode(&frame).unwrap();
assert_eq!(codec.decode(&wire).unwrap(), frame);
}
#[test]
fn rejects_truncated_length_prefix() {
let codec = FrameCodec::default();
assert_eq!(
codec.decode(&[0u8, 1]).unwrap_err(),
CodecError::TruncatedLengthPrefix { available: 2 }
);
}
#[test]
fn rejects_truncated_payload_with_declared_and_available() {
let declared: u32 = 64;
let mut wire = declared.to_be_bytes().to_vec();
wire.extend_from_slice(b"part!"); assert_eq!(
decode_frame(&wire, DEFAULT_MAX_FRAME_BYTES).unwrap_err(),
CodecError::TruncatedPayload {
declared: 64,
available: 5
}
);
}
#[test]
fn rejects_truncated_payload_when_prefix_declares_exactly_one_byte_more() {
let frame = Frame::Cancel {
id: OperationId::from("op-1"),
};
let wire = encode_frame(&frame).unwrap();
assert_eq!(
decode_frame(&wire[..wire.len() - 1], DEFAULT_MAX_FRAME_BYTES).unwrap_err(),
CodecError::TruncatedPayload {
declared: wire.len() - LENGTH_PREFIX_BYTES,
available: wire.len() - LENGTH_PREFIX_BYTES - 1
}
);
}
#[test]
fn rejects_oversized_frame() {
let codec = FrameCodec::new(4);
let frame = Frame::Cancel {
id: OperationId::from("op-1"),
};
let wire = encode_frame(&frame).unwrap();
assert_eq!(
codec.decode(&wire).unwrap_err(),
CodecError::FrameTooLarge {
declared: wire.len() - LENGTH_PREFIX_BYTES,
max: 4
}
);
}
#[test]
fn rejects_unknown_frame_kind() {
let payload = br#"{"kind":"ping"}"#;
assert_eq!(
decode_payload(payload).unwrap_err(),
CodecError::UnknownFrameKind("ping".to_string())
);
}
#[test]
fn rejects_non_json_payload() {
let payload = b"not json";
match decode_payload(payload).unwrap_err() {
CodecError::InvalidJson(_) => {}
other => panic!("expected InvalidJson, got {other:?}"),
}
}
#[test]
fn rejects_zero_length_payload() {
let codec = FrameCodec::default();
match codec.decode(&0u32.to_be_bytes()).unwrap_err() {
CodecError::InvalidJson(_) => {}
other => panic!("expected InvalidJson, got {other:?}"),
}
}
#[test]
fn rejects_missing_required_field() {
let payload = br#"{"kind":"cancel"}"#;
match decode_payload(payload).unwrap_err() {
CodecError::InvalidFields { kind, .. } => assert_eq!(kind, "cancel"),
other => panic!("expected InvalidFields, got {other:?}"),
}
}
#[test]
fn rejects_unknown_top_level_field() {
let payload = br#"{"kind":"cancel","id":"op-1","unexpected":true}"#;
match decode_payload(payload).unwrap_err() {
CodecError::InvalidFields { kind, detail } => {
assert_eq!(kind, "cancel");
assert!(
detail.contains("unexpected"),
"detail should name the offending field: {detail}"
);
}
other => panic!("expected InvalidFields, got {other:?}"),
}
}
#[test]
fn rejects_unknown_field_alongside_optional_fields() {
let payload =
br#"{"kind":"subscribe","id":"op-2","topic":"a.b","resume_cursor":1,"extra":{"nested":true}}"#;
match decode_payload(payload).unwrap_err() {
CodecError::InvalidFields { kind, detail } => {
assert_eq!(kind, "subscribe");
assert!(detail.contains("extra"), "detail: {detail}");
}
other => panic!("expected InvalidFields, got {other:?}"),
}
}
#[test]
fn opaque_payload_values_do_not_reject_unknown_keys() {
let payload = br#"{"kind":"event","topic":"a.b","cursor":1,"occurred_at":"2026-08-04T11:00:00Z","payload":{"anything":{"goes":true}}}"#;
decode_payload(payload).unwrap();
}
#[test]
fn every_frame_kind_round_trips_with_strict_decoding() {
let frames = sample_frames();
assert_eq!(frames.len(), crate::frame::FRAME_KINDS.len());
for frame in frames {
let wire = encode_frame(&frame).unwrap();
let decoded = decode_frame(&wire, DEFAULT_MAX_FRAME_BYTES)
.unwrap_or_else(|e| panic!("kind {:?} failed to decode: {e}", frame.kind()));
assert_eq!(decoded, frame, "kind {:?} did not round-trip", frame.kind());
}
}
#[test]
fn rejects_connection_terminal_error_carrying_an_id() {
let payload =
br#"{"kind":"error","id":"op-1","code":"frame_too_large","message":"too big"}"#;
match decode_payload(payload).unwrap_err() {
CodecError::InconsistentErrorScope { detail } => {
assert!(detail.contains("frame_too_large"), "detail: {detail}");
assert!(detail.contains("op-1"), "detail: {detail}");
}
other => panic!("expected InconsistentErrorScope, got {other:?}"),
}
}
#[test]
fn rejects_request_terminal_error_without_an_id() {
let payload = br#"{"kind":"error","code":"cancelled","message":"cancelled"}"#;
match decode_payload(payload).unwrap_err() {
CodecError::InconsistentErrorScope { detail } => {
assert!(detail.contains("cancelled"), "detail: {detail}");
}
other => panic!("expected InconsistentErrorScope, got {other:?}"),
}
}
#[test]
fn accepts_connection_terminal_error_without_an_id() {
let payload =
br#"{"kind":"error","code":"unsupported_version","message":"no common version"}"#;
let frame = decode_payload(payload).unwrap();
assert!(matches!(frame, Frame::Error { id: None, .. }));
}
#[test]
fn accepts_request_terminal_error_with_an_id() {
let payload =
br#"{"kind":"error","id":"op-9","code":"deadline_exceeded","message":"too slow"}"#;
let frame = decode_payload(payload).unwrap();
match frame {
Frame::Error { id, code, .. } => {
assert_eq!(id, Some(OperationId::from("op-9")));
assert_eq!(code, crate::error::WireErrorCode::DeadlineExceeded);
}
other => panic!("expected an error frame, got {other:?}"),
}
}
#[test]
fn rejects_empty_operation_id() {
let payload = br#"{"kind":"cancel","id":""}"#;
match decode_payload(payload).unwrap_err() {
CodecError::InvalidFields { kind, detail } => {
assert_eq!(kind, "cancel");
assert!(detail.contains("non-empty"), "detail: {detail}");
}
other => panic!("expected InvalidFields, got {other:?}"),
}
}
#[test]
fn rejects_valid_json_that_is_not_an_object() {
for payload in [
b"[1, 2, 3]".as_slice(),
b"\"cancel\"".as_slice(),
b"42".as_slice(),
b"null".as_slice(),
] {
assert_eq!(
decode_payload(payload).unwrap_err(),
CodecError::MissingKind,
"payload {payload:?}"
);
}
}
#[test]
fn rejects_payload_with_no_kind_field() {
let payload = br#"{"id":"op-1"}"#;
assert_eq!(
decode_payload(payload).unwrap_err(),
CodecError::MissingKind
);
}
#[test]
fn rejects_non_string_kind() {
for payload in [
br#"{"kind":42}"#.as_slice(),
br#"{"kind":null}"#.as_slice(),
br#"{"kind":["cancel"]}"#.as_slice(),
] {
assert_eq!(
decode_payload(payload).unwrap_err(),
CodecError::MissingKind,
"payload {payload:?}"
);
}
}
#[test]
fn rejects_wrong_typed_field() {
let payload = br#"{"kind":"cancel","id":7}"#;
match decode_payload(payload).unwrap_err() {
CodecError::InvalidFields { kind, detail } => {
assert_eq!(kind, "cancel");
assert!(detail.contains("id"), "detail: {detail}");
}
other => panic!("expected InvalidFields, got {other:?}"),
}
let payload = br#"{"kind":"subscribe_ack","id":"op-2","topic":"a.b","start_cursor":"42"}"#;
match decode_payload(payload).unwrap_err() {
CodecError::InvalidFields { kind, .. } => assert_eq!(kind, "subscribe_ack"),
other => panic!("expected InvalidFields, got {other:?}"),
}
}
#[test]
fn opaque_payloads_are_preserved_semantically_not_byte_for_byte() {
let raw = br#"{"kind":"response","id":"op-1","result":{"zeta":1,"alpha":{"n":9007199254740993},"huge":18446744073709551616,"neg":-9007199254740993}}"#;
let frame = decode_payload(raw).unwrap();
let Frame::Response { result, .. } = &frame else {
panic!("expected a response frame");
};
assert_eq!(result["alpha"]["n"], serde_json::json!(9007199254740993u64));
assert_eq!(result["neg"], serde_json::json!(-9007199254740993i64));
assert!(result["huge"].is_f64());
assert!(!result["huge"].is_u64());
assert_eq!(result["huge"].as_f64().unwrap(), 2.0f64.powi(64));
let wire = encode_frame(&frame).unwrap();
assert_eq!(decode_payload(&wire[LENGTH_PREFIX_BYTES..]).unwrap(), frame);
let payload = std::str::from_utf8(&wire[LENGTH_PREFIX_BYTES..]).unwrap();
let alpha = payload.find("\"alpha\"").unwrap();
let zeta = payload.find("\"zeta\"").unwrap();
assert!(
alpha < zeta,
"expected sorted key order in re-encoded payload: {payload}"
);
}
#[test]
fn encode_rejects_empty_operation_id() {
let frame = Frame::Cancel {
id: OperationId::from(""),
};
match encode_frame(&frame).unwrap_err() {
CodecError::InvalidFields { kind, detail } => {
assert_eq!(kind, "cancel");
assert!(detail.contains("non-empty"), "detail: {detail}");
}
other => panic!("expected InvalidFields, got {other:?}"),
}
}
#[test]
fn encode_rejects_empty_operation_id_in_every_id_field() {
let frames = [
Frame::Request {
id: OperationId::from(""),
ops: "stats()".to_string(),
deadline_ms: None,
namespace: None,
actor_id: None,
visible_namespaces: None,
},
Frame::Response {
id: OperationId::from(""),
result: serde_json::json!({}),
},
Frame::Subscribe {
id: OperationId::from(""),
topic: "a.b".to_string(),
resume_cursor: None,
},
Frame::SubscribeAck {
id: OperationId::from(""),
topic: "a.b".to_string(),
start_cursor: 0,
},
Frame::Unsubscribe {
id: OperationId::from(""),
topic: "a.b".to_string(),
},
Frame::UnsubscribeAck {
id: OperationId::from(""),
topic: "a.b".to_string(),
},
Frame::Error {
id: Some(OperationId::from("")),
code: crate::error::WireErrorCode::Internal,
message: "failure".to_string(),
unrecognized_code: None,
},
];
for frame in frames {
match encode_frame(&frame).unwrap_err() {
CodecError::InvalidFields { detail, .. } => {
assert!(detail.contains("non-empty"), "detail: {detail}");
}
other => panic!(
"kind {:?}: expected InvalidFields, got {other:?}",
frame.kind()
),
}
}
}
#[test]
fn encode_rejects_zero_protocol_version_in_both_handshake_kinds() {
for frame in [
Frame::Handshake {
version: crate::version::ProtocolVersion::new(0),
},
Frame::HandshakeAck {
version: crate::version::ProtocolVersion::new(0),
},
] {
match encode_frame(&frame).unwrap_err() {
CodecError::InvalidFields { kind, detail } => {
assert_eq!(kind, frame.kind());
assert!(detail.contains("version 0"), "detail: {detail}");
}
other => panic!(
"kind {:?}: expected InvalidFields, got {other:?}",
frame.kind()
),
}
}
}
#[test]
fn encode_rejects_connection_terminal_error_carrying_an_id() {
let frame = Frame::Error {
id: Some(OperationId::from("op-1")),
code: crate::error::WireErrorCode::FrameTooLarge,
message: "too big".to_string(),
unrecognized_code: None,
};
match encode_frame(&frame).unwrap_err() {
CodecError::InconsistentErrorScope { detail } => {
assert!(detail.contains("frame_too_large"), "detail: {detail}");
assert!(detail.contains("op-1"), "detail: {detail}");
}
other => panic!("expected InconsistentErrorScope, got {other:?}"),
}
}
#[test]
fn encode_rejects_request_terminal_error_without_an_id() {
let frame = Frame::Error {
id: None,
code: crate::error::WireErrorCode::Cancelled,
message: "cancelled".to_string(),
unrecognized_code: None,
};
match encode_frame(&frame).unwrap_err() {
CodecError::InconsistentErrorScope { detail } => {
assert!(detail.contains("cancelled"), "detail: {detail}");
}
other => panic!("expected InconsistentErrorScope, got {other:?}"),
}
}
#[test]
fn encode_accepts_both_consistent_error_scopes() {
let consistent = [
Frame::Error {
id: None,
code: crate::error::WireErrorCode::MalformedFrame,
message: "connection scope".to_string(),
unrecognized_code: None,
},
Frame::Error {
id: Some(OperationId::from("op-1")),
code: crate::error::WireErrorCode::DeadlineExceeded,
message: "request scope".to_string(),
unrecognized_code: None,
},
];
for frame in consistent {
let wire = encode_frame(&frame).expect("consistent scope must encode");
assert_eq!(decode_frame(&wire, DEFAULT_MAX_FRAME_BYTES).unwrap(), frame);
}
}
#[test]
fn unknown_code_fallback_preserves_the_raw_string_with_an_id() {
let payload =
br#"{"kind":"error","id":"op-7","code":"future_code_xyz","message":"from newer peer"}"#;
let frame = decode_payload(payload).unwrap();
match frame {
Frame::Error {
id,
code,
unrecognized_code,
..
} => {
assert_eq!(code, crate::error::WireErrorCode::Internal);
assert_eq!(id, Some(OperationId::from("op-7")));
assert_eq!(unrecognized_code.as_deref(), Some("future_code_xyz"));
}
other => panic!("expected an error frame, got {other:?}"),
}
}
#[test]
fn unknown_code_fallback_preserves_the_raw_string_without_an_id() {
let payload = br#"{"kind":"error","code":"future_code_xyz","message":"from newer peer"}"#;
match decode_payload(payload).unwrap() {
Frame::Error {
id,
code,
unrecognized_code,
..
} => {
assert_eq!(code, crate::error::WireErrorCode::Internal);
assert!(id.is_none());
assert_eq!(unrecognized_code.as_deref(), Some("future_code_xyz"));
}
other => panic!("expected an error frame, got {other:?}"),
}
}
#[test]
fn closed_set_code_never_populates_unrecognized_code() {
let payload =
br#"{"kind":"error","id":"op-9","code":"deadline_exceeded","message":"too slow"}"#;
match decode_payload(payload).unwrap() {
Frame::Error {
unrecognized_code, ..
} => assert!(unrecognized_code.is_none()),
other => panic!("expected an error frame, got {other:?}"),
}
}
#[test]
fn decode_scope_rejection_detail_is_reclassified_without_the_shared_prefix() {
let payload =
br#"{"kind":"error","id":"op-1","code":"frame_too_large","message":"too big"}"#;
match decode_payload(payload).unwrap_err() {
CodecError::InconsistentErrorScope { detail } => {
assert!(
!detail.contains(INCONSISTENT_SCOPE_ERROR_PREFIX),
"detail must not repeat the shared prefix: {detail}"
);
assert!(detail.starts_with("connection-terminal code"));
}
other => panic!("expected InconsistentErrorScope, got {other:?}"),
}
}
#[test]
fn encode_rejects_a_fallback_frame_with_an_id() {
let payload =
br#"{"kind":"error","id":"op-7","code":"future_code_xyz","message":"from newer peer"}"#;
let frame = decode_payload(payload).unwrap();
match encode_frame(&frame).unwrap_err() {
CodecError::FallbackFrameNotEncodable { code } => {
assert_eq!(code, "future_code_xyz");
}
other => panic!("expected FallbackFrameNotEncodable, got {other:?}"),
}
}
#[test]
fn encode_rejects_a_fallback_frame_without_an_id() {
let payload = br#"{"kind":"error","code":"future_code_xyz","message":"from newer peer"}"#;
let frame = decode_payload(payload).unwrap();
match encode_frame(&frame).unwrap_err() {
CodecError::FallbackFrameNotEncodable { code } => {
assert_eq!(code, "future_code_xyz");
}
other => panic!("expected FallbackFrameNotEncodable, got {other:?}"),
}
}
#[test]
fn encode_accepts_an_honest_internal_error_frame() {
let frame = Frame::Error {
id: Some(OperationId::from("op-1")),
code: crate::error::WireErrorCode::Internal,
message: "boom".to_string(),
unrecognized_code: None,
};
let wire = encode_frame(&frame).expect("honest internal must encode");
let decoded = decode_frame(&wire, DEFAULT_MAX_FRAME_BYTES).unwrap();
assert_eq!(decoded, frame);
}
#[test]
fn decode_with_consumed_reports_consumed_length_on_a_two_frame_buffer() {
let first = Frame::Cancel {
id: OperationId::from("op-1"),
};
let second = Frame::Handshake {
version: crate::version::CURRENT_VERSION,
};
let codec = FrameCodec::default();
let mut buf = codec.encode(&first).unwrap();
let first_len = buf.len();
buf.extend_from_slice(&codec.encode(&second).unwrap());
let (decoded, consumed) = codec.decode_with_consumed(&buf).unwrap();
assert_eq!(decoded, first);
assert_eq!(consumed, first_len);
let (decoded2, consumed2) = codec.decode_with_consumed(&buf[consumed..]).unwrap();
assert_eq!(decoded2, second);
assert_eq!(consumed + consumed2, buf.len());
assert_eq!(codec.decode(&buf).unwrap(), first);
}
#[test]
fn decode_with_consumed_errors_without_reporting_length_on_truncation() {
let codec = FrameCodec::default();
let wire = codec
.encode(&Frame::Cancel {
id: OperationId::from("op-1"),
})
.unwrap();
let err = codec
.decode_with_consumed(&wire[..wire.len() - 1])
.unwrap_err();
assert!(matches!(err, CodecError::TruncatedPayload { .. }));
}
#[test]
fn serde_json_feature_posture_is_default_keys_serialize_sorted() {
let mut map = serde_json::Map::new();
map.insert("zeta".to_string(), serde_json::json!(1));
map.insert("alpha".to_string(), serde_json::json!(2));
map.insert("mid".to_string(), serde_json::json!(3));
let wire = serde_json::to_string(&serde_json::Value::Object(map)).unwrap();
assert_eq!(
wire, r#"{"alpha":2,"mid":3,"zeta":1}"#,
"serde_json keys did not serialize in sorted order — \
`preserve_order` may have been enabled workspace-wide"
);
}
#[test]
fn codec_errors_map_to_their_wire_error_codes() {
use crate::error::WireErrorCode;
let cases: &[(CodecError, WireErrorCode)] = &[
(
CodecError::TruncatedLengthPrefix { available: 2 },
WireErrorCode::MalformedFrame,
),
(
CodecError::TruncatedPayload {
declared: 8,
available: 3,
},
WireErrorCode::MalformedFrame,
),
(
CodecError::FrameTooLarge {
declared: 10,
max: 4,
},
WireErrorCode::FrameTooLarge,
),
(
CodecError::U32PrefixLimitExceeded {
declared: 10,
max: 4,
},
WireErrorCode::FrameTooLarge,
),
(
CodecError::InvalidJson("x".to_string()),
WireErrorCode::MalformedFrame,
),
(CodecError::MissingKind, WireErrorCode::MalformedFrame),
(
CodecError::UnknownFrameKind("ping".to_string()),
WireErrorCode::MalformedFrame,
),
(
CodecError::InvalidFields {
kind: "cancel".to_string(),
detail: "x".to_string(),
},
WireErrorCode::MalformedFrame,
),
(
CodecError::InconsistentErrorScope {
detail: "x".to_string(),
},
WireErrorCode::MalformedFrame,
),
(
CodecError::FallbackFrameNotEncodable {
code: "future_code_xyz".to_string(),
},
WireErrorCode::Internal,
),
];
for (err, expected) in cases {
assert_eq!(
err.wire_code(),
*expected,
"{err:?} mapped to {:?}, expected {expected:?}",
err.wire_code()
);
assert_eq!(
WireErrorCode::from(err),
*expected,
"From<&CodecError> disagrees"
);
}
}
}
#[cfg(test)]
mod configured_max_tests {
use super::*;
use crate::frame::Frame;
fn oversized_frame() -> Frame {
Frame::Request {
id: crate::frame::OperationId("x".into()),
ops: "a".repeat(4096),
deadline_ms: None,
namespace: None,
actor_id: None,
visible_namespaces: None,
}
}
#[test]
fn codec_encode_honors_its_own_maximum_not_the_default() {
let codec = FrameCodec::new(1024);
let err = codec.encode(&oversized_frame()).unwrap_err();
match err {
CodecError::FrameTooLarge { max, .. } => assert_eq!(max, 1024),
other => panic!("expected FrameTooLarge with the configured max, got {other:?}"),
}
}
#[test]
fn codec_encode_accepts_a_frame_within_its_own_maximum() {
let codec = FrameCodec::new(1024 * 1024);
let wire = codec
.encode(&oversized_frame())
.expect("within configured max");
assert_eq!(codec.decode(&wire).unwrap(), oversized_frame());
}
#[test]
fn free_function_encode_still_uses_the_default_maximum() {
let wire = encode_frame(&oversized_frame()).expect("well under 8 MiB");
assert!(wire.len() > 4096);
}
#[test]
fn codec_configured_above_u32_capacity_still_encodes_normal_frames() {
assert!(encode_frame_with_max(&oversized_frame(), usize::MAX).is_ok());
}
#[test]
fn encode_size_guard_reports_the_configured_max_when_it_is_binding() {
let err = check_encode_payload_len(32, 16).unwrap_err();
assert_eq!(
err,
CodecError::FrameTooLarge {
declared: 32,
max: 16
}
);
}
#[cfg(target_pointer_width = "64")]
#[test]
fn encode_size_guard_reports_the_u32_prefix_capacity_when_it_is_binding() {
let err = check_encode_payload_len(u32::MAX as usize + 1, usize::MAX).unwrap_err();
assert_eq!(
err,
CodecError::U32PrefixLimitExceeded {
declared: u32::MAX as usize + 1,
max: u32::MAX as usize
}
);
assert!(check_encode_payload_len(u32::MAX as usize, usize::MAX).is_ok());
}
}