1pub 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
18pub type Bucket = u64;
20
21#[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 Utf8(String),
32 Int32(i32),
34 Int64(i64),
36 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#[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 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 pub fn components(&self) -> &[EntityValue] {
112 &self.components
113 }
114}
115
116#[derive(Debug, Clone, PartialEq, Eq, Snafu)]
118pub enum EntityIdentityError {
119 #[snafu(display("entity identity must contain at least one component"))]
121 Empty,
122}
123
124#[derive(Debug, Clone, Default, PartialEq, Eq)]
126pub struct Coverage {
127 present: RoaringTreemap,
128}
129
130impl Coverage {
131 pub fn empty() -> Self {
133 Self::default()
134 }
135
136 pub fn from_treemap(present: RoaringTreemap) -> Self {
138 Self { present }
139 }
140
141 pub fn present(&self) -> &RoaringTreemap {
143 &self.present
144 }
145
146 pub fn into_treemap(self) -> RoaringTreemap {
148 self.present
149 }
150
151 pub fn union(&self, other: &Self) -> Self {
153 Self::from_treemap(&self.present | &other.present)
154 }
155
156 pub fn union_inplace(&mut self, other: &Self) {
158 self.present |= other.present();
159 }
160
161 pub fn intersect(&self, other: &Self) -> Self {
163 Self::from_treemap(&self.present & &other.present)
164 }
165
166 pub fn intersection_cardinality(&self, other: &Self) -> u64 {
168 self.present.intersection_len(&other.present)
169 }
170
171 pub fn cardinality(&self) -> u64 {
173 self.present.len()
174 }
175
176 pub fn is_empty(&self) -> bool {
178 self.present.is_empty()
179 }
180
181 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 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 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 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 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 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 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#[derive(Debug, Clone, Default, PartialEq, Eq)]
337pub struct EntityCoverage {
338 by_identity: BTreeMap<EntityIdentity, Coverage>,
339}
340
341impl EntityCoverage {
342 pub fn empty() -> Self {
344 Self::default()
345 }
346
347 pub fn get(&self, identity: &EntityIdentity) -> Option<&Coverage> {
352 self.by_identity.get(identity)
353 }
354
355 pub fn iter(&self) -> btree_map::Iter<'_, EntityIdentity, Coverage> {
357 self.by_identity.iter()
358 }
359
360 pub fn identity_count(&self) -> usize {
362 self.by_identity.len()
363 }
364
365 pub fn is_empty(&self) -> bool {
367 self.by_identity.is_empty()
368 }
369
370 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 pub fn union(&self, other: &Self) -> Self {
380 let mut union = self.clone();
381 union.union_inplace(other);
382 union
383 }
384
385 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 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 pub fn cardinality(&self) -> u128 {
409 self.by_identity
410 .values()
411 .map(|coverage| u128::from(coverage.cardinality()))
412 .sum()
413 }
414
415 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 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}