feat(cluster): RFC 010 c2 — owned wire envelope
Frame enum per the RFC inventory; hand-rolled encode/decode with u32 LE length prefix + u8 tag; strings u16-prefixed, payload blobs u32-prefixed; MAX_FRAME_LEN cap (control plane never carries bulk, §5). Streaming decode: Ok(None) = need more bytes, every Err = corruption. postcard confined to encode_payload/decode_payload — the single codec seam (§2); postcard gains the alloc feature for to_allocvec (still no_std-aligned, no default features). DownReason/RejectReason travel as single tag bytes; Down tag 5 is reserved for c11's Disconnected. Tests: per-frame roundtrip, back-to-back frames, golden heartbeat bytes, zero-length payload, every-prefix incomplete, unknown frame/enum tags, length prefix lying long (with and without bytes present) and short, truncation mid-string, adversarial lengths (u32::MAX, cap+1, zero), UTF-8 corruption, payload seam roundtrip through a real Send frame.
This commit is contained in:
+4
-1
@@ -51,13 +51,16 @@ cc = "1"
|
||||
libc = "0.2"
|
||||
# RFC 010 §2 — only compiled under `--features cluster`.
|
||||
serde = { version = "1", default-features = false, optional = true }
|
||||
postcard = { version = "1", default-features = false, optional = true }
|
||||
# `alloc` (not `std`): the seam serializes to Vec; postcard stays no_std-aligned.
|
||||
postcard = { version = "1", default-features = false, features = ["alloc"], optional = true }
|
||||
|
||||
[target.'cfg(loom)'.dependencies]
|
||||
loom = "0.7"
|
||||
|
||||
[dev-dependencies]
|
||||
libc = "0.2"
|
||||
# derive + std for cluster envelope tests only; the lib itself never needs them
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
tokio = { version = "1", features = ["rt", "rt-multi-thread", "macros", "sync", "time"] }
|
||||
|
||||
[profile.dev]
|
||||
|
||||
@@ -4,3 +4,5 @@
|
||||
//! trait (c3), and everything above them land in later chunks. This module is
|
||||
//! intentionally empty so the `cluster` feature's default-build invariance is
|
||||
//! reviewable in isolation.
|
||||
|
||||
pub mod envelope;
|
||||
|
||||
@@ -0,0 +1,506 @@
|
||||
//! RFC 010 c2 — the owned wire envelope.
|
||||
//!
|
||||
//! Every control-plane frame is `u32` little-endian length prefix (of tag +
|
||||
//! body), `u8` tag, hand-encoded body. postcard appears in exactly one place:
|
||||
//! the payload blob inside `Send`/`SendNamed`, via [`encode_payload`] /
|
||||
//! [`decode_payload`] — the seam where a codec swap would land (RFC 010 §2).
|
||||
//! Everything else is hand-rolled and wholly owned.
|
||||
//!
|
||||
//! Integers are little-endian. Strings are `u16` length + UTF-8 bytes.
|
||||
//! Payload blobs are `u32` length + bytes. Enum-shaped fields
|
||||
//! ([`RejectReason`], [`DownReason`]) are a single tag byte.
|
||||
|
||||
use crate::monitor::DownReason;
|
||||
use crate::pg::Incarnation;
|
||||
|
||||
/// Wire protocol version, checked in the handshake (c5).
|
||||
pub const PROTO_VERSION: u32 = 1;
|
||||
|
||||
/// Hard cap on the length prefix. The control plane never carries bulk data
|
||||
/// (RFC 010 §5 — that is the jarred rkyv plane), so anything larger is
|
||||
/// corruption or an attack, not a legitimate frame.
|
||||
pub const MAX_FRAME_LEN: usize = 16 * 1024 * 1024;
|
||||
|
||||
/// Per-node metadata exchanged in the handshake (RFC 010 §1: not identity).
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct NodeMeta {
|
||||
pub role: String,
|
||||
pub region: String,
|
||||
}
|
||||
|
||||
/// Why a `Hello` was rejected.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum RejectReason {
|
||||
/// Build hashes differ — not the same binary.
|
||||
HashMismatch,
|
||||
/// The offered node name is already claimed by a live peer.
|
||||
NameTaken,
|
||||
/// Wire protocol version mismatch.
|
||||
ProtoVersion,
|
||||
}
|
||||
|
||||
/// The control-plane frame inventory (RFC 010, *Implementation details*).
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub enum Frame {
|
||||
Hello {
|
||||
proto_version: u32,
|
||||
build_hash: u64,
|
||||
node_name: String,
|
||||
incarnation: Incarnation,
|
||||
meta: NodeMeta,
|
||||
},
|
||||
HelloAck {
|
||||
node_name: String,
|
||||
incarnation: Incarnation,
|
||||
meta: NodeMeta,
|
||||
},
|
||||
HelloReject {
|
||||
reason: RejectReason,
|
||||
},
|
||||
Heartbeat,
|
||||
Send {
|
||||
/// Target slot index (node is implicit in the connection, incarnation
|
||||
/// is bound at handshake — RFC 010 §3).
|
||||
index: u32,
|
||||
generation: u32,
|
||||
type_hash: u64,
|
||||
payload: Vec<u8>,
|
||||
},
|
||||
SendNamed {
|
||||
name: String,
|
||||
type_hash: u64,
|
||||
payload: Vec<u8>,
|
||||
},
|
||||
Monitor {
|
||||
monitor_id: u64,
|
||||
index: u32,
|
||||
generation: u32,
|
||||
},
|
||||
Demonitor {
|
||||
monitor_id: u64,
|
||||
},
|
||||
Down {
|
||||
monitor_id: u64,
|
||||
reason: DownReason,
|
||||
},
|
||||
}
|
||||
|
||||
// Frame tags. 0 is deliberately unassigned so an all-zero buffer never parses.
|
||||
const TAG_HELLO: u8 = 1;
|
||||
const TAG_HELLO_ACK: u8 = 2;
|
||||
const TAG_HELLO_REJECT: u8 = 3;
|
||||
const TAG_HEARTBEAT: u8 = 4;
|
||||
const TAG_SEND: u8 = 5;
|
||||
const TAG_SEND_NAMED: u8 = 6;
|
||||
const TAG_MONITOR: u8 = 7;
|
||||
const TAG_DEMONITOR: u8 = 8;
|
||||
const TAG_DOWN: u8 = 9;
|
||||
|
||||
// RejectReason tags.
|
||||
const REJ_HASH_MISMATCH: u8 = 1;
|
||||
const REJ_NAME_TAKEN: u8 = 2;
|
||||
const REJ_PROTO_VERSION: u8 = 3;
|
||||
|
||||
// DownReason tags. c11 adds `Disconnected = 5`; do not reuse tags.
|
||||
const DR_EXIT: u8 = 1;
|
||||
const DR_PANIC: u8 = 2;
|
||||
const DR_STOPPED: u8 = 3;
|
||||
const DR_NOPROC: u8 = 4;
|
||||
|
||||
/// Frame could not be encoded. The output buffer is left exactly as it was.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum EncodeError {
|
||||
/// tag + body exceed [`MAX_FRAME_LEN`].
|
||||
FrameTooLarge { len: usize },
|
||||
/// A string field exceeds `u16::MAX` bytes.
|
||||
StringTooLong { len: usize },
|
||||
}
|
||||
|
||||
impl core::fmt::Display for EncodeError {
|
||||
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
|
||||
match self {
|
||||
Self::FrameTooLarge { len } => {
|
||||
write!(f, "frame body of {len} bytes exceeds MAX_FRAME_LEN")
|
||||
}
|
||||
Self::StringTooLong { len } => {
|
||||
write!(f, "string field of {len} bytes exceeds u16::MAX")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl std::error::Error for EncodeError {}
|
||||
|
||||
/// Frame could not be decoded. Everything here is *corruption* — "not enough
|
||||
/// bytes yet" is the `Ok(None)` streaming case, never an error.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum DecodeError {
|
||||
/// The length prefix exceeds [`MAX_FRAME_LEN`].
|
||||
FrameTooLarge { declared: usize },
|
||||
/// The length prefix is zero — there is no tag byte.
|
||||
EmptyFrame,
|
||||
/// Unknown frame tag.
|
||||
UnknownTag(u8),
|
||||
/// Unknown tag for an enum-shaped field.
|
||||
UnknownEnumTag { what: &'static str, tag: u8 },
|
||||
/// A field ran past the declared frame end (the length prefix lied long,
|
||||
/// or a length-carrying field inside the body lied).
|
||||
Truncated,
|
||||
/// Bytes were left over after the body (the length prefix lied short).
|
||||
Trailing { extra: usize },
|
||||
/// A string field was not valid UTF-8.
|
||||
Utf8,
|
||||
}
|
||||
|
||||
impl core::fmt::Display for DecodeError {
|
||||
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
|
||||
match self {
|
||||
Self::FrameTooLarge { declared } => {
|
||||
write!(f, "declared frame length {declared} exceeds MAX_FRAME_LEN")
|
||||
}
|
||||
Self::EmptyFrame => write!(f, "zero-length frame (no tag byte)"),
|
||||
Self::UnknownTag(t) => write!(f, "unknown frame tag {t}"),
|
||||
Self::UnknownEnumTag { what, tag } => write!(f, "unknown {what} tag {tag}"),
|
||||
Self::Truncated => write!(f, "frame body truncated mid-field"),
|
||||
Self::Trailing { extra } => write!(f, "{extra} trailing bytes after frame body"),
|
||||
Self::Utf8 => write!(f, "string field is not valid UTF-8"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl std::error::Error for DecodeError {}
|
||||
|
||||
impl Frame {
|
||||
/// Append this frame, length-prefixed, to `out`.
|
||||
///
|
||||
/// On error `out` is left untouched.
|
||||
pub fn encode(&self, out: &mut Vec<u8>) -> Result<(), EncodeError> {
|
||||
let start = out.len();
|
||||
out.extend_from_slice(&[0u8; 4]); // length placeholder, patched below
|
||||
let result = self.encode_body(out);
|
||||
match result {
|
||||
Ok(()) => {
|
||||
let frame_len = out.len() - start - 4;
|
||||
if frame_len > MAX_FRAME_LEN {
|
||||
out.truncate(start);
|
||||
return Err(EncodeError::FrameTooLarge { len: frame_len });
|
||||
}
|
||||
// Cast is lossless: MAX_FRAME_LEN < u32::MAX, checked above.
|
||||
let len32 = frame_len as u32;
|
||||
out[start..start + 4].copy_from_slice(&len32.to_le_bytes());
|
||||
Ok(())
|
||||
}
|
||||
Err(e) => {
|
||||
out.truncate(start);
|
||||
Err(e)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn encode_body(&self, out: &mut Vec<u8>) -> Result<(), EncodeError> {
|
||||
match self {
|
||||
Frame::Hello {
|
||||
proto_version,
|
||||
build_hash,
|
||||
node_name,
|
||||
incarnation,
|
||||
meta,
|
||||
} => {
|
||||
out.push(TAG_HELLO);
|
||||
put_u32(out, *proto_version);
|
||||
put_u64(out, *build_hash);
|
||||
put_str(out, node_name)?;
|
||||
put_u32(out, incarnation.get());
|
||||
put_meta(out, meta)?;
|
||||
}
|
||||
Frame::HelloAck {
|
||||
node_name,
|
||||
incarnation,
|
||||
meta,
|
||||
} => {
|
||||
out.push(TAG_HELLO_ACK);
|
||||
put_str(out, node_name)?;
|
||||
put_u32(out, incarnation.get());
|
||||
put_meta(out, meta)?;
|
||||
}
|
||||
Frame::HelloReject { reason } => {
|
||||
out.push(TAG_HELLO_REJECT);
|
||||
out.push(match reason {
|
||||
RejectReason::HashMismatch => REJ_HASH_MISMATCH,
|
||||
RejectReason::NameTaken => REJ_NAME_TAKEN,
|
||||
RejectReason::ProtoVersion => REJ_PROTO_VERSION,
|
||||
});
|
||||
}
|
||||
Frame::Heartbeat => out.push(TAG_HEARTBEAT),
|
||||
Frame::Send {
|
||||
index,
|
||||
generation,
|
||||
type_hash,
|
||||
payload,
|
||||
} => {
|
||||
out.push(TAG_SEND);
|
||||
put_u32(out, *index);
|
||||
put_u32(out, *generation);
|
||||
put_u64(out, *type_hash);
|
||||
put_blob(out, payload)?;
|
||||
}
|
||||
Frame::SendNamed {
|
||||
name,
|
||||
type_hash,
|
||||
payload,
|
||||
} => {
|
||||
out.push(TAG_SEND_NAMED);
|
||||
put_str(out, name)?;
|
||||
put_u64(out, *type_hash);
|
||||
put_blob(out, payload)?;
|
||||
}
|
||||
Frame::Monitor {
|
||||
monitor_id,
|
||||
index,
|
||||
generation,
|
||||
} => {
|
||||
out.push(TAG_MONITOR);
|
||||
put_u64(out, *monitor_id);
|
||||
put_u32(out, *index);
|
||||
put_u32(out, *generation);
|
||||
}
|
||||
Frame::Demonitor { monitor_id } => {
|
||||
out.push(TAG_DEMONITOR);
|
||||
put_u64(out, *monitor_id);
|
||||
}
|
||||
Frame::Down { monitor_id, reason } => {
|
||||
out.push(TAG_DOWN);
|
||||
put_u64(out, *monitor_id);
|
||||
out.push(match reason {
|
||||
DownReason::Exit => DR_EXIT,
|
||||
DownReason::Panic => DR_PANIC,
|
||||
DownReason::Stopped => DR_STOPPED,
|
||||
DownReason::NoProc => DR_NOPROC,
|
||||
});
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Try to decode one frame from the start of `buf`.
|
||||
///
|
||||
/// `Ok(Some((frame, consumed)))` — a full frame; the caller advances by
|
||||
/// `consumed`. `Ok(None)` — not enough bytes yet (streaming); read more
|
||||
/// and retry. `Err(_)` — the bytes are corrupt; the connection is dead.
|
||||
pub fn decode(buf: &[u8]) -> Result<Option<(Frame, usize)>, DecodeError> {
|
||||
let Some(prefix) = buf.get(0..4) else {
|
||||
return Ok(None);
|
||||
};
|
||||
let mut len4 = [0u8; 4];
|
||||
len4.copy_from_slice(prefix);
|
||||
let declared = u32::from_le_bytes(len4) as usize;
|
||||
if declared > MAX_FRAME_LEN {
|
||||
return Err(DecodeError::FrameTooLarge { declared });
|
||||
}
|
||||
if declared == 0 {
|
||||
return Err(DecodeError::EmptyFrame);
|
||||
}
|
||||
let Some(body) = buf.get(4..4 + declared) else {
|
||||
return Ok(None);
|
||||
};
|
||||
let mut r = Reader { buf: body, pos: 0 };
|
||||
let frame = Self::decode_body(&mut r)?;
|
||||
if r.pos != body.len() {
|
||||
return Err(DecodeError::Trailing {
|
||||
extra: body.len() - r.pos,
|
||||
});
|
||||
}
|
||||
Ok(Some((frame, 4 + declared)))
|
||||
}
|
||||
|
||||
fn decode_body(r: &mut Reader<'_>) -> Result<Frame, DecodeError> {
|
||||
let tag = r.u8()?;
|
||||
let frame = match tag {
|
||||
TAG_HELLO => Frame::Hello {
|
||||
proto_version: r.u32()?,
|
||||
build_hash: r.u64()?,
|
||||
node_name: r.string()?,
|
||||
incarnation: Incarnation::new(r.u32()?),
|
||||
meta: r.meta()?,
|
||||
},
|
||||
TAG_HELLO_ACK => Frame::HelloAck {
|
||||
node_name: r.string()?,
|
||||
incarnation: Incarnation::new(r.u32()?),
|
||||
meta: r.meta()?,
|
||||
},
|
||||
TAG_HELLO_REJECT => Frame::HelloReject {
|
||||
reason: match r.u8()? {
|
||||
REJ_HASH_MISMATCH => RejectReason::HashMismatch,
|
||||
REJ_NAME_TAKEN => RejectReason::NameTaken,
|
||||
REJ_PROTO_VERSION => RejectReason::ProtoVersion,
|
||||
t => {
|
||||
return Err(DecodeError::UnknownEnumTag {
|
||||
what: "RejectReason",
|
||||
tag: t,
|
||||
})
|
||||
}
|
||||
},
|
||||
},
|
||||
TAG_HEARTBEAT => Frame::Heartbeat,
|
||||
TAG_SEND => Frame::Send {
|
||||
index: r.u32()?,
|
||||
generation: r.u32()?,
|
||||
type_hash: r.u64()?,
|
||||
payload: r.blob()?,
|
||||
},
|
||||
TAG_SEND_NAMED => Frame::SendNamed {
|
||||
name: r.string()?,
|
||||
type_hash: r.u64()?,
|
||||
payload: r.blob()?,
|
||||
},
|
||||
TAG_MONITOR => Frame::Monitor {
|
||||
monitor_id: r.u64()?,
|
||||
index: r.u32()?,
|
||||
generation: r.u32()?,
|
||||
},
|
||||
TAG_DEMONITOR => Frame::Demonitor {
|
||||
monitor_id: r.u64()?,
|
||||
},
|
||||
TAG_DOWN => Frame::Down {
|
||||
monitor_id: r.u64()?,
|
||||
reason: match r.u8()? {
|
||||
DR_EXIT => DownReason::Exit,
|
||||
DR_PANIC => DownReason::Panic,
|
||||
DR_STOPPED => DownReason::Stopped,
|
||||
DR_NOPROC => DownReason::NoProc,
|
||||
t => {
|
||||
return Err(DecodeError::UnknownEnumTag {
|
||||
what: "DownReason",
|
||||
tag: t,
|
||||
})
|
||||
}
|
||||
},
|
||||
},
|
||||
t => return Err(DecodeError::UnknownTag(t)),
|
||||
};
|
||||
Ok(frame)
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Body writers
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
fn put_u32(out: &mut Vec<u8>, v: u32) {
|
||||
out.extend_from_slice(&v.to_le_bytes());
|
||||
}
|
||||
|
||||
fn put_u64(out: &mut Vec<u8>, v: u64) {
|
||||
out.extend_from_slice(&v.to_le_bytes());
|
||||
}
|
||||
|
||||
fn put_str(out: &mut Vec<u8>, s: &str) -> Result<(), EncodeError> {
|
||||
let Ok(len) = u16::try_from(s.len()) else {
|
||||
return Err(EncodeError::StringTooLong { len: s.len() });
|
||||
};
|
||||
out.extend_from_slice(&len.to_le_bytes());
|
||||
out.extend_from_slice(s.as_bytes());
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn put_blob(out: &mut Vec<u8>, b: &[u8]) -> Result<(), EncodeError> {
|
||||
let Ok(len) = u32::try_from(b.len()) else {
|
||||
return Err(EncodeError::FrameTooLarge { len: b.len() });
|
||||
};
|
||||
out.extend_from_slice(&len.to_le_bytes());
|
||||
out.extend_from_slice(b);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn put_meta(out: &mut Vec<u8>, m: &NodeMeta) -> Result<(), EncodeError> {
|
||||
put_str(out, &m.role)?;
|
||||
put_str(out, &m.region)
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Body reader
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
struct Reader<'a> {
|
||||
buf: &'a [u8],
|
||||
pos: usize,
|
||||
}
|
||||
|
||||
impl Reader<'_> {
|
||||
fn take(&mut self, n: usize) -> Result<&[u8], DecodeError> {
|
||||
let end = self.pos.checked_add(n).ok_or(DecodeError::Truncated)?;
|
||||
let s = self.buf.get(self.pos..end).ok_or(DecodeError::Truncated)?;
|
||||
self.pos = end;
|
||||
Ok(s)
|
||||
}
|
||||
|
||||
fn u8(&mut self) -> Result<u8, DecodeError> {
|
||||
Ok(self.take(1)?[0])
|
||||
}
|
||||
|
||||
fn u16(&mut self) -> Result<u16, DecodeError> {
|
||||
let mut b = [0u8; 2];
|
||||
b.copy_from_slice(self.take(2)?);
|
||||
Ok(u16::from_le_bytes(b))
|
||||
}
|
||||
|
||||
fn u32(&mut self) -> Result<u32, DecodeError> {
|
||||
let mut b = [0u8; 4];
|
||||
b.copy_from_slice(self.take(4)?);
|
||||
Ok(u32::from_le_bytes(b))
|
||||
}
|
||||
|
||||
fn u64(&mut self) -> Result<u64, DecodeError> {
|
||||
let mut b = [0u8; 8];
|
||||
b.copy_from_slice(self.take(8)?);
|
||||
Ok(u64::from_le_bytes(b))
|
||||
}
|
||||
|
||||
fn string(&mut self) -> Result<String, DecodeError> {
|
||||
let len = self.u16()? as usize;
|
||||
let bytes = self.take(len)?;
|
||||
match core::str::from_utf8(bytes) {
|
||||
Ok(s) => Ok(s.to_owned()),
|
||||
Err(_) => Err(DecodeError::Utf8),
|
||||
}
|
||||
}
|
||||
|
||||
fn blob(&mut self) -> Result<Vec<u8>, DecodeError> {
|
||||
let len = self.u32()? as usize;
|
||||
Ok(self.take(len)?.to_vec())
|
||||
}
|
||||
|
||||
fn meta(&mut self) -> Result<NodeMeta, DecodeError> {
|
||||
Ok(NodeMeta {
|
||||
role: self.string()?,
|
||||
region: self.string()?,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// The postcard seam (RFC 010 §2) — the ONLY place payload bytes are produced
|
||||
// or consumed. A codec swap lands here and nowhere else.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
/// Payload (de)serialization failed at the codec seam.
|
||||
#[derive(Debug)]
|
||||
pub struct PayloadError(String);
|
||||
|
||||
impl core::fmt::Display for PayloadError {
|
||||
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
|
||||
write!(f, "payload codec: {}", self.0)
|
||||
}
|
||||
}
|
||||
|
||||
impl std::error::Error for PayloadError {}
|
||||
|
||||
/// Serialize a payload value to the wire blob.
|
||||
pub fn encode_payload<T: serde::Serialize + ?Sized>(value: &T) -> Result<Vec<u8>, PayloadError> {
|
||||
postcard::to_allocvec(value).map_err(|e| PayloadError(e.to_string()))
|
||||
}
|
||||
|
||||
/// Deserialize a payload value from the wire blob.
|
||||
pub fn decode_payload<T: serde::de::DeserializeOwned>(bytes: &[u8]) -> Result<T, PayloadError> {
|
||||
postcard::from_bytes(bytes).map_err(|e| PayloadError(e.to_string()))
|
||||
}
|
||||
@@ -0,0 +1,271 @@
|
||||
//! RFC 010 c2 — owned envelope tests (roadmap: per-frame roundtrip,
|
||||
//! truncation mid-field, unknown tag, length prefix lying long and short,
|
||||
//! zero-length payload, adversarial lengths).
|
||||
#![cfg(feature = "cluster")]
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
use smarm::cluster::envelope::{
|
||||
decode_payload, encode_payload, DecodeError, Frame, NodeMeta, RejectReason, MAX_FRAME_LEN,
|
||||
PROTO_VERSION,
|
||||
};
|
||||
use smarm::monitor::DownReason;
|
||||
use smarm::pg::Incarnation;
|
||||
|
||||
fn meta() -> NodeMeta {
|
||||
NodeMeta {
|
||||
role: "worker".into(),
|
||||
region: "eu-west".into(),
|
||||
}
|
||||
}
|
||||
|
||||
fn all_frames() -> Vec<Frame> {
|
||||
vec![
|
||||
Frame::Hello {
|
||||
proto_version: PROTO_VERSION,
|
||||
build_hash: 0xDEAD_BEEF_CAFE_F00D,
|
||||
node_name: "alpha".into(),
|
||||
incarnation: Incarnation::new(7),
|
||||
meta: meta(),
|
||||
},
|
||||
Frame::HelloAck {
|
||||
node_name: "beta".into(),
|
||||
incarnation: Incarnation::new(9),
|
||||
meta: meta(),
|
||||
},
|
||||
Frame::HelloReject {
|
||||
reason: RejectReason::NameTaken,
|
||||
},
|
||||
Frame::Heartbeat,
|
||||
Frame::Send {
|
||||
index: 42,
|
||||
generation: 3,
|
||||
type_hash: 0x1234_5678_9ABC_DEF0,
|
||||
payload: vec![1, 2, 3, 4, 5],
|
||||
},
|
||||
Frame::SendNamed {
|
||||
name: "the_counter".into(),
|
||||
type_hash: 0xFFFF_0000_FFFF_0000,
|
||||
payload: vec![],
|
||||
},
|
||||
Frame::Monitor {
|
||||
monitor_id: 77,
|
||||
index: 42,
|
||||
generation: 3,
|
||||
},
|
||||
Frame::Demonitor { monitor_id: 77 },
|
||||
Frame::Down {
|
||||
monitor_id: 77,
|
||||
reason: DownReason::Panic,
|
||||
},
|
||||
]
|
||||
}
|
||||
|
||||
fn encode_one(f: &Frame) -> Vec<u8> {
|
||||
let mut buf = Vec::new();
|
||||
f.encode(&mut buf).unwrap();
|
||||
buf
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn per_frame_roundtrip() {
|
||||
for f in all_frames() {
|
||||
let buf = encode_one(&f);
|
||||
let (decoded, consumed) = Frame::decode(&buf).unwrap().unwrap();
|
||||
assert_eq!(decoded, f, "roundtrip mismatch");
|
||||
assert_eq!(consumed, buf.len(), "consumed != buffer length for {f:?}");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn back_to_back_frames_decode_sequentially() {
|
||||
let mut buf = Vec::new();
|
||||
for f in all_frames() {
|
||||
f.encode(&mut buf).unwrap();
|
||||
}
|
||||
let mut off = 0;
|
||||
let mut decoded = Vec::new();
|
||||
while off < buf.len() {
|
||||
let (f, n) = Frame::decode(&buf[off..]).unwrap().unwrap();
|
||||
decoded.push(f);
|
||||
off += n;
|
||||
}
|
||||
assert_eq!(decoded, all_frames());
|
||||
assert_eq!(off, buf.len());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn heartbeat_golden_bytes() {
|
||||
// Locks the layout: u32 LE length prefix, then the tag byte.
|
||||
let buf = encode_one(&Frame::Heartbeat);
|
||||
assert_eq!(buf, vec![1, 0, 0, 0, 4]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn zero_length_payload_roundtrips() {
|
||||
let f = Frame::Send {
|
||||
index: 0,
|
||||
generation: 0,
|
||||
type_hash: 0,
|
||||
payload: vec![],
|
||||
};
|
||||
let buf = encode_one(&f);
|
||||
let (decoded, consumed) = Frame::decode(&buf).unwrap().unwrap();
|
||||
assert_eq!(decoded, f);
|
||||
assert_eq!(consumed, buf.len());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn incomplete_is_none_not_error() {
|
||||
let buf = encode_one(&all_frames()[0]);
|
||||
// Every strict prefix short of the full frame must report "need more".
|
||||
for cut in 0..buf.len() {
|
||||
assert_eq!(
|
||||
Frame::decode(&buf[..cut]).unwrap(),
|
||||
None,
|
||||
"cut at {cut} should be incomplete"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unknown_frame_tag() {
|
||||
let buf = vec![1, 0, 0, 0, 250];
|
||||
assert_eq!(Frame::decode(&buf), Err(DecodeError::UnknownTag(250)));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unknown_enum_tags() {
|
||||
// HelloReject with a bogus reason tag.
|
||||
let buf = vec![2, 0, 0, 0, 3, 99];
|
||||
assert_eq!(
|
||||
Frame::decode(&buf),
|
||||
Err(DecodeError::UnknownEnumTag {
|
||||
what: "RejectReason",
|
||||
tag: 99
|
||||
})
|
||||
);
|
||||
// Down with a bogus reason tag (id = 0u64).
|
||||
let mut buf = vec![10, 0, 0, 0, 9];
|
||||
buf.extend_from_slice(&0u64.to_le_bytes());
|
||||
buf.push(200);
|
||||
assert_eq!(
|
||||
Frame::decode(&buf),
|
||||
Err(DecodeError::UnknownEnumTag {
|
||||
what: "DownReason",
|
||||
tag: 200
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn length_prefix_lying_long_with_bytes_present_is_trailing() {
|
||||
let mut buf = encode_one(&Frame::Heartbeat);
|
||||
// Declare 3 extra body bytes and actually supply them.
|
||||
let declared = u32::from_le_bytes([buf[0], buf[1], buf[2], buf[3]]) + 3;
|
||||
buf[0..4].copy_from_slice(&declared.to_le_bytes());
|
||||
buf.extend_from_slice(&[0xAA, 0xBB, 0xCC]);
|
||||
assert_eq!(Frame::decode(&buf), Err(DecodeError::Trailing { extra: 3 }));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn length_prefix_lying_long_without_bytes_is_incomplete() {
|
||||
// Indistinguishable from a partial read — must be None, not an error.
|
||||
let mut buf = encode_one(&Frame::Heartbeat);
|
||||
let declared = u32::from_le_bytes([buf[0], buf[1], buf[2], buf[3]]) + 3;
|
||||
buf[0..4].copy_from_slice(&declared.to_le_bytes());
|
||||
assert_eq!(Frame::decode(&buf).unwrap(), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn length_prefix_lying_short_truncates_a_field() {
|
||||
let f = &all_frames()[0]; // Hello: plenty of fields to cut into
|
||||
let mut buf = encode_one(f);
|
||||
let declared = u32::from_le_bytes([buf[0], buf[1], buf[2], buf[3]]);
|
||||
let lie = declared - 4; // cut mid-field
|
||||
buf[0..4].copy_from_slice(&lie.to_le_bytes());
|
||||
assert_eq!(Frame::decode(&buf), Err(DecodeError::Truncated));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn truncation_mid_string_field() {
|
||||
// A frame whose declared length is intact but whose inner string length
|
||||
// runs past the body: SendNamed claiming a 1000-byte name in a tiny body.
|
||||
let mut body = vec![6u8]; // TAG_SEND_NAMED
|
||||
body.extend_from_slice(&1000u16.to_le_bytes());
|
||||
body.extend_from_slice(b"short");
|
||||
let mut buf = (body.len() as u32).to_le_bytes().to_vec();
|
||||
buf.extend_from_slice(&body);
|
||||
assert_eq!(Frame::decode(&buf), Err(DecodeError::Truncated));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn adversarial_lengths() {
|
||||
// Length prefix of u32::MAX: reject as oversized, do not wait for 4 GiB.
|
||||
let buf = [0xFF, 0xFF, 0xFF, 0xFF, 0];
|
||||
assert_eq!(
|
||||
Frame::decode(&buf),
|
||||
Err(DecodeError::FrameTooLarge {
|
||||
declared: u32::MAX as usize
|
||||
})
|
||||
);
|
||||
// Just over the cap: also rejected.
|
||||
let over = (MAX_FRAME_LEN as u32 + 1).to_le_bytes();
|
||||
assert!(matches!(
|
||||
Frame::decode(&over),
|
||||
Err(DecodeError::FrameTooLarge { .. })
|
||||
));
|
||||
// Zero-length frame: there is no tag byte; corrupt, not incomplete.
|
||||
let buf = [0, 0, 0, 0];
|
||||
assert_eq!(Frame::decode(&buf), Err(DecodeError::EmptyFrame));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn invalid_utf8_in_string_field() {
|
||||
let mut buf = encode_one(&Frame::SendNamed {
|
||||
name: "abcd".into(),
|
||||
type_hash: 0,
|
||||
payload: vec![],
|
||||
});
|
||||
// name bytes start after: 4 (len) + 1 (tag) + 2 (str len) = offset 7
|
||||
buf[7] = 0xFF;
|
||||
assert_eq!(Frame::decode(&buf), Err(DecodeError::Utf8));
|
||||
}
|
||||
|
||||
#[derive(Debug, PartialEq, Serialize, Deserialize)]
|
||||
struct Ping {
|
||||
seq: u64,
|
||||
label: String,
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn payload_seam_roundtrip() {
|
||||
let ping = Ping {
|
||||
seq: 31337,
|
||||
label: "hello".into(),
|
||||
};
|
||||
let blob = encode_payload(&ping).unwrap();
|
||||
// Carry it through a real frame, as it will travel in c9.
|
||||
let f = Frame::Send {
|
||||
index: 1,
|
||||
generation: 1,
|
||||
type_hash: 0xABCD,
|
||||
payload: blob,
|
||||
};
|
||||
let buf = encode_one(&f);
|
||||
let (decoded, _) = Frame::decode(&buf).unwrap().unwrap();
|
||||
let Frame::Send { payload, .. } = decoded else {
|
||||
panic!("wrong frame");
|
||||
};
|
||||
let back: Ping = decode_payload(&payload).unwrap();
|
||||
assert_eq!(back, ping);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn payload_seam_rejects_truncated_blob() {
|
||||
let blob = encode_payload(&Ping {
|
||||
seq: 1,
|
||||
label: "x".into(),
|
||||
})
|
||||
.unwrap();
|
||||
assert!(decode_payload::<Ping>(&blob[..blob.len() - 1]).is_err());
|
||||
}
|
||||
Reference in New Issue
Block a user