1use crate::frame::Frame;
13
14pub const LENGTH_PREFIX_BYTES: usize = 4;
16
17pub(crate) const INCONSISTENT_SCOPE_ERROR_PREFIX: &str = "error frame violates the id/scope rule: ";
26
27pub const DEFAULT_MAX_FRAME_BYTES: usize = 8 * 1024 * 1024;
37
38#[derive(Debug, thiserror::Error, PartialEq, Eq)]
48pub enum CodecError {
49 #[error("truncated length prefix: got {available} of {LENGTH_PREFIX_BYTES} bytes")]
52 TruncatedLengthPrefix { available: usize },
53
54 #[error("truncated payload: declared {declared} bytes, got {available}")]
56 TruncatedPayload { declared: usize, available: usize },
57
58 #[error("frame of {declared} bytes exceeds the {max} byte maximum")]
62 FrameTooLarge { declared: usize, max: usize },
63
64 #[error("frame of {declared} bytes exceeds the u32 length prefix's {max} byte capacity")]
71 U32PrefixLimitExceeded { declared: usize, max: usize },
72
73 #[error("payload is not valid JSON: {0}")]
78 InvalidJson(String),
79
80 #[error("payload has no string \"kind\" discriminant field")]
83 MissingKind,
84
85 #[error("unknown frame kind: {0:?}")]
88 UnknownFrameKind(String),
89
90 #[error("frame kind {kind:?}: {detail}")]
97 InvalidFields { kind: String, detail: String },
98
99 #[error("error frame violates the id/scope rule: {detail}")]
109 InconsistentErrorScope { detail: String },
110
111 #[error(
121 "fallback error frame (unrecognized code {code:?}) is not re-encodable: \
122 re-encoding would emit \"internal\" and discard the newer wire code"
123 )]
124 FallbackFrameNotEncodable { code: String },
125}
126
127impl CodecError {
128 pub const fn wire_code(&self) -> crate::error::WireErrorCode {
143 match self {
144 CodecError::FrameTooLarge { .. } | CodecError::U32PrefixLimitExceeded { .. } => {
145 crate::error::WireErrorCode::FrameTooLarge
146 }
147 CodecError::TruncatedLengthPrefix { .. }
148 | CodecError::TruncatedPayload { .. }
149 | CodecError::InvalidJson(_)
150 | CodecError::MissingKind
151 | CodecError::UnknownFrameKind(_)
152 | CodecError::InvalidFields { .. }
153 | CodecError::InconsistentErrorScope { .. } => {
154 crate::error::WireErrorCode::MalformedFrame
155 }
156 CodecError::FallbackFrameNotEncodable { .. } => crate::error::WireErrorCode::Internal,
157 }
158 }
159}
160
161impl From<&CodecError> for crate::error::WireErrorCode {
162 fn from(err: &CodecError) -> Self {
163 err.wire_code()
164 }
165}
166
167#[derive(Debug, Clone, Copy, PartialEq, Eq)]
174pub struct FrameCodec {
175 max_frame_bytes: usize,
176}
177
178impl FrameCodec {
179 pub const fn new(max_frame_bytes: usize) -> Self {
182 Self { max_frame_bytes }
183 }
184
185 pub const fn max_frame_bytes(&self) -> usize {
186 self.max_frame_bytes
187 }
188
189 pub fn encode(&self, frame: &Frame) -> Result<Vec<u8>, CodecError> {
193 encode_frame_with_max(frame, self.max_frame_bytes)
194 }
195
196 pub fn decode(&self, buf: &[u8]) -> Result<Frame, CodecError> {
206 decode_frame(buf, self.max_frame_bytes)
207 }
208
209 pub fn decode_with_consumed(&self, buf: &[u8]) -> Result<(Frame, usize), CodecError> {
216 decode_frame_with_consumed(buf, self.max_frame_bytes)
217 }
218}
219
220impl Default for FrameCodec {
221 fn default() -> Self {
222 Self::new(DEFAULT_MAX_FRAME_BYTES)
223 }
224}
225
226pub fn encode_frame(frame: &Frame) -> Result<Vec<u8>, CodecError> {
233 encode_frame_with_max(frame, DEFAULT_MAX_FRAME_BYTES)
234}
235
236pub fn encode_frame_with_max(frame: &Frame, max_frame_bytes: usize) -> Result<Vec<u8>, CodecError> {
262 validate_frame_for_wire(frame)?;
263 let payload = serde_json::to_vec(frame).map_err(|e| CodecError::InvalidJson(e.to_string()))?;
264 check_encode_payload_len(payload.len(), max_frame_bytes)?;
265 let mut buf = Vec::with_capacity(LENGTH_PREFIX_BYTES + payload.len());
266 buf.extend_from_slice(&(payload.len() as u32).to_be_bytes());
267 buf.extend_from_slice(&payload);
268 Ok(buf)
269}
270
271fn validate_frame_for_wire(frame: &Frame) -> Result<(), CodecError> {
297 use crate::error::TerminalScope;
298
299 fn check_id(kind: &str, id: &crate::frame::OperationId) -> Result<(), CodecError> {
300 if id.0.is_empty() {
301 return Err(CodecError::InvalidFields {
302 kind: kind.to_string(),
303 detail: "operation id must be a non-empty string".to_string(),
304 });
305 }
306 Ok(())
307 }
308
309 match frame {
310 Frame::Handshake { version } => {
311 if version.get() == 0 {
312 return Err(CodecError::InvalidFields {
313 kind: "handshake".to_string(),
314 detail: "protocol version 0 does not exist".to_string(),
315 });
316 }
317 }
318 Frame::HandshakeAck { version } => {
319 if version.get() == 0 {
320 return Err(CodecError::InvalidFields {
321 kind: "handshake_ack".to_string(),
322 detail: "protocol version 0 does not exist".to_string(),
323 });
324 }
325 }
326 Frame::Event { .. } => {}
327 Frame::Request { id, .. } => check_id("request", id)?,
328 Frame::Response { id, .. } => check_id("response", id)?,
329 Frame::Cancel { id } => check_id("cancel", id)?,
330 Frame::Subscribe { id, .. } => check_id("subscribe", id)?,
331 Frame::SubscribeAck { id, .. } => check_id("subscribe_ack", id)?,
332 Frame::Unsubscribe { id, .. } => check_id("unsubscribe", id)?,
333 Frame::UnsubscribeAck { id, .. } => check_id("unsubscribe_ack", id)?,
334 Frame::Error {
335 id,
336 code,
337 unrecognized_code,
338 ..
339 } => {
340 if let Some(raw_code) = unrecognized_code {
343 return Err(CodecError::FallbackFrameNotEncodable {
344 code: raw_code.clone(),
345 });
346 }
347 if let Some(id) = id {
348 check_id("error", id)?;
349 }
350 match (code.terminal_scope(), id) {
351 (TerminalScope::Connection, Some(id)) => {
352 return Err(CodecError::InconsistentErrorScope {
353 detail: format!(
354 "connection-terminal code {code} must not carry an operation id, got {id}"
355 ),
356 });
357 }
358 (TerminalScope::Request, None) => {
359 return Err(CodecError::InconsistentErrorScope {
360 detail: format!(
361 "request-terminal code {code} must echo the operation id it terminates"
362 ),
363 });
364 }
365 (TerminalScope::Connection, None) | (TerminalScope::Request, Some(_)) => {}
366 }
367 }
368 }
369 Ok(())
370}
371
372fn check_encode_payload_len(payload_len: usize, max_frame_bytes: usize) -> Result<(), CodecError> {
376 if payload_len > max_frame_bytes {
377 return Err(CodecError::FrameTooLarge {
378 declared: payload_len,
379 max: max_frame_bytes,
380 });
381 }
382 if payload_len > u32::MAX as usize {
390 return Err(CodecError::U32PrefixLimitExceeded {
391 declared: payload_len,
392 max: u32::MAX as usize,
393 });
394 }
395 Ok(())
396}
397
398pub fn decode_frame(buf: &[u8], max_frame_bytes: usize) -> Result<Frame, CodecError> {
405 decode_frame_with_consumed(buf, max_frame_bytes).map(|(frame, _)| frame)
406}
407
408pub fn decode_frame_with_consumed(
424 buf: &[u8],
425 max_frame_bytes: usize,
426) -> Result<(Frame, usize), CodecError> {
427 if buf.len() < LENGTH_PREFIX_BYTES {
428 return Err(CodecError::TruncatedLengthPrefix {
429 available: buf.len(),
430 });
431 }
432 let mut len_bytes = [0u8; LENGTH_PREFIX_BYTES];
433 len_bytes.copy_from_slice(&buf[..LENGTH_PREFIX_BYTES]);
434 let declared = u32::from_be_bytes(len_bytes) as usize;
435
436 if declared > max_frame_bytes {
437 return Err(CodecError::FrameTooLarge {
438 declared,
439 max: max_frame_bytes,
440 });
441 }
442
443 let available = buf.len() - LENGTH_PREFIX_BYTES;
444 if available < declared {
445 return Err(CodecError::TruncatedPayload {
446 declared,
447 available,
448 });
449 }
450
451 let payload = &buf[LENGTH_PREFIX_BYTES..LENGTH_PREFIX_BYTES + declared];
452 let frame = decode_payload(payload)?;
453 Ok((frame, LENGTH_PREFIX_BYTES + declared))
454}
455
456pub(crate) fn decode_payload(payload: &[u8]) -> Result<Frame, CodecError> {
479 let value: serde_json::Value =
480 serde_json::from_slice(payload).map_err(|e| CodecError::InvalidJson(e.to_string()))?;
481
482 let kind = value
483 .as_object()
484 .and_then(|obj| obj.get("kind"))
485 .and_then(|k| k.as_str())
486 .ok_or(CodecError::MissingKind)?
487 .to_string();
488
489 if !crate::frame::FRAME_KINDS.contains(&kind.as_str()) {
490 return Err(CodecError::UnknownFrameKind(kind));
491 }
492
493 serde_json::from_value(value).map_err(|e| {
494 let detail = e.to_string();
498 match detail.strip_prefix(INCONSISTENT_SCOPE_ERROR_PREFIX) {
499 Some(scope_detail) => CodecError::InconsistentErrorScope {
500 detail: scope_detail.to_string(),
501 },
502 None => CodecError::InvalidFields { kind, detail },
503 }
504 })
505}
506
507#[cfg(test)]
508mod tests {
509 use super::*;
510 use crate::frame::OperationId;
511
512 fn sample_frames() -> Vec<Frame> {
515 vec![
516 Frame::Handshake {
517 version: crate::version::CURRENT_VERSION,
518 },
519 Frame::HandshakeAck {
520 version: crate::version::CURRENT_VERSION,
521 },
522 Frame::Request {
523 id: OperationId::from("op-1"),
524 ops: "stats()".to_string(),
525 deadline_ms: Some(5000),
526 namespace: Some("research".to_string()),
527 actor_id: Some("lambda".to_string()),
528 visible_namespaces: Some(vec!["research".to_string(), "ops".to_string()]),
529 },
530 Frame::Response {
531 id: OperationId::from("op-1"),
532 result: serde_json::json!({"ok": true, "tool": "stats", "result": {"entities": 3}}),
533 },
534 Frame::Error {
535 id: Some(OperationId::from("op-1")),
536 code: crate::error::WireErrorCode::PeerClassDenied,
537 message: "denied".to_string(),
538 unrecognized_code: None,
539 },
540 Frame::Cancel {
541 id: OperationId::from("op-1"),
542 },
543 Frame::Subscribe {
544 id: OperationId::from("op-2"),
545 topic: "comm.message_created".to_string(),
546 resume_cursor: Some(42),
547 },
548 Frame::SubscribeAck {
549 id: OperationId::from("op-2"),
550 topic: "comm.message_created".to_string(),
551 start_cursor: 42,
552 },
553 Frame::Unsubscribe {
554 id: OperationId::from("op-3"),
555 topic: "comm.message_created".to_string(),
556 },
557 Frame::UnsubscribeAck {
558 id: OperationId::from("op-3"),
559 topic: "comm.message_created".to_string(),
560 },
561 Frame::Event {
562 topic: "comm.message_created".to_string(),
563 cursor: 43,
564 occurred_at: "2026-08-04T11:00:00Z".to_string(),
565 payload: serde_json::json!({"message_id": "m-1"}),
566 },
567 ]
568 }
569
570 #[test]
571 fn round_trips_a_cancel_frame() {
572 let frame = Frame::Cancel {
573 id: OperationId::from("op-1"),
574 };
575 let codec = FrameCodec::default();
576 let wire = codec.encode(&frame).unwrap();
577 assert_eq!(codec.decode(&wire).unwrap(), frame);
578 }
579
580 #[test]
581 fn rejects_truncated_length_prefix() {
582 let codec = FrameCodec::default();
583 assert_eq!(
584 codec.decode(&[0u8, 1]).unwrap_err(),
585 CodecError::TruncatedLengthPrefix { available: 2 }
586 );
587 }
588
589 #[test]
590 fn rejects_truncated_payload_with_declared_and_available() {
591 let declared: u32 = 64;
595 let mut wire = declared.to_be_bytes().to_vec();
596 wire.extend_from_slice(b"part!"); assert_eq!(
598 decode_frame(&wire, DEFAULT_MAX_FRAME_BYTES).unwrap_err(),
599 CodecError::TruncatedPayload {
600 declared: 64,
601 available: 5
602 }
603 );
604 }
605
606 #[test]
607 fn rejects_truncated_payload_when_prefix_declares_exactly_one_byte_more() {
608 let frame = Frame::Cancel {
610 id: OperationId::from("op-1"),
611 };
612 let wire = encode_frame(&frame).unwrap();
613 assert_eq!(
614 decode_frame(&wire[..wire.len() - 1], DEFAULT_MAX_FRAME_BYTES).unwrap_err(),
615 CodecError::TruncatedPayload {
616 declared: wire.len() - LENGTH_PREFIX_BYTES,
617 available: wire.len() - LENGTH_PREFIX_BYTES - 1
618 }
619 );
620 }
621
622 #[test]
623 fn rejects_oversized_frame() {
624 let codec = FrameCodec::new(4);
625 let frame = Frame::Cancel {
626 id: OperationId::from("op-1"),
627 };
628 let wire = encode_frame(&frame).unwrap();
629 assert_eq!(
630 codec.decode(&wire).unwrap_err(),
631 CodecError::FrameTooLarge {
632 declared: wire.len() - LENGTH_PREFIX_BYTES,
633 max: 4
634 }
635 );
636 }
637
638 #[test]
639 fn rejects_unknown_frame_kind() {
640 let payload = br#"{"kind":"ping"}"#;
641 assert_eq!(
642 decode_payload(payload).unwrap_err(),
643 CodecError::UnknownFrameKind("ping".to_string())
644 );
645 }
646
647 #[test]
648 fn rejects_non_json_payload() {
649 let payload = b"not json";
650 match decode_payload(payload).unwrap_err() {
651 CodecError::InvalidJson(_) => {}
652 other => panic!("expected InvalidJson, got {other:?}"),
653 }
654 }
655
656 #[test]
657 fn rejects_zero_length_payload() {
658 let codec = FrameCodec::default();
659 match codec.decode(&0u32.to_be_bytes()).unwrap_err() {
660 CodecError::InvalidJson(_) => {}
661 other => panic!("expected InvalidJson, got {other:?}"),
662 }
663 }
664
665 #[test]
666 fn rejects_missing_required_field() {
667 let payload = br#"{"kind":"cancel"}"#;
668 match decode_payload(payload).unwrap_err() {
669 CodecError::InvalidFields { kind, .. } => assert_eq!(kind, "cancel"),
670 other => panic!("expected InvalidFields, got {other:?}"),
671 }
672 }
673
674 #[test]
675 fn rejects_unknown_top_level_field() {
676 let payload = br#"{"kind":"cancel","id":"op-1","unexpected":true}"#;
679 match decode_payload(payload).unwrap_err() {
680 CodecError::InvalidFields { kind, detail } => {
681 assert_eq!(kind, "cancel");
682 assert!(
683 detail.contains("unexpected"),
684 "detail should name the offending field: {detail}"
685 );
686 }
687 other => panic!("expected InvalidFields, got {other:?}"),
688 }
689 }
690
691 #[test]
692 fn rejects_unknown_field_alongside_optional_fields() {
693 let payload =
696 br#"{"kind":"subscribe","id":"op-2","topic":"a.b","resume_cursor":1,"extra":{"nested":true}}"#;
697 match decode_payload(payload).unwrap_err() {
698 CodecError::InvalidFields { kind, detail } => {
699 assert_eq!(kind, "subscribe");
700 assert!(detail.contains("extra"), "detail: {detail}");
701 }
702 other => panic!("expected InvalidFields, got {other:?}"),
703 }
704 }
705
706 #[test]
707 fn opaque_payload_values_do_not_reject_unknown_keys() {
708 let payload = br#"{"kind":"event","topic":"a.b","cursor":1,"occurred_at":"2026-08-04T11:00:00Z","payload":{"anything":{"goes":true}}}"#;
712 decode_payload(payload).unwrap();
713 }
714
715 #[test]
716 fn every_frame_kind_round_trips_with_strict_decoding() {
717 let frames = sample_frames();
721 assert_eq!(frames.len(), crate::frame::FRAME_KINDS.len());
722 for frame in frames {
723 let wire = encode_frame(&frame).unwrap();
724 let decoded = decode_frame(&wire, DEFAULT_MAX_FRAME_BYTES)
725 .unwrap_or_else(|e| panic!("kind {:?} failed to decode: {e}", frame.kind()));
726 assert_eq!(decoded, frame, "kind {:?} did not round-trip", frame.kind());
727 }
728 }
729
730 #[test]
731 fn rejects_connection_terminal_error_carrying_an_id() {
732 let payload =
735 br#"{"kind":"error","id":"op-1","code":"frame_too_large","message":"too big"}"#;
736 match decode_payload(payload).unwrap_err() {
737 CodecError::InconsistentErrorScope { detail } => {
738 assert!(detail.contains("frame_too_large"), "detail: {detail}");
739 assert!(detail.contains("op-1"), "detail: {detail}");
740 }
741 other => panic!("expected InconsistentErrorScope, got {other:?}"),
742 }
743 }
744
745 #[test]
746 fn rejects_request_terminal_error_without_an_id() {
747 let payload = br#"{"kind":"error","code":"cancelled","message":"cancelled"}"#;
750 match decode_payload(payload).unwrap_err() {
751 CodecError::InconsistentErrorScope { detail } => {
752 assert!(detail.contains("cancelled"), "detail: {detail}");
753 }
754 other => panic!("expected InconsistentErrorScope, got {other:?}"),
755 }
756 }
757
758 #[test]
759 fn accepts_connection_terminal_error_without_an_id() {
760 let payload =
762 br#"{"kind":"error","code":"unsupported_version","message":"no common version"}"#;
763 let frame = decode_payload(payload).unwrap();
764 assert!(matches!(frame, Frame::Error { id: None, .. }));
765 }
766
767 #[test]
768 fn accepts_request_terminal_error_with_an_id() {
769 let payload =
771 br#"{"kind":"error","id":"op-9","code":"deadline_exceeded","message":"too slow"}"#;
772 let frame = decode_payload(payload).unwrap();
773 match frame {
774 Frame::Error { id, code, .. } => {
775 assert_eq!(id, Some(OperationId::from("op-9")));
776 assert_eq!(code, crate::error::WireErrorCode::DeadlineExceeded);
777 }
778 other => panic!("expected an error frame, got {other:?}"),
779 }
780 }
781
782 #[test]
783 fn rejects_empty_operation_id() {
784 let payload = br#"{"kind":"cancel","id":""}"#;
786 match decode_payload(payload).unwrap_err() {
787 CodecError::InvalidFields { kind, detail } => {
788 assert_eq!(kind, "cancel");
789 assert!(detail.contains("non-empty"), "detail: {detail}");
790 }
791 other => panic!("expected InvalidFields, got {other:?}"),
792 }
793 }
794
795 #[test]
796 fn rejects_valid_json_that_is_not_an_object() {
797 for payload in [
798 b"[1, 2, 3]".as_slice(),
799 b"\"cancel\"".as_slice(),
800 b"42".as_slice(),
801 b"null".as_slice(),
802 ] {
803 assert_eq!(
804 decode_payload(payload).unwrap_err(),
805 CodecError::MissingKind,
806 "payload {payload:?}"
807 );
808 }
809 }
810
811 #[test]
812 fn rejects_payload_with_no_kind_field() {
813 let payload = br#"{"id":"op-1"}"#;
814 assert_eq!(
815 decode_payload(payload).unwrap_err(),
816 CodecError::MissingKind
817 );
818 }
819
820 #[test]
821 fn rejects_non_string_kind() {
822 for payload in [
823 br#"{"kind":42}"#.as_slice(),
824 br#"{"kind":null}"#.as_slice(),
825 br#"{"kind":["cancel"]}"#.as_slice(),
826 ] {
827 assert_eq!(
828 decode_payload(payload).unwrap_err(),
829 CodecError::MissingKind,
830 "payload {payload:?}"
831 );
832 }
833 }
834
835 #[test]
836 fn rejects_wrong_typed_field() {
837 let payload = br#"{"kind":"cancel","id":7}"#;
839 match decode_payload(payload).unwrap_err() {
840 CodecError::InvalidFields { kind, detail } => {
841 assert_eq!(kind, "cancel");
842 assert!(detail.contains("id"), "detail: {detail}");
843 }
844 other => panic!("expected InvalidFields, got {other:?}"),
845 }
846
847 let payload = br#"{"kind":"subscribe_ack","id":"op-2","topic":"a.b","start_cursor":"42"}"#;
849 match decode_payload(payload).unwrap_err() {
850 CodecError::InvalidFields { kind, .. } => assert_eq!(kind, "subscribe_ack"),
851 other => panic!("expected InvalidFields, got {other:?}"),
852 }
853 }
854
855 #[test]
856 fn opaque_payloads_are_preserved_semantically_not_byte_for_byte() {
857 let raw = br#"{"kind":"response","id":"op-1","result":{"zeta":1,"alpha":{"n":9007199254740993},"huge":18446744073709551616,"neg":-9007199254740993}}"#;
870 let frame = decode_payload(raw).unwrap();
871 let Frame::Response { result, .. } = &frame else {
872 panic!("expected a response frame");
873 };
874
875 assert_eq!(result["alpha"]["n"], serde_json::json!(9007199254740993u64));
877 assert_eq!(result["neg"], serde_json::json!(-9007199254740993i64));
878
879 assert!(result["huge"].is_f64());
883 assert!(!result["huge"].is_u64());
884 assert_eq!(result["huge"].as_f64().unwrap(), 2.0f64.powi(64));
885
886 let wire = encode_frame(&frame).unwrap();
889 assert_eq!(decode_payload(&wire[LENGTH_PREFIX_BYTES..]).unwrap(), frame);
890 let payload = std::str::from_utf8(&wire[LENGTH_PREFIX_BYTES..]).unwrap();
891 let alpha = payload.find("\"alpha\"").unwrap();
892 let zeta = payload.find("\"zeta\"").unwrap();
893 assert!(
894 alpha < zeta,
895 "expected sorted key order in re-encoded payload: {payload}"
896 );
897 }
898
899 #[test]
902 fn encode_rejects_empty_operation_id() {
903 let frame = Frame::Cancel {
907 id: OperationId::from(""),
908 };
909 match encode_frame(&frame).unwrap_err() {
910 CodecError::InvalidFields { kind, detail } => {
911 assert_eq!(kind, "cancel");
912 assert!(detail.contains("non-empty"), "detail: {detail}");
913 }
914 other => panic!("expected InvalidFields, got {other:?}"),
915 }
916 }
917
918 #[test]
919 fn encode_rejects_empty_operation_id_in_every_id_field() {
920 let frames = [
923 Frame::Request {
924 id: OperationId::from(""),
925 ops: "stats()".to_string(),
926 deadline_ms: None,
927 namespace: None,
928 actor_id: None,
929 visible_namespaces: None,
930 },
931 Frame::Response {
932 id: OperationId::from(""),
933 result: serde_json::json!({}),
934 },
935 Frame::Subscribe {
936 id: OperationId::from(""),
937 topic: "a.b".to_string(),
938 resume_cursor: None,
939 },
940 Frame::SubscribeAck {
941 id: OperationId::from(""),
942 topic: "a.b".to_string(),
943 start_cursor: 0,
944 },
945 Frame::Unsubscribe {
946 id: OperationId::from(""),
947 topic: "a.b".to_string(),
948 },
949 Frame::UnsubscribeAck {
950 id: OperationId::from(""),
951 topic: "a.b".to_string(),
952 },
953 Frame::Error {
954 id: Some(OperationId::from("")),
955 code: crate::error::WireErrorCode::Internal,
956 message: "failure".to_string(),
957 unrecognized_code: None,
958 },
959 ];
960 for frame in frames {
961 match encode_frame(&frame).unwrap_err() {
962 CodecError::InvalidFields { detail, .. } => {
963 assert!(detail.contains("non-empty"), "detail: {detail}");
964 }
965 other => panic!(
966 "kind {:?}: expected InvalidFields, got {other:?}",
967 frame.kind()
968 ),
969 }
970 }
971 }
972
973 #[test]
974 fn encode_rejects_zero_protocol_version_in_both_handshake_kinds() {
975 for frame in [
976 Frame::Handshake {
977 version: crate::version::ProtocolVersion::new(0),
978 },
979 Frame::HandshakeAck {
980 version: crate::version::ProtocolVersion::new(0),
981 },
982 ] {
983 match encode_frame(&frame).unwrap_err() {
984 CodecError::InvalidFields { kind, detail } => {
985 assert_eq!(kind, frame.kind());
986 assert!(detail.contains("version 0"), "detail: {detail}");
987 }
988 other => panic!(
989 "kind {:?}: expected InvalidFields, got {other:?}",
990 frame.kind()
991 ),
992 }
993 }
994 }
995
996 #[test]
997 fn encode_rejects_connection_terminal_error_carrying_an_id() {
998 let frame = Frame::Error {
1002 id: Some(OperationId::from("op-1")),
1003 code: crate::error::WireErrorCode::FrameTooLarge,
1004 message: "too big".to_string(),
1005 unrecognized_code: None,
1006 };
1007 match encode_frame(&frame).unwrap_err() {
1008 CodecError::InconsistentErrorScope { detail } => {
1009 assert!(detail.contains("frame_too_large"), "detail: {detail}");
1010 assert!(detail.contains("op-1"), "detail: {detail}");
1011 }
1012 other => panic!("expected InconsistentErrorScope, got {other:?}"),
1013 }
1014 }
1015
1016 #[test]
1017 fn encode_rejects_request_terminal_error_without_an_id() {
1018 let frame = Frame::Error {
1022 id: None,
1023 code: crate::error::WireErrorCode::Cancelled,
1024 message: "cancelled".to_string(),
1025 unrecognized_code: None,
1026 };
1027 match encode_frame(&frame).unwrap_err() {
1028 CodecError::InconsistentErrorScope { detail } => {
1029 assert!(detail.contains("cancelled"), "detail: {detail}");
1030 }
1031 other => panic!("expected InconsistentErrorScope, got {other:?}"),
1032 }
1033 }
1034
1035 #[test]
1036 fn encode_accepts_both_consistent_error_scopes() {
1037 let consistent = [
1040 Frame::Error {
1041 id: None,
1042 code: crate::error::WireErrorCode::MalformedFrame,
1043 message: "connection scope".to_string(),
1044 unrecognized_code: None,
1045 },
1046 Frame::Error {
1047 id: Some(OperationId::from("op-1")),
1048 code: crate::error::WireErrorCode::DeadlineExceeded,
1049 message: "request scope".to_string(),
1050 unrecognized_code: None,
1051 },
1052 ];
1053 for frame in consistent {
1054 let wire = encode_frame(&frame).expect("consistent scope must encode");
1055 assert_eq!(decode_frame(&wire, DEFAULT_MAX_FRAME_BYTES).unwrap(), frame);
1056 }
1057 }
1058
1059 #[test]
1062 fn unknown_code_fallback_preserves_the_raw_string_with_an_id() {
1063 let payload =
1069 br#"{"kind":"error","id":"op-7","code":"future_code_xyz","message":"from newer peer"}"#;
1070 let frame = decode_payload(payload).unwrap();
1071 match frame {
1072 Frame::Error {
1073 id,
1074 code,
1075 unrecognized_code,
1076 ..
1077 } => {
1078 assert_eq!(code, crate::error::WireErrorCode::Internal);
1079 assert_eq!(id, Some(OperationId::from("op-7")));
1080 assert_eq!(unrecognized_code.as_deref(), Some("future_code_xyz"));
1081 }
1082 other => panic!("expected an error frame, got {other:?}"),
1083 }
1084 }
1085
1086 #[test]
1087 fn unknown_code_fallback_preserves_the_raw_string_without_an_id() {
1088 let payload = br#"{"kind":"error","code":"future_code_xyz","message":"from newer peer"}"#;
1089 match decode_payload(payload).unwrap() {
1090 Frame::Error {
1091 id,
1092 code,
1093 unrecognized_code,
1094 ..
1095 } => {
1096 assert_eq!(code, crate::error::WireErrorCode::Internal);
1097 assert!(id.is_none());
1098 assert_eq!(unrecognized_code.as_deref(), Some("future_code_xyz"));
1099 }
1100 other => panic!("expected an error frame, got {other:?}"),
1101 }
1102 }
1103
1104 #[test]
1105 fn closed_set_code_never_populates_unrecognized_code() {
1106 let payload =
1109 br#"{"kind":"error","id":"op-9","code":"deadline_exceeded","message":"too slow"}"#;
1110 match decode_payload(payload).unwrap() {
1111 Frame::Error {
1112 unrecognized_code, ..
1113 } => assert!(unrecognized_code.is_none()),
1114 other => panic!("expected an error frame, got {other:?}"),
1115 }
1116 }
1117
1118 #[test]
1119 fn decode_scope_rejection_detail_is_reclassified_without_the_shared_prefix() {
1120 let payload =
1124 br#"{"kind":"error","id":"op-1","code":"frame_too_large","message":"too big"}"#;
1125 match decode_payload(payload).unwrap_err() {
1126 CodecError::InconsistentErrorScope { detail } => {
1127 assert!(
1128 !detail.contains(INCONSISTENT_SCOPE_ERROR_PREFIX),
1129 "detail must not repeat the shared prefix: {detail}"
1130 );
1131 assert!(detail.starts_with("connection-terminal code"));
1132 }
1133 other => panic!("expected InconsistentErrorScope, got {other:?}"),
1134 }
1135 }
1136
1137 #[test]
1140 fn encode_rejects_a_fallback_frame_with_an_id() {
1141 let payload =
1147 br#"{"kind":"error","id":"op-7","code":"future_code_xyz","message":"from newer peer"}"#;
1148 let frame = decode_payload(payload).unwrap();
1149 match encode_frame(&frame).unwrap_err() {
1150 CodecError::FallbackFrameNotEncodable { code } => {
1151 assert_eq!(code, "future_code_xyz");
1152 }
1153 other => panic!("expected FallbackFrameNotEncodable, got {other:?}"),
1154 }
1155 }
1156
1157 #[test]
1158 fn encode_rejects_a_fallback_frame_without_an_id() {
1159 let payload = br#"{"kind":"error","code":"future_code_xyz","message":"from newer peer"}"#;
1162 let frame = decode_payload(payload).unwrap();
1163 match encode_frame(&frame).unwrap_err() {
1164 CodecError::FallbackFrameNotEncodable { code } => {
1165 assert_eq!(code, "future_code_xyz");
1166 }
1167 other => panic!("expected FallbackFrameNotEncodable, got {other:?}"),
1168 }
1169 }
1170
1171 #[test]
1172 fn encode_accepts_an_honest_internal_error_frame() {
1173 let frame = Frame::Error {
1177 id: Some(OperationId::from("op-1")),
1178 code: crate::error::WireErrorCode::Internal,
1179 message: "boom".to_string(),
1180 unrecognized_code: None,
1181 };
1182 let wire = encode_frame(&frame).expect("honest internal must encode");
1183 let decoded = decode_frame(&wire, DEFAULT_MAX_FRAME_BYTES).unwrap();
1184 assert_eq!(decoded, frame);
1185 }
1186
1187 #[test]
1190 fn decode_with_consumed_reports_consumed_length_on_a_two_frame_buffer() {
1191 let first = Frame::Cancel {
1195 id: OperationId::from("op-1"),
1196 };
1197 let second = Frame::Handshake {
1198 version: crate::version::CURRENT_VERSION,
1199 };
1200 let codec = FrameCodec::default();
1201 let mut buf = codec.encode(&first).unwrap();
1202 let first_len = buf.len();
1203 buf.extend_from_slice(&codec.encode(&second).unwrap());
1204
1205 let (decoded, consumed) = codec.decode_with_consumed(&buf).unwrap();
1206 assert_eq!(decoded, first);
1207 assert_eq!(consumed, first_len);
1208
1209 let (decoded2, consumed2) = codec.decode_with_consumed(&buf[consumed..]).unwrap();
1210 assert_eq!(decoded2, second);
1211 assert_eq!(consumed + consumed2, buf.len());
1212
1213 assert_eq!(codec.decode(&buf).unwrap(), first);
1215 }
1216
1217 #[test]
1218 fn decode_with_consumed_errors_without_reporting_length_on_truncation() {
1219 let codec = FrameCodec::default();
1220 let wire = codec
1221 .encode(&Frame::Cancel {
1222 id: OperationId::from("op-1"),
1223 })
1224 .unwrap();
1225 let err = codec
1226 .decode_with_consumed(&wire[..wire.len() - 1])
1227 .unwrap_err();
1228 assert!(matches!(err, CodecError::TruncatedPayload { .. }));
1229 }
1230
1231 #[test]
1234 fn serde_json_feature_posture_is_default_keys_serialize_sorted() {
1235 let mut map = serde_json::Map::new();
1248 map.insert("zeta".to_string(), serde_json::json!(1));
1249 map.insert("alpha".to_string(), serde_json::json!(2));
1250 map.insert("mid".to_string(), serde_json::json!(3));
1251 let wire = serde_json::to_string(&serde_json::Value::Object(map)).unwrap();
1252 assert_eq!(
1253 wire, r#"{"alpha":2,"mid":3,"zeta":1}"#,
1254 "serde_json keys did not serialize in sorted order — \
1255 `preserve_order` may have been enabled workspace-wide"
1256 );
1257 }
1258
1259 #[test]
1262 fn codec_errors_map_to_their_wire_error_codes() {
1263 use crate::error::WireErrorCode;
1264
1265 let cases: &[(CodecError, WireErrorCode)] = &[
1266 (
1267 CodecError::TruncatedLengthPrefix { available: 2 },
1268 WireErrorCode::MalformedFrame,
1269 ),
1270 (
1271 CodecError::TruncatedPayload {
1272 declared: 8,
1273 available: 3,
1274 },
1275 WireErrorCode::MalformedFrame,
1276 ),
1277 (
1278 CodecError::FrameTooLarge {
1279 declared: 10,
1280 max: 4,
1281 },
1282 WireErrorCode::FrameTooLarge,
1283 ),
1284 (
1285 CodecError::U32PrefixLimitExceeded {
1286 declared: 10,
1287 max: 4,
1288 },
1289 WireErrorCode::FrameTooLarge,
1290 ),
1291 (
1292 CodecError::InvalidJson("x".to_string()),
1293 WireErrorCode::MalformedFrame,
1294 ),
1295 (CodecError::MissingKind, WireErrorCode::MalformedFrame),
1296 (
1297 CodecError::UnknownFrameKind("ping".to_string()),
1298 WireErrorCode::MalformedFrame,
1299 ),
1300 (
1301 CodecError::InvalidFields {
1302 kind: "cancel".to_string(),
1303 detail: "x".to_string(),
1304 },
1305 WireErrorCode::MalformedFrame,
1306 ),
1307 (
1308 CodecError::InconsistentErrorScope {
1309 detail: "x".to_string(),
1310 },
1311 WireErrorCode::MalformedFrame,
1312 ),
1313 (
1314 CodecError::FallbackFrameNotEncodable {
1315 code: "future_code_xyz".to_string(),
1316 },
1317 WireErrorCode::Internal,
1318 ),
1319 ];
1320 for (err, expected) in cases {
1321 assert_eq!(
1322 err.wire_code(),
1323 *expected,
1324 "{err:?} mapped to {:?}, expected {expected:?}",
1325 err.wire_code()
1326 );
1327 assert_eq!(
1328 WireErrorCode::from(err),
1329 *expected,
1330 "From<&CodecError> disagrees"
1331 );
1332 }
1333 }
1334}
1335
1336#[cfg(test)]
1337mod configured_max_tests {
1338 use super::*;
1339 use crate::frame::Frame;
1340
1341 fn oversized_frame() -> Frame {
1342 Frame::Request {
1343 id: crate::frame::OperationId("x".into()),
1344 ops: "a".repeat(4096),
1345 deadline_ms: None,
1346 namespace: None,
1347 actor_id: None,
1348 visible_namespaces: None,
1349 }
1350 }
1351
1352 #[test]
1353 fn codec_encode_honors_its_own_maximum_not_the_default() {
1354 let codec = FrameCodec::new(1024);
1355 let err = codec.encode(&oversized_frame()).unwrap_err();
1356 match err {
1357 CodecError::FrameTooLarge { max, .. } => assert_eq!(max, 1024),
1358 other => panic!("expected FrameTooLarge with the configured max, got {other:?}"),
1359 }
1360 }
1361
1362 #[test]
1363 fn codec_encode_accepts_a_frame_within_its_own_maximum() {
1364 let codec = FrameCodec::new(1024 * 1024);
1365 let wire = codec
1366 .encode(&oversized_frame())
1367 .expect("within configured max");
1368 assert_eq!(codec.decode(&wire).unwrap(), oversized_frame());
1369 }
1370
1371 #[test]
1372 fn free_function_encode_still_uses_the_default_maximum() {
1373 let wire = encode_frame(&oversized_frame()).expect("well under 8 MiB");
1374 assert!(wire.len() > 4096);
1375 }
1376
1377 #[test]
1378 fn codec_configured_above_u32_capacity_still_encodes_normal_frames() {
1379 assert!(encode_frame_with_max(&oversized_frame(), usize::MAX).is_ok());
1385 }
1386
1387 #[test]
1388 fn encode_size_guard_reports_the_configured_max_when_it_is_binding() {
1389 let err = check_encode_payload_len(32, 16).unwrap_err();
1390 assert_eq!(
1391 err,
1392 CodecError::FrameTooLarge {
1393 declared: 32,
1394 max: 16
1395 }
1396 );
1397 }
1398
1399 #[cfg(target_pointer_width = "64")]
1400 #[test]
1401 fn encode_size_guard_reports_the_u32_prefix_capacity_when_it_is_binding() {
1402 let err = check_encode_payload_len(u32::MAX as usize + 1, usize::MAX).unwrap_err();
1409 assert_eq!(
1410 err,
1411 CodecError::U32PrefixLimitExceeded {
1412 declared: u32::MAX as usize + 1,
1413 max: u32::MAX as usize
1414 }
1415 );
1416 assert!(check_encode_payload_len(u32::MAX as usize, usize::MAX).is_ok());
1418 }
1419}