1use snafu::Snafu;
12use uuid::Uuid;
13
14use crate::metadata::index::{IndexKind, IndexSpec, TimeIndexGranularity};
15
16pub const SEGMENT_COVERAGE_DIR: &str = "_coverage/segments";
18pub const TABLE_SNAPSHOT_DIR: &str = "_coverage/table";
20pub const COVERAGE_EXT: &str = "roar";
22
23#[derive(Debug, Snafu)]
25#[non_exhaustive]
26pub enum CoverageLayoutError {
27 #[snafu(display("Invalid coverage id: {coverage_id}"))]
29 InvalidCoverageId {
30 coverage_id: String,
32 },
33}
34
35pub fn validate_coverage_id(coverage_id: &str) -> Result<(), CoverageLayoutError> {
42 if coverage_id.is_empty() || coverage_id.len() > 128 {
43 return Err(CoverageLayoutError::InvalidCoverageId {
44 coverage_id: coverage_id.to_string(),
45 });
46 }
47
48 if !coverage_id.chars().any(|c| c.is_ascii_alphanumeric()) {
50 return Err(CoverageLayoutError::InvalidCoverageId {
51 coverage_id: coverage_id.to_string(),
52 });
53 }
54
55 if coverage_id.starts_with('.') {
57 return Err(CoverageLayoutError::InvalidCoverageId {
58 coverage_id: coverage_id.to_string(),
59 });
60 }
61
62 if coverage_id.contains('/') || coverage_id.contains('\\') || coverage_id.contains("..") {
64 return Err(CoverageLayoutError::InvalidCoverageId {
65 coverage_id: coverage_id.to_string(),
66 });
67 }
68
69 let ok = coverage_id
71 .chars()
72 .all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | '-'));
73
74 if !ok {
75 return Err(CoverageLayoutError::InvalidCoverageId {
76 coverage_id: coverage_id.to_string(),
77 });
78 }
79
80 Ok(())
81}
82
83pub fn segment_coverage_key(coverage_id: &str) -> Result<String, CoverageLayoutError> {
85 validate_coverage_id(coverage_id)?;
86 Ok(format!(
87 "{SEGMENT_COVERAGE_DIR}/{coverage_id}.{COVERAGE_EXT}"
88 ))
89}
90
91pub fn table_snapshot_key(version: u64, snapshot_id: &str) -> Result<String, CoverageLayoutError> {
93 validate_coverage_id(snapshot_id)?;
94 Ok(format!(
95 "{TABLE_SNAPSHOT_DIR}/{version}-{snapshot_id}.{COVERAGE_EXT}"
96 ))
97}
98
99fn coverage_id_v2(
100 domain_prefix: &[u8],
101 output_prefix: &str,
102 index: &IndexSpec,
103 coverage_bytes: &[u8],
104) -> String {
105 let mut h = blake3::Hasher::new();
106
107 h.update(domain_prefix);
109 h.update(b"\0");
110
111 h.update(index.column.as_bytes());
112 h.update(b"\0");
113
114 match &index.kind {
115 IndexKind::Timestamp {
116 index_granularity,
117 timezone,
118 } => {
119 h.update(b"T");
120 hash_time_index_granularity(&mut h, index_granularity);
121 h.update(b"\0");
122 match timezone {
123 Some(timezone) => {
124 h.update(b"S");
125 h.update(timezone.as_bytes());
126 }
127 None => {
128 h.update(b"N");
129 }
130 }
131 }
132 IndexKind::Int64 { index_granularity } => {
133 h.update(b"I");
134 h.update(&index_granularity.get().to_le_bytes());
135 }
136 IndexKind::UInt64 { index_granularity } => {
137 h.update(b"U");
138 h.update(&index_granularity.get().to_le_bytes());
139 }
140 }
141
142 h.update(b"\0");
143 h.update(coverage_bytes);
144
145 let hex = h.finalize().to_hex();
146 format!("{output_prefix}-{}", &hex[..32])
147}
148
149fn entity_coverage_id_v1(
150 domain_prefix: &[u8],
151 output_prefix: &str,
152 index: &IndexSpec,
153 coverage_bytes: &[u8],
154) -> String {
155 let mut h = blake3::Hasher::new();
156 h.update(domain_prefix);
157 h.update(b"\0");
158
159 h.update(b"C");
160 hash_len_prefixed(&mut h, index.column.as_bytes());
161 h.update(b"E");
162 hash_usize(&mut h, index.entity_columns.len());
163 for column in &index.entity_columns {
164 hash_len_prefixed(&mut h, column.as_bytes());
165 }
166 h.update(b"K");
167 match &index.kind {
168 IndexKind::Timestamp {
169 index_granularity,
170 timezone,
171 } => {
172 h.update(b"T");
173 hash_time_index_granularity(&mut h, index_granularity);
174 match timezone {
175 Some(timezone) => {
176 h.update(b"S");
177 hash_len_prefixed(&mut h, timezone.as_bytes());
178 }
179 None => {
180 h.update(b"N");
181 }
182 }
183 }
184 IndexKind::Int64 { index_granularity } => {
185 h.update(b"I");
186 h.update(&index_granularity.get().to_le_bytes());
187 }
188 IndexKind::UInt64 { index_granularity } => {
189 h.update(b"U");
190 h.update(&index_granularity.get().to_le_bytes());
191 }
192 }
193 h.update(b"\0");
194 h.update(coverage_bytes);
195
196 let hex = h.finalize().to_hex();
197 format!("{output_prefix}-{}", &hex[..32])
198}
199
200fn hash_len_prefixed(hasher: &mut blake3::Hasher, bytes: &[u8]) {
201 hash_usize(hasher, bytes.len());
202 hasher.update(bytes);
203}
204
205fn hash_usize(hasher: &mut blake3::Hasher, value: usize) {
206 hasher.update(value.to_string().as_bytes());
207 hasher.update(b":");
208}
209
210fn hash_time_index_granularity(
211 hasher: &mut blake3::Hasher,
212 index_granularity: &TimeIndexGranularity,
213) {
214 match index_granularity {
215 TimeIndexGranularity::Seconds(n) => {
216 hasher.update(b"S");
217 hasher.update(&n.to_le_bytes());
218 }
219 TimeIndexGranularity::Minutes(n) => {
220 hasher.update(b"M");
221 hasher.update(&n.to_le_bytes());
222 }
223 TimeIndexGranularity::Hours(n) => {
224 hasher.update(b"H");
225 hasher.update(&n.to_le_bytes());
226 }
227 TimeIndexGranularity::Days(n) => {
228 hasher.update(b"D");
229 hasher.update(&n.to_le_bytes());
230 }
231 }
232}
233
234pub fn segment_coverage_id_v2(index: &IndexSpec, coverage_bytes: &[u8]) -> String {
236 coverage_id_v2(b"segcov-v2", "segcov", index, coverage_bytes)
237}
238
239pub fn table_coverage_id_v2(index: &IndexSpec, coverage_bytes: &[u8]) -> String {
241 coverage_id_v2(b"tblcov-v2", "tblcov", index, coverage_bytes)
242}
243
244pub(crate) fn segment_entity_coverage_id_v1(index: &IndexSpec, coverage_bytes: &[u8]) -> String {
246 entity_coverage_id_v1(b"entity-segcov-v1", "segcov", index, coverage_bytes)
247}
248
249pub(crate) fn table_entity_coverage_id_v1(index: &IndexSpec, coverage_bytes: &[u8]) -> String {
251 entity_coverage_id_v1(b"entity-tblcov-v1", "tblcov", index, coverage_bytes)
252}
253
254pub(crate) fn coverage_file_id_for_attempt(content_id: &str, attempt_id: &Uuid) -> String {
256 format!("{content_id}-{attempt_id}")
257}
258
259#[cfg(test)]
260mod tests {
261 use std::num::NonZeroU64;
262
263 use super::*;
264
265 fn timestamp_index(column: &str, index_granularity: TimeIndexGranularity) -> IndexSpec {
266 IndexSpec {
267 column: column.to_string(),
268 entity_columns: Vec::new(),
269 kind: IndexKind::Timestamp {
270 index_granularity,
271 timezone: None,
272 },
273 }
274 }
275
276 #[test]
277 fn validate_coverage_id_accepts_valid_ids() {
278 let long = "a".repeat(128);
279 let valid_ids = ["abc", "A_B-1.2", long.as_str()];
280
281 for id in valid_ids {
282 validate_coverage_id(id).expect("valid id should pass");
283 }
284 }
285
286 #[test]
287 fn validate_coverage_id_rejects_empty_or_too_long() {
288 let too_long = "x".repeat(129);
289 assert!(validate_coverage_id("").is_err());
290 assert!(validate_coverage_id(&too_long).is_err());
291 }
292
293 #[test]
294 fn validate_coverage_id_rejects_path_components() {
295 for id in ["a/b", "a\\b", "a..b", "..", "../etc"] {
296 assert!(validate_coverage_id(id).is_err(), "id `{id}` should fail");
297 }
298 }
299
300 #[test]
301 fn validate_coverage_id_rejects_disallowed_chars() {
302 for id in ["space id", "id*", "id@", "id$", "id:"] {
303 assert!(validate_coverage_id(id).is_err(), "id `{id}` should fail");
304 }
305 }
306
307 #[test]
308 fn segment_coverage_key_formats_and_validates() {
309 let id = "seg-001";
310 let key = segment_coverage_key(id).expect("valid id");
311 assert_eq!(key, "_coverage/segments/seg-001.roar");
312
313 assert!(segment_coverage_key("bad/id").is_err());
315 }
316
317 #[test]
318 fn table_snapshot_key_formats() {
319 let key = table_snapshot_key(42, "snap-001").expect("valid snapshot id");
320 assert_eq!(key, "_coverage/table/42-snap-001.roar");
321 }
322
323 #[test]
324 fn segment_coverage_id_matches_golden_value_and_is_valid() {
325 let index = timestamp_index("ts", TimeIndexGranularity::Minutes(1));
326 let bytes = b"bitmap-bytes";
327
328 let id1 = segment_coverage_id_v2(&index, bytes);
329 let id2 = segment_coverage_id_v2(&index, bytes);
330
331 assert_eq!(id1, "segcov-00720d0b60b246ef53e757b286681cc0");
332 assert_eq!(id1, id2, "same inputs must produce stable id");
333 assert!(id1.starts_with("segcov-"));
334 assert_eq!(id1.len(), "segcov-".len() + 32, "prefix + 32 hex chars");
335 validate_coverage_id(&id1).expect("derived id should be valid");
336 }
337
338 #[test]
339 fn segment_coverage_id_changes_with_inputs() {
340 let bytes = b"bytes";
341
342 let base_index = timestamp_index("ts", TimeIndexGranularity::Seconds(5));
343 let base = segment_coverage_id_v2(&base_index, bytes);
344 let different_granularity = segment_coverage_id_v2(
345 ×tamp_index("ts", TimeIndexGranularity::Hours(5)),
346 bytes,
347 );
348 let different_column = segment_coverage_id_v2(
349 ×tamp_index("event_time", TimeIndexGranularity::Seconds(5)),
350 bytes,
351 );
352 let different_kind = segment_coverage_id_v2(
353 &IndexSpec {
354 column: "ts".to_string(),
355 entity_columns: Vec::new(),
356 kind: IndexKind::UInt64 {
357 index_granularity: NonZeroU64::new(5).unwrap(),
358 },
359 },
360 bytes,
361 );
362 let different_integer_domain = segment_coverage_id_v2(
363 &IndexSpec {
364 column: "ts".to_string(),
365 entity_columns: Vec::new(),
366 kind: IndexKind::Int64 {
367 index_granularity: NonZeroU64::new(5).unwrap(),
368 },
369 },
370 bytes,
371 );
372 let different_integer_granularity = segment_coverage_id_v2(
373 &IndexSpec {
374 column: "ts".to_string(),
375 entity_columns: Vec::new(),
376 kind: IndexKind::UInt64 {
377 index_granularity: NonZeroU64::new(6).unwrap(),
378 },
379 },
380 bytes,
381 );
382 let different_bytes = segment_coverage_id_v2(&base_index, b"other");
383
384 assert_ne!(
385 base, different_granularity,
386 "index granularity should affect id"
387 );
388 assert_ne!(base, different_column, "index column should affect id");
389 assert_ne!(base, different_kind, "index kind should affect id");
390 assert_ne!(different_kind, different_integer_domain);
391 assert_ne!(different_kind, different_integer_granularity);
392 assert_ne!(base, different_bytes, "coverage bytes should affect id");
393 }
394
395 #[test]
396 fn table_coverage_id_matches_golden_value_and_is_valid() {
397 let index = timestamp_index("ts", TimeIndexGranularity::Hours(1));
398 let bytes = b"table-bitmap";
399
400 let id1 = table_coverage_id_v2(&index, bytes);
401 let id2 = table_coverage_id_v2(&index, bytes);
402
403 assert_eq!(id1, "tblcov-38f0aa9c3e526d0cdabf234af8fb0fd3");
404 assert_eq!(id1, id2, "same inputs must produce stable id");
405 assert!(id1.starts_with("tblcov-"));
406 assert_eq!(id1.len(), "tblcov-".len() + 32, "prefix + 32 hex chars");
407 validate_coverage_id(&id1).expect("derived id should be valid");
408 }
409
410 #[test]
411 fn table_coverage_id_changes_with_inputs() {
412 let bytes = b"bytes";
413
414 let base_index = timestamp_index("ts", TimeIndexGranularity::Minutes(15));
415 let base = table_coverage_id_v2(&base_index, bytes);
416 let different_granularity =
417 table_coverage_id_v2(×tamp_index("ts", TimeIndexGranularity::Days(1)), bytes);
418 let different_column = table_coverage_id_v2(
419 ×tamp_index("event_time", TimeIndexGranularity::Minutes(15)),
420 bytes,
421 );
422 let different_bytes = table_coverage_id_v2(&base_index, b"other");
423
424 assert_ne!(
425 base, different_granularity,
426 "index granularity should affect id"
427 );
428 assert_ne!(base, different_column, "index column should affect id");
429 assert_ne!(base, different_bytes, "coverage bytes should affect id");
430 }
431
432 #[test]
433 fn entity_coverage_ids_match_golden_values_and_include_ordered_columns() {
434 let index = IndexSpec {
435 column: "ts".to_string(),
436 entity_columns: vec!["symbol".to_string(), "venue".to_string()],
437 kind: IndexKind::Timestamp {
438 index_granularity: TimeIndexGranularity::Minutes(1),
439 timezone: None,
440 },
441 };
442 let mut renamed = index.clone();
443 renamed.entity_columns[0] = "device".to_string();
444 let mut reordered = index.clone();
445 reordered.entity_columns.reverse();
446 let bytes = b"entity-coverage-bytes";
447
448 let segment = segment_entity_coverage_id_v1(&index, bytes);
449 assert_eq!(segment, "segcov-67c0022aad0d9f5bf5ea813e9ef88119");
450 assert_ne!(segment, segment_entity_coverage_id_v1(&renamed, bytes));
451 assert_ne!(segment, segment_entity_coverage_id_v1(&reordered, bytes));
452
453 let table = table_entity_coverage_id_v1(&index, bytes);
454 assert_eq!(table, "tblcov-9c54647467c3a0e89e60675e00b7c75b");
455 assert_ne!(table, table_entity_coverage_id_v1(&renamed, bytes));
456 assert_ne!(table, table_entity_coverage_id_v1(&reordered, bytes));
457 }
458
459 #[test]
460 fn coverage_file_ids_are_owned_by_the_append_attempt() {
461 let content_id = "segcov-0123456789abcdef0123456789abcdef";
462 let first = coverage_file_id_for_attempt(content_id, &Uuid::from_u128(1));
463 let second = coverage_file_id_for_attempt(content_id, &Uuid::from_u128(2));
464
465 assert_ne!(first, second);
466 validate_coverage_id(&first).expect("first id should be valid");
467 validate_coverage_id(&second).expect("second id should be valid");
468 }
469}