1use std::{collections::BTreeMap, path::Path};
11
12use log::warn;
13
14use crate::{
15 coverage::{
16 Coverage, EntityCoverage, EntityIdentity, EntityValue,
17 io::{CoverageError, read_coverage_sidecar, read_entity_coverage_sidecar},
18 },
19 metadata::schema_compat::{ensure_entity_identity_matches_schema, require_table_schema},
20 transaction_log::table_state::TableCoveragePointer,
21};
22
23use super::{TimeSeriesTable, error::TableError};
24
25fn ensure_entity_coverage_identity_schema(
26 coverage: &EntityCoverage,
27 table: &TimeSeriesTable,
28) -> Result<(), CoverageError> {
29 let schema = require_table_schema(&table.state().table_meta)
30 .map_err(|source| CoverageError::EntityIdentitySchema { source })?;
31 for (identity, _) in coverage.iter() {
32 ensure_entity_identity_matches_schema(schema, table.index_spec(), identity)
33 .map_err(|source| CoverageError::EntityIdentitySchema { source })?;
34 }
35 Ok(())
36}
37
38impl TimeSeriesTable {
39 fn ensure_global_coverage_query(&self) -> Result<(), TableError> {
40 if self.index_spec().entity_columns.is_empty() {
41 Ok(())
42 } else {
43 Err(TableError::EntityIdentityRequired {
44 entity_columns: self.index_spec().entity_columns.clone(),
45 })
46 }
47 }
48
49 fn resolve_entity_identity(
50 &self,
51 components: &[(&str, EntityValue)],
52 ) -> Result<EntityIdentity, TableError> {
53 let entity_columns = &self.index_spec().entity_columns;
54 if entity_columns.is_empty() {
55 return Err(TableError::EntityIdentityNotConfigured);
56 }
57
58 let mut provided = BTreeMap::new();
59 for (column, value) in components {
60 if !entity_columns.iter().any(|expected| expected == *column) {
61 return Err(TableError::UnexpectedEntityIdentityColumn {
62 column: (*column).to_string(),
63 });
64 }
65 if provided.insert(*column, value).is_some() {
66 return Err(TableError::DuplicateEntityIdentityColumn {
67 column: (*column).to_string(),
68 });
69 }
70 }
71
72 let ordered = entity_columns
73 .iter()
74 .map(|column| {
75 provided
76 .get(column.as_str())
77 .map(|value| (**value).clone())
78 .ok_or_else(|| TableError::MissingEntityIdentityColumn {
79 column: column.clone(),
80 })
81 })
82 .collect::<Result<Vec<_>, _>>()?;
83
84 let identity = EntityIdentity::try_new(ordered)
85 .map_err(|_| TableError::EntityIdentityNotConfigured)?;
86 let schema = self.state().table_meta.logical_schema.as_ref().ok_or(
87 TableError::MissingCanonicalSchema {
88 version: self.state().version,
89 },
90 )?;
91 ensure_entity_identity_matches_schema(schema, self.index_spec(), &identity)
92 .map_err(|source| TableError::SchemaCompatibility { source })?;
93 Ok(identity)
94 }
95
96 async fn read_validated_entity_coverage_sidecar(
97 &self,
98 path: &Path,
99 ) -> Result<EntityCoverage, CoverageError> {
100 let coverage = read_entity_coverage_sidecar(self.location(), path).await?;
101 ensure_entity_coverage_identity_schema(&coverage, self)?;
102 Ok(coverage)
103 }
104
105 pub(crate) async fn recover_table_coverage_from_segments(
110 &self,
111 ) -> Result<Coverage, TableError> {
112 let mut acc = Coverage::empty();
113
114 for seg in self.state().segments.values() {
115 let path = seg.coverage_path.as_ref().ok_or_else(|| {
116 TableError::ExistingSegmentMissingCoverage {
117 path: seg.path.clone(),
118 }
119 })?;
120
121 let cov = read_coverage_sidecar(self.location(), Path::new(path))
122 .await
123 .map_err(|source| TableError::SegmentCoverageSidecarRead {
124 path: seg.path.clone(),
125 coverage_path: path.clone(),
126 source: Box::new(source),
127 })?;
128
129 acc.union_inplace(&cov);
131 }
132
133 Ok(acc)
134 }
135
136 pub(crate) async fn recover_table_entity_coverage_from_segments(
137 &self,
138 ) -> Result<EntityCoverage, TableError> {
139 let mut acc = EntityCoverage::empty();
140
141 for seg in self.state().segments.values() {
142 let path = seg.coverage_path.as_ref().ok_or_else(|| {
143 TableError::ExistingSegmentMissingCoverage {
144 path: seg.path.clone(),
145 }
146 })?;
147
148 let coverage = self
149 .read_validated_entity_coverage_sidecar(Path::new(path))
150 .await
151 .map_err(|source| TableError::SegmentCoverageSidecarRead {
152 path: seg.path.clone(),
153 coverage_path: path.clone(),
154 source: Box::new(source),
155 })?;
156 acc.union_inplace(&coverage);
157 }
158
159 Ok(acc)
160 }
161
162 fn ensure_table_coverage_index_matches(
163 &self,
164 ptr: &TableCoveragePointer,
165 ) -> Result<(), TableError> {
166 let expected = self.index_spec().kind.clone();
167 if ptr.index_kind != expected {
168 return Err(TableError::TableCoverageIndexKindMismatch {
169 expected,
170 actual: ptr.index_kind.clone(),
171 pointer_version: ptr.version,
172 });
173 }
174 Ok(())
175 }
176
177 pub async fn load_table_coverage_snapshot_only(&self) -> Result<Coverage, TableError> {
184 match &self.state().table_coverage {
185 None => {
186 if self.state().segments.is_empty() {
187 return Ok(Coverage::empty());
188 }
189 Err(TableError::MissingTableCoveragePointer)
190 }
191 Some(ptr) => {
192 self.ensure_table_coverage_index_matches(ptr)?;
193 read_coverage_sidecar(self.location(), Path::new(&ptr.coverage_path))
194 .await
195 .map_err(|source| TableError::CoverageSidecar { source })
196 }
197 }
198 }
199
200 pub(crate) async fn load_table_entity_coverage_snapshot_only(
201 &self,
202 ) -> Result<EntityCoverage, TableError> {
203 match &self.state().table_coverage {
204 None => {
205 if self.state().segments.is_empty() {
206 return Ok(EntityCoverage::empty());
207 }
208 Err(TableError::MissingTableCoveragePointer)
209 }
210 Some(ptr) => {
211 self.ensure_table_coverage_index_matches(ptr)?;
212 self.read_validated_entity_coverage_sidecar(Path::new(&ptr.coverage_path))
213 .await
214 .map_err(|source| TableError::CoverageSidecar { source })
215 }
216 }
217 }
218
219 pub(crate) async fn load_table_snapshot_coverage_readonly(
227 &self,
228 ) -> Result<Coverage, TableError> {
229 match &self.state().table_coverage {
230 None => {
231 if self.state().segments.is_empty() {
232 return Ok(Coverage::empty());
233 }
234 self.recover_table_coverage_from_segments().await
235 }
236 Some(ptr) => {
237 self.ensure_table_coverage_index_matches(ptr)?;
238
239 match read_coverage_sidecar(self.location(), Path::new(&ptr.coverage_path)).await {
240 Ok(cov) => Ok(cov),
241 Err(snapshot_err) => {
242 warn!(
243 "Failed to read table coverage snapshot at {} (version {}): {:?}. \
244 Attempting recovery from segment sidecars (readonly).",
245 ptr.coverage_path, ptr.version, snapshot_err
246 );
247 self.recover_table_coverage_from_segments().await
248 }
249 }
250 }
251 }
252 }
253
254 pub(crate) async fn load_table_entity_snapshot_coverage_readonly(
255 &self,
256 ) -> Result<EntityCoverage, TableError> {
257 match &self.state().table_coverage {
258 None => {
259 if self.state().segments.is_empty() {
260 return Ok(EntityCoverage::empty());
261 }
262 self.recover_table_entity_coverage_from_segments().await
263 }
264 Some(ptr) => {
265 self.ensure_table_coverage_index_matches(ptr)?;
266 match self.load_table_entity_coverage_snapshot_only().await {
267 Ok(coverage) => Ok(coverage),
268 Err(snapshot_err) => {
269 warn!(
270 "Failed to read entity coverage snapshot at {} (version {}): {:?}. \
271 Attempting recovery from segment sidecars (readonly).",
272 ptr.coverage_path, ptr.version, snapshot_err
273 );
274 self.recover_table_entity_coverage_from_segments().await
275 }
276 }
277 }
278 }
279 }
280}
281use std::ops::RangeInclusive;
288
289use snafu::ResultExt;
290
291use crate::{
292 coverage::Bucket,
293 coverage::bucket::{bucket_for_exclusive_end, bucket_range},
294 metadata::table_metadata::{IndexValue, validate_index_range},
295 table::error::{CoverageBucketSnafu, InvalidRangeSnafu},
296};
297
298impl TimeSeriesTable {
299 fn bucket_range_for_index_range<S, E>(
300 &self,
301 start: S,
302 end: E,
303 ) -> Result<RangeInclusive<Bucket>, TableError>
304 where
305 S: Into<IndexValue>,
306 E: Into<IndexValue>,
307 {
308 let start = start.into();
309 let end = end.into();
310 validate_index_range(&self.index_spec().kind, &start, &end).context(InvalidRangeSnafu)?;
311 bucket_range(&self.index_spec().kind, &start, &end).context(CoverageBucketSnafu)
312 }
313
314 pub async fn coverage_ratio_for_range<S, E>(&self, start: S, end: E) -> Result<f64, TableError>
341 where
342 S: Into<IndexValue>,
343 E: Into<IndexValue>,
344 {
345 self.ensure_global_coverage_query()?;
346 let range = self.bucket_range_for_index_range(start, end)?;
347 let cov = self.load_table_snapshot_coverage_readonly().await?;
348 Ok(cov.coverage_ratio(&range))
349 }
350
351 pub async fn coverage_ratio_for_entity_range<S, E>(
382 &self,
383 entity: &[(&str, EntityValue)],
384 start: S,
385 end: E,
386 ) -> Result<f64, TableError>
387 where
388 S: Into<IndexValue>,
389 E: Into<IndexValue>,
390 {
391 let identity = self.resolve_entity_identity(entity)?;
392 let range = self.bucket_range_for_index_range(start, end)?;
393 let coverage = self.load_table_entity_snapshot_coverage_readonly().await?;
394 Ok(coverage
395 .get(&identity)
396 .map_or(0.0, |coverage| coverage.coverage_ratio(&range)))
397 }
398
399 pub async fn max_gap_len_for_range<S, E>(&self, start: S, end: E) -> Result<u128, TableError>
423 where
424 S: Into<IndexValue>,
425 E: Into<IndexValue>,
426 {
427 self.ensure_global_coverage_query()?;
428 let range = self.bucket_range_for_index_range(start, end)?;
429 let cov = self.load_table_snapshot_coverage_readonly().await?;
430 Ok(cov.max_gap_len(&range))
431 }
432
433 pub async fn max_gap_len_for_entity_range<S, E>(
446 &self,
447 entity: &[(&str, EntityValue)],
448 start: S,
449 end: E,
450 ) -> Result<u128, TableError>
451 where
452 S: Into<IndexValue>,
453 E: Into<IndexValue>,
454 {
455 let identity = self.resolve_entity_identity(entity)?;
456 let range = self.bucket_range_for_index_range(start, end)?;
457 let coverage = self.load_table_entity_snapshot_coverage_readonly().await?;
458 Ok(coverage.get(&identity).map_or_else(
459 || Coverage::range_cardinality(&range),
460 |c| c.max_gap_len(&range),
461 ))
462 }
463
464 pub async fn last_fully_covered_window<E>(
490 &self,
491 end: E,
492 window_len_buckets: u64,
493 ) -> Result<Option<RangeInclusive<Bucket>>, TableError>
494 where
495 E: Into<IndexValue>,
496 {
497 self.ensure_global_coverage_query()?;
498 let end = end.into();
499 end.validate_kind(&self.index_spec().kind)
500 .context(InvalidRangeSnafu)?;
501 if window_len_buckets == 0 {
502 return Ok(None);
503 }
504
505 let end_bucket =
506 bucket_for_exclusive_end(&self.index_spec().kind, &end).context(CoverageBucketSnafu)?;
507 let cov = self.load_table_snapshot_coverage_readonly().await?;
508 Ok(cov.last_window_at_or_before(end_bucket, window_len_buckets))
509 }
510
511 pub async fn last_fully_covered_window_for_entity<E>(
525 &self,
526 entity: &[(&str, EntityValue)],
527 end: E,
528 window_len_buckets: u64,
529 ) -> Result<Option<RangeInclusive<Bucket>>, TableError>
530 where
531 E: Into<IndexValue>,
532 {
533 let identity = self.resolve_entity_identity(entity)?;
534 let end = end.into();
535 end.validate_kind(&self.index_spec().kind)
536 .context(InvalidRangeSnafu)?;
537 if window_len_buckets == 0 {
538 return Ok(None);
539 }
540
541 let end_bucket =
542 bucket_for_exclusive_end(&self.index_spec().kind, &end).context(CoverageBucketSnafu)?;
543 let coverage = self.load_table_entity_snapshot_coverage_readonly().await?;
544 Ok(coverage
545 .get(&identity)
546 .and_then(|coverage| coverage.last_window_at_or_before(end_bucket, window_len_buckets)))
547 }
548}
549
550#[cfg(test)]
551mod tests {
552 use std::{num::NonZeroU64, path::Path};
553
554 use super::*;
555 use crate::{
556 coverage::{
557 Coverage, EntityCoverage, EntityIdentity, EntityValue,
558 bucket::{BucketError, bucket_id},
559 io::{write_coverage_sidecar_atomic, write_coverage_sidecar_new_bytes},
560 serde::entity_coverage_to_bytes,
561 },
562 metadata::logical_schema::{LogicalDataType, LogicalField, LogicalSchema},
563 metadata::schema_compat::SchemaCompatibilityError,
564 metadata::table_metadata::{
565 IndexKind, IndexSpec, IndexValueError, TableKind, TableMeta, TimeBucket,
566 },
567 storage::TableLocation,
568 table::test_util::{
569 TestResult, TestRow, make_basic_table_meta, make_int32_entity_table_meta, utc_datetime,
570 write_int32_entity_parquet, write_test_parquet,
571 },
572 };
573 use chrono::{DateTime, TimeZone, Utc};
574 use tempfile::TempDir;
575
576 type HelperResult<T> = Result<T, Box<dyn std::error::Error>>;
577
578 fn ts_from_secs(secs: i64) -> DateTime<Utc> {
579 Utc.timestamp_opt(secs, 0)
580 .single()
581 .expect("valid timestamp")
582 }
583
584 async fn make_table() -> HelperResult<(TempDir, TimeSeriesTable)> {
585 let tmp = TempDir::new()?;
586 let location = TableLocation::local(tmp.path());
587 let mut meta = make_basic_table_meta();
588 let TableKind::TimeSeries(index) = &mut meta.kind else {
589 unreachable!("test metadata is time series")
590 };
591 index.entity_columns.clear();
592 let table = TimeSeriesTable::create(location, meta).await?;
593 Ok((tmp, table))
594 }
595
596 async fn table_with_index_coverage(
597 kind: IndexKind,
598 coverage: Coverage,
599 ) -> HelperResult<(TempDir, TimeSeriesTable)> {
600 let tmp = TempDir::new()?;
601 let mut table = TimeSeriesTable::create(
602 TableLocation::local(tmp.path()),
603 TableMeta::new_time_series(IndexSpec {
604 column: "index".to_string(),
605 entity_columns: Vec::new(),
606 kind: kind.clone(),
607 }),
608 )
609 .await?;
610 let coverage_path = "_coverage/table/query-test.roar";
611 write_coverage_sidecar_atomic(table.location(), Path::new(coverage_path), &coverage)
612 .await?;
613 let version = table.state().version;
614 table.state_mut().table_coverage = Some(TableCoveragePointer {
615 index_kind: kind,
616 coverage_path: coverage_path.to_string(),
617 version,
618 });
619 Ok((tmp, table))
620 }
621
622 async fn append_segment(
623 table: &mut TimeSeriesTable,
624 tmp: &TempDir,
625 rel_path: &str,
626 rows: &[TestRow],
627 ) -> HelperResult<()> {
628 let abs = tmp.path().join(rel_path);
629 write_test_parquet(&abs, true, false, rows)?;
630 table.append_parquet_segment(rel_path).await?;
631 Ok(())
632 }
633
634 #[tokio::test]
635 async fn entity_sidecar_identity_schema_is_validated_for_snapshots_and_recovery() -> TestResult
636 {
637 let tmp = TempDir::new()?;
638 let location = TableLocation::local(tmp.path());
639 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
640 let mut wrong_arity = EntityCoverage::empty();
641 wrong_arity.union_coverage(
642 EntityIdentity::try_new(vec!["A".into(), "X".into()])?,
643 Coverage::from_iter([0]),
644 );
645 let wrong_arity_bytes = entity_coverage_to_bytes(&wrong_arity)?;
646 let snapshot_path = "_coverage/table/wrong-arity.roar";
647 write_coverage_sidecar_new_bytes(&location, Path::new(snapshot_path), &wrong_arity_bytes)
648 .await?;
649 let version = table.state().version;
650 let index_kind = table.index_spec().kind.clone();
651 table.state_mut().table_coverage = Some(TableCoveragePointer {
652 index_kind,
653 coverage_path: snapshot_path.to_string(),
654 version,
655 });
656
657 let snapshot_error = table
658 .load_table_entity_coverage_snapshot_only()
659 .await
660 .expect_err("snapshot identity arity must match the table");
661 assert!(matches!(
662 snapshot_error,
663 TableError::CoverageSidecar {
664 source: CoverageError::EntityIdentitySchema {
665 source: SchemaCompatibilityError::EntityIdentityArityMismatch {
666 expected: 1,
667 actual: 2,
668 },
669 },
670 }
671 ));
672
673 let mut wrong_type = EntityCoverage::empty();
674 wrong_type.union_coverage(
675 EntityIdentity::try_new(vec![EntityValue::Int32(1)])?,
676 Coverage::from_iter([0]),
677 );
678 let wrong_type_bytes = entity_coverage_to_bytes(&wrong_type)?;
679 table.state_mut().table_coverage = None;
680 let segment_path = "data/entity-arity.parquet";
681 append_segment(
682 &mut table,
683 &tmp,
684 segment_path,
685 &[TestRow {
686 ts_millis: 1_000,
687 symbol: "A",
688 price: 1.0,
689 }],
690 )
691 .await?;
692 let segment_coverage_path = table
693 .state()
694 .segments
695 .get(segment_path)
696 .and_then(|segment| segment.coverage_path.clone())
697 .expect("segment coverage path");
698 tokio::fs::write(tmp.path().join(&segment_coverage_path), &wrong_type_bytes).await?;
699
700 let recovery_error = table
701 .recover_table_entity_coverage_from_segments()
702 .await
703 .expect_err("segment identity type must match the table schema");
704 assert!(matches!(
705 recovery_error,
706 TableError::SegmentCoverageSidecarRead {
707 path,
708 coverage_path,
709 source,
710 } if path == segment_path
711 && coverage_path == segment_coverage_path
712 && matches!(
713 &*source,
714 CoverageError::EntityIdentitySchema {
715 source: SchemaCompatibilityError::EntityIdentityTypeMismatch {
716 column,
717 expected: crate::metadata::logical_schema::LogicalDataType::Utf8,
718 actual: "int32",
719 },
720 }
721 if column == "symbol"
722 )
723 ));
724 Ok(())
725 }
726
727 #[tokio::test]
728 async fn entity_coverage_ratio_selects_only_the_requested_identity() -> TestResult {
729 let tmp = TempDir::new()?;
730 let location = TableLocation::local(tmp.path());
731 let mut table = TimeSeriesTable::create(location, make_basic_table_meta()).await?;
732 append_segment(
733 &mut table,
734 &tmp,
735 "data/entity-ratio.parquet",
736 &[
737 TestRow {
738 ts_millis: 1_000,
739 symbol: "A",
740 price: 1.0,
741 },
742 TestRow {
743 ts_millis: 61_000,
744 symbol: "B",
745 price: 2.0,
746 },
747 ],
748 )
749 .await?;
750 let state_before = table.state().clone();
751 let start = ts_from_secs(0);
752 let end = ts_from_secs(120);
753
754 assert_eq!(
755 table
756 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end)
757 .await?,
758 0.5
759 );
760 assert_eq!(
761 table
762 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("B"))], start, end)
763 .await?,
764 0.5
765 );
766 assert_eq!(
767 table
768 .coverage_ratio_for_entity_range(
769 &[("symbol", EntityValue::from("unseen"))],
770 start,
771 end,
772 )
773 .await?,
774 0.0
775 );
776 assert!(matches!(
777 table.coverage_ratio_for_range(start, end).await,
778 Err(TableError::EntityIdentityRequired { entity_columns })
779 if entity_columns == ["symbol"]
780 ));
781 assert_eq!(table.state(), &state_before);
782 Ok(())
783 }
784
785 #[tokio::test]
786 async fn numeric_entity_queries_require_the_exact_scalar_type() -> TestResult {
787 let tmp = TempDir::new()?;
788 let mut table = TimeSeriesTable::create(
789 TableLocation::local(tmp.path()),
790 make_int32_entity_table_meta(),
791 )
792 .await?;
793 let path = "data/numeric-coverage.parquet";
794 write_int32_entity_parquet(
795 &tmp.path().join(path),
796 &[1_000, 61_000],
797 &[-1, i32::MAX],
798 &[10.0, 20.0],
799 )?;
800 table.append_parquet_segment(path).await?;
801 let start = ts_from_secs(0);
802 let end = ts_from_secs(120);
803
804 assert_eq!(
805 table
806 .coverage_ratio_for_entity_range(
807 &[("device_id", EntityValue::Int32(-1))],
808 start,
809 end,
810 )
811 .await?,
812 0.5
813 );
814 assert_eq!(
815 table
816 .coverage_ratio_for_entity_range(
817 &[("device_id", EntityValue::Int32(42))],
818 start,
819 end,
820 )
821 .await?,
822 0.0
823 );
824 assert!(matches!(
825 table
826 .coverage_ratio_for_entity_range(
827 &[("device_id", EntityValue::Int64(-1))],
828 start,
829 end,
830 )
831 .await,
832 Err(TableError::SchemaCompatibility {
833 source: SchemaCompatibilityError::EntityIdentityTypeMismatch {
834 column,
835 expected: LogicalDataType::Int32,
836 actual: "int64",
837 },
838 }) if column == "device_id"
839 ));
840 assert!(matches!(
841 table
842 .coverage_ratio_for_entity_range(
843 &[("device_id", EntityValue::from("-1"))],
844 start,
845 end,
846 )
847 .await,
848 Err(TableError::SchemaCompatibility {
849 source: SchemaCompatibilityError::EntityIdentityTypeMismatch {
850 column,
851 expected: LogicalDataType::Int32,
852 actual: "utf8",
853 },
854 }) if column == "device_id"
855 ));
856 Ok(())
857 }
858
859 #[tokio::test]
860 async fn entity_gap_and_window_queries_are_isolated_and_recover_readonly() -> TestResult {
861 let tmp = TempDir::new()?;
862 let mut table =
863 TimeSeriesTable::create(TableLocation::local(tmp.path()), make_basic_table_meta())
864 .await?;
865 append_segment(
866 &mut table,
867 &tmp,
868 "data/entity-gaps.parquet",
869 &[
870 TestRow {
871 ts_millis: 1_000,
872 symbol: "A",
873 price: 1.0,
874 },
875 TestRow {
876 ts_millis: 181_000,
877 symbol: "A",
878 price: 2.0,
879 },
880 TestRow {
881 ts_millis: 1_000,
882 symbol: "B",
883 price: 3.0,
884 },
885 TestRow {
886 ts_millis: 61_000,
887 symbol: "B",
888 price: 4.0,
889 },
890 TestRow {
891 ts_millis: 121_000,
892 symbol: "B",
893 price: 5.0,
894 },
895 ],
896 )
897 .await?;
898 let start = ts_from_secs(0);
899 let end = ts_from_secs(240);
900
901 assert_eq!(
902 table
903 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end)
904 .await?,
905 0.5
906 );
907 assert_eq!(
908 table
909 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("B"))], start, end)
910 .await?,
911 0.75
912 );
913 assert_eq!(
914 table
915 .max_gap_len_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end)
916 .await?,
917 2
918 );
919 assert_eq!(
920 table
921 .max_gap_len_for_entity_range(&[("symbol", EntityValue::from("B"))], start, end)
922 .await?,
923 1
924 );
925 assert_eq!(
926 table
927 .last_fully_covered_window_for_entity(
928 &[("symbol", EntityValue::from("A"))],
929 end,
930 2,
931 )
932 .await?,
933 None
934 );
935 assert_eq!(
936 table
937 .last_fully_covered_window_for_entity(
938 &[("symbol", EntityValue::from("B"))],
939 end,
940 2,
941 )
942 .await?,
943 Some(0x8000_0000_0000_0001..=0x8000_0000_0000_0002)
944 );
945
946 let snapshot_path = table
947 .state()
948 .table_coverage
949 .as_ref()
950 .expect("snapshot pointer")
951 .coverage_path
952 .clone();
953 let state_before = table.state().clone();
954 tokio::fs::remove_file(tmp.path().join(snapshot_path)).await?;
955
956 assert_eq!(
957 table
958 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end)
959 .await?,
960 0.5
961 );
962 assert_eq!(
963 table
964 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("B"))], start, end)
965 .await?,
966 0.75
967 );
968 assert_eq!(
969 table
970 .max_gap_len_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end)
971 .await?,
972 2
973 );
974 assert_eq!(
975 table
976 .last_fully_covered_window_for_entity(
977 &[("symbol", EntityValue::from("B"))],
978 end,
979 2,
980 )
981 .await?,
982 Some(0x8000_0000_0000_0001..=0x8000_0000_0000_0002)
983 );
984 assert_eq!(
985 table
986 .max_gap_len_for_entity_range(
987 &[("symbol", EntityValue::from("unseen"))],
988 start,
989 end,
990 )
991 .await?,
992 4
993 );
994 assert_eq!(
995 table
996 .last_fully_covered_window_for_entity(
997 &[("symbol", EntityValue::from("unseen"))],
998 end,
999 1,
1000 )
1001 .await?,
1002 None
1003 );
1004 assert_eq!(
1005 table
1006 .last_fully_covered_window_for_entity(
1007 &[("symbol", EntityValue::from("unseen"))],
1008 end,
1009 0,
1010 )
1011 .await?,
1012 None
1013 );
1014
1015 assert!(matches!(
1016 table.coverage_ratio_for_range(start, end).await,
1017 Err(TableError::EntityIdentityRequired { entity_columns })
1018 if entity_columns == ["symbol"]
1019 ));
1020 assert!(matches!(
1021 table.max_gap_len_for_range(start, end).await,
1022 Err(TableError::EntityIdentityRequired { entity_columns })
1023 if entity_columns == ["symbol"]
1024 ));
1025 assert!(matches!(
1026 table.last_fully_covered_window(end, 0).await,
1027 Err(TableError::EntityIdentityRequired { entity_columns })
1028 if entity_columns == ["symbol"]
1029 ));
1030 assert_eq!(table.state(), &state_before);
1031 Ok(())
1032 }
1033
1034 #[tokio::test]
1035 async fn entity_identity_input_is_validated_and_canonicalized() -> TestResult {
1036 let tmp = TempDir::new()?;
1037 let mut meta = make_basic_table_meta();
1038 let TableKind::TimeSeries(index) = &mut meta.kind else {
1039 unreachable!("test metadata is time series")
1040 };
1041 index.entity_columns = vec!["symbol".to_string(), "venue".to_string()];
1042 let mut fields = meta
1043 .logical_schema
1044 .as_ref()
1045 .expect("test schema")
1046 .columns()
1047 .to_vec();
1048 fields.push(LogicalField {
1049 name: "venue".to_string(),
1050 data_type: LogicalDataType::Utf8,
1051 nullable: false,
1052 });
1053 meta.logical_schema = Some(LogicalSchema::new(fields)?);
1054 let table = TimeSeriesTable::create(TableLocation::local(tmp.path()), meta).await?;
1055
1056 let identity = table.resolve_entity_identity(&[
1057 ("venue", EntityValue::from("X")),
1058 ("symbol", EntityValue::from("A")),
1059 ])?;
1060 assert_eq!(
1061 identity.components(),
1062 [EntityValue::from("A"), EntityValue::from("X")]
1063 );
1064 let start = ts_from_secs(0);
1065 let end = ts_from_secs(60);
1066 assert!(matches!(
1067 table
1068 .coverage_ratio_for_entity_range(&[("venue", EntityValue::from("X"))], start, end)
1069 .await,
1070 Err(TableError::MissingEntityIdentityColumn { column }) if column == "symbol"
1071 ));
1072 assert!(matches!(
1073 table
1074 .coverage_ratio_for_entity_range(
1075 &[
1076 ("device", EntityValue::from("A")),
1077 ("venue", EntityValue::from("X")),
1078 ],
1079 start,
1080 end,
1081 )
1082 .await,
1083 Err(TableError::UnexpectedEntityIdentityColumn { column }) if column == "device"
1084 ));
1085 assert!(matches!(
1086 table
1087 .coverage_ratio_for_entity_range(
1088 &[
1089 ("symbol", EntityValue::from("A")),
1090 ("symbol", EntityValue::from("B")),
1091 ("venue", EntityValue::from("X")),
1092 ],
1093 start,
1094 end,
1095 )
1096 .await,
1097 Err(TableError::DuplicateEntityIdentityColumn { column }) if column == "symbol"
1098 ));
1099
1100 let (_tmp, global_table) = make_table().await?;
1101 assert!(matches!(
1102 global_table
1103 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end)
1104 .await,
1105 Err(TableError::EntityIdentityNotConfigured)
1106 ));
1107 Ok(())
1108 }
1109
1110 async fn table_with_sparse_coverage() -> HelperResult<(TempDir, TimeSeriesTable)> {
1111 let (tmp, mut table) = make_table().await?;
1113 append_segment(
1114 &mut table,
1115 &tmp,
1116 "data/sparse.parquet",
1117 &[
1118 TestRow {
1119 ts_millis: 1_000,
1120 symbol: "A",
1121 price: 1.0,
1122 },
1123 TestRow {
1124 ts_millis: 61_000,
1125 symbol: "A",
1126 price: 2.0,
1127 },
1128 TestRow {
1129 ts_millis: 180_000,
1130 symbol: "A",
1131 price: 3.0,
1132 },
1133 ],
1134 )
1135 .await?;
1136 Ok((tmp, table))
1137 }
1138
1139 async fn table_with_contiguous_run() -> HelperResult<(TempDir, TimeSeriesTable)> {
1140 let (tmp, mut table) = make_table().await?;
1142 append_segment(
1143 &mut table,
1144 &tmp,
1145 "data/window.parquet",
1146 &[
1147 TestRow {
1148 ts_millis: 240_000,
1149 symbol: "A",
1150 price: 1.0,
1151 },
1152 TestRow {
1153 ts_millis: 300_000,
1154 symbol: "A",
1155 price: 2.0,
1156 },
1157 ],
1158 )
1159 .await?;
1160 Ok((tmp, table))
1161 }
1162
1163 #[tokio::test]
1164 async fn bucket_range_rejects_invalid_range() -> TestResult {
1165 let (_tmp, table) = make_table().await?;
1166 let ts = utc_datetime(2024, 1, 1, 0, 0, 0);
1167
1168 let err = table
1169 .bucket_range_for_index_range(ts, ts)
1170 .expect_err("start >= end should be invalid");
1171 assert!(matches!(err, TableError::InvalidRange { .. }));
1172 Ok(())
1173 }
1174
1175 #[tokio::test]
1176 async fn bucket_range_uses_64_bit_signed_timestamp_mapping() -> TestResult {
1177 let (_tmp, table) = make_table().await?;
1178 let start = ts_from_secs(0);
1179 let end = ts_from_secs(180); let range = table.bucket_range_for_index_range(start, end)?;
1182 assert_eq!(range, 0x8000_0000_0000_0000..=0x8000_0000_0000_0002);
1183 Ok(())
1184 }
1185
1186 #[tokio::test]
1187 async fn signed_coverage_queries_handle_gaps_extremes_and_last_window() -> TestResult {
1188 let kind = IndexKind::Int64 {
1189 bucket_width: NonZeroU64::new(10).unwrap(),
1190 };
1191 let coverage: Coverage = [-10i64, 0, 10]
1192 .into_iter()
1193 .map(|value| bucket_id(&kind, &value.into()).unwrap())
1194 .collect();
1195 let huge_gap = u128::from(
1196 bucket_id(&kind, &(-10i64).into()).unwrap()
1197 - bucket_id(&kind, &i64::MIN.into()).unwrap(),
1198 )
1199 .max(u128::from(
1200 bucket_for_exclusive_end(&kind, &i64::MAX.into()).unwrap()
1201 - bucket_id(&kind, &10i64.into()).unwrap(),
1202 ));
1203 let (_tmp, table) = table_with_index_coverage(kind, coverage).await?;
1204
1205 assert_eq!(table.coverage_ratio_for_range(-20i64, 30i64).await?, 0.6);
1206 assert_eq!(table.max_gap_len_for_range(-20i64, 30i64).await?, 1);
1207 assert_eq!(table.max_gap_len_for_range(-10i64, 0i64).await?, 0);
1208 assert_eq!(table.max_gap_len_for_range(-50i64, -20i64).await?, 3);
1209 assert_eq!(
1210 table.max_gap_len_for_range(i64::MIN, i64::MAX).await?,
1211 huge_gap
1212 );
1213
1214 let window = table
1215 .last_fully_covered_window(10i64, 2)
1216 .await?
1217 .expect("signed window across zero");
1218 assert_eq!(
1219 window,
1220 bucket_id(&table.index_spec().kind, &(-10i64).into()).unwrap()
1221 ..=bucket_id(&table.index_spec().kind, &0i64.into()).unwrap()
1222 );
1223 Ok(())
1224 }
1225
1226 #[tokio::test]
1227 async fn unsigned_coverage_queries_preserve_large_values_and_boundaries() -> TestResult {
1228 let kind = IndexKind::UInt64 {
1229 bucket_width: NonZeroU64::new(1).unwrap(),
1230 };
1231 let start = i64::MAX as u64 + 1;
1232 let coverage: Coverage = [start, start + 1, u64::MAX - 2, u64::MAX - 1]
1233 .into_iter()
1234 .collect();
1235 let (_tmp, table) = table_with_index_coverage(kind, coverage).await?;
1236
1237 let requested = u128::from(u64::MAX) - u128::from(start);
1238 let ratio = table.coverage_ratio_for_range(start, u64::MAX).await?;
1239 assert!((ratio - 4.0 / requested as f64).abs() < f64::EPSILON);
1240 assert_eq!(
1241 table.max_gap_len_for_range(start, u64::MAX).await?,
1242 u128::from(u64::MAX) - u128::from(start) - 4
1243 );
1244 assert_eq!(
1245 table.last_fully_covered_window(start + 2, 2).await?,
1246 Some(start..=start + 1)
1247 );
1248 assert_eq!(
1249 table.last_fully_covered_window(u64::MAX, 2).await?,
1250 Some(u64::MAX - 2..=u64::MAX - 1)
1251 );
1252 Ok(())
1253 }
1254
1255 #[tokio::test]
1256 async fn coverage_ratio_uses_snapshot_when_present() -> TestResult {
1257 let (_tmp, table) = table_with_sparse_coverage().await?;
1258 let start = ts_from_secs(0);
1259 let end = ts_from_secs(240); let ratio = table.coverage_ratio_for_range(start, end).await?;
1262 assert!((ratio - 0.75).abs() < 1e-12);
1263 Ok(())
1264 }
1265
1266 #[tokio::test]
1267 async fn coverage_ratio_recovers_when_snapshot_missing() -> TestResult {
1268 let (_tmp, mut table) = table_with_sparse_coverage().await?;
1269 table.state_mut().table_coverage = None;
1270
1271 let ratio = table
1272 .coverage_ratio_for_range(ts_from_secs(0), ts_from_secs(240))
1273 .await?;
1274 assert!((ratio - 0.75).abs() < 1e-12);
1275 Ok(())
1276 }
1277
1278 #[tokio::test]
1279 async fn coverage_ratio_errors_when_recovery_missing_segment_coverage_path() -> TestResult {
1280 let (_tmp, mut table) = table_with_sparse_coverage().await?;
1281 table.state_mut().table_coverage = None;
1282 let segment = table
1283 .state_mut()
1284 .segments
1285 .values_mut()
1286 .next()
1287 .expect("segment present");
1288 let segment_path = segment.path.clone();
1289 segment.coverage_path = None;
1290
1291 let err = table
1292 .coverage_ratio_for_range(ts_from_secs(0), ts_from_secs(240))
1293 .await
1294 .expect_err("missing segment coverage_path should bubble up");
1295 assert!(matches!(
1296 err,
1297 TableError::ExistingSegmentMissingCoverage { path } if path == segment_path
1298 ));
1299 Ok(())
1300 }
1301
1302 #[tokio::test]
1303 async fn coverage_ratio_errors_on_bucket_mismatch() -> TestResult {
1304 let (_tmp, mut table) = table_with_sparse_coverage().await?;
1305 let mut ptr = table
1306 .state()
1307 .table_coverage
1308 .clone()
1309 .expect("snapshot pointer present");
1310 ptr.index_kind = IndexKind::Timestamp {
1311 bucket: TimeBucket::Hours(1),
1312 timezone: None,
1313 };
1314 table.state_mut().table_coverage = Some(ptr.clone());
1315
1316 let err = table
1317 .coverage_ratio_for_range(ts_from_secs(0), ts_from_secs(240))
1318 .await
1319 .expect_err("mismatched bucket spec should error");
1320
1321 match err {
1322 TableError::TableCoverageIndexKindMismatch {
1323 expected, actual, ..
1324 } => {
1325 assert_eq!(expected, table.index_spec().kind);
1326 assert_eq!(actual, ptr.index_kind);
1327 }
1328 other => panic!("unexpected error: {other:?}"),
1329 }
1330 Ok(())
1331 }
1332
1333 #[tokio::test]
1334 async fn coverage_ratio_handles_empty_table() -> TestResult {
1335 let (_tmp, table) = make_table().await?;
1336 let ratio = table
1337 .coverage_ratio_for_range(ts_from_secs(0), ts_from_secs(60))
1338 .await?;
1339 assert_eq!(ratio, 0.0);
1340 Ok(())
1341 }
1342
1343 #[tokio::test]
1344 async fn coverage_ratio_handles_bucket_ids_above_u32() -> TestResult {
1345 let (_tmp, table) = make_table().await?;
1346 let start = ts_from_secs(0);
1347 let end = ts_from_secs(((u32::MAX as i64) + 3) * 60);
1348
1349 let ratio = table.coverage_ratio_for_range(start, end).await?;
1350 assert_eq!(ratio, 0.0);
1351 Ok(())
1352 }
1353
1354 #[tokio::test]
1355 async fn max_gap_len_reports_missing_run() -> TestResult {
1356 let (_tmp, table) = table_with_sparse_coverage().await?;
1357 let gap = table
1358 .max_gap_len_for_range(ts_from_secs(0), ts_from_secs(240))
1359 .await?;
1360 assert_eq!(gap, 1);
1361 Ok(())
1362 }
1363
1364 #[tokio::test]
1365 async fn last_window_returns_none_for_zero_length() -> TestResult {
1366 let (_tmp, table) = make_table().await?;
1367 let res = table.last_fully_covered_window(ts_from_secs(0), 0).await?;
1368 assert!(res.is_none());
1369 Ok(())
1370 }
1371
1372 #[tokio::test]
1373 async fn last_window_respects_half_open_end_and_run_length() -> TestResult {
1374 let (_tmp, table) = table_with_contiguous_run().await?;
1375 let ts_end = ts_from_secs(360); let win = table
1378 .last_fully_covered_window(ts_end, 2)
1379 .await?
1380 .expect("window should be present");
1381 assert_eq!(win, 0x8000_0000_0000_0004..=0x8000_0000_0000_0005);
1382
1383 let none = table.last_fully_covered_window(ts_end, 3).await?;
1384 assert!(none.is_none());
1385 Ok(())
1386 }
1387
1388 #[tokio::test]
1389 async fn coverage_queries_validate_before_reading_coverage() -> TestResult {
1390 let tmp = TempDir::new()?;
1391 let kind = IndexKind::UInt64 {
1392 bucket_width: NonZeroU64::new(1).unwrap(),
1393 };
1394 let meta = TableMeta::new_time_series(IndexSpec {
1395 column: "offset".to_string(),
1396 entity_columns: Vec::new(),
1397 kind: kind.clone(),
1398 });
1399 let mut table = TimeSeriesTable::create(TableLocation::local(tmp.path()), meta).await?;
1400 let version = table.state().version;
1401 table.state_mut().table_coverage = Some(TableCoveragePointer {
1402 index_kind: kind,
1403 coverage_path: "_coverage/table/missing.roar".to_string(),
1404 version,
1405 });
1406
1407 let error = table
1408 .last_fully_covered_window(ts_from_secs(1), 1)
1409 .await
1410 .expect_err("endpoint domain must match the table index");
1411 assert!(matches!(
1412 error,
1413 TableError::InvalidRange {
1414 source: IndexValueError::KindMismatch {
1415 expected: "uint64",
1416 actual: "timestamp"
1417 }
1418 }
1419 ));
1420
1421 assert!(matches!(
1422 table.coverage_ratio_for_range(1u64, 1u64).await,
1423 Err(TableError::InvalidRange {
1424 source: IndexValueError::InvalidRange { .. }
1425 })
1426 ));
1427 assert!(matches!(
1428 table.max_gap_len_for_range(2u64, 1u64).await,
1429 Err(TableError::InvalidRange {
1430 source: IndexValueError::InvalidRange { .. }
1431 })
1432 ));
1433 assert!(matches!(
1434 table.coverage_ratio_for_range(0u64, 1i64).await,
1435 Err(TableError::InvalidRange {
1436 source: IndexValueError::KindMismatch { .. }
1437 })
1438 ));
1439 assert_eq!(table.last_fully_covered_window(0u64, 0).await?, None);
1440 assert!(matches!(
1441 table.last_fully_covered_window(0u64, 1).await,
1442 Err(TableError::CoverageBucket {
1443 source: BucketError::RangeEndUnderflow { .. }
1444 })
1445 ));
1446 Ok(())
1447 }
1448
1449 #[tokio::test]
1450 async fn last_window_errors_when_recovery_fails() -> TestResult {
1451 let (_tmp, mut table) = table_with_contiguous_run().await?;
1452 table.state_mut().table_coverage = None;
1453 let segment = table
1454 .state_mut()
1455 .segments
1456 .values_mut()
1457 .next()
1458 .expect("segment present");
1459 let segment_path = segment.path.clone();
1460 segment.coverage_path = None;
1461
1462 let err = table
1463 .last_fully_covered_window(ts_from_secs(360), 1)
1464 .await
1465 .expect_err("missing coverage_path should bubble up");
1466 assert!(matches!(
1467 err,
1468 TableError::ExistingSegmentMissingCoverage { path } if path == segment_path
1469 ));
1470 Ok(())
1471 }
1472}