use parking_lot::RwLock;
use std::collections::HashMap;
use std::sync::atomic::{AtomicBool, Ordering};
pub(crate) const VMETA_PREFIX: &str = "vmeta:";
#[derive(Clone, Copy, Debug)]
struct Interval {
first: u64,
last: u64,
ts: u64,
}
pub(crate) struct CopyMarkers {
by_table: RwLock<HashMap<String, Vec<Interval>>>,
any: AtomicBool,
}
impl CopyMarkers {
pub(crate) fn new() -> Self {
Self {
by_table: RwLock::new(HashMap::new()),
any: AtomicBool::new(false),
}
}
#[inline]
pub(crate) fn is_empty(&self) -> bool {
!self.any.load(Ordering::Relaxed)
}
pub(crate) fn len(&self) -> usize {
if self.is_empty() {
return 0;
}
self.by_table.read().values().map(|v| v.len()).sum()
}
pub(crate) fn marker_key(table: &str, first: u64, last: u64) -> String {
format!("{VMETA_PREFIX}{table}:{first:020}:{last:020}")
}
fn parse_marker(key: &[u8], value: &[u8]) -> Option<(String, u64, u64, u64)> {
let key = std::str::from_utf8(key).ok()?;
let rest = key.strip_prefix(VMETA_PREFIX)?;
let (rest, last) = rest.rsplit_once(':')?;
let (table, first) = rest.rsplit_once(':')?;
let first: u64 = first.parse().ok()?;
let last: u64 = last.parse().ok()?;
if value.len() < 8 {
return None;
}
let ts = u64::from_be_bytes(value.get(0..8)?.try_into().ok()?);
Some((table.to_string(), first, last, ts))
}
pub(crate) fn insert(&self, table: &str, first: u64, last: u64, ts: u64) {
let mut g = self.by_table.write();
g.entry(table.to_string())
.or_default()
.push(Interval { first, last, ts });
self.any.store(true, Ordering::Relaxed);
}
pub(crate) fn load_record(&self, key: &[u8], value: &[u8]) {
if let Some((table, first, last, ts)) = Self::parse_marker(key, value) {
self.insert(&table, first, last, ts);
}
}
pub(crate) fn table_has_markers(&self, table: &str) -> bool {
if self.is_empty() {
return false;
}
self.by_table.read().get(table).is_some_and(|v| !v.is_empty())
}
pub(crate) fn covering_ts(&self, table: &str, row_id: u64) -> Option<u64> {
if self.is_empty() {
return None;
}
let g = self.by_table.read();
let intervals = g.get(table)?;
intervals
.iter()
.find(|iv| row_id >= iv.first && row_id <= iv.last)
.map(|iv| iv.ts)
}
}
impl Default for CopyMarkers {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn key_roundtrips_including_colon_in_table() {
for (t, f, l, ts) in [
("t", 1u64, 100u64, 42u64),
("weird:name", 5, 9, 7),
("x", 0, u64::MAX, 1),
] {
let k = CopyMarkers::marker_key(t, f, l);
let v = ts.to_be_bytes();
let (pt, pf, pl, pts) = CopyMarkers::parse_marker(k.as_bytes(), &v).unwrap();
assert_eq!((pt.as_str(), pf, pl, pts), (t, f, l, ts));
}
}
#[test]
fn covering_ts_finds_and_misses() {
let m = CopyMarkers::new();
assert!(m.is_empty());
assert_eq!(m.covering_ts("t", 5), None);
m.insert("t", 10, 20, 100);
m.insert("t", 30, 40, 200);
assert!(!m.is_empty());
assert_eq!(m.covering_ts("t", 9), None);
assert_eq!(m.covering_ts("t", 10), Some(100));
assert_eq!(m.covering_ts("t", 15), Some(100));
assert_eq!(m.covering_ts("t", 20), Some(100));
assert_eq!(m.covering_ts("t", 25), None);
assert_eq!(m.covering_ts("t", 35), Some(200));
assert_eq!(m.covering_ts("other", 15), None);
}
#[test]
fn load_record_skips_garbage() {
let m = CopyMarkers::new();
m.load_record(
b"vmeta:t:00000000000000000001:00000000000000000009",
&50u64.to_be_bytes(),
);
assert_eq!(m.covering_ts("t", 5), Some(50));
m.load_record(b"vmeta:garbage", b"x");
m.load_record(b"not-a-marker", b"");
assert_eq!(m.covering_ts("t", 5), Some(50));
}
#[test]
fn key_ordering_is_numeric() {
let a = CopyMarkers::marker_key("t", 2, 9);
let b = CopyMarkers::marker_key("t", 10, 20);
assert!(a < b, "{a} should sort before {b}");
}
}