use snafu::Snafu;
use uuid::Uuid;
use crate::metadata::index::{IndexKind, IndexSpec, TimeIndexGranularity};
pub const SEGMENT_COVERAGE_DIR: &str = "_coverage/segments";
pub const TABLE_SNAPSHOT_DIR: &str = "_coverage/table";
pub const COVERAGE_EXT: &str = "roar";
#[derive(Debug, Snafu)]
#[non_exhaustive]
pub enum CoverageLayoutError {
#[snafu(display("Invalid coverage id: {coverage_id}"))]
InvalidCoverageId {
coverage_id: String,
},
}
pub fn validate_coverage_id(coverage_id: &str) -> Result<(), CoverageLayoutError> {
if coverage_id.is_empty() || coverage_id.len() > 128 {
return Err(CoverageLayoutError::InvalidCoverageId {
coverage_id: coverage_id.to_string(),
});
}
if !coverage_id.chars().any(|c| c.is_ascii_alphanumeric()) {
return Err(CoverageLayoutError::InvalidCoverageId {
coverage_id: coverage_id.to_string(),
});
}
if coverage_id.starts_with('.') {
return Err(CoverageLayoutError::InvalidCoverageId {
coverage_id: coverage_id.to_string(),
});
}
if coverage_id.contains('/') || coverage_id.contains('\\') || coverage_id.contains("..") {
return Err(CoverageLayoutError::InvalidCoverageId {
coverage_id: coverage_id.to_string(),
});
}
let ok = coverage_id
.chars()
.all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-'));
if !ok {
return Err(CoverageLayoutError::InvalidCoverageId {
coverage_id: coverage_id.to_string(),
});
}
Ok(())
}
pub fn segment_coverage_key(coverage_id: &str) -> Result<String, CoverageLayoutError> {
validate_coverage_id(coverage_id)?;
Ok(format!(
"{SEGMENT_COVERAGE_DIR}/{coverage_id}.{COVERAGE_EXT}"
))
}
pub fn table_snapshot_key(version: u64, snapshot_id: &str) -> Result<String, CoverageLayoutError> {
validate_coverage_id(snapshot_id)?;
Ok(format!(
"{TABLE_SNAPSHOT_DIR}/{version}-{snapshot_id}.{COVERAGE_EXT}"
))
}
fn coverage_id_v2(
domain_prefix: &[u8],
output_prefix: &str,
index: &IndexSpec,
coverage_bytes: &[u8],
) -> String {
let mut h = blake3::Hasher::new();
h.update(domain_prefix);
h.update(b"\0");
h.update(index.column.as_bytes());
h.update(b"\0");
match &index.kind {
IndexKind::Timestamp {
index_granularity,
timezone,
} => {
h.update(b"T");
hash_time_index_granularity(&mut h, index_granularity);
h.update(b"\0");
match timezone {
Some(timezone) => {
h.update(b"S");
h.update(timezone.as_bytes());
}
None => {
h.update(b"N");
}
}
}
IndexKind::Int64 { index_granularity } => {
h.update(b"I");
h.update(&index_granularity.get().to_le_bytes());
}
IndexKind::UInt64 { index_granularity } => {
h.update(b"U");
h.update(&index_granularity.get().to_le_bytes());
}
}
h.update(b"\0");
h.update(coverage_bytes);
let hex = h.finalize().to_hex();
format!("{output_prefix}-{}", &hex[..32])
}
fn entity_coverage_id_v1(
domain_prefix: &[u8],
output_prefix: &str,
index: &IndexSpec,
coverage_bytes: &[u8],
) -> String {
let mut h = blake3::Hasher::new();
h.update(domain_prefix);
h.update(b"\0");
h.update(b"C");
hash_len_prefixed(&mut h, index.column.as_bytes());
h.update(b"E");
hash_usize(&mut h, index.entity_columns.len());
for column in &index.entity_columns {
hash_len_prefixed(&mut h, column.as_bytes());
}
h.update(b"K");
match &index.kind {
IndexKind::Timestamp {
index_granularity,
timezone,
} => {
h.update(b"T");
hash_time_index_granularity(&mut h, index_granularity);
match timezone {
Some(timezone) => {
h.update(b"S");
hash_len_prefixed(&mut h, timezone.as_bytes());
}
None => {
h.update(b"N");
}
}
}
IndexKind::Int64 { index_granularity } => {
h.update(b"I");
h.update(&index_granularity.get().to_le_bytes());
}
IndexKind::UInt64 { index_granularity } => {
h.update(b"U");
h.update(&index_granularity.get().to_le_bytes());
}
}
h.update(b"\0");
h.update(coverage_bytes);
let hex = h.finalize().to_hex();
format!("{output_prefix}-{}", &hex[..32])
}
fn hash_len_prefixed(hasher: &mut blake3::Hasher, bytes: &[u8]) {
hash_usize(hasher, bytes.len());
hasher.update(bytes);
}
fn hash_usize(hasher: &mut blake3::Hasher, value: usize) {
hasher.update(value.to_string().as_bytes());
hasher.update(b":");
}
fn hash_time_index_granularity(
hasher: &mut blake3::Hasher,
index_granularity: &TimeIndexGranularity,
) {
match index_granularity {
TimeIndexGranularity::Seconds(n) => {
hasher.update(b"S");
hasher.update(&n.to_le_bytes());
}
TimeIndexGranularity::Minutes(n) => {
hasher.update(b"M");
hasher.update(&n.to_le_bytes());
}
TimeIndexGranularity::Hours(n) => {
hasher.update(b"H");
hasher.update(&n.to_le_bytes());
}
TimeIndexGranularity::Days(n) => {
hasher.update(b"D");
hasher.update(&n.to_le_bytes());
}
}
}
pub fn segment_coverage_id_v2(index: &IndexSpec, coverage_bytes: &[u8]) -> String {
coverage_id_v2(b"segcov-v2", "segcov", index, coverage_bytes)
}
pub fn table_coverage_id_v2(index: &IndexSpec, coverage_bytes: &[u8]) -> String {
coverage_id_v2(b"tblcov-v2", "tblcov", index, coverage_bytes)
}
pub(crate) fn segment_entity_coverage_id_v1(index: &IndexSpec, coverage_bytes: &[u8]) -> String {
entity_coverage_id_v1(b"entity-segcov-v1", "segcov", index, coverage_bytes)
}
pub(crate) fn table_entity_coverage_id_v1(index: &IndexSpec, coverage_bytes: &[u8]) -> String {
entity_coverage_id_v1(b"entity-tblcov-v1", "tblcov", index, coverage_bytes)
}
pub(crate) fn coverage_file_id_for_attempt(content_id: &str, attempt_id: &Uuid) -> String {
format!("{content_id}-{attempt_id}")
}
#[cfg(test)]
mod tests {
use std::num::NonZeroU64;
use super::*;
fn timestamp_index(column: &str, index_granularity: TimeIndexGranularity) -> IndexSpec {
IndexSpec {
column: column.to_string(),
entity_columns: Vec::new(),
kind: IndexKind::Timestamp {
index_granularity,
timezone: None,
},
}
}
#[test]
fn validate_coverage_id_accepts_valid_ids() {
let long = "a".repeat(128);
let valid_ids = ["abc", "A_B-1.2", long.as_str()];
for id in valid_ids {
validate_coverage_id(id).expect("valid id should pass");
}
}
#[test]
fn validate_coverage_id_rejects_empty_or_too_long() {
let too_long = "x".repeat(129);
assert!(validate_coverage_id("").is_err());
assert!(validate_coverage_id(&too_long).is_err());
}
#[test]
fn validate_coverage_id_rejects_path_components() {
for id in ["a/b", "a\\b", "a..b", "..", "../etc"] {
assert!(validate_coverage_id(id).is_err(), "id `{id}` should fail");
}
}
#[test]
fn validate_coverage_id_rejects_disallowed_chars() {
for id in ["space id", "id*", "id@", "id$", "id:"] {
assert!(validate_coverage_id(id).is_err(), "id `{id}` should fail");
}
}
#[test]
fn segment_coverage_key_formats_and_validates() {
let id = "seg-001";
let key = segment_coverage_key(id).expect("valid id");
assert_eq!(key, "_coverage/segments/seg-001.roar");
assert!(segment_coverage_key("bad/id").is_err());
}
#[test]
fn table_snapshot_key_formats() {
let key = table_snapshot_key(42, "snap-001").expect("valid snapshot id");
assert_eq!(key, "_coverage/table/42-snap-001.roar");
}
#[test]
fn segment_coverage_id_matches_golden_value_and_is_valid() {
let index = timestamp_index("ts", TimeIndexGranularity::Minutes(1));
let bytes = b"bitmap-bytes";
let id1 = segment_coverage_id_v2(&index, bytes);
let id2 = segment_coverage_id_v2(&index, bytes);
assert_eq!(id1, "segcov-00720d0b60b246ef53e757b286681cc0");
assert_eq!(id1, id2, "same inputs must produce stable id");
assert!(id1.starts_with("segcov-"));
assert_eq!(id1.len(), "segcov-".len() + 32, "prefix + 32 hex chars");
validate_coverage_id(&id1).expect("derived id should be valid");
}
#[test]
fn segment_coverage_id_changes_with_inputs() {
let bytes = b"bytes";
let base_index = timestamp_index("ts", TimeIndexGranularity::Seconds(5));
let base = segment_coverage_id_v2(&base_index, bytes);
let different_granularity = segment_coverage_id_v2(
×tamp_index("ts", TimeIndexGranularity::Hours(5)),
bytes,
);
let different_column = segment_coverage_id_v2(
×tamp_index("event_time", TimeIndexGranularity::Seconds(5)),
bytes,
);
let different_kind = segment_coverage_id_v2(
&IndexSpec {
column: "ts".to_string(),
entity_columns: Vec::new(),
kind: IndexKind::UInt64 {
index_granularity: NonZeroU64::new(5).unwrap(),
},
},
bytes,
);
let different_integer_domain = segment_coverage_id_v2(
&IndexSpec {
column: "ts".to_string(),
entity_columns: Vec::new(),
kind: IndexKind::Int64 {
index_granularity: NonZeroU64::new(5).unwrap(),
},
},
bytes,
);
let different_integer_granularity = segment_coverage_id_v2(
&IndexSpec {
column: "ts".to_string(),
entity_columns: Vec::new(),
kind: IndexKind::UInt64 {
index_granularity: NonZeroU64::new(6).unwrap(),
},
},
bytes,
);
let different_bytes = segment_coverage_id_v2(&base_index, b"other");
assert_ne!(
base, different_granularity,
"index granularity should affect id"
);
assert_ne!(base, different_column, "index column should affect id");
assert_ne!(base, different_kind, "index kind should affect id");
assert_ne!(different_kind, different_integer_domain);
assert_ne!(different_kind, different_integer_granularity);
assert_ne!(base, different_bytes, "coverage bytes should affect id");
}
#[test]
fn table_coverage_id_matches_golden_value_and_is_valid() {
let index = timestamp_index("ts", TimeIndexGranularity::Hours(1));
let bytes = b"table-bitmap";
let id1 = table_coverage_id_v2(&index, bytes);
let id2 = table_coverage_id_v2(&index, bytes);
assert_eq!(id1, "tblcov-38f0aa9c3e526d0cdabf234af8fb0fd3");
assert_eq!(id1, id2, "same inputs must produce stable id");
assert!(id1.starts_with("tblcov-"));
assert_eq!(id1.len(), "tblcov-".len() + 32, "prefix + 32 hex chars");
validate_coverage_id(&id1).expect("derived id should be valid");
}
#[test]
fn table_coverage_id_changes_with_inputs() {
let bytes = b"bytes";
let base_index = timestamp_index("ts", TimeIndexGranularity::Minutes(15));
let base = table_coverage_id_v2(&base_index, bytes);
let different_granularity =
table_coverage_id_v2(×tamp_index("ts", TimeIndexGranularity::Days(1)), bytes);
let different_column = table_coverage_id_v2(
×tamp_index("event_time", TimeIndexGranularity::Minutes(15)),
bytes,
);
let different_bytes = table_coverage_id_v2(&base_index, b"other");
assert_ne!(
base, different_granularity,
"index granularity should affect id"
);
assert_ne!(base, different_column, "index column should affect id");
assert_ne!(base, different_bytes, "coverage bytes should affect id");
}
#[test]
fn entity_coverage_ids_match_golden_values_and_include_ordered_columns() {
let index = IndexSpec {
column: "ts".to_string(),
entity_columns: vec!["symbol".to_string(), "venue".to_string()],
kind: IndexKind::Timestamp {
index_granularity: TimeIndexGranularity::Minutes(1),
timezone: None,
},
};
let mut renamed = index.clone();
renamed.entity_columns[0] = "device".to_string();
let mut reordered = index.clone();
reordered.entity_columns.reverse();
let bytes = b"entity-coverage-bytes";
let segment = segment_entity_coverage_id_v1(&index, bytes);
assert_eq!(segment, "segcov-67c0022aad0d9f5bf5ea813e9ef88119");
assert_ne!(segment, segment_entity_coverage_id_v1(&renamed, bytes));
assert_ne!(segment, segment_entity_coverage_id_v1(&reordered, bytes));
let table = table_entity_coverage_id_v1(&index, bytes);
assert_eq!(table, "tblcov-9c54647467c3a0e89e60675e00b7c75b");
assert_ne!(table, table_entity_coverage_id_v1(&renamed, bytes));
assert_ne!(table, table_entity_coverage_id_v1(&reordered, bytes));
}
#[test]
fn coverage_file_ids_are_owned_by_the_append_attempt() {
let content_id = "segcov-0123456789abcdef0123456789abcdef";
let first = coverage_file_id_for_attempt(content_id, &Uuid::from_u128(1));
let second = coverage_file_id_for_attempt(content_id, &Uuid::from_u128(2));
assert_ne!(first, second);
validate_coverage_id(&first).expect("first id should be valid");
validate_coverage_id(&second).expect("second id should be valid");
}
}