Decode concatenated row payloads into `(doc_id, msgpack_data)` pairs. Also used by the Control Plane sync layer to filter snapshot documents by a shape predicate before sending them to subscribers. Input: zero or more msgpack arrays back-to-back. Elements may be either: - raw scan rows from `encode_raw_document_rows` with `{id, data}` wrappers - plain msgpack rows from aggregate/join paths seria
(bytes: &[u8])
| 55 | /// For wrapped scan rows, the `data` field's raw bytes are extracted. For |
| 56 | /// plain rows, the entire row value is returned as `msgpack_data`. |
| 57 | pub(crate) fn decode_raw_scan_to_docs(bytes: &[u8]) -> Vec<(String, Vec<u8>)> { |
| 58 | use nodedb_query::msgpack_scan; |
| 59 | |
| 60 | let mut results = Vec::new(); |
| 61 | let mut pos = 0; |
| 62 | |
| 63 | while pos < bytes.len() { |
| 64 | let first = bytes[pos]; |
| 65 | let (count, hdr_len) = if (0x90..=0x9f).contains(&first) { |
| 66 | ((first & 0x0f) as usize, 1) |
| 67 | } else if first == 0xdc && pos + 3 <= bytes.len() { |
| 68 | ( |
| 69 | u16::from_be_bytes([bytes[pos + 1], bytes[pos + 2]]) as usize, |
| 70 | 3, |
| 71 | ) |
| 72 | } else if first == 0xdd && pos + 5 <= bytes.len() { |
| 73 | ( |
| 74 | u32::from_be_bytes([ |
| 75 | bytes[pos + 1], |
| 76 | bytes[pos + 2], |
| 77 | bytes[pos + 3], |
| 78 | bytes[pos + 4], |
| 79 | ]) as usize, |
| 80 | 5, |
| 81 | ) |
| 82 | } else { |
| 83 | break; |
| 84 | }; |
| 85 | |
| 86 | let mut inner = pos + hdr_len; |
| 87 | for _ in 0..count { |
| 88 | if inner >= bytes.len() { |
| 89 | break; |
| 90 | } |
| 91 | |
| 92 | let elem_start = inner; |
| 93 | let elem_end = msgpack_scan::skip_value(bytes, inner).unwrap_or(bytes.len()); |
| 94 | |
| 95 | let id = msgpack_scan::extract_field(bytes, elem_start, "id") |
| 96 | .and_then(|(s, _e)| msgpack_scan::read_value(bytes, s)) |
| 97 | .and_then(|v| match v { |
| 98 | nodedb_types::Value::String(s) => Some(s), |
| 99 | _ => None, |
| 100 | }) |
| 101 | .unwrap_or_default(); |
| 102 | |
| 103 | let data = msgpack_scan::extract_field(bytes, elem_start, "data") |
| 104 | .map(|(s, e)| bytes[s..e].to_vec()) |
| 105 | .unwrap_or_else(|| bytes[elem_start..elem_end].to_vec()); |
| 106 | |
| 107 | results.push((id, data)); |
| 108 | |
| 109 | inner = elem_end; |
| 110 | } |
| 111 | pos = inner; |
| 112 | } |
| 113 | |
| 114 | results |