1mod error;
11
12pub use error::CoverageQueryError;
13
14use std::{collections::BTreeMap, path::Path};
15
16use snafu::{Backtrace, OptionExt, ResultExt};
17
18use crate::{
19 coverage::{
20 Coverage, EntityCoverage, EntityIdentity, EntityValue,
21 io::{CoverageSidecarError, read_coverage_sidecar, read_entity_coverage_sidecar},
22 },
23 metadata::schema_compat::{ensure_entity_identity_matches_schema, require_table_schema},
24 table::{AppendError, TableError, TimeSeriesTable, error::CoverageQuerySnafu},
25};
26
27use self::error as query_error;
28
29pub(crate) trait CoverageRecoveryError: Sized {
30 fn missing_segment_coverage_path(segment_path: String) -> Self;
31
32 fn segment_coverage_read_failed(
33 segment_path: String,
34 coverage_path: String,
35 source: CoverageSidecarError,
36 ) -> Self;
37}
38
39impl CoverageRecoveryError for AppendError {
40 fn missing_segment_coverage_path(segment_path: String) -> Self {
41 Self::ExistingSegmentMissingCoverageMetadata { segment_path }
42 }
43
44 fn segment_coverage_read_failed(
45 segment_path: String,
46 coverage_path: String,
47 source: CoverageSidecarError,
48 ) -> Self {
49 Self::ExistingSegmentCoverageSidecarRead {
50 segment_path,
51 coverage_path,
52 source: Box::new(source),
53 }
54 }
55}
56
57impl CoverageRecoveryError for CoverageQueryError {
58 fn missing_segment_coverage_path(segment_path: String) -> Self {
59 Self::ExistingSegmentMissingCoverageMetadata {
60 segment_path,
61 backtrace: Box::new(Backtrace::capture()),
62 }
63 }
64
65 fn segment_coverage_read_failed(
66 segment_path: String,
67 coverage_path: String,
68 source: CoverageSidecarError,
69 ) -> Self {
70 Self::SegmentCoverageSidecarRead {
71 segment_path,
72 coverage_path,
73 source: Box::new(source),
74 }
75 }
76}
77
78fn ensure_entity_coverage_identity_schema(
79 coverage: &EntityCoverage,
80 table: &TimeSeriesTable,
81) -> Result<(), CoverageSidecarError> {
82 let schema =
83 require_table_schema(&table.state().table_meta).map_err(CoverageSidecarError::from)?;
84 for (identity, _) in coverage.iter() {
85 ensure_entity_identity_matches_schema(schema, table.index_spec(), identity)
86 .map_err(CoverageSidecarError::from)?;
87 }
88 Ok(())
89}
90
91impl TimeSeriesTable {
92 fn ensure_global_coverage_query(&self) -> Result<(), CoverageQueryError> {
93 if self.index_spec().entity_columns.is_empty() {
94 Ok(())
95 } else {
96 query_error::EntityIdentityRequiredSnafu {
97 entity_columns: self.index_spec().entity_columns.clone(),
98 }
99 .fail()
100 }
101 }
102
103 fn resolve_entity_identity(
104 &self,
105 components: &[(&str, EntityValue)],
106 ) -> Result<EntityIdentity, CoverageQueryError> {
107 let entity_columns = &self.index_spec().entity_columns;
108 if entity_columns.is_empty() {
109 return query_error::EntityIdentityNotConfiguredSnafu.fail();
110 }
111
112 let mut provided = BTreeMap::new();
113 for (column, value) in components {
114 if !entity_columns.iter().any(|expected| expected == *column) {
115 return query_error::UnexpectedEntityIdentityColumnSnafu {
116 column: (*column).to_string(),
117 }
118 .fail();
119 }
120 if provided.insert(*column, value).is_some() {
121 return query_error::DuplicateEntityIdentityColumnSnafu {
122 column: (*column).to_string(),
123 }
124 .fail();
125 }
126 }
127
128 let ordered = entity_columns
129 .iter()
130 .map(|column| {
131 provided
132 .get(column.as_str())
133 .map(|value| (**value).clone())
134 .context(query_error::MissingEntityIdentityColumnSnafu {
135 column: column.clone(),
136 })
137 })
138 .collect::<Result<Vec<_>, _>>()?;
139
140 let identity =
141 EntityIdentity::try_new(ordered).context(query_error::InvalidEntityIdentitySnafu)?;
142 let schema = require_table_schema(&self.state().table_meta)
143 .context(query_error::SchemaCompatibilitySnafu)?;
144 ensure_entity_identity_matches_schema(schema, self.index_spec(), &identity)
145 .context(query_error::SchemaCompatibilitySnafu)?;
146 Ok(identity)
147 }
148
149 async fn read_and_validate_entity_coverage_sidecar(
150 &self,
151 path: &Path,
152 ) -> Result<EntityCoverage, CoverageSidecarError> {
153 let coverage = read_entity_coverage_sidecar(self.location(), path).await?;
154 ensure_entity_coverage_identity_schema(&coverage, self)?;
155 Ok(coverage)
156 }
157
158 pub(crate) async fn recover_global_coverage_from_segments<E>(&self) -> Result<Coverage, E>
163 where
164 E: CoverageRecoveryError,
165 {
166 let mut acc = Coverage::empty();
167
168 for seg in self.state().segments.values() {
169 let path = seg
170 .coverage_path
171 .as_ref()
172 .ok_or_else(|| E::missing_segment_coverage_path(seg.path.clone()))?;
173
174 let cov = read_coverage_sidecar(self.location(), Path::new(path))
175 .await
176 .map_err(|source| {
177 E::segment_coverage_read_failed(seg.path.clone(), path.clone(), source)
178 })?;
179
180 acc.union_inplace(&cov);
182 }
183
184 Ok(acc)
185 }
186
187 pub(crate) async fn recover_entity_coverage_from_segments<E>(&self) -> Result<EntityCoverage, E>
188 where
189 E: CoverageRecoveryError,
190 {
191 let mut acc = EntityCoverage::empty();
192
193 for seg in self.state().segments.values() {
194 let path = seg
195 .coverage_path
196 .as_ref()
197 .ok_or_else(|| E::missing_segment_coverage_path(seg.path.clone()))?;
198
199 let coverage = self
200 .read_and_validate_entity_coverage_sidecar(Path::new(path))
201 .await
202 .map_err(|source| {
203 E::segment_coverage_read_failed(seg.path.clone(), path.clone(), source)
204 })?;
205 acc.union_inplace(&coverage);
206 }
207
208 Ok(acc)
209 }
210
211 #[cfg(test)]
212 async fn load_global_coverage_snapshot_for_query(
213 &self,
214 ) -> Result<Coverage, CoverageQueryError> {
215 match &self.state().table_coverage {
216 None => {
217 if self.state().segments.is_empty() {
218 return Ok(Coverage::empty());
219 }
220 query_error::MissingTableCoveragePointerSnafu.fail()
221 }
222 Some(ptr) => read_coverage_sidecar(self.location(), Path::new(&ptr.coverage_path))
223 .await
224 .context(query_error::CoverageSnapshotReadSnafu {
225 coverage_path: ptr.coverage_path.clone(),
226 }),
227 }
228 }
229
230 #[cfg(test)]
237 pub(crate) async fn load_table_coverage_snapshot_only(&self) -> Result<Coverage, TableError> {
238 self.load_global_coverage_snapshot_for_query()
239 .await
240 .context(CoverageQuerySnafu)
241 }
242
243 pub(crate) async fn load_global_coverage_with_recovery<E>(&self) -> Result<Coverage, E>
251 where
252 E: CoverageRecoveryError,
253 {
254 match &self.state().table_coverage {
255 None => {
256 if self.state().segments.is_empty() {
257 return Ok(Coverage::empty());
258 }
259 self.recover_global_coverage_from_segments::<E>().await
260 }
261 Some(ptr) => {
262 match read_coverage_sidecar(self.location(), Path::new(&ptr.coverage_path)).await {
263 Ok(cov) => Ok(cov),
264 Err(snapshot_err) => {
265 tracing::warn!(
266 name: "coverage.recover",
267 target: "timeseries_table_format::table::coverage",
268 coverage_mode = "global",
269 snapshot_version = ptr.version,
270 coverage_path = %ptr.coverage_path,
271 error = %snapshot_err,
272 recovery_source = "segment_sidecars",
273 "Failed to read table coverage snapshot; attempting read-only recovery from segment sidecars"
274 );
275 let coverage = self.recover_global_coverage_from_segments::<E>().await?;
276 tracing::debug!(
277 name: "coverage.recover",
278 target: "timeseries_table_format::table::coverage",
279 coverage_mode = "global",
280 snapshot_version = ptr.version,
281 coverage_path = %ptr.coverage_path,
282 recovery_source = "segment_sidecars",
283 covered_index_interval_count = coverage.cardinality(),
284 outcome = "succeeded",
285 "Recovered table coverage from segment sidecars"
286 );
287 Ok(coverage)
288 }
289 }
290 }
291 }
292 }
293
294 pub(crate) async fn load_entity_coverage_with_recovery<E>(&self) -> Result<EntityCoverage, E>
295 where
296 E: CoverageRecoveryError,
297 {
298 match &self.state().table_coverage {
299 None => {
300 if self.state().segments.is_empty() {
301 return Ok(EntityCoverage::empty());
302 }
303 self.recover_entity_coverage_from_segments::<E>().await
304 }
305 Some(ptr) => {
306 match self
307 .read_and_validate_entity_coverage_sidecar(Path::new(&ptr.coverage_path))
308 .await
309 {
310 Ok(coverage) => Ok(coverage),
311 Err(snapshot_err) => {
312 tracing::warn!(
313 name: "coverage.recover",
314 target: "timeseries_table_format::table::coverage",
315 coverage_mode = "entity",
316 snapshot_version = ptr.version,
317 coverage_path = %ptr.coverage_path,
318 error = %snapshot_err,
319 recovery_source = "segment_sidecars",
320 "Failed to read entity coverage snapshot; attempting read-only recovery from segment sidecars"
321 );
322 let coverage = self.recover_entity_coverage_from_segments::<E>().await?;
323 tracing::debug!(
324 name: "coverage.recover",
325 target: "timeseries_table_format::table::coverage",
326 coverage_mode = "entity",
327 snapshot_version = ptr.version,
328 coverage_path = %ptr.coverage_path,
329 recovery_source = "segment_sidecars",
330 coverage_identity_count = coverage.identity_count(),
331 outcome = "succeeded",
332 "Recovered entity coverage from segment sidecars"
333 );
334 Ok(coverage)
335 }
336 }
337 }
338 }
339 }
340}
341
342use std::ops::RangeInclusive;
349
350use crate::{
351 coverage::IndexIntervalId,
352 coverage::index_interval::{index_interval_id_for_exclusive_end, index_interval_id_range},
353 metadata::index::{IndexValue, validate_index_range},
354};
355
356impl TimeSeriesTable {
357 fn interval_ids_for_query_range<S, E>(
358 &self,
359 start: S,
360 end: E,
361 ) -> Result<RangeInclusive<IndexIntervalId>, CoverageQueryError>
362 where
363 S: Into<IndexValue>,
364 E: Into<IndexValue>,
365 {
366 let start = start.into();
367 let end = end.into();
368 validate_index_range(&self.index_spec().kind, &start, &end)
369 .context(query_error::InvalidRangeSnafu)?;
370 index_interval_id_range(&self.index_spec().kind, &start, &end)
371 .context(query_error::IndexIntervalMappingSnafu)
372 }
373
374 pub async fn coverage_ratio_for_range<S, E>(&self, start: S, end: E) -> Result<f64, TableError>
403 where
404 S: Into<IndexValue>,
405 E: Into<IndexValue>,
406 {
407 async {
408 self.ensure_global_coverage_query()?;
409 let range = self.interval_ids_for_query_range(start, end)?;
410 let cov = self
411 .load_global_coverage_with_recovery::<CoverageQueryError>()
412 .await?;
413 Ok(cov.coverage_ratio(&range))
414 }
415 .await
416 .context(CoverageQuerySnafu)
417 }
418
419 pub async fn coverage_ratio_for_entity_range<S, E>(
450 &self,
451 entity: &[(&str, EntityValue)],
452 start: S,
453 end: E,
454 ) -> Result<f64, TableError>
455 where
456 S: Into<IndexValue>,
457 E: Into<IndexValue>,
458 {
459 async {
460 let identity = self.resolve_entity_identity(entity)?;
461 let range = self.interval_ids_for_query_range(start, end)?;
462 let coverage = self
463 .load_entity_coverage_with_recovery::<CoverageQueryError>()
464 .await?;
465 Ok(coverage
466 .get(&identity)
467 .map_or(0.0, |coverage| coverage.coverage_ratio(&range)))
468 }
469 .await
470 .context(CoverageQuerySnafu)
471 }
472
473 pub async fn max_gap_len_for_range<S, E>(&self, start: S, end: E) -> Result<u128, TableError>
499 where
500 S: Into<IndexValue>,
501 E: Into<IndexValue>,
502 {
503 async {
504 self.ensure_global_coverage_query()?;
505 let range = self.interval_ids_for_query_range(start, end)?;
506 let cov = self
507 .load_global_coverage_with_recovery::<CoverageQueryError>()
508 .await?;
509 Ok(cov.max_gap_len(&range))
510 }
511 .await
512 .context(CoverageQuerySnafu)
513 }
514
515 pub async fn max_gap_len_for_entity_range<S, E>(
529 &self,
530 entity: &[(&str, EntityValue)],
531 start: S,
532 end: E,
533 ) -> Result<u128, TableError>
534 where
535 S: Into<IndexValue>,
536 E: Into<IndexValue>,
537 {
538 async {
539 let identity = self.resolve_entity_identity(entity)?;
540 let range = self.interval_ids_for_query_range(start, end)?;
541 let coverage = self
542 .load_entity_coverage_with_recovery::<CoverageQueryError>()
543 .await?;
544 Ok(coverage.get(&identity).map_or_else(
545 || Coverage::range_cardinality(&range),
546 |c| c.max_gap_len(&range),
547 ))
548 }
549 .await
550 .context(CoverageQuerySnafu)
551 }
552
553 pub async fn last_fully_covered_window<E>(
580 &self,
581 end: E,
582 window_len_intervals: u64,
583 ) -> Result<Option<RangeInclusive<IndexIntervalId>>, TableError>
584 where
585 E: Into<IndexValue>,
586 {
587 async {
588 self.ensure_global_coverage_query()?;
589 let end = end.into();
590 end.validate_kind(&self.index_spec().kind)
591 .context(query_error::InvalidRangeSnafu)?;
592 if window_len_intervals == 0 {
593 return Ok(None);
594 }
595
596 let end_index_interval_id =
597 index_interval_id_for_exclusive_end(&self.index_spec().kind, &end)
598 .context(query_error::IndexIntervalMappingSnafu)?;
599 let cov = self
600 .load_global_coverage_with_recovery::<CoverageQueryError>()
601 .await?;
602 Ok(cov.last_window_at_or_before(end_index_interval_id, window_len_intervals))
603 }
604 .await
605 .context(CoverageQuerySnafu)
606 }
607
608 pub async fn last_fully_covered_window_for_entity<E>(
623 &self,
624 entity: &[(&str, EntityValue)],
625 end: E,
626 window_len_intervals: u64,
627 ) -> Result<Option<RangeInclusive<IndexIntervalId>>, TableError>
628 where
629 E: Into<IndexValue>,
630 {
631 async {
632 let identity = self.resolve_entity_identity(entity)?;
633 let end = end.into();
634 end.validate_kind(&self.index_spec().kind)
635 .context(query_error::InvalidRangeSnafu)?;
636 if window_len_intervals == 0 {
637 return Ok(None);
638 }
639
640 let end_index_interval_id =
641 index_interval_id_for_exclusive_end(&self.index_spec().kind, &end)
642 .context(query_error::IndexIntervalMappingSnafu)?;
643 let coverage = self
644 .load_entity_coverage_with_recovery::<CoverageQueryError>()
645 .await?;
646 Ok(coverage.get(&identity).and_then(|coverage| {
647 coverage.last_window_at_or_before(end_index_interval_id, window_len_intervals)
648 }))
649 }
650 .await
651 .context(CoverageQuerySnafu)
652 }
653}
654
655#[cfg(test)]
656mod tests {
657 use std::{collections::BTreeSet, error::Error as _, num::NonZeroU64, path::Path};
658
659 use super::*;
660 use crate::{
661 coverage::{
662 Coverage, EntityCoverage, EntityIdentity, EntityValue,
663 index_interval::{IndexIntervalMappingError, index_interval_id_for_value},
664 io::{
665 CoverageSidecarError, write_coverage_sidecar_atomic,
666 write_coverage_sidecar_new_bytes,
667 },
668 serde::{CoverageCodecError, coverage_to_bytes, entity_coverage_to_bytes},
669 },
670 metadata::logical_schema::{LogicalDataType, LogicalField, LogicalSchema},
671 metadata::schema_compat::SchemaCompatibilityError,
672 metadata::{
673 index::{IndexKind, IndexSpec, IndexValueError},
674 table::{TableKind, TableMeta},
675 },
676 storage::{StorageError, TableLocation},
677 table::test_util::{
678 TestResult, TestRow, TraceCapture, append_parquet_fixture, make_basic_table_meta,
679 make_int32_entity_table_meta, utc_datetime, write_int32_entity_parquet,
680 write_test_parquet,
681 },
682 transaction_log::table_state::TableCoveragePointer,
683 };
684 use chrono::{DateTime, TimeZone, Utc};
685 use snafu::ErrorCompat;
686 use tempfile::TempDir;
687
688 type HelperResult<T> = Result<T, Box<dyn std::error::Error>>;
689
690 fn ts_from_secs(secs: i64) -> DateTime<Utc> {
691 Utc.timestamp_opt(secs, 0)
692 .single()
693 .expect("valid timestamp")
694 }
695
696 fn coverage_query_source(error: &TableError) -> &CoverageQueryError {
697 error
698 .source()
699 .and_then(|source| source.downcast_ref::<CoverageQueryError>())
700 .expect("coverage query source")
701 }
702
703 fn coverage_sidecar_source(error: &CoverageQueryError) -> &CoverageSidecarError {
704 error
705 .source()
706 .and_then(|source| source.downcast_ref::<Box<CoverageSidecarError>>())
707 .map(Box::as_ref)
708 .expect("coverage sidecar source")
709 }
710
711 fn single_segment_and_coverage_paths(table: &TimeSeriesTable) -> (String, String) {
712 assert_eq!(table.state().segments.len(), 1, "expected one segment");
713 let segment = table
714 .state()
715 .segments
716 .values()
717 .next()
718 .expect("segment present");
719 (
720 segment.path.clone(),
721 segment
722 .coverage_path
723 .clone()
724 .expect("segment coverage path"),
725 )
726 }
727
728 async fn make_table() -> HelperResult<(TempDir, TimeSeriesTable)> {
729 let tmp = TempDir::new()?;
730 let location = TableLocation::local(tmp.path());
731 let mut meta = make_basic_table_meta();
732 let TableKind::TimeSeries(index) = &mut meta.kind else {
733 unreachable!("test metadata is time series")
734 };
735 index.entity_columns.clear();
736 let table = TimeSeriesTable::create(location, meta).await?;
737 Ok((tmp, table))
738 }
739
740 #[tokio::test]
741 async fn strict_snapshot_missing_preserves_storage_source_and_backtrace() -> TestResult {
742 let (_tmp, mut table) = make_table().await?;
743 let coverage_path = "_coverage/table/missing.roar";
744 table.state_mut().table_coverage = Some(TableCoveragePointer {
745 index_kind: table.index_spec().kind.clone(),
746 coverage_path: coverage_path.to_string(),
747 version: table.state().version,
748 });
749
750 let error = table
751 .load_table_coverage_snapshot_only()
752 .await
753 .expect_err("missing snapshot must fail");
754 let query = coverage_query_source(&error);
755 assert!(matches!(
756 query,
757 CoverageQueryError::CoverageSnapshotRead {
758 coverage_path: path,
759 ..
760 } if path == coverage_path
761 ));
762 let sidecar = coverage_sidecar_source(query);
763 let storage = sidecar
764 .source()
765 .and_then(|source| source.downcast_ref::<StorageError>())
766 .expect("storage source");
767 assert!(matches!(storage, StorageError::NotFound { .. }));
768 assert!(std::ptr::eq(
769 ErrorCompat::backtrace(&error).expect("table backtrace"),
770 ErrorCompat::backtrace(storage).expect("storage backtrace"),
771 ));
772 Ok(())
773 }
774
775 #[tokio::test]
776 async fn strict_snapshot_corruption_and_truncation_preserve_codec_source_and_backtrace()
777 -> TestResult {
778 let (tmp, mut table) = make_table().await?;
779 let valid = coverage_to_bytes(&Coverage::from_iter([1u64]))?;
780 let cases = [
781 ("corrupt", b"not a bitmap".to_vec()),
782 ("truncated", valid[..valid.len() - 1].to_vec()),
783 ];
784
785 for (name, bytes) in cases {
786 let coverage_path = format!("_coverage/table/{name}.roar");
787 let absolute = tmp.path().join(&coverage_path);
788 std::fs::create_dir_all(absolute.parent().expect("coverage parent"))?;
789 std::fs::write(&absolute, bytes)?;
790 table.state_mut().table_coverage = Some(TableCoveragePointer {
791 index_kind: table.index_spec().kind.clone(),
792 coverage_path,
793 version: table.state().version,
794 });
795
796 let error = table
797 .load_table_coverage_snapshot_only()
798 .await
799 .expect_err("invalid snapshot must fail");
800 let query = coverage_query_source(&error);
801 let sidecar = coverage_sidecar_source(query);
802 let codec = sidecar
803 .source()
804 .and_then(|source| source.downcast_ref::<CoverageCodecError>())
805 .expect("codec source");
806 assert!(matches!(
807 codec,
808 CoverageCodecError::BitmapDeserialization { .. }
809 ));
810 assert!(std::ptr::eq(
811 ErrorCompat::backtrace(&error).expect("table backtrace"),
812 ErrorCompat::backtrace(codec).expect("codec backtrace"),
813 ));
814 }
815 Ok(())
816 }
817
818 #[tokio::test]
819 async fn strict_snapshot_requires_a_pointer_when_segments_exist() -> TestResult {
820 let (_tmp, mut table) = table_with_sparse_coverage().await?;
821 table.state_mut().table_coverage = None;
822
823 assert!(matches!(
824 table.load_table_coverage_snapshot_only().await,
825 Err(TableError::CoverageQuery {
826 source: CoverageQueryError::MissingTableCoveragePointer { .. }
827 })
828 ));
829 Ok(())
830 }
831
832 #[cfg(unix)]
833 #[tokio::test]
834 async fn strict_snapshot_permission_denied_preserves_storage_source() -> TestResult {
835 use std::os::unix::fs::PermissionsExt;
836
837 let (tmp, mut table) = make_table().await?;
838 let coverage_path = "_coverage/table/denied.roar";
839 let absolute = tmp.path().join(coverage_path);
840 write_coverage_sidecar_atomic(
841 table.location(),
842 Path::new(coverage_path),
843 &Coverage::empty(),
844 )
845 .await?;
846 table.state_mut().table_coverage = Some(TableCoveragePointer {
847 index_kind: table.index_spec().kind.clone(),
848 coverage_path: coverage_path.to_string(),
849 version: table.state().version,
850 });
851 let original_permissions = std::fs::metadata(&absolute)?.permissions();
852 let mut denied_permissions = original_permissions.clone();
853 denied_permissions.set_mode(0o0);
854 std::fs::set_permissions(&absolute, denied_permissions)?;
855
856 let error = table
857 .load_table_coverage_snapshot_only()
858 .await
859 .expect_err("permission-denied snapshot must fail");
860 std::fs::set_permissions(&absolute, original_permissions)?;
861
862 let query = coverage_query_source(&error);
863 let sidecar = coverage_sidecar_source(query);
864 let storage = sidecar
865 .source()
866 .and_then(|source| source.downcast_ref::<StorageError>())
867 .expect("storage source");
868 assert!(matches!(storage, StorageError::OtherIo { .. }));
869 assert!(std::ptr::eq(
870 ErrorCompat::backtrace(&error).expect("table backtrace"),
871 ErrorCompat::backtrace(storage).expect("storage backtrace"),
872 ));
873 Ok(())
874 }
875
876 async fn table_with_index_coverage(
877 kind: IndexKind,
878 coverage: Coverage,
879 ) -> HelperResult<(TempDir, TimeSeriesTable)> {
880 let tmp = TempDir::new()?;
881 let mut table = TimeSeriesTable::create(
882 TableLocation::local(tmp.path()),
883 TableMeta::new_time_series(IndexSpec {
884 column: "index".to_string(),
885 entity_columns: Vec::new(),
886 kind: kind.clone(),
887 }),
888 )
889 .await?;
890 let coverage_path = "_coverage/table/query-test.roar";
891 write_coverage_sidecar_atomic(table.location(), Path::new(coverage_path), &coverage)
892 .await?;
893 let version = table.state().version;
894 table.state_mut().table_coverage = Some(TableCoveragePointer {
895 index_kind: kind,
896 coverage_path: coverage_path.to_string(),
897 version,
898 });
899 Ok((tmp, table))
900 }
901
902 async fn append_segment(
903 table: &mut TimeSeriesTable,
904 tmp: &TempDir,
905 rel_path: &str,
906 rows: &[TestRow],
907 ) -> HelperResult<String> {
908 let abs = tmp.path().join(rel_path);
909 write_test_parquet(&abs, true, false, rows)?;
910 let existing_paths = table
911 .state()
912 .segments
913 .keys()
914 .cloned()
915 .collect::<BTreeSet<_>>();
916 append_parquet_fixture(table, rel_path).await?;
917 table
918 .state()
919 .segments
920 .keys()
921 .find(|path| !existing_paths.contains(*path))
922 .cloned()
923 .ok_or_else(|| "append did not add a segment".into())
924 }
925
926 #[tokio::test]
927 async fn entity_sidecar_identity_schema_is_validated_for_snapshots_and_recovery() -> TestResult
928 {
929 let tmp = TempDir::new()?;
930 let location = TableLocation::local(tmp.path());
931 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
932 let mut wrong_arity = EntityCoverage::empty();
933 wrong_arity.union_coverage(
934 EntityIdentity::try_new(vec!["A".into(), "X".into()])?,
935 Coverage::from_iter([0]),
936 );
937 let wrong_arity_bytes = entity_coverage_to_bytes(&wrong_arity)?;
938 let snapshot_path = "_coverage/table/wrong-arity.roar";
939 write_coverage_sidecar_new_bytes(&location, Path::new(snapshot_path), &wrong_arity_bytes)
940 .await?;
941 let version = table.state().version;
942 let index_kind = table.index_spec().kind.clone();
943 table.state_mut().table_coverage = Some(TableCoveragePointer {
944 index_kind,
945 coverage_path: snapshot_path.to_string(),
946 version,
947 });
948
949 let snapshot_error = table
950 .read_and_validate_entity_coverage_sidecar(Path::new(snapshot_path))
951 .await
952 .expect_err("snapshot identity arity must match the table");
953 assert!(matches!(
954 snapshot_error,
955 CoverageSidecarError::EntityIdentitySchema { source, .. } if matches!(
956 source.as_ref(),
957 SchemaCompatibilityError::EntityIdentityArityMismatch {
958 expected: 1,
959 actual: 2,
960 }
961 )
962 ));
963
964 let mut wrong_type = EntityCoverage::empty();
965 wrong_type.union_coverage(
966 EntityIdentity::try_new(vec![EntityValue::Int32(1)])?,
967 Coverage::from_iter([0]),
968 );
969 let wrong_type_bytes = entity_coverage_to_bytes(&wrong_type)?;
970 table.state_mut().table_coverage = None;
971 let segment_path = "data/entity-arity.parquet";
972 let committed_segment_path = append_segment(
973 &mut table,
974 &tmp,
975 segment_path,
976 &[TestRow {
977 ts_millis: 1_000,
978 symbol: "A",
979 price: 1.0,
980 }],
981 )
982 .await?;
983 let segment_coverage_path = table
984 .state()
985 .segments
986 .get(&committed_segment_path)
987 .and_then(|segment| segment.coverage_path.clone())
988 .expect("segment coverage path");
989 tokio::fs::write(tmp.path().join(&segment_coverage_path), &wrong_type_bytes).await?;
990
991 let recovery_error = table
992 .recover_entity_coverage_from_segments::<CoverageQueryError>()
993 .await
994 .expect_err("segment identity type must match the table schema");
995 assert!(matches!(
996 recovery_error,
997 CoverageQueryError::SegmentCoverageSidecarRead {
998 segment_path: path,
999 coverage_path,
1000 source,
1001 } if path == committed_segment_path
1002 && coverage_path == segment_coverage_path
1003 && matches!(
1004 &*source,
1005 CoverageSidecarError::EntityIdentitySchema {
1006 source,
1007 ..
1008 }
1009 if matches!(
1010 source.as_ref(),
1011 SchemaCompatibilityError::EntityIdentityTypeMismatch {
1012 column,
1013 expected: crate::metadata::logical_schema::LogicalDataType::Utf8,
1014 actual: "int32",
1015 } if column == "symbol"
1016 )
1017 )
1018 ));
1019 Ok(())
1020 }
1021
1022 #[tokio::test]
1023 async fn entity_coverage_ratio_selects_only_the_requested_identity() -> TestResult {
1024 let tmp = TempDir::new()?;
1025 let location = TableLocation::local(tmp.path());
1026 let mut table = TimeSeriesTable::create(location, make_basic_table_meta()).await?;
1027 append_segment(
1028 &mut table,
1029 &tmp,
1030 "data/entity-ratio.parquet",
1031 &[
1032 TestRow {
1033 ts_millis: 1_000,
1034 symbol: "A",
1035 price: 1.0,
1036 },
1037 TestRow {
1038 ts_millis: 61_000,
1039 symbol: "B",
1040 price: 2.0,
1041 },
1042 ],
1043 )
1044 .await?;
1045 let state_before = table.state().clone();
1046 let start = ts_from_secs(0);
1047 let end = ts_from_secs(120);
1048
1049 assert_eq!(
1050 table
1051 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end)
1052 .await?,
1053 0.5
1054 );
1055 assert_eq!(
1056 table
1057 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("B"))], start, end)
1058 .await?,
1059 0.5
1060 );
1061 assert_eq!(
1062 table
1063 .coverage_ratio_for_entity_range(
1064 &[("symbol", EntityValue::from("unseen"))],
1065 start,
1066 end,
1067 )
1068 .await?,
1069 0.0
1070 );
1071 assert!(matches!(
1072 table.coverage_ratio_for_range(start, end).await,
1073 Err(TableError::CoverageQuery {
1074 source: CoverageQueryError::EntityIdentityRequired { entity_columns, .. }
1075 })
1076 if entity_columns == ["symbol"]
1077 ));
1078 assert_eq!(table.state(), &state_before);
1079 Ok(())
1080 }
1081
1082 #[tokio::test]
1083 async fn numeric_entity_queries_require_the_exact_scalar_type() -> TestResult {
1084 let tmp = TempDir::new()?;
1085 let mut table = TimeSeriesTable::create(
1086 TableLocation::local(tmp.path()),
1087 make_int32_entity_table_meta(),
1088 )
1089 .await?;
1090 let path = "data/numeric-coverage.parquet";
1091 write_int32_entity_parquet(
1092 &tmp.path().join(path),
1093 &[1_000, 61_000],
1094 &[-1, i32::MAX],
1095 &[10.0, 20.0],
1096 )?;
1097 append_parquet_fixture(&mut table, path).await?;
1098 let start = ts_from_secs(0);
1099 let end = ts_from_secs(120);
1100
1101 assert_eq!(
1102 table
1103 .coverage_ratio_for_entity_range(
1104 &[("device_id", EntityValue::Int32(-1))],
1105 start,
1106 end,
1107 )
1108 .await?,
1109 0.5
1110 );
1111 assert_eq!(
1112 table
1113 .coverage_ratio_for_entity_range(
1114 &[("device_id", EntityValue::Int32(42))],
1115 start,
1116 end,
1117 )
1118 .await?,
1119 0.0
1120 );
1121 assert!(matches!(
1122 table
1123 .coverage_ratio_for_entity_range(
1124 &[("device_id", EntityValue::Int64(-1))],
1125 start,
1126 end,
1127 )
1128 .await,
1129 Err(TableError::CoverageQuery {
1130 source: CoverageQueryError::SchemaCompatibility { source, .. }
1131 }) if matches!(
1132 source.as_ref(),
1133 SchemaCompatibilityError::EntityIdentityTypeMismatch {
1134 column,
1135 expected: LogicalDataType::Int32,
1136 actual: "int64",
1137 } if column == "device_id"
1138 )
1139 ));
1140 assert!(matches!(
1141 table
1142 .coverage_ratio_for_entity_range(
1143 &[("device_id", EntityValue::from("-1"))],
1144 start,
1145 end,
1146 )
1147 .await,
1148 Err(TableError::CoverageQuery {
1149 source: CoverageQueryError::SchemaCompatibility { source, .. }
1150 }) if matches!(
1151 source.as_ref(),
1152 SchemaCompatibilityError::EntityIdentityTypeMismatch {
1153 column,
1154 expected: LogicalDataType::Int32,
1155 actual: "utf8",
1156 } if column == "device_id"
1157 )
1158 ));
1159 Ok(())
1160 }
1161
1162 #[tokio::test]
1163 async fn entity_gap_and_window_queries_are_isolated_and_recover_readonly() -> TestResult {
1164 let tmp = TempDir::new()?;
1165 let mut table =
1166 TimeSeriesTable::create(TableLocation::local(tmp.path()), make_basic_table_meta())
1167 .await?;
1168 append_segment(
1169 &mut table,
1170 &tmp,
1171 "data/entity-gaps.parquet",
1172 &[
1173 TestRow {
1174 ts_millis: 1_000,
1175 symbol: "A",
1176 price: 1.0,
1177 },
1178 TestRow {
1179 ts_millis: 181_000,
1180 symbol: "A",
1181 price: 2.0,
1182 },
1183 TestRow {
1184 ts_millis: 1_000,
1185 symbol: "B",
1186 price: 3.0,
1187 },
1188 TestRow {
1189 ts_millis: 61_000,
1190 symbol: "B",
1191 price: 4.0,
1192 },
1193 TestRow {
1194 ts_millis: 121_000,
1195 symbol: "B",
1196 price: 5.0,
1197 },
1198 ],
1199 )
1200 .await?;
1201 let start = ts_from_secs(0);
1202 let end = ts_from_secs(240);
1203
1204 assert_eq!(
1205 table
1206 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end)
1207 .await?,
1208 0.5
1209 );
1210 assert_eq!(
1211 table
1212 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("B"))], start, end)
1213 .await?,
1214 0.75
1215 );
1216 assert_eq!(
1217 table
1218 .max_gap_len_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end)
1219 .await?,
1220 2
1221 );
1222 assert_eq!(
1223 table
1224 .max_gap_len_for_entity_range(&[("symbol", EntityValue::from("B"))], start, end)
1225 .await?,
1226 1
1227 );
1228 assert_eq!(
1229 table
1230 .last_fully_covered_window_for_entity(
1231 &[("symbol", EntityValue::from("A"))],
1232 end,
1233 2,
1234 )
1235 .await?,
1236 None
1237 );
1238 assert_eq!(
1239 table
1240 .last_fully_covered_window_for_entity(
1241 &[("symbol", EntityValue::from("B"))],
1242 end,
1243 2,
1244 )
1245 .await?,
1246 Some(0x8000_0000_0000_0001..=0x8000_0000_0000_0002)
1247 );
1248
1249 let snapshot_path = table
1250 .state()
1251 .table_coverage
1252 .as_ref()
1253 .expect("snapshot pointer")
1254 .coverage_path
1255 .clone();
1256 let state_before = table.state().clone();
1257 tokio::fs::remove_file(tmp.path().join(snapshot_path)).await?;
1258
1259 assert_eq!(
1260 table
1261 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end)
1262 .await?,
1263 0.5
1264 );
1265 assert_eq!(
1266 table
1267 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("B"))], start, end)
1268 .await?,
1269 0.75
1270 );
1271 assert_eq!(
1272 table
1273 .max_gap_len_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end)
1274 .await?,
1275 2
1276 );
1277 assert_eq!(
1278 table
1279 .last_fully_covered_window_for_entity(
1280 &[("symbol", EntityValue::from("B"))],
1281 end,
1282 2,
1283 )
1284 .await?,
1285 Some(0x8000_0000_0000_0001..=0x8000_0000_0000_0002)
1286 );
1287 assert_eq!(
1288 table
1289 .max_gap_len_for_entity_range(
1290 &[("symbol", EntityValue::from("unseen"))],
1291 start,
1292 end,
1293 )
1294 .await?,
1295 4
1296 );
1297 assert_eq!(
1298 table
1299 .last_fully_covered_window_for_entity(
1300 &[("symbol", EntityValue::from("unseen"))],
1301 end,
1302 1,
1303 )
1304 .await?,
1305 None
1306 );
1307 assert_eq!(
1308 table
1309 .last_fully_covered_window_for_entity(
1310 &[("symbol", EntityValue::from("unseen"))],
1311 end,
1312 0,
1313 )
1314 .await?,
1315 None
1316 );
1317
1318 assert!(matches!(
1319 table.coverage_ratio_for_range(start, end).await,
1320 Err(TableError::CoverageQuery {
1321 source: CoverageQueryError::EntityIdentityRequired { entity_columns, .. }
1322 })
1323 if entity_columns == ["symbol"]
1324 ));
1325 assert!(matches!(
1326 table.max_gap_len_for_range(start, end).await,
1327 Err(TableError::CoverageQuery {
1328 source: CoverageQueryError::EntityIdentityRequired { entity_columns, .. }
1329 })
1330 if entity_columns == ["symbol"]
1331 ));
1332 assert!(matches!(
1333 table.last_fully_covered_window(end, 0).await,
1334 Err(TableError::CoverageQuery {
1335 source: CoverageQueryError::EntityIdentityRequired { entity_columns, .. }
1336 })
1337 if entity_columns == ["symbol"]
1338 ));
1339 assert_eq!(table.state(), &state_before);
1340 Ok(())
1341 }
1342
1343 #[tokio::test]
1344 async fn entity_identity_input_is_validated_and_canonicalized() -> TestResult {
1345 let tmp = TempDir::new()?;
1346 let mut meta = make_basic_table_meta();
1347 let TableKind::TimeSeries(index) = &mut meta.kind else {
1348 unreachable!("test metadata is time series")
1349 };
1350 index.entity_columns = vec!["symbol".to_string(), "venue".to_string()];
1351 let mut fields = meta
1352 .logical_schema
1353 .as_ref()
1354 .expect("test schema")
1355 .columns()
1356 .to_vec();
1357 fields.push(LogicalField {
1358 name: "venue".to_string(),
1359 data_type: LogicalDataType::Utf8,
1360 nullable: false,
1361 });
1362 meta.logical_schema = Some(LogicalSchema::new(fields)?);
1363 let table = TimeSeriesTable::create(TableLocation::local(tmp.path()), meta).await?;
1364
1365 let identity = table.resolve_entity_identity(&[
1366 ("venue", EntityValue::from("X")),
1367 ("symbol", EntityValue::from("A")),
1368 ])?;
1369 assert_eq!(
1370 identity.components(),
1371 [EntityValue::from("A"), EntityValue::from("X")]
1372 );
1373 let start = ts_from_secs(0);
1374 let end = ts_from_secs(60);
1375 assert!(matches!(
1376 table
1377 .coverage_ratio_for_entity_range(&[("venue", EntityValue::from("X"))], start, end)
1378 .await,
1379 Err(TableError::CoverageQuery {
1380 source: CoverageQueryError::MissingEntityIdentityColumn { column, .. }
1381 }) if column == "symbol"
1382 ));
1383 assert!(matches!(
1384 table
1385 .coverage_ratio_for_entity_range(
1386 &[
1387 ("device", EntityValue::from("A")),
1388 ("venue", EntityValue::from("X")),
1389 ],
1390 start,
1391 end,
1392 )
1393 .await,
1394 Err(TableError::CoverageQuery {
1395 source: CoverageQueryError::UnexpectedEntityIdentityColumn { column, .. }
1396 }) if column == "device"
1397 ));
1398 assert!(matches!(
1399 table
1400 .coverage_ratio_for_entity_range(
1401 &[
1402 ("symbol", EntityValue::from("A")),
1403 ("symbol", EntityValue::from("B")),
1404 ("venue", EntityValue::from("X")),
1405 ],
1406 start,
1407 end,
1408 )
1409 .await,
1410 Err(TableError::CoverageQuery {
1411 source: CoverageQueryError::DuplicateEntityIdentityColumn { column, .. }
1412 }) if column == "symbol"
1413 ));
1414
1415 let (_tmp, global_table) = make_table().await?;
1416 assert!(matches!(
1417 global_table
1418 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end)
1419 .await,
1420 Err(TableError::CoverageQuery {
1421 source: CoverageQueryError::EntityIdentityNotConfigured { .. }
1422 })
1423 ));
1424 Ok(())
1425 }
1426
1427 async fn table_with_sparse_coverage() -> HelperResult<(TempDir, TimeSeriesTable)> {
1428 let (tmp, mut table) = make_table().await?;
1430 append_segment(
1431 &mut table,
1432 &tmp,
1433 "data/sparse.parquet",
1434 &[
1435 TestRow {
1436 ts_millis: 1_000,
1437 symbol: "A",
1438 price: 1.0,
1439 },
1440 TestRow {
1441 ts_millis: 61_000,
1442 symbol: "A",
1443 price: 2.0,
1444 },
1445 TestRow {
1446 ts_millis: 180_000,
1447 symbol: "A",
1448 price: 3.0,
1449 },
1450 ],
1451 )
1452 .await?;
1453 Ok((tmp, table))
1454 }
1455
1456 async fn table_with_contiguous_run() -> HelperResult<(TempDir, TimeSeriesTable)> {
1457 let (tmp, mut table) = make_table().await?;
1459 append_segment(
1460 &mut table,
1461 &tmp,
1462 "data/window.parquet",
1463 &[
1464 TestRow {
1465 ts_millis: 240_000,
1466 symbol: "A",
1467 price: 1.0,
1468 },
1469 TestRow {
1470 ts_millis: 300_000,
1471 symbol: "A",
1472 price: 2.0,
1473 },
1474 ],
1475 )
1476 .await?;
1477 Ok((tmp, table))
1478 }
1479
1480 #[tokio::test]
1481 async fn index_interval_id_range_rejects_invalid_index_range() -> TestResult {
1482 let (_tmp, table) = make_table().await?;
1483 let ts = utc_datetime(2024, 1, 1, 0, 0, 0);
1484
1485 let err = table
1486 .interval_ids_for_query_range(ts, ts)
1487 .expect_err("start >= end should be invalid");
1488 assert!(matches!(err, CoverageQueryError::InvalidRange { .. }));
1489 Ok(())
1490 }
1491
1492 #[tokio::test]
1493 async fn index_interval_id_range_uses_64_bit_signed_timestamp_mapping() -> TestResult {
1494 let (_tmp, table) = make_table().await?;
1495 let start = ts_from_secs(0);
1496 let end = ts_from_secs(180); let range = table.interval_ids_for_query_range(start, end)?;
1499 assert_eq!(range, 0x8000_0000_0000_0000..=0x8000_0000_0000_0002);
1500 Ok(())
1501 }
1502
1503 #[tokio::test]
1504 async fn signed_coverage_queries_handle_gaps_extremes_and_last_window() -> TestResult {
1505 let kind = IndexKind::Int64 {
1506 index_granularity: NonZeroU64::new(10).unwrap(),
1507 };
1508 let coverage: Coverage = [-10i64, 0, 10]
1509 .into_iter()
1510 .map(|value| index_interval_id_for_value(&kind, &value.into()).unwrap())
1511 .collect();
1512 let huge_gap = u128::from(
1513 index_interval_id_for_value(&kind, &(-10i64).into()).unwrap()
1514 - index_interval_id_for_value(&kind, &i64::MIN.into()).unwrap(),
1515 )
1516 .max(u128::from(
1517 index_interval_id_for_exclusive_end(&kind, &i64::MAX.into()).unwrap()
1518 - index_interval_id_for_value(&kind, &10i64.into()).unwrap(),
1519 ));
1520 let (_tmp, table) = table_with_index_coverage(kind, coverage).await?;
1521
1522 assert_eq!(table.coverage_ratio_for_range(-20i64, 30i64).await?, 0.6);
1523 assert_eq!(table.max_gap_len_for_range(-20i64, 30i64).await?, 1);
1524 assert_eq!(table.max_gap_len_for_range(-10i64, 0i64).await?, 0);
1525 assert_eq!(table.max_gap_len_for_range(-50i64, -20i64).await?, 3);
1526 assert_eq!(
1527 table.max_gap_len_for_range(i64::MIN, i64::MAX).await?,
1528 huge_gap
1529 );
1530
1531 let window = table
1532 .last_fully_covered_window(10i64, 2)
1533 .await?
1534 .expect("signed window across zero");
1535 assert_eq!(
1536 window,
1537 index_interval_id_for_value(&table.index_spec().kind, &(-10i64).into()).unwrap()
1538 ..=index_interval_id_for_value(&table.index_spec().kind, &0i64.into()).unwrap()
1539 );
1540 Ok(())
1541 }
1542
1543 #[tokio::test]
1544 async fn unsigned_coverage_queries_preserve_large_values_and_boundaries() -> TestResult {
1545 let kind = IndexKind::UInt64 {
1546 index_granularity: NonZeroU64::new(1).unwrap(),
1547 };
1548 let start = i64::MAX as u64 + 1;
1549 let coverage: Coverage = [start, start + 1, u64::MAX - 2, u64::MAX - 1]
1550 .into_iter()
1551 .collect();
1552 let (_tmp, table) = table_with_index_coverage(kind, coverage).await?;
1553
1554 let requested = u128::from(u64::MAX) - u128::from(start);
1555 let ratio = table.coverage_ratio_for_range(start, u64::MAX).await?;
1556 assert!((ratio - 4.0 / requested as f64).abs() < f64::EPSILON);
1557 assert_eq!(
1558 table.max_gap_len_for_range(start, u64::MAX).await?,
1559 u128::from(u64::MAX) - u128::from(start) - 4
1560 );
1561 assert_eq!(
1562 table.last_fully_covered_window(start + 2, 2).await?,
1563 Some(start..=start + 1)
1564 );
1565 assert_eq!(
1566 table.last_fully_covered_window(u64::MAX, 2).await?,
1567 Some(u64::MAX - 2..=u64::MAX - 1)
1568 );
1569 Ok(())
1570 }
1571
1572 #[tokio::test]
1573 async fn coverage_ratio_uses_snapshot_when_present() -> TestResult {
1574 let (_tmp, table) = table_with_sparse_coverage().await?;
1575 let start = ts_from_secs(0);
1576 let end = ts_from_secs(240); let ratio = table.coverage_ratio_for_range(start, end).await?;
1579 assert!((ratio - 0.75).abs() < 1e-12);
1580 Ok(())
1581 }
1582
1583 fn assert_recovery_events(
1584 capture: &TraceCapture,
1585 pointer: &TableCoveragePointer,
1586 mode: &str,
1587 count_field: &str,
1588 count: &str,
1589 forbidden_values: &[&str],
1590 ) {
1591 let recovery_events: Vec<_> = capture
1592 .events()
1593 .into_iter()
1594 .filter(|event| event.name == "coverage.recover")
1595 .collect();
1596 assert_eq!(recovery_events.len(), 2);
1597 assert!(
1598 recovery_events
1599 .iter()
1600 .all(|event| { event.target == "timeseries_table_format::table::coverage" })
1601 );
1602
1603 let warning = recovery_events
1604 .iter()
1605 .find(|event| event.level == tracing::Level::WARN)
1606 .expect("recovery warning");
1607 assert_eq!(
1608 warning.fields.get("coverage_mode").map(String::as_str),
1609 Some(mode)
1610 );
1611 assert_eq!(
1612 warning.fields.get("snapshot_version"),
1613 Some(&pointer.version.to_string())
1614 );
1615 assert_eq!(
1616 warning.fields.get("coverage_path"),
1617 Some(&pointer.coverage_path)
1618 );
1619 assert_eq!(
1620 warning.fields.get("recovery_source").map(String::as_str),
1621 Some("segment_sidecars")
1622 );
1623 assert!(
1624 warning
1625 .fields
1626 .get("message")
1627 .is_some_and(|message| message.contains("attempting read-only recovery"))
1628 );
1629 assert!(
1630 warning
1631 .fields
1632 .get("error")
1633 .is_some_and(|error| error.contains(&pointer.coverage_path))
1634 );
1635
1636 let completion = recovery_events
1637 .iter()
1638 .find(|event| event.level == tracing::Level::DEBUG)
1639 .expect("recovery completion");
1640 assert_eq!(
1641 completion.fields.get("outcome").map(String::as_str),
1642 Some("succeeded")
1643 );
1644 assert_eq!(
1645 completion.fields.get(count_field).map(String::as_str),
1646 Some(count)
1647 );
1648
1649 for value in recovery_events
1650 .iter()
1651 .flat_map(|event| event.fields.values())
1652 {
1653 for forbidden in forbidden_values
1654 .iter()
1655 .copied()
1656 .chain(["LogicalSchema", "RecordBatch"])
1657 {
1658 assert!(
1659 !value.contains(forbidden),
1660 "diagnostic value contains sensitive data '{forbidden}': {value}"
1661 );
1662 }
1663 }
1664 }
1665
1666 #[tokio::test]
1667 async fn global_coverage_snapshot_recovery_emits_safe_structured_events() -> TestResult {
1668 let (tmp, table) = table_with_sparse_coverage().await?;
1669 let pointer = table
1670 .state()
1671 .table_coverage
1672 .as_ref()
1673 .expect("snapshot pointer")
1674 .clone();
1675 tokio::fs::remove_file(tmp.path().join(&pointer.coverage_path)).await?;
1676 let capture = TraceCapture::default();
1677
1678 let ratio = capture
1679 .run(table.coverage_ratio_for_range(ts_from_secs(0), ts_from_secs(240)))
1680 .await?;
1681
1682 assert!((ratio - 0.75).abs() < 1e-12);
1683 let table_root = tmp.path().display().to_string();
1684 assert_recovery_events(
1685 &capture,
1686 &pointer,
1687 "global",
1688 "covered_index_interval_count",
1689 "3",
1690 &[&table_root],
1691 );
1692 Ok(())
1693 }
1694
1695 #[tokio::test]
1696 async fn entity_coverage_snapshot_recovery_emits_safe_structured_events() -> TestResult {
1697 const SENSITIVE_ENTITY: &str = "sensitive-entity-value";
1698 let tmp = TempDir::new()?;
1699 let mut table =
1700 TimeSeriesTable::create(TableLocation::local(tmp.path()), make_basic_table_meta())
1701 .await?;
1702 append_segment(
1703 &mut table,
1704 &tmp,
1705 "data/sensitive.parquet",
1706 &[TestRow {
1707 ts_millis: 1_000,
1708 symbol: SENSITIVE_ENTITY,
1709 price: 987_654.25,
1710 }],
1711 )
1712 .await?;
1713 let pointer = table
1714 .state()
1715 .table_coverage
1716 .as_ref()
1717 .expect("snapshot pointer")
1718 .clone();
1719 tokio::fs::remove_file(tmp.path().join(&pointer.coverage_path)).await?;
1720 let capture = TraceCapture::default();
1721
1722 let ratio = capture
1723 .run(table.coverage_ratio_for_entity_range(
1724 &[("symbol", EntityValue::from(SENSITIVE_ENTITY))],
1725 ts_from_secs(0),
1726 ts_from_secs(60),
1727 ))
1728 .await?;
1729
1730 assert_eq!(ratio, 1.0);
1731 let table_root = tmp.path().display().to_string();
1732 assert_recovery_events(
1733 &capture,
1734 &pointer,
1735 "entity",
1736 "coverage_identity_count",
1737 "1",
1738 &[&table_root, SENSITIVE_ENTITY, "987654.25"],
1739 );
1740 Ok(())
1741 }
1742
1743 #[tokio::test]
1744 async fn coverage_ratio_recovers_when_snapshot_missing() -> TestResult {
1745 let (_tmp, mut table) = table_with_sparse_coverage().await?;
1746 table.state_mut().table_coverage = None;
1747
1748 let ratio = table
1749 .coverage_ratio_for_range(ts_from_secs(0), ts_from_secs(240))
1750 .await?;
1751 assert!((ratio - 0.75).abs() < 1e-12);
1752 Ok(())
1753 }
1754
1755 #[tokio::test]
1756 async fn coverage_ratio_errors_when_recovery_missing_segment_coverage_path() -> TestResult {
1757 let (_tmp, mut table) = table_with_sparse_coverage().await?;
1758 table.state_mut().table_coverage = None;
1759 let segment = table
1760 .state_mut()
1761 .segments
1762 .values_mut()
1763 .next()
1764 .expect("segment present");
1765 let segment_path = segment.path.clone();
1766 segment.coverage_path = None;
1767
1768 let err = table
1769 .coverage_ratio_for_range(ts_from_secs(0), ts_from_secs(240))
1770 .await
1771 .expect_err("missing segment coverage_path should bubble up");
1772 assert!(matches!(
1773 err,
1774 TableError::CoverageQuery {
1775 source: CoverageQueryError::ExistingSegmentMissingCoverageMetadata {
1776 segment_path: path,
1777 ..
1778 }
1779 } if path == segment_path
1780 ));
1781 Ok(())
1782 }
1783
1784 #[tokio::test]
1785 async fn recovery_missing_segment_sidecar_preserves_storage_source_and_backtrace() -> TestResult
1786 {
1787 let (tmp, mut table) = table_with_sparse_coverage().await?;
1788 table.state_mut().table_coverage = None;
1789 let (segment_path, coverage_path) = single_segment_and_coverage_paths(&table);
1790 tokio::fs::remove_file(tmp.path().join(&coverage_path)).await?;
1791
1792 let error = table
1793 .coverage_ratio_for_range(ts_from_secs(0), ts_from_secs(240))
1794 .await
1795 .expect_err("missing segment sidecar must fail recovery");
1796 let query = coverage_query_source(&error);
1797 assert!(matches!(
1798 query,
1799 CoverageQueryError::SegmentCoverageSidecarRead {
1800 segment_path: actual_segment_path,
1801 coverage_path: actual_coverage_path,
1802 ..
1803 } if actual_segment_path == &segment_path && actual_coverage_path == &coverage_path
1804 ));
1805 let sidecar = coverage_sidecar_source(query);
1806 let storage = sidecar
1807 .source()
1808 .and_then(|source| source.downcast_ref::<StorageError>())
1809 .expect("storage source");
1810 assert!(matches!(storage, StorageError::NotFound { .. }));
1811 assert!(std::ptr::eq(
1812 ErrorCompat::backtrace(&error).expect("table backtrace"),
1813 ErrorCompat::backtrace(storage).expect("storage backtrace"),
1814 ));
1815 Ok(())
1816 }
1817
1818 #[tokio::test]
1819 async fn recovery_corrupt_segment_sidecar_preserves_codec_source_and_backtrace() -> TestResult {
1820 let (tmp, mut table) = table_with_sparse_coverage().await?;
1821 table.state_mut().table_coverage = None;
1822 let (segment_path, coverage_path) = single_segment_and_coverage_paths(&table);
1823 tokio::fs::write(tmp.path().join(&coverage_path), b"not a coverage bitmap").await?;
1824
1825 let error = table
1826 .coverage_ratio_for_range(ts_from_secs(0), ts_from_secs(240))
1827 .await
1828 .expect_err("corrupt segment sidecar must fail recovery");
1829 let query = coverage_query_source(&error);
1830 assert!(matches!(
1831 query,
1832 CoverageQueryError::SegmentCoverageSidecarRead {
1833 segment_path: actual_segment_path,
1834 coverage_path: actual_coverage_path,
1835 ..
1836 } if actual_segment_path == &segment_path && actual_coverage_path == &coverage_path
1837 ));
1838 let sidecar = coverage_sidecar_source(query);
1839 let codec = sidecar
1840 .source()
1841 .and_then(|source| source.downcast_ref::<CoverageCodecError>())
1842 .expect("codec source");
1843 assert!(matches!(
1844 codec,
1845 CoverageCodecError::BitmapDeserialization { .. }
1846 ));
1847 assert!(std::ptr::eq(
1848 ErrorCompat::backtrace(&error).expect("table backtrace"),
1849 ErrorCompat::backtrace(codec).expect("codec backtrace"),
1850 ));
1851 Ok(())
1852 }
1853
1854 #[tokio::test]
1855 async fn coverage_ratio_handles_empty_table() -> TestResult {
1856 let (_tmp, table) = make_table().await?;
1857 let ratio = table
1858 .coverage_ratio_for_range(ts_from_secs(0), ts_from_secs(60))
1859 .await?;
1860 assert_eq!(ratio, 0.0);
1861 Ok(())
1862 }
1863
1864 #[tokio::test]
1865 async fn coverage_ratio_handles_interval_ids_above_u32() -> TestResult {
1866 let (_tmp, table) = make_table().await?;
1867 let start = ts_from_secs(0);
1868 let end = ts_from_secs(((u32::MAX as i64) + 3) * 60);
1869
1870 let ratio = table.coverage_ratio_for_range(start, end).await?;
1871 assert_eq!(ratio, 0.0);
1872 Ok(())
1873 }
1874
1875 #[tokio::test]
1876 async fn max_gap_len_reports_missing_run() -> TestResult {
1877 let (_tmp, table) = table_with_sparse_coverage().await?;
1878 let gap = table
1879 .max_gap_len_for_range(ts_from_secs(0), ts_from_secs(240))
1880 .await?;
1881 assert_eq!(gap, 1);
1882 Ok(())
1883 }
1884
1885 #[tokio::test]
1886 async fn last_window_returns_none_for_zero_length() -> TestResult {
1887 let (_tmp, table) = make_table().await?;
1888 let res = table.last_fully_covered_window(ts_from_secs(0), 0).await?;
1889 assert!(res.is_none());
1890 Ok(())
1891 }
1892
1893 #[tokio::test]
1894 async fn last_window_respects_half_open_end_and_run_length() -> TestResult {
1895 let (_tmp, table) = table_with_contiguous_run().await?;
1896 let ts_end = ts_from_secs(360); let win = table
1899 .last_fully_covered_window(ts_end, 2)
1900 .await?
1901 .expect("window should be present");
1902 assert_eq!(win, 0x8000_0000_0000_0004..=0x8000_0000_0000_0005);
1903
1904 let none = table.last_fully_covered_window(ts_end, 3).await?;
1905 assert!(none.is_none());
1906 Ok(())
1907 }
1908
1909 #[tokio::test]
1910 async fn coverage_queries_validate_before_reading_coverage() -> TestResult {
1911 let (_tmp, timestamp_table) = make_table().await?;
1912 let timestamp = ts_from_secs(1);
1913 assert!(matches!(
1914 timestamp_table
1915 .coverage_ratio_for_range(timestamp, timestamp)
1916 .await,
1917 Err(TableError::CoverageQuery {
1918 source: CoverageQueryError::InvalidRange {
1919 source: IndexValueError::InvalidRange { .. },
1920 ..
1921 }
1922 })
1923 ));
1924
1925 let signed_kind = IndexKind::Int64 {
1926 index_granularity: NonZeroU64::new(1).unwrap(),
1927 };
1928 let (_tmp, signed_table) =
1929 table_with_index_coverage(signed_kind, Coverage::empty()).await?;
1930 assert!(matches!(
1931 signed_table.max_gap_len_for_range(1i64, 0i64).await,
1932 Err(TableError::CoverageQuery {
1933 source: CoverageQueryError::InvalidRange {
1934 source: IndexValueError::InvalidRange { .. },
1935 ..
1936 }
1937 })
1938 ));
1939
1940 let tmp = TempDir::new()?;
1941 let kind = IndexKind::UInt64 {
1942 index_granularity: NonZeroU64::new(1).unwrap(),
1943 };
1944 let meta = TableMeta::new_time_series(IndexSpec {
1945 column: "offset".to_string(),
1946 entity_columns: Vec::new(),
1947 kind: kind.clone(),
1948 });
1949 let mut table = TimeSeriesTable::create(TableLocation::local(tmp.path()), meta).await?;
1950 let version = table.state().version;
1951 table.state_mut().table_coverage = Some(TableCoveragePointer {
1952 index_kind: kind,
1953 coverage_path: "_coverage/table/missing.roar".to_string(),
1954 version,
1955 });
1956
1957 let error = table
1958 .last_fully_covered_window(ts_from_secs(1), 1)
1959 .await
1960 .expect_err("endpoint domain must match the table index");
1961 let query = coverage_query_source(&error);
1962 assert!(matches!(
1963 query,
1964 CoverageQueryError::InvalidRange {
1965 source: IndexValueError::KindMismatch {
1966 expected: "uint64",
1967 actual: "timestamp"
1968 },
1969 ..
1970 }
1971 ));
1972 assert!(matches!(
1973 query.source(),
1974 Some(source) if source.downcast_ref::<IndexValueError>().is_some()
1975 ));
1976 assert!(std::ptr::eq(
1977 ErrorCompat::backtrace(&error).expect("table backtrace"),
1978 ErrorCompat::backtrace(query).expect("coverage query backtrace"),
1979 ));
1980
1981 assert!(matches!(
1982 table.coverage_ratio_for_range(1u64, 1u64).await,
1983 Err(TableError::CoverageQuery {
1984 source: CoverageQueryError::InvalidRange {
1985 source: IndexValueError::InvalidRange { .. },
1986 ..
1987 }
1988 })
1989 ));
1990 assert!(matches!(
1991 table.max_gap_len_for_range(2u64, 1u64).await,
1992 Err(TableError::CoverageQuery {
1993 source: CoverageQueryError::InvalidRange {
1994 source: IndexValueError::InvalidRange { .. },
1995 ..
1996 }
1997 })
1998 ));
1999 assert!(matches!(
2000 table.coverage_ratio_for_range(0u64, 1i64).await,
2001 Err(TableError::CoverageQuery {
2002 source: CoverageQueryError::InvalidRange {
2003 source: IndexValueError::KindMismatch { .. },
2004 ..
2005 }
2006 })
2007 ));
2008 assert_eq!(table.last_fully_covered_window(0u64, 0).await?, None);
2009 assert!(matches!(
2010 table.last_fully_covered_window(0u64, 1).await,
2011 Err(TableError::CoverageQuery {
2012 source: CoverageQueryError::IndexIntervalMapping {
2013 source: IndexIntervalMappingError::RangeEndUnderflow { .. },
2014 ..
2015 }
2016 })
2017 ));
2018 Ok(())
2019 }
2020
2021 #[tokio::test]
2022 async fn last_window_errors_when_recovery_fails() -> TestResult {
2023 let (_tmp, mut table) = table_with_contiguous_run().await?;
2024 table.state_mut().table_coverage = None;
2025 let segment = table
2026 .state_mut()
2027 .segments
2028 .values_mut()
2029 .next()
2030 .expect("segment present");
2031 let segment_path = segment.path.clone();
2032 segment.coverage_path = None;
2033
2034 let err = table
2035 .last_fully_covered_window(ts_from_secs(360), 1)
2036 .await
2037 .expect_err("missing coverage_path should bubble up");
2038 assert!(matches!(
2039 err,
2040 TableError::CoverageQuery {
2041 source: CoverageQueryError::ExistingSegmentMissingCoverageMetadata {
2042 segment_path: path,
2043 ..
2044 }
2045 } if path == segment_path
2046 ));
2047 Ok(())
2048 }
2049}