1pub mod index_interval;
4pub(crate) mod io;
5pub(crate) mod layout;
6pub(crate) mod serde;
7
8pub use crate::formats::parquet::SegmentCoverageError;
9pub use io::CoverageSidecarError;
10pub use layout::CoverageLayoutError;
11pub use serde::CoverageCodecError;
12
13use std::{
14 collections::{BTreeMap, btree_map},
15 ops::RangeInclusive,
16};
17
18use ::serde::{Deserialize, Deserializer, Serialize, Serializer, de::Error as _};
19use snafu::Snafu;
20
21pub use roaring::RoaringTreemap;
22
23pub type IndexIntervalId = u64;
25
26#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
28#[serde(
29 tag = "type",
30 content = "value",
31 rename_all = "lowercase",
32 deny_unknown_fields
33)]
34pub enum EntityValue {
35 Utf8(String),
37 Int32(i32),
39 Int64(i64),
41 UInt64(u64),
43}
44
45impl From<String> for EntityValue {
46 fn from(value: String) -> Self {
47 Self::Utf8(value)
48 }
49}
50
51impl From<&str> for EntityValue {
52 fn from(value: &str) -> Self {
53 Self::Utf8(value.to_string())
54 }
55}
56
57impl From<i32> for EntityValue {
58 fn from(value: i32) -> Self {
59 Self::Int32(value)
60 }
61}
62
63impl From<i64> for EntityValue {
64 fn from(value: i64) -> Self {
65 Self::Int64(value)
66 }
67}
68
69impl From<u64> for EntityValue {
70 fn from(value: u64) -> Self {
71 Self::UInt64(value)
72 }
73}
74
75#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
80pub struct EntityIdentity {
81 components: Vec<EntityValue>,
82}
83
84impl Serialize for EntityIdentity {
85 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
86 where
87 S: Serializer,
88 {
89 self.components.serialize(serializer)
90 }
91}
92
93impl<'de> Deserialize<'de> for EntityIdentity {
94 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
95 where
96 D: Deserializer<'de>,
97 {
98 let components = Vec::<EntityValue>::deserialize(deserializer)?;
99 Self::try_new(components).map_err(D::Error::custom)
100 }
101}
102
103impl EntityIdentity {
104 pub fn try_new(components: Vec<EntityValue>) -> Result<Self, EntityIdentityError> {
109 if components.is_empty() {
110 return Err(EntityIdentityError::Empty);
111 }
112 Ok(Self { components })
113 }
114
115 pub fn components(&self) -> &[EntityValue] {
117 &self.components
118 }
119}
120
121#[derive(Debug, Clone, PartialEq, Eq, Snafu)]
123#[non_exhaustive]
124pub enum EntityIdentityError {
125 #[snafu(display("entity identity must contain at least one component"))]
127 Empty,
128}
129
130#[derive(Debug, Clone, Default, PartialEq, Eq)]
132pub struct Coverage {
133 present: RoaringTreemap,
134}
135
136impl Coverage {
137 pub fn empty() -> Self {
139 Self::default()
140 }
141
142 pub fn from_treemap(present: RoaringTreemap) -> Self {
144 Self { present }
145 }
146
147 pub fn present(&self) -> &RoaringTreemap {
149 &self.present
150 }
151
152 pub fn into_treemap(self) -> RoaringTreemap {
154 self.present
155 }
156
157 pub fn union(&self, other: &Self) -> Self {
159 Self::from_treemap(&self.present | &other.present)
160 }
161
162 pub fn union_inplace(&mut self, other: &Self) {
164 self.present |= other.present();
165 }
166
167 pub fn intersect(&self, other: &Self) -> Self {
169 Self::from_treemap(&self.present & &other.present)
170 }
171
172 pub fn intersection_cardinality(&self, other: &Self) -> u64 {
174 self.present.intersection_len(&other.present)
175 }
176
177 pub fn cardinality(&self) -> u64 {
179 self.present.len()
180 }
181
182 pub fn is_empty(&self) -> bool {
184 self.present.is_empty()
185 }
186
187 pub fn range_cardinality(range: &RangeInclusive<IndexIntervalId>) -> u128 {
189 if range.is_empty() {
190 return 0;
191 }
192 u128::from(*range.end()) - u128::from(*range.start()) + 1
193 }
194
195 pub fn covered_cardinality(&self, range: &RangeInclusive<IndexIntervalId>) -> u64 {
197 if range.is_empty() {
198 0
199 } else {
200 self.present.range_cardinality(range.clone())
201 }
202 }
203
204 pub fn missing_runs(
210 &self,
211 range: &RangeInclusive<IndexIntervalId>,
212 max_run_len: Option<u64>,
213 ) -> Vec<RangeInclusive<IndexIntervalId>> {
214 if range.is_empty() || max_run_len == Some(0) {
215 return Vec::new();
216 }
217
218 let start = *range.start();
219 let end = *range.end();
220 let mut cursor = Some(start);
221 let mut runs = Vec::new();
222 let mut present = self.present.iter();
223 present.advance_to(start);
224
225 for index_interval_id in present {
226 if index_interval_id > end {
227 break;
228 }
229 let Some(missing_start) = cursor else {
230 break;
231 };
232 if missing_start < index_interval_id {
233 runs.push(missing_start..=index_interval_id - 1);
234 }
235 cursor = index_interval_id.checked_add(1);
236 }
237
238 if let Some(missing_start) = cursor.filter(|value| *value <= end) {
239 runs.push(missing_start..=end);
240 }
241
242 match max_run_len {
243 Some(max_len) => split_runs_by_len(runs, max_len),
244 None => runs,
245 }
246 }
247
248 pub fn last_run_with_min_len(
250 &self,
251 range: &RangeInclusive<IndexIntervalId>,
252 min_len: u64,
253 ) -> Option<RangeInclusive<IndexIntervalId>> {
254 if range.is_empty() || min_len == 0 {
255 return None;
256 }
257
258 let mut iter = self.present.iter();
259 iter.advance_to(*range.start());
260 iter.advance_back_to(*range.end());
261
262 let mut current: Option<(IndexIntervalId, IndexIntervalId)> = None;
263 let mut last = None;
264 for index_interval_id in iter {
265 match current {
266 Some((start, end)) if end.checked_add(1) == Some(index_interval_id) => {
267 current = Some((start, index_interval_id));
268 }
269 Some((start, end)) => {
270 if inclusive_len(start, end) >= u128::from(min_len) {
271 last = Some(start..=end);
272 }
273 current = Some((index_interval_id, index_interval_id));
274 }
275 None => current = Some((index_interval_id, index_interval_id)),
276 }
277 }
278
279 if let Some((start, end)) = current
280 && inclusive_len(start, end) >= u128::from(min_len)
281 {
282 last = Some(start..=end);
283 }
284 last
285 }
286
287 pub fn coverage_ratio(&self, range: &RangeInclusive<IndexIntervalId>) -> f64 {
289 let expected = Self::range_cardinality(range);
290 if expected == 0 {
291 return 1.0;
292 }
293 self.covered_cardinality(range) as f64 / expected as f64
294 }
295
296 pub fn max_gap_len(&self, range: &RangeInclusive<IndexIntervalId>) -> u128 {
298 self.missing_runs(range, None)
299 .into_iter()
300 .map(|run| inclusive_len(*run.start(), *run.end()))
301 .max()
302 .unwrap_or(0)
303 }
304
305 pub fn last_window_at_or_before(
307 &self,
308 end_index_interval_id: IndexIntervalId,
309 len: u64,
310 ) -> Option<RangeInclusive<IndexIntervalId>> {
311 if len == 0 {
312 return None;
313 }
314
315 let mut iter = self.present.iter();
316 iter.advance_back_to(end_index_interval_id);
317 let mut run_end = None;
318 let mut previous: Option<IndexIntervalId> = None;
319 let mut run_len = 0u64;
320
321 for index_interval_id in iter.rev() {
322 if previous.and_then(|value| value.checked_sub(1)) == Some(index_interval_id) {
323 run_len += 1;
324 } else {
325 run_end = Some(index_interval_id);
326 run_len = 1;
327 }
328 previous = Some(index_interval_id);
329
330 if run_len >= len {
331 return run_end.map(|end| index_interval_id..=end);
332 }
333 }
334 None
335 }
336}
337
338#[derive(Debug, Clone, Default, PartialEq, Eq)]
343pub struct EntityCoverage {
344 by_identity: BTreeMap<EntityIdentity, Coverage>,
345}
346
347impl EntityCoverage {
348 pub fn empty() -> Self {
350 Self::default()
351 }
352
353 pub fn get(&self, identity: &EntityIdentity) -> Option<&Coverage> {
358 self.by_identity.get(identity)
359 }
360
361 pub fn iter(&self) -> btree_map::Iter<'_, EntityIdentity, Coverage> {
363 self.by_identity.iter()
364 }
365
366 pub fn identity_count(&self) -> usize {
368 self.by_identity.len()
369 }
370
371 pub fn is_empty(&self) -> bool {
373 self.by_identity.is_empty()
374 }
375
376 pub fn union_coverage(&mut self, identity: EntityIdentity, coverage: Coverage) {
378 self.by_identity
379 .entry(identity)
380 .and_modify(|current| current.union_inplace(&coverage))
381 .or_insert(coverage);
382 }
383
384 pub fn union(&self, other: &Self) -> Self {
386 let mut union = self.clone();
387 union.union_inplace(other);
388 union
389 }
390
391 pub fn union_inplace(&mut self, other: &Self) {
393 for (identity, coverage) in other.iter() {
394 if let Some(current) = self.by_identity.get_mut(identity) {
395 current.union_inplace(coverage);
396 } else {
397 self.by_identity.insert(identity.clone(), coverage.clone());
398 }
399 }
400 }
401
402 pub fn intersect(&self, other: &Self) -> Self {
404 let mut intersection = Self::empty();
405 for (identity, coverage) in self.iter() {
406 if let Some(other_coverage) = other.get(identity) {
407 intersection.union_coverage(identity.clone(), coverage.intersect(other_coverage));
408 }
409 }
410 intersection
411 }
412
413 pub fn cardinality(&self) -> u128 {
415 self.by_identity
416 .values()
417 .map(|coverage| u128::from(coverage.cardinality()))
418 .sum()
419 }
420
421 pub fn intersection_cardinality(&self, other: &Self) -> u128 {
423 self.iter()
424 .filter_map(|(identity, coverage)| {
425 other.get(identity).map(|other_coverage| {
426 u128::from(coverage.intersection_cardinality(other_coverage))
427 })
428 })
429 .sum()
430 }
431
432 pub fn first_overlapping_identity_and_interval_id<'a>(
434 &'a self,
435 other: &Self,
436 ) -> Option<(&'a EntityIdentity, IndexIntervalId)> {
437 self.iter().find_map(|(identity, coverage)| {
438 let other_coverage = other.get(identity)?;
439 if coverage.present().is_disjoint(other_coverage.present()) {
440 return None;
441 }
442 coverage
443 .intersect(other_coverage)
444 .present()
445 .min()
446 .map(|index_interval_id| (identity, index_interval_id))
447 })
448 }
449}
450
451impl FromIterator<IndexIntervalId> for Coverage {
452 fn from_iter<I>(iter: I) -> Self
453 where
454 I: IntoIterator<Item = IndexIntervalId>,
455 {
456 Self::from_treemap(iter.into_iter().collect())
457 }
458}
459
460fn inclusive_len(start: IndexIntervalId, end: IndexIntervalId) -> u128 {
461 u128::from(end) - u128::from(start) + 1
462}
463
464fn split_runs_by_len(
465 runs: Vec<RangeInclusive<IndexIntervalId>>,
466 max_len: u64,
467) -> Vec<RangeInclusive<IndexIntervalId>> {
468 if max_len == 0 {
469 return Vec::new();
470 }
471
472 let mut out = Vec::new();
473 for range in runs {
474 let mut start = *range.start();
475 let end = *range.end();
476 loop {
477 let chunk_end = start.saturating_add(max_len - 1).min(end);
478 out.push(start..=chunk_end);
479 if chunk_end == end {
480 break;
481 }
482 start = chunk_end + 1;
483 }
484 }
485 out
486}
487
488#[cfg(test)]
489mod tests {
490 use super::*;
491
492 #[test]
493 fn basic_set_operations_use_u64_domain() {
494 let a: Coverage = [0, u64::from(u32::MAX) + 1, u64::MAX].into_iter().collect();
495 let b: Coverage = [u64::from(u32::MAX) + 1, 7].into_iter().collect();
496
497 assert_eq!(a.cardinality(), 3);
498 assert_eq!(a.intersection_cardinality(&b), 1);
499 assert_eq!(a.intersect(&b).cardinality(), 1);
500 assert_eq!(a.union(&b).cardinality(), 4);
501 }
502
503 #[test]
504 fn sparse_huge_ranges_do_not_require_expected_bitmap() {
505 let coverage: Coverage = [0, 2, u64::MAX].into_iter().collect();
506 let range = 0..=u64::MAX;
507
508 assert_eq!(Coverage::range_cardinality(&range), 1u128 << 64);
509 assert_eq!(coverage.covered_cardinality(&range), 3);
510 assert_eq!(
511 coverage.missing_runs(&range, None),
512 vec![1..=1, 3..=u64::MAX - 1]
513 );
514 assert_eq!(coverage.max_gap_len(&range), u128::from(u64::MAX) - 3);
515 }
516
517 #[test]
518 fn missing_runs_split_without_enumerating_missing_points() {
519 let coverage: Coverage = [2, 7].into_iter().collect();
520 assert_eq!(
521 coverage.missing_runs(&(0..=9), Some(2)),
522 vec![0..=1, 3..=4, 5..=6, 8..=9]
523 );
524 }
525
526 #[test]
527 fn range_queries_respect_requested_bounds() {
528 let coverage: Coverage = [2, 3, 4, 7, 8, 9].into_iter().collect();
529 let range = 1..=8;
530
531 assert_eq!(coverage.covered_cardinality(&range), 5);
532 assert_eq!(coverage.missing_runs(&range, None), vec![1..=1, 5..=6]);
533 assert_eq!(coverage.max_gap_len(&range), 2);
534 assert_eq!(coverage.last_run_with_min_len(&range, 2), Some(7..=8));
535 assert_eq!(coverage.coverage_ratio(&range), 5.0 / 8.0);
536 }
537
538 #[test]
539 fn last_window_handles_zero_and_u64_max() {
540 let coverage: Coverage = [0, 1, u64::MAX - 2, u64::MAX - 1, u64::MAX]
541 .into_iter()
542 .collect();
543
544 assert_eq!(
545 coverage.last_window_at_or_before(u64::MAX, 3),
546 Some(u64::MAX - 2..=u64::MAX)
547 );
548 assert_eq!(coverage.last_window_at_or_before(1, 2), Some(0..=1));
549 assert_eq!(coverage.last_window_at_or_before(u64::MAX, 0), None);
550 }
551
552 fn identity(components: &[&str]) -> EntityIdentity {
553 EntityIdentity::try_new(
554 components
555 .iter()
556 .map(|component| EntityValue::from(*component))
557 .collect(),
558 )
559 .unwrap()
560 }
561
562 #[test]
563 fn entity_identity_preserves_component_order() {
564 assert_eq!(
565 identity(&["venue", "symbol"]).components(),
566 &[EntityValue::from("venue"), EntityValue::from("symbol")]
567 );
568 assert!(identity(&["A", "Z"]) < identity(&["B", "A"]));
569 assert!(identity(&["A", "A"]) < identity(&["A", "Z"]));
570 assert_ne!(identity(&["A", "A"]), identity(&["A", "Z"]));
571 assert_eq!(
572 EntityIdentity::try_new(Vec::new()),
573 Err(EntityIdentityError::Empty)
574 );
575 }
576
577 #[test]
578 fn entity_identity_preserves_scalar_types_in_json() {
579 let identity = EntityIdentity::try_new(vec![
580 EntityValue::from("sensor"),
581 EntityValue::Int32(-1),
582 EntityValue::Int64(i64::MIN),
583 EntityValue::UInt64(u64::MAX),
584 ])
585 .unwrap();
586
587 let json = serde_json::to_value(&identity).unwrap();
588 assert_eq!(
589 json,
590 serde_json::json!([
591 { "type": "utf8", "value": "sensor" },
592 { "type": "int32", "value": -1 },
593 { "type": "int64", "value": i64::MIN },
594 { "type": "uint64", "value": u64::MAX },
595 ])
596 );
597 assert_eq!(
598 serde_json::from_value::<EntityIdentity>(json).unwrap(),
599 identity
600 );
601 assert_ne!(EntityValue::Int32(1), EntityValue::Int64(1));
602 assert_ne!(EntityValue::Int64(1), EntityValue::UInt64(1));
603 }
604
605 #[test]
606 fn entity_coverage_unions_only_matching_identities() {
607 let a = identity(&["A"]);
608 let b = identity(&["B"]);
609 let c = identity(&["C"]);
610 let mut left = EntityCoverage::empty();
611 left.union_coverage(a.clone(), [1, 2].into_iter().collect());
612 left.union_coverage(b.clone(), [1].into_iter().collect());
613
614 let mut right = EntityCoverage::empty();
615 right.union_coverage(a.clone(), [2, 3].into_iter().collect());
616 right.union_coverage(c.clone(), [4].into_iter().collect());
617
618 let union = left.union(&right);
619 assert_eq!(union.identity_count(), 3);
620 assert_eq!(union.cardinality(), 5);
621 assert_eq!(
622 union.get(&a).unwrap().present().iter().collect::<Vec<_>>(),
623 vec![1, 2, 3]
624 );
625 assert_eq!(
626 union.get(&b).unwrap().present().iter().collect::<Vec<_>>(),
627 vec![1]
628 );
629 assert_eq!(
630 union.get(&c).unwrap().present().iter().collect::<Vec<_>>(),
631 vec![4]
632 );
633 }
634
635 #[test]
636 fn entity_coverage_intersection_requires_identity_and_interval() {
637 let a = identity(&["A"]);
638 let b = identity(&["B"]);
639 let mut left = EntityCoverage::empty();
640 left.union_coverage(a.clone(), [1, 2].into_iter().collect());
641 left.union_coverage(b.clone(), [7].into_iter().collect());
642
643 let mut right = EntityCoverage::empty();
644 right.union_coverage(a.clone(), [2, 7].into_iter().collect());
645 right.union_coverage(b.clone(), [1].into_iter().collect());
646
647 let intersection = left.intersect(&right);
648 assert_eq!(intersection.identity_count(), 2);
649 assert_eq!(intersection.cardinality(), 1);
650 assert_eq!(left.intersection_cardinality(&right), 1);
651 assert_eq!(
652 intersection
653 .get(&a)
654 .unwrap()
655 .present()
656 .iter()
657 .collect::<Vec<_>>(),
658 vec![2]
659 );
660 assert!(intersection.get(&b).unwrap().is_empty());
661 }
662
663 #[test]
664 fn entity_coverage_counts_same_interval_once_per_identity() {
665 let mut coverage = EntityCoverage::empty();
666 coverage.union_coverage(identity(&["A"]), [u64::MAX].into_iter().collect());
667 coverage.union_coverage(identity(&["B"]), [u64::MAX].into_iter().collect());
668
669 assert_eq!(coverage.cardinality(), 2);
670 }
671
672 #[test]
673 fn first_overlapping_identity_and_interval_id_is_deterministic() {
674 let first = identity(&["A", "one"]);
675 let later = identity(&["B", "one"]);
676 let mut left = EntityCoverage::empty();
677 left.union_coverage(later.clone(), [1].into_iter().collect());
678 left.union_coverage(first.clone(), [9, 3].into_iter().collect());
679
680 let mut right = EntityCoverage::empty();
681 right.union_coverage(later, [1].into_iter().collect());
682 right.union_coverage(first.clone(), [3, 9].into_iter().collect());
683
684 assert_eq!(
685 left.first_overlapping_identity_and_interval_id(&right),
686 Some((&first, 3))
687 );
688 }
689
690 #[test]
691 fn first_overlapping_identity_and_interval_id_handles_dense_intervals() {
692 let entity = identity(&["dense"]);
693 let last = u64::from(u32::MAX);
694 let mut dense = RoaringTreemap::new();
695 dense.insert_range(0..=last);
696
697 let mut left = EntityCoverage::empty();
698 left.union_coverage(entity.clone(), Coverage::from_treemap(dense));
699 let mut right = EntityCoverage::empty();
700 right.union_coverage(entity.clone(), [last].into_iter().collect());
701
702 assert_eq!(
703 left.first_overlapping_identity_and_interval_id(&right),
704 Some((&entity, last))
705 );
706 }
707
708 #[test]
709 fn explicit_empty_entity_coverage_is_preserved() {
710 let entity = identity(&["empty"]);
711 let mut coverage = EntityCoverage::empty();
712 coverage.union_coverage(entity.clone(), Coverage::empty());
713
714 assert!(!coverage.is_empty());
715 assert!(coverage.get(&entity).unwrap().is_empty());
716 assert!(coverage.get(&identity(&["absent"])).is_none());
717 assert_eq!(coverage.identity_count(), 1);
718 }
719}