use crate::{Error, Result};
pub const SUPPORTED_COMMITLOG_VERSION: i32 = 7;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct CommitLogVersionGates {
pub version: i32,
}
impl CommitLogVersionGates {
pub fn from_version(version: i32) -> Result<Self> {
if version != SUPPORTED_COMMITLOG_VERSION {
return Err(Error::UnsupportedCommitLogVersion {
version,
floor: SUPPORTED_COMMITLOG_VERSION,
ceiling: SUPPORTED_COMMITLOG_VERSION,
});
}
Ok(Self { version })
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CommitLogDescriptor {
pub version: i32,
pub id: i64,
pub params_json: String,
pub compression_class: Option<String>,
pub encrypted: bool,
pub header_len: usize,
}
impl CommitLogDescriptor {
pub fn parse(bytes: &[u8]) -> Result<Self> {
if bytes.len() < 18 {
return Err(Error::CorruptCommitLogFrame(format!(
"segment too short for a descriptor header: {} bytes",
bytes.len()
)));
}
let version = read_i32_be(bytes, 0);
let gates = CommitLogVersionGates::from_version(version)?;
let id = read_i64_be(bytes, 4);
let params_len = read_u16_be(bytes, 12) as usize;
let params_start = 14;
let params_end = params_start + params_len;
let crc_end = params_end + 4;
if bytes.len() < crc_end {
return Err(Error::CorruptCommitLogFrame(format!(
"descriptor header claims {} param bytes but segment is only {} bytes",
params_len,
bytes.len()
)));
}
let params_bytes = &bytes[params_start..params_end];
let stored_crc = read_u32_be(bytes, params_end);
let mut hasher = crc32fast::Hasher::new();
hasher.update(&version.to_be_bytes());
let id_u = id as u64;
hasher.update(&((id_u & 0xFFFF_FFFF) as u32).to_be_bytes());
hasher.update(&((id_u >> 32) as u32).to_be_bytes());
hasher.update(&(params_len as u32).to_be_bytes());
hasher.update(params_bytes);
let computed_crc = hasher.finalize();
if computed_crc != stored_crc {
return Err(Error::CorruptCommitLogFrame(format!(
"descriptor CRC mismatch: stored={stored_crc:#010x} computed={computed_crc:#010x}"
)));
}
let params_json = String::from_utf8_lossy(params_bytes).into_owned();
let (compression_class, encrypted) = parse_params(¶ms_json);
Ok(Self {
version: gates.version,
id,
params_json,
compression_class,
encrypted,
header_len: crc_end,
})
}
pub fn is_unsupported_payload(&self) -> bool {
self.compression_class.is_some() || self.encrypted
}
}
fn parse_params(json: &str) -> (Option<String>, bool) {
let value: serde_json::Value = match serde_json::from_str(json) {
Ok(v) => v,
Err(_) => return (None, false),
};
let obj = match value.as_object() {
Some(o) => o,
None => return (None, false),
};
let compression_class = obj
.get("compressionClass")
.and_then(|c| c.as_str())
.filter(|s| !s.is_empty())
.map(|s| s.to_string());
let encrypted = ["encCipher", "encKeyAlias", "encIV"]
.iter()
.any(|k| obj.get(*k).is_some_and(|v| !v.is_null()));
(compression_class, encrypted)
}
#[inline]
fn read_u16_be(b: &[u8], off: usize) -> u16 {
u16::from_be_bytes([b[off], b[off + 1]])
}
#[inline]
fn read_i32_be(b: &[u8], off: usize) -> i32 {
i32::from_be_bytes([b[off], b[off + 1], b[off + 2], b[off + 3]])
}
#[inline]
fn read_u32_be(b: &[u8], off: usize) -> u32 {
u32::from_be_bytes([b[off], b[off + 1], b[off + 2], b[off + 3]])
}
#[inline]
fn read_i64_be(b: &[u8], off: usize) -> i64 {
let mut a = [0u8; 8];
a.copy_from_slice(&b[off..off + 8]);
i64::from_be_bytes(a)
}
#[cfg(test)]
pub(crate) mod tests {
use super::*;
pub(crate) fn build_header(version: i32, id: i64, params: &str) -> Vec<u8> {
let params_bytes = params.as_bytes();
let mut out = Vec::new();
out.extend_from_slice(&version.to_be_bytes());
out.extend_from_slice(&id.to_be_bytes());
out.extend_from_slice(&(params_bytes.len() as u16).to_be_bytes());
out.extend_from_slice(params_bytes);
let mut hasher = crc32fast::Hasher::new();
hasher.update(&version.to_be_bytes());
let id_u = id as u64;
hasher.update(&((id_u & 0xFFFF_FFFF) as u32).to_be_bytes());
hasher.update(&((id_u >> 32) as u32).to_be_bytes());
hasher.update(&(params_bytes.len() as u32).to_be_bytes());
hasher.update(params_bytes);
out.extend_from_slice(&hasher.finalize().to_be_bytes());
out
}
#[test]
fn version_gate_accepts_only_v7() {
assert!(CommitLogVersionGates::from_version(7).is_ok());
for v in [-1, 0, 5, 6, 8, 100] {
assert!(matches!(
CommitLogVersionGates::from_version(v),
Err(Error::UnsupportedCommitLogVersion { .. })
));
}
}
#[test]
fn parses_well_formed_uncompressed_header() {
let id = 1_689_012_345_678_i64;
let bytes = build_header(7, id, "{}");
let desc = CommitLogDescriptor::parse(&bytes).expect("parse");
assert_eq!(desc.version, 7);
assert_eq!(desc.id, id);
assert_eq!(desc.header_len, bytes.len());
assert!(!desc.is_unsupported_payload());
assert!(desc.compression_class.is_none());
}
#[test]
fn rejects_unsupported_version_before_body() {
let bytes = build_header(6, 42, "{}");
assert!(matches!(
CommitLogDescriptor::parse(&bytes),
Err(Error::UnsupportedCommitLogVersion { version: 6, .. })
));
}
#[test]
fn rejects_crc_mismatch() {
let mut bytes = build_header(7, 42, "{}");
let last = bytes.len() - 1;
bytes[last] ^= 0xFF;
assert!(matches!(
CommitLogDescriptor::parse(&bytes),
Err(Error::CorruptCommitLogFrame(_))
));
}
#[test]
fn rejects_short_header() {
assert!(matches!(
CommitLogDescriptor::parse(&[0u8; 4]),
Err(Error::CorruptCommitLogFrame(_))
));
}
#[test]
fn detects_compression_class_from_params() {
let params = r#"{"compressionClass":"LZ4Compressor","compressionParameters":{}}"#;
let bytes = build_header(7, 7, params);
let desc = CommitLogDescriptor::parse(&bytes).expect("parse");
assert_eq!(desc.compression_class.as_deref(), Some("LZ4Compressor"));
assert!(desc.is_unsupported_payload());
}
#[test]
fn compression_null_parses_as_uncompressed_not_unknown() {
let params = r#"{"compressionClass":null}"#;
let bytes = build_header(7, 9, params);
let desc = CommitLogDescriptor::parse(&bytes).expect("parse");
assert_eq!(desc.compression_class, None);
assert!(!desc.is_unsupported_payload());
}
#[test]
fn encryption_key_null_parses_as_unencrypted() {
let params = r#"{"compressionClass":null,"encCipher":null}"#;
let bytes = build_header(7, 11, params);
let desc = CommitLogDescriptor::parse(&bytes).expect("parse");
assert_eq!(desc.compression_class, None);
assert!(!desc.is_unsupported_payload());
}
#[test]
fn detects_encryption_from_real_context_keys() {
let params = r#"{"encCipher":"AES/CBC/PKCS5Padding","encKeyAlias":"testing:1"}"#;
let bytes = build_header(7, 13, params);
let desc = CommitLogDescriptor::parse(&bytes).expect("parse");
assert!(desc.encrypted);
assert!(desc.compression_class.is_none());
assert!(desc.is_unsupported_payload());
}
}