use nodedb_query::msgpack_scan;
pub(crate) fn extract_msgpack_elements(payload: &[u8]) -> Vec<Vec<u8>> {
if payload.is_empty() {
return Vec::new();
}
let Some((count, mut pos)) = msgpack_scan::array_header(payload, 0) else {
tracing::warn!(
payload_len = payload.len(),
"payload_merge: payload is not a msgpack array; treating as single element"
);
return vec![payload.to_vec()];
};
let mut rows = Vec::with_capacity(count);
for _ in 0..count {
if pos >= payload.len() {
break;
}
let start = pos;
match msgpack_scan::skip_value(payload, pos) {
Some(next) => {
rows.push(payload[start..next].to_vec());
pos = next;
}
None => {
tracing::warn!(
pos,
payload_len = payload.len(),
"payload_merge: could not skip msgpack element; stopping early"
);
break;
}
}
}
rows
}
pub(crate) fn encode_msgpack_array(rows: &[Vec<u8>]) -> Vec<u8> {
let total_data: usize = rows.iter().map(|r| r.len()).sum();
let mut out = Vec::with_capacity(total_data + 5);
let n = rows.len();
if n < 16 {
out.push(0x90 | n as u8);
} else if n <= u16::MAX as usize {
out.push(0xdc);
out.extend_from_slice(&(n as u16).to_be_bytes());
} else {
out.push(0xdd);
out.extend_from_slice(&(n as u32).to_be_bytes());
}
for row in rows {
out.extend_from_slice(row);
}
out
}
pub(crate) fn merge_msgpack_arrays(payloads: &[Vec<u8>]) -> Vec<u8> {
let mut elements: Vec<Vec<u8>> = Vec::new();
for payload in payloads {
elements.extend(extract_msgpack_elements(payload));
}
encode_msgpack_array(&elements)
}
#[cfg(test)]
mod tests {
use super::*;
fn int_array(start: usize, n: usize) -> Vec<u8> {
let rows: Vec<Vec<u8>> = (start..start + n).map(|i| vec![(i % 128) as u8]).collect();
encode_msgpack_array(&rows)
}
#[test]
fn merge_concatenates_all_elements() {
let chunks = vec![
int_array(0, 1000),
int_array(1000, 1000),
int_array(2000, 500),
];
let merged = merge_msgpack_arrays(&chunks);
let elements = extract_msgpack_elements(&merged);
assert_eq!(
elements.len(),
2500,
"merging three array chunks must yield every element, not just the first chunk"
);
}
#[test]
fn single_array_round_trips() {
let one = int_array(0, 3);
let merged = merge_msgpack_arrays(std::slice::from_ref(&one));
assert_eq!(extract_msgpack_elements(&merged).len(), 3);
}
#[test]
fn empty_input_is_empty_array() {
let merged = merge_msgpack_arrays(&[]);
assert_eq!(extract_msgpack_elements(&merged).len(), 0);
}
}