Skip to main content

timeseries_table_format/
coverage.rs

1//! In-memory coverage and gap analysis over the 64-bit index interval ID domain.
2
3pub 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
23/// Ordered 64-bit index interval ID.
24pub type IndexIntervalId = u64;
25
26/// Exact scalar value in an entity identity.
27#[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    /// UTF-8 text from an Arrow `Utf8` or `LargeUtf8` column.
36    Utf8(String),
37    /// Signed 32-bit integer.
38    Int32(i32),
39    /// Signed 64-bit integer.
40    Int64(i64),
41    /// Unsigned 64-bit integer.
42    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/// Ordered composite entity identity.
76///
77/// Component positions correspond to the table's configured entity-column
78/// order. Column names are deliberately not repeated in every identity.
79#[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    /// Construct an identity from at least one ordered component.
105    ///
106    /// # Errors
107    /// Returns [`EntityIdentityError::Empty`] when `components` is empty.
108    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    /// Borrow components in configured entity-column order.
116    pub fn components(&self) -> &[EntityValue] {
117        &self.components
118    }
119}
120
121/// Invalid entity identity construction.
122#[derive(Debug, Clone, PartialEq, Eq, Snafu)]
123#[non_exhaustive]
124pub enum EntityIdentityError {
125    /// An entity-aware table requires at least one identity component.
126    #[snafu(display("entity identity must contain at least one component"))]
127    Empty,
128}
129
130/// In-memory coverage over a discrete set of index interval IDs.
131#[derive(Debug, Clone, Default, PartialEq, Eq)]
132pub struct Coverage {
133    present: RoaringTreemap,
134}
135
136impl Coverage {
137    /// Construct an empty coverage set.
138    pub fn empty() -> Self {
139        Self::default()
140    }
141
142    /// Wrap an existing treemap.
143    pub fn from_treemap(present: RoaringTreemap) -> Self {
144        Self { present }
145    }
146
147    /// Borrow the present index interval IDs.
148    pub fn present(&self) -> &RoaringTreemap {
149        &self.present
150    }
151
152    /// Consume the coverage and return its treemap.
153    pub fn into_treemap(self) -> RoaringTreemap {
154        self.present
155    }
156
157    /// Return the union of two coverage sets.
158    pub fn union(&self, other: &Self) -> Self {
159        Self::from_treemap(&self.present | &other.present)
160    }
161
162    /// Merge another coverage set into this one.
163    pub fn union_inplace(&mut self, other: &Self) {
164        self.present |= other.present();
165    }
166
167    /// Return the intersection of two coverage sets.
168    pub fn intersect(&self, other: &Self) -> Self {
169        Self::from_treemap(&self.present & &other.present)
170    }
171
172    /// Count index interval IDs present in both sets without materializing them.
173    pub fn intersection_cardinality(&self, other: &Self) -> u64 {
174        self.present.intersection_len(&other.present)
175    }
176
177    /// Number of present index interval IDs.
178    pub fn cardinality(&self) -> u64 {
179        self.present.len()
180    }
181
182    /// Whether no index interval IDs are present.
183    pub fn is_empty(&self) -> bool {
184        self.present.is_empty()
185    }
186
187    /// Number of index interval IDs in an inclusive range.
188    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    /// Count present index interval IDs in a range without materializing it.
196    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    /// Return missing contiguous runs in an inclusive requested range.
205    ///
206    /// Long runs are optionally split into chunks of at most `max_run_len`.
207    /// Work is proportional to present index interval IDs and returned runs, not the size
208    /// of the requested range.
209    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    /// Return the last covered contiguous run of at least `min_len` IDs.
249    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    /// Coverage ratio in `[0.0, 1.0]` for an inclusive range.
288    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    /// Length of the largest missing run in an inclusive range.
297    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    /// Return the last fully-covered window ending at or before an index interval ID.
306    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/// Independent index interval IDs for each ordered entity identity.
339///
340/// Explicit identities with empty coverage are preserved. An absent identity
341/// is still treated as empty and returned as `None` by [`EntityCoverage::get`].
342#[derive(Debug, Clone, Default, PartialEq, Eq)]
343pub struct EntityCoverage {
344    by_identity: BTreeMap<EntityIdentity, Coverage>,
345}
346
347impl EntityCoverage {
348    /// Construct empty entity-scoped coverage.
349    pub fn empty() -> Self {
350        Self::default()
351    }
352
353    /// Borrow one identity's coverage.
354    ///
355    /// `None` means the identity is absent. Set operations treat absence as
356    /// empty coverage.
357    pub fn get(&self, identity: &EntityIdentity) -> Option<&Coverage> {
358        self.by_identity.get(identity)
359    }
360
361    /// Iterate identities and their coverage in canonical order.
362    pub fn iter(&self) -> btree_map::Iter<'_, EntityIdentity, Coverage> {
363        self.by_identity.iter()
364    }
365
366    /// Number of stored identities.
367    pub fn identity_count(&self) -> usize {
368        self.by_identity.len()
369    }
370
371    /// Whether no identities are stored.
372    pub fn is_empty(&self) -> bool {
373        self.by_identity.is_empty()
374    }
375
376    /// Merge coverage for one identity.
377    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    /// Return the union of two entity-scoped coverage values.
385    pub fn union(&self, other: &Self) -> Self {
386        let mut union = self.clone();
387        union.union_inplace(other);
388        union
389    }
390
391    /// Merge another entity-scoped coverage value into this one.
392    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    /// Return overlap only where both identity and index interval ID match.
403    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    /// Count covered `(entity identity, index interval ID)` pairs.
414    pub fn cardinality(&self) -> u128 {
415        self.by_identity
416            .values()
417            .map(|coverage| u128::from(coverage.cardinality()))
418            .sum()
419    }
420
421    /// Count overlapping `(entity identity, index interval ID)` pairs.
422    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    /// Return the first canonical identity and smallest overlapping index interval ID.
433    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}