Skip to main content

timeseries_table_format/
coverage.rs

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