use std::path::Path;
use super::super::Result;
use super::COORD_SLOT_SIZE;
pub(super) enum BlobLocation {
Worker {
worker_id: u32,
byte_offset: u64,
byte_length: u64,
},
Straddler {
bytes: Vec<u8>,
},
Empty,
}
#[derive(Clone, Copy, Debug)]
pub(super) enum StraddlerSide {
Left,
Right,
}
enum StraddlerPartial {
Left(Vec<u8>),
Right(Vec<u8>),
}
pub(super) struct ConcurrentBlobLocationRouter {
worker_files: Vec<std::sync::Arc<std::fs::File>>,
slots: Vec<std::sync::Mutex<Option<BlobLocation>>>,
straddler_partials: std::sync::Mutex<rustc_hash::FxHashMap<usize, StraddlerPartial>>,
notify: std::sync::Condvar,
notify_mu: std::sync::Mutex<()>,
aborted: std::sync::atomic::AtomicBool,
producer_done: std::sync::atomic::AtomicBool,
abort_error: std::sync::Mutex<Option<String>>,
wait_gauge: super::StallGauge,
pub(super) stats: std::sync::Mutex<ConcurrentRouterStats>,
}
#[derive(Default)]
pub(super) struct ConcurrentRouterStats {
pub num_worker: u64,
pub num_straddlers: u64,
pub num_empty: u64,
pub worker_bytes: u64,
pub straddler_bytes: u64,
pub straddler_encode_ns: u64,
}
impl ConcurrentBlobLocationRouter {
pub(super) fn new(
per_way_rcs: &PerWayRcs,
worker_files: Vec<std::sync::Arc<std::fs::File>>,
) -> Result<Self> {
let num_way_blobs = per_way_rcs.num_blobs();
let mut slots: Vec<std::sync::Mutex<Option<BlobLocation>>> =
Vec::with_capacity(num_way_blobs);
let mut num_empty: u64 = 0;
for blob_idx in 0..num_way_blobs {
if per_way_rcs.blob_has_nonzero_refs(blob_idx)? {
slots.push(std::sync::Mutex::new(None));
} else {
slots.push(std::sync::Mutex::new(Some(BlobLocation::Empty)));
num_empty += 1;
}
}
let straddler_partials: std::sync::Mutex<rustc_hash::FxHashMap<usize, StraddlerPartial>> =
std::sync::Mutex::new(rustc_hash::FxHashMap::default());
let stats = ConcurrentRouterStats {
num_empty,
..Default::default()
};
Ok(Self {
worker_files,
slots,
straddler_partials,
notify: std::sync::Condvar::new(),
notify_mu: std::sync::Mutex::new(()),
aborted: std::sync::atomic::AtomicBool::new(false),
producer_done: std::sync::atomic::AtomicBool::new(false),
abort_error: std::sync::Mutex::new(None),
wait_gauge: super::StallGauge::new("WAIT_S4_ROUTER_START", "WAIT_S4_ROUTER_END"),
stats: std::sync::Mutex::new(stats),
})
}
pub(super) fn num_blobs(&self) -> usize {
self.slots.len()
}
pub(super) fn worker_file(&self, worker_id: usize) -> &std::sync::Arc<std::fs::File> {
&self.worker_files[worker_id]
}
pub(super) fn publish_worker(
&self,
blob_idx: usize,
worker_id: u32,
byte_offset: u64,
byte_length: u64,
) -> Result<()> {
let mut guard = self.slots[blob_idx]
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if guard.is_some() {
return Err(format!(
"router publish_worker: blob {blob_idx} already has a location \
(likely duplicate emission across workers)"
)
.into());
}
*guard = Some(BlobLocation::Worker {
worker_id,
byte_offset,
byte_length,
});
drop(guard);
{
let mut s = self
.stats
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
s.num_worker += 1;
s.worker_bytes += byte_length;
}
self.notify.notify_all();
Ok(())
}
pub(super) fn publish_straddler_half(
&self,
blob_idx: usize,
side: StraddlerSide,
raw_bytes: Vec<u8>,
per_way_rcs: &PerWayRcs,
inject_prepass: bool,
encode_scratch: &mut Vec<u8>,
) -> Result<()> {
let mut map = self
.straddler_partials
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let (left_bytes, right_bytes) = match (map.remove(&blob_idx), side) {
(None, StraddlerSide::Left) => {
map.insert(blob_idx, StraddlerPartial::Left(raw_bytes));
return Ok(());
}
(None, StraddlerSide::Right) => {
map.insert(blob_idx, StraddlerPartial::Right(raw_bytes));
return Ok(());
}
(Some(StraddlerPartial::Left(left)), StraddlerSide::Right) => (left, raw_bytes),
(Some(StraddlerPartial::Right(right)), StraddlerSide::Left) => (raw_bytes, right),
(Some(StraddlerPartial::Left(_)), StraddlerSide::Left) => {
return Err(
format!("router straddler blob {blob_idx}: duplicate left half").into(),
);
}
(Some(StraddlerPartial::Right(_)), StraddlerSide::Right) => {
return Err(
format!("router straddler blob {blob_idx}: duplicate right half").into(),
);
}
};
drop(map);
let t_enc = std::time::Instant::now();
let mut coord_bytes = left_bytes;
coord_bytes.extend_from_slice(&right_bytes);
encode_scratch.clear();
encode_blob_payload_from_record(
&coord_bytes,
per_way_rcs.blob_record(blob_idx),
blob_idx,
inject_prepass,
encode_scratch,
)
.map_err(|e| format!("router straddler encode blob {blob_idx}: {e}"))?;
#[allow(clippy::cast_possible_truncation)]
let encode_ns = t_enc.elapsed().as_nanos() as u64;
let bytes = std::mem::take(encode_scratch);
let byte_len = bytes.len() as u64;
let mut slot_guard = self.slots[blob_idx]
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if slot_guard.is_some() {
return Err(format!(
"router straddler blob {blob_idx}: slot already populated at encode time"
)
.into());
}
*slot_guard = Some(BlobLocation::Straddler { bytes });
drop(slot_guard);
{
let mut s = self
.stats
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
s.num_straddlers += 1;
s.straddler_bytes += byte_len;
s.straddler_encode_ns += encode_ns;
}
self.notify.notify_all();
Ok(())
}
pub(super) fn mark_producer_done(&self) {
self.producer_done
.store(true, std::sync::atomic::Ordering::SeqCst);
self.notify.notify_all();
}
pub(super) fn abort(&self, msg: String) {
{
let mut guard = self
.abort_error
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if guard.is_none() {
*guard = Some(msg);
}
}
self.aborted
.store(true, std::sync::atomic::Ordering::SeqCst);
self.notify.notify_all();
}
pub(super) fn is_aborted(&self) -> bool {
self.aborted.load(std::sync::atomic::Ordering::SeqCst)
}
pub(super) fn wait_ready(&self, blob_idx: usize) -> Result<()> {
{
let guard = self.slots[blob_idx]
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if guard.is_some() {
return Ok(());
}
}
let _stall = self.wait_gauge.track();
loop {
let mu_guard = self
.notify_mu
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if self.aborted.load(std::sync::atomic::Ordering::SeqCst) {
drop(mu_guard);
let err = self
.abort_error
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
.unwrap_or_else(|| "router aborted (no message recorded)".to_string());
return Err(err.into());
}
{
let guard = self.slots[blob_idx]
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if guard.is_some() {
return Ok(());
}
if self.producer_done.load(std::sync::atomic::Ordering::SeqCst) {
return Err(format!(
"router: no publication for blob {blob_idx} after producer finished"
)
.into());
}
}
drop(
self.notify
.wait(mu_guard)
.unwrap_or_else(std::sync::PoisonError::into_inner),
);
}
}
pub(super) fn pread_blob_payload(&self, blob_idx: usize, buf: &mut Vec<u8>) -> Result<()> {
self.wait_ready(blob_idx)?;
use std::os::unix::fs::FileExt as _;
let loc = {
let guard = self.slots[blob_idx]
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
match &*guard {
Some(loc) => loc.clone(),
None => {
return Err(format!(
"router pread_blob_payload: slot {blob_idx} empty after wait_ready"
)
.into());
}
}
};
match loc {
BlobLocation::Worker {
worker_id,
byte_offset,
byte_length,
} => {
#[allow(clippy::cast_possible_truncation)]
let len = byte_length as usize;
buf.resize(len, 0);
if len > 0 {
self.worker_files[worker_id as usize]
.read_exact_at(buf, byte_offset)
.map_err(|e| {
format!("router pread worker {worker_id} blob {blob_idx}: {e}")
})?;
}
}
BlobLocation::Straddler { bytes } => {
buf.clear();
buf.extend_from_slice(&bytes);
}
BlobLocation::Empty => {
buf.clear();
}
}
Ok(())
}
}
impl Clone for BlobLocation {
fn clone(&self) -> Self {
match self {
Self::Worker {
worker_id,
byte_offset,
byte_length,
} => Self::Worker {
worker_id: *worker_id,
byte_offset: *byte_offset,
byte_length: *byte_length,
},
Self::Straddler { bytes } => Self::Straddler {
bytes: bytes.clone(),
},
Self::Empty => Self::Empty,
}
}
}
pub(super) struct AbortOnDrop<'a> {
router: &'a ConcurrentBlobLocationRouter,
label: &'static str,
armed: std::cell::Cell<bool>,
}
impl<'a> AbortOnDrop<'a> {
pub(super) fn new(router: &'a ConcurrentBlobLocationRouter, label: &'static str) -> Self {
Self {
router,
label,
armed: std::cell::Cell::new(true),
}
}
pub(super) fn disarm(&self) {
self.armed.set(false);
}
}
impl Drop for AbortOnDrop<'_> {
fn drop(&mut self) {
if self.armed.get() {
self.router
.abort(format!("{} panicked (AbortOnDrop guard fired)", self.label));
}
}
}
pub(super) struct PerWayRcs {
data: Vec<u8>,
offsets: Vec<usize>, }
impl PerWayRcs {
pub(super) fn blob_record(&self, blob_idx: usize) -> &[u8] {
&self.data[self.offsets[blob_idx]..self.offsets[blob_idx + 1]]
}
pub(super) fn num_blobs(&self) -> usize {
self.offsets.len().saturating_sub(1)
}
#[cfg(test)]
pub(super) fn decode_blob_into<'a>(
&self,
blob_idx: usize,
scratch: &'a mut Vec<u32>,
) -> Result<&'a [u32]> {
decode_blob_record_into(self.blob_record(blob_idx), blob_idx, scratch)?;
Ok(scratch.as_slice())
}
pub(super) fn blob_has_nonzero_refs(&self, blob_idx: usize) -> Result<bool> {
blob_record_has_nonzero_refs(self.blob_record(blob_idx), blob_idx)
}
}
fn scan_blob_record(cursor: &mut protohoggr::Cursor<'_>, blob_idx: usize) -> Result<()> {
let num_ways = cursor
.read_varint()
.map_err(|e| format!("per-way sidecar blob {blob_idx} num_ways: {e}"))?;
#[allow(clippy::cast_possible_truncation)]
let num_ways_usize = num_ways as usize;
for way_idx in 0..num_ways_usize {
cursor
.read_varint()
.map_err(|e| format!("per-way sidecar blob {blob_idx} way {way_idx}: {e}"))?;
}
Ok(())
}
#[cfg(test)]
fn decode_blob_record_into(record: &[u8], blob_idx: usize, scratch: &mut Vec<u32>) -> Result<()> {
let mut cursor = protohoggr::Cursor::new(record);
let num_ways = cursor
.read_varint()
.map_err(|e| format!("per-way sidecar blob {blob_idx} num_ways: {e}"))?;
#[allow(clippy::cast_possible_truncation)]
let num_ways_usize = num_ways as usize;
scratch.clear();
scratch.reserve(num_ways_usize);
for way_idx in 0..num_ways_usize {
let rc = cursor
.read_varint()
.map_err(|e| format!("per-way sidecar blob {blob_idx} way {way_idx}: {e}"))?;
#[allow(clippy::cast_possible_truncation)]
scratch.push(rc as u32);
}
if cursor.remaining() != 0 {
return Err(format!(
"per-way sidecar blob {blob_idx} has {} trailing bytes",
cursor.remaining()
)
.into());
}
Ok(())
}
fn blob_record_has_nonzero_refs(record: &[u8], blob_idx: usize) -> Result<bool> {
let mut cursor = protohoggr::Cursor::new(record);
let num_ways = cursor
.read_varint()
.map_err(|e| format!("per-way sidecar blob {blob_idx} num_ways: {e}"))?;
#[allow(clippy::cast_possible_truncation)]
let num_ways_usize = num_ways as usize;
for way_idx in 0..num_ways_usize {
let rc = cursor
.read_varint()
.map_err(|e| format!("per-way sidecar blob {blob_idx} way {way_idx}: {e}"))?;
if rc != 0 {
return Ok(true);
}
}
if cursor.remaining() != 0 {
return Err(format!(
"per-way sidecar blob {blob_idx} has {} trailing bytes",
cursor.remaining()
)
.into());
}
Ok(false)
}
fn build_per_way_refcount_index(data: Vec<u8>, num_way_blobs: usize) -> Result<PerWayRcs> {
let mut cursor = protohoggr::Cursor::new(&data);
let mut offsets: Vec<usize> = Vec::with_capacity(num_way_blobs + 1);
offsets.push(0);
let data_len = data.len();
for blob_idx in 0..num_way_blobs {
scan_blob_record(&mut cursor, blob_idx)?;
offsets.push(data_len - cursor.remaining());
}
if cursor.remaining() != 0 {
return Err(format!(
"per-way refcount sidecar has {} trailing bytes",
cursor.remaining()
)
.into());
}
Ok(PerWayRcs { data, offsets })
}
#[cfg(test)]
pub(super) fn parse_per_way_refcount_sidecar_bytes(
data: &[u8],
num_way_blobs: usize,
) -> Result<PerWayRcs> {
build_per_way_refcount_index(data.to_vec(), num_way_blobs)
}
pub(super) fn load_per_way_refcount_sidecar_indexed(
path: &Path,
num_way_blobs: usize,
) -> Result<PerWayRcs> {
let data = std::fs::read(path).map_err(|e| format!("read per-way refcount sidecar: {e}"))?;
build_per_way_refcount_index(data, num_way_blobs)
}
const _: () = {
assert!(COORD_SLOT_SIZE == 8);
};
#[cfg(test)]
pub(super) fn encode_blob_payload(
coord_bytes: &[u8],
per_way_rcs: &[u32],
output: &mut Vec<u8>,
) -> std::result::Result<(), String> {
let expected_bytes: u64 = per_way_rcs.iter().map(|&r| u64::from(r)).sum::<u64>() * 8;
if coord_bytes.len() as u64 != expected_bytes {
return Err(format!(
"coord_bytes length mismatch: got {} bytes, expected {} (8 * sum(per_way_rcs))",
coord_bytes.len(),
expected_bytes
));
}
let mut cursor: usize = 0;
for &rc in per_way_rcs {
let mut last_lat: i32 = 0;
let mut last_lon: i32 = 0;
for _ in 0..rc {
let off = cursor;
cursor += COORD_SLOT_SIZE;
let lat = i32::from_le_bytes([
coord_bytes[off],
coord_bytes[off + 1],
coord_bytes[off + 2],
coord_bytes[off + 3],
]);
let lon = i32::from_le_bytes([
coord_bytes[off + 4],
coord_bytes[off + 5],
coord_bytes[off + 6],
coord_bytes[off + 7],
]);
let dlat = i64::from(lat) - i64::from(last_lat);
let dlon = i64::from(lon) - i64::from(last_lon);
protohoggr::encode_varint(output, protohoggr::zigzag_encode_64(dlat));
protohoggr::encode_varint(output, protohoggr::zigzag_encode_64(dlon));
last_lat = lat;
last_lon = lon;
}
}
Ok(())
}
#[hotpath::measure]
pub(super) fn encode_blob_payload_from_record(
coord_bytes: &[u8],
record: &[u8],
blob_idx: usize,
inject_prepass: bool,
output: &mut Vec<u8>,
) -> std::result::Result<(), String> {
let mut cursor = protohoggr::Cursor::new(record);
let num_ways = cursor
.read_varint()
.map_err(|e| format!("per-way sidecar blob {blob_idx} num_ways: {e}"))?;
#[allow(clippy::cast_possible_truncation)]
let num_ways_usize = num_ways as usize;
let mut coord_cursor: usize = 0;
for way_idx in 0..num_ways_usize {
let rc = cursor
.read_varint()
.map_err(|e| format!("per-way sidecar blob {blob_idx} way {way_idx}: {e}"))?;
#[allow(clippy::cast_possible_truncation)]
let rc_usize = rc as usize;
let way_bytes = rc_usize
.checked_mul(COORD_SLOT_SIZE)
.ok_or_else(|| format!("blob {blob_idx} way {way_idx}: refcount byte size overflow"))?;
if coord_cursor + way_bytes > coord_bytes.len() {
return Err(format!(
"coord_bytes length mismatch for blob {blob_idx}: got {} bytes, \
need at least {} bytes by way {way_idx}",
coord_bytes.len(),
coord_cursor + way_bytes,
));
}
let mut last_lat: i32 = 0;
let mut last_lon: i32 = 0;
let mut pins = if inject_prepass {
vec![0_u8; rc_usize.div_ceil(8)]
} else {
Vec::new()
};
for i in 0..rc_usize {
let off = coord_cursor;
coord_cursor += COORD_SLOT_SIZE;
let packed_lat = i32::from_le_bytes([
coord_bytes[off],
coord_bytes[off + 1],
coord_bytes[off + 2],
coord_bytes[off + 3],
]);
let lon = i32::from_le_bytes([
coord_bytes[off + 4],
coord_bytes[off + 5],
coord_bytes[off + 6],
coord_bytes[off + 7],
]);
let lat = if inject_prepass {
if packed_lat & 1 != 0 {
pins[i / 8] |= 1 << (i % 8);
}
packed_lat >> 1
} else {
packed_lat
};
let dlat = i64::from(lat) - i64::from(last_lat);
let dlon = i64::from(lon) - i64::from(last_lon);
protohoggr::encode_varint(output, protohoggr::zigzag_encode_64(dlat));
protohoggr::encode_varint(output, protohoggr::zigzag_encode_64(dlon));
last_lat = lat;
last_lon = lon;
}
if inject_prepass {
output.extend_from_slice(&pins);
}
}
if cursor.remaining() != 0 {
return Err(format!(
"per-way sidecar blob {blob_idx} has {} trailing bytes",
cursor.remaining()
));
}
if coord_cursor != coord_bytes.len() {
return Err(format!(
"coord_bytes length mismatch for blob {blob_idx}: got {} bytes, expected {}",
coord_bytes.len(),
coord_cursor,
));
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
fn make_coord_bytes(coords: &[(i32, i32)]) -> Vec<u8> {
let mut buf = Vec::with_capacity(coords.len() * 8);
for &(lat, lon) in coords {
buf.extend_from_slice(&lat.to_le_bytes());
buf.extend_from_slice(&lon.to_le_bytes());
}
buf
}
fn decode_zigzag_varints(data: &[u8], count: usize) -> Vec<i64> {
let mut cursor = protohoggr::Cursor::new(data);
let mut out = Vec::with_capacity(count);
for _ in 0..count {
let v = cursor.read_varint().expect("read varint");
out.push(protohoggr::zigzag_decode_64(v));
}
out
}
fn reconstruct_coords(output: &[u8], per_way_rcs: &[u32]) -> Vec<(i32, i32)> {
let total_refs: usize = per_way_rcs.iter().map(|&r| r as usize).sum();
let deltas = decode_zigzag_varints(output, total_refs * 2);
let mut result = Vec::with_capacity(total_refs);
let mut delta_idx = 0;
for &rc in per_way_rcs {
let mut last_lat: i64 = 0;
let mut last_lon: i64 = 0;
for _ in 0..rc {
let lat = last_lat + deltas[delta_idx];
let lon = last_lon + deltas[delta_idx + 1];
delta_idx += 2;
#[allow(clippy::cast_possible_truncation)]
result.push((lat as i32, lon as i32));
last_lat = lat;
last_lon = lon;
}
}
result
}
#[test]
fn encode_blob_payload_single_way_single_ref() {
let coords = [(12345_i32, 67890_i32)];
let cb = make_coord_bytes(&coords);
let rcs = [1u32];
let mut out = Vec::new();
encode_blob_payload(&cb, &rcs, &mut out).expect("encode");
let reconstructed = reconstruct_coords(&out, &rcs);
assert_eq!(reconstructed, coords);
}
#[test]
fn encode_blob_payload_single_way_multi_ref() {
let coords = [(100_i32, 200_i32), (150_i32, 250_i32), (90_i32, 180_i32)];
let cb = make_coord_bytes(&coords);
let rcs = [3u32];
let mut out = Vec::new();
encode_blob_payload(&cb, &rcs, &mut out).expect("encode");
let reconstructed = reconstruct_coords(&out, &rcs);
assert_eq!(reconstructed, coords);
}
#[test]
fn encode_blob_payload_multiple_ways() {
let coords = [
(10_i32, 20_i32),
(30_i32, 40_i32), (500_i32, 600_i32),
(510_i32, 610_i32),
(490_i32, 590_i32), (9999_i32, -9999_i32), ];
let cb = make_coord_bytes(&coords);
let rcs = [2u32, 3u32, 1u32];
let mut out = Vec::new();
encode_blob_payload(&cb, &rcs, &mut out).expect("encode");
let reconstructed = reconstruct_coords(&out, &rcs);
assert_eq!(reconstructed, coords);
}
#[test]
fn encode_blob_payload_empty_blob() {
let cb: Vec<u8> = Vec::new();
let rcs: &[u32] = &[];
let mut out = vec![0xAAu8, 0xBBu8];
encode_blob_payload(&cb, rcs, &mut out).expect("encode");
assert_eq!(out, [0xAAu8, 0xBBu8]);
}
#[test]
fn encode_blob_payload_zero_coords() {
let coords = [(0_i32, 0_i32), (0_i32, 0_i32)];
let cb = make_coord_bytes(&coords);
let rcs = [2u32];
let mut out = Vec::new();
encode_blob_payload(&cb, &rcs, &mut out).expect("encode");
let reconstructed = reconstruct_coords(&out, &rcs);
assert_eq!(reconstructed, coords);
}
#[test]
fn encode_blob_payload_negative_coords() {
let coords = [(-1_000_000_i32, -1_000_000_i32), (0_i32, 0_i32)];
let cb = make_coord_bytes(&coords);
let rcs = [2u32];
let mut out = Vec::new();
encode_blob_payload(&cb, &rcs, &mut out).expect("encode");
let reconstructed = reconstruct_coords(&out, &rcs);
assert_eq!(reconstructed, coords);
}
#[test]
fn encode_blob_payload_length_mismatch() {
let cb = vec![0u8; 15];
let rcs = [2u32];
let result = encode_blob_payload(&cb, &rcs, &mut Vec::new());
assert!(result.is_err(), "expected Err for length mismatch");
}
fn make_sidecar_bytes(blobs: &[&[u32]]) -> Vec<u8> {
let mut out: Vec<u8> = Vec::new();
for blob_rcs in blobs {
protohoggr::encode_varint(&mut out, blob_rcs.len() as u64);
for &rc in *blob_rcs {
protohoggr::encode_varint(&mut out, rc as u64);
}
}
out
}
#[test]
fn parse_per_way_refcount_sidecar_bytes_basic() {
let sidecar = make_sidecar_bytes(&[&[3, 1], &[2]]);
let pwr = parse_per_way_refcount_sidecar_bytes(&sidecar, 2).expect("parse");
let mut scratch: Vec<u32> = Vec::new();
assert_eq!(pwr.num_blobs(), 2);
assert_eq!(
pwr.decode_blob_into(0, &mut scratch).expect("decode 0"),
&[3u32, 1u32]
);
assert_eq!(
pwr.decode_blob_into(1, &mut scratch).expect("decode 1"),
&[2u32]
);
}
#[test]
fn parse_per_way_refcount_sidecar_bytes_empty_blob() {
let sidecar = make_sidecar_bytes(&[&[], &[1]]);
let pwr = parse_per_way_refcount_sidecar_bytes(&sidecar, 2).expect("parse");
let mut scratch: Vec<u32> = Vec::new();
assert_eq!(
pwr.decode_blob_into(0, &mut scratch).expect("decode 0"),
&[] as &[u32]
);
assert_eq!(
pwr.decode_blob_into(1, &mut scratch).expect("decode 1"),
&[1u32]
);
}
#[test]
fn encode_blob_payload_from_record_matches_decoded_path() {
let coords = [
(10_i32, 20_i32),
(30_i32, 40_i32),
(500_i32, 600_i32),
(510_i32, 610_i32),
(490_i32, 590_i32),
(9999_i32, -9999_i32),
];
let cb = make_coord_bytes(&coords);
let pwr = make_per_way_rcs(&[&[2u32, 3u32, 1u32]]);
let mut from_record = Vec::new();
let mut from_decoded = Vec::new();
let mut scratch: Vec<u32> = Vec::new();
encode_blob_payload_from_record(&cb, pwr.blob_record(0), 0, false, &mut from_record)
.expect("encode from record");
let decoded = pwr
.decode_blob_into(0, &mut scratch)
.expect("decode 0")
.to_vec();
encode_blob_payload(&cb, &decoded, &mut from_decoded).expect("encode decoded");
assert_eq!(from_record, from_decoded);
}
fn make_per_way_rcs(blobs: &[&[u32]]) -> PerWayRcs {
let sidecar = make_sidecar_bytes(blobs);
parse_per_way_refcount_sidecar_bytes(&sidecar, blobs.len()).expect("make_per_way_rcs")
}
}