1use std::collections::HashMap;
10
11#[cfg(feature = "test-counters")]
12use std::cell::Cell;
13
14#[cfg(feature = "test-counters")]
15thread_local! {
16 static REBUILD_TABLE_STATE_COUNT: Cell<usize> = const { Cell::new(0) };
17}
18
19#[cfg(feature = "test-counters")]
20pub fn rebuild_table_state_count() -> usize {
22 REBUILD_TABLE_STATE_COUNT.with(|c| c.get())
23}
24
25#[cfg(feature = "test-counters")]
26pub fn reset_rebuild_table_state_count() {
28 REBUILD_TABLE_STATE_COUNT.with(|c| c.set(0));
29}
30
31use crate::{
32 metadata::{
33 schema_compat::{ensure_entity_identity_matches_schema, ensure_index_spec_matches_schema},
34 segments::sort_segment_meta_by_index,
35 },
36 storage::ensure_canonical_relative_storage_path,
37 transaction_log::*,
38};
39
40fn validate_persisted_storage_path(path: &str, description: &str) -> Result<(), CommitError> {
41 ensure_canonical_relative_storage_path(path).map_err(|source| {
42 CommitError::InvalidPersistedPath {
43 description: description.to_string(),
44 path: path.to_string(),
45 source: Box::new(source),
46 }
47 })
48}
49
50#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct TableCoveragePointer {
53 pub index_kind: IndexKind,
55 pub coverage_path: String,
57 pub version: u64,
59}
60
61#[derive(Debug, Clone, PartialEq, Eq)]
68pub struct TableState {
69 pub version: u64,
71 pub table_meta: TableMeta,
73 pub segments: HashMap<String, SegmentMeta>,
75
76 pub table_coverage: Option<TableCoveragePointer>,
78}
79
80impl TableState {
81 pub fn segments_sorted_by_index(
86 &self,
87 ) -> Result<Vec<&SegmentMeta>, crate::metadata::index::IndexValueError> {
88 let mut v: Vec<&SegmentMeta> = self.segments.values().collect();
89 sort_segment_meta_by_index(&mut v)?;
90 Ok(v)
91 }
92}
93
94impl TransactionLogStore {
95 pub async fn rebuild_table_state(&self) -> Result<TableState, CommitError> {
103 self.replay_table_state(|_, _| {}).await
104 }
105
106 pub(crate) async fn replay_table_state<F>(
108 &self,
109 mut observe: F,
110 ) -> Result<TableState, CommitError>
111 where
112 F: FnMut(u64, &LogAction),
113 {
114 #[cfg(feature = "test-counters")]
115 REBUILD_TABLE_STATE_COUNT.with(|c| c.set(c.get() + 1));
116
117 let current_version = self.load_current_version().await?;
118
119 if current_version == 0 {
120 return Err(CommitError::UninitializedTableState {
121 backtrace: snafu::Backtrace::capture(),
122 });
123 }
124
125 let mut table_meta: Option<TableMeta> = None;
126 let mut segments: HashMap<String, SegmentMeta> = HashMap::new();
127 let mut persisted_segment_layouts = Vec::new();
128
129 let mut table_coverage: Option<TableCoveragePointer> = None;
130
131 for v in 1..=current_version {
133 let commit = self.load_commit(v).await?;
134
135 if commit.version != v {
137 return Err(CommitError::CommitVersionMismatch {
138 expected: v,
139 found: commit.version,
140 backtrace: snafu::Backtrace::capture(),
141 });
142 }
143
144 for action in commit.actions {
145 observe(v, &action);
146 match action {
147 LogAction::AddSegment(meta) => {
148 validate_persisted_storage_path(&meta.path, "segment path")?;
149 if let Some(coverage_path) = &meta.coverage_path {
150 validate_persisted_storage_path(
151 coverage_path,
152 "segment coverage path",
153 )?;
154 }
155 if segments.contains_key(&meta.path) {
156 return Err(CommitError::DuplicateLiveSegmentPath {
157 path: meta.path,
158 backtrace: snafu::Backtrace::capture(),
159 });
160 }
161 persisted_segment_layouts
162 .push((meta.path.clone(), meta.entity_layout.clone()));
163 segments.insert(meta.path.clone(), meta);
164 }
165 LogAction::RemoveSegment { path } => {
166 validate_persisted_storage_path(&path, "segment path")?;
167 segments.remove(&path);
168 }
169 LogAction::UpdateTableMeta(delta) => {
170 if let Some(previous) = &table_meta {
171 previous
172 .ensure_valid_transition_to(&delta)
173 .map_err(CommitError::from)?;
174 }
175 table_meta = Some(delta);
177 }
178 LogAction::UpdateTableCoverage {
179 index_kind,
180 coverage_path,
181 } => {
182 validate_persisted_storage_path(&coverage_path, "table coverage path")?;
183 table_coverage = Some(TableCoveragePointer {
184 index_kind,
185 coverage_path,
186 version: v,
187 })
188 }
189 }
190 }
191 }
192
193 let table_meta = table_meta.ok_or_else(|| CommitError::MissingTableMetadata {
194 current_version,
195 backtrace: snafu::Backtrace::capture(),
196 })?;
197 table_meta
198 .ensure_read_compatible()
199 .map_err(CommitError::from)?;
200
201 if let TableKind::TimeSeries(index) = &table_meta.kind {
202 index
203 .validate()
204 .map_err(|source| CommitError::InvalidIndexSpec {
205 source,
206 backtrace: snafu::Backtrace::capture(),
207 })?;
208 if let Some(pointer) = &table_coverage
209 && pointer.index_kind != index.kind
210 {
211 return Err(CommitError::CoverageIndexKindMismatch {
212 expected: index.kind.clone(),
213 actual: pointer.index_kind.clone(),
214 pointer_version: pointer.version,
215 backtrace: Box::new(snafu::Backtrace::capture()),
216 });
217 }
218 let schema = table_meta.logical_schema.as_ref();
219 if let Some(schema) = schema {
220 ensure_index_spec_matches_schema(schema, index).map_err(|source| {
221 CommitError::TableSchemaCompatibility {
222 source: Box::new(source),
223 }
224 })?;
225 }
226 if schema.is_none() && !persisted_segment_layouts.is_empty() {
227 return Err(CommitError::MissingLogicalSchemaForSegments {
228 backtrace: snafu::Backtrace::capture(),
229 });
230 }
231 let entity_column_count = index.entity_columns.len();
232 for (path, layout) in persisted_segment_layouts {
233 match (&layout, entity_column_count) {
234 (SegmentEntityLayout::NotApplicable, 0) | (SegmentEntityLayout::Mixed, 1..) => {
235 }
236 (SegmentEntityLayout::Single(identity), 1..) => {
237 let Some(schema) = schema else {
238 return Err(CommitError::MissingLogicalSchemaForSegments {
239 backtrace: snafu::Backtrace::capture(),
240 });
241 };
242 ensure_entity_identity_matches_schema(schema, index, identity).map_err(
243 |source| CommitError::SegmentEntityIdentitySchema {
244 path: path.clone(),
245 source: Box::new(source),
246 },
247 )?;
248 }
249 _ => {
250 return Err(CommitError::InvalidSegmentEntityLayout {
251 path,
252 entity_column_count,
253 layout,
254 backtrace: snafu::Backtrace::capture(),
255 });
256 }
257 }
258 }
259 for segment in segments.values() {
260 segment.validate_bounds(&index.kind).map_err(|source| {
261 CommitError::SegmentMetadata {
262 source: Box::new(source),
263 }
264 })?;
265 }
266 }
267
268 Ok(TableState {
269 version: current_version,
270 table_meta,
271 segments,
272 table_coverage,
273 })
274 }
275}
276
277#[cfg(test)]
278mod tests {
279 use super::*;
280 use crate::coverage::EntityIdentity;
281 use crate::metadata::{
282 logical_schema::{LogicalDataType, LogicalField, LogicalSchema, LogicalTimestampUnit},
283 protocol::TABLE_PROTOCOL_VERSION,
284 };
285 use crate::storage::layout;
286 use crate::storage::{StorageError, TableLocation};
287 use crate::transaction_log::{
288 FileFormat, IndexKind, IndexSpec, LogAction, SegmentEntityLayout, SegmentMeta, TableKind,
289 TableMeta, TimeIndexGranularity, TransactionLogStore,
290 };
291 use chrono::TimeZone;
292 use tempfile::TempDir;
293
294 type TestResult = Result<(), Box<dyn std::error::Error>>;
295
296 fn create_test_log_store() -> (TempDir, TransactionLogStore) {
297 let tmp = TempDir::new().expect("create temp dir");
298 let location = TableLocation::local(tmp.path());
299 let store = TransactionLogStore::new(location);
300 (tmp, store)
301 }
302
303 fn sample_table_meta() -> TableMeta {
304 let entity_columns = vec!["symbol".to_string()];
305 TableMeta {
306 kind: TableKind::TimeSeries(IndexSpec {
307 column: "ts".to_string(),
308 entity_columns: entity_columns.clone(),
309 kind: IndexKind::Timestamp {
310 index_granularity: TimeIndexGranularity::Minutes(1),
311 timezone: None,
312 },
313 }),
314 logical_schema: Some(schema_for_entities(&entity_columns)),
315 created_at: chrono::Utc
316 .with_ymd_and_hms(2025, 1, 1, 0, 0, 0)
317 .single()
318 .expect("valid sample table metadata timestamp"),
319 protocol_version: TABLE_PROTOCOL_VERSION,
320 required_reader_features: Default::default(),
321 required_writer_features: Default::default(),
322 }
323 }
324
325 fn schema_for_entities(entity_columns: &[String]) -> LogicalSchema {
326 let mut fields = vec![LogicalField {
327 name: "ts".to_string(),
328 data_type: LogicalDataType::Timestamp {
329 unit: LogicalTimestampUnit::Millis,
330 timezone: None,
331 },
332 nullable: false,
333 }];
334 fields.extend(entity_columns.iter().map(|column| LogicalField {
335 name: column.clone(),
336 data_type: LogicalDataType::Utf8,
337 nullable: false,
338 }));
339 LogicalSchema::new(fields).expect("valid test schema")
340 }
341
342 fn sample_segment(id: &str) -> SegmentMeta {
343 SegmentMeta {
344 path: format!("data/{id}.parquet"),
345 format: FileFormat::Parquet,
346 entity_layout: SegmentEntityLayout::Single(
347 EntityIdentity::try_new(vec!["A".into()]).expect("valid sample identity"),
348 ),
349 index_min: IndexValue::Timestamp(
350 chrono::Utc
351 .with_ymd_and_hms(2025, 1, 1, 0, 0, 0)
352 .single()
353 .expect("valid sample segment index_min"),
354 ),
355 index_max: IndexValue::Timestamp(
356 chrono::Utc
357 .with_ymd_and_hms(2025, 1, 1, 1, 0, 0)
358 .single()
359 .expect("valid sample segment index_max"),
360 ),
361 row_count: 42,
362 file_size: None,
363 coverage_path: None,
364 }
365 }
366
367 async fn prepend_raw_action(tmp: &TempDir, action: serde_json::Value) -> TestResult {
368 let commit_path = tmp.path().join(layout::commit_rel_path(1));
369 let mut commit: serde_json::Value =
370 serde_json::from_slice(&tokio::fs::read(&commit_path).await?)?;
371 commit["actions"]
372 .as_array_mut()
373 .expect("valid committed actions")
374 .insert(0, action);
375 tokio::fs::write(commit_path, serde_json::to_vec(&commit)?).await?;
376 Ok(())
377 }
378
379 fn segment_with_ts(id: &str, ts_min: i64, ts_max: i64) -> SegmentMeta {
380 SegmentMeta {
381 path: format!("data/{id}.parquet"),
382 format: FileFormat::Parquet,
383 entity_layout: SegmentEntityLayout::Single(
384 EntityIdentity::try_new(vec!["A".into()]).expect("valid sample identity"),
385 ),
386 index_min: (chrono::Utc.timestamp_opt(ts_min, 0).single().unwrap()).into(),
387 index_max: (chrono::Utc.timestamp_opt(ts_max, 0).single().unwrap()).into(),
388 row_count: 1,
389 file_size: None,
390 coverage_path: None,
391 }
392 }
393
394 #[test]
395 fn segments_sorted_by_index_orders_hashmap_deterministically() {
396 let mut segments = HashMap::new();
397 let seg_c = segment_with_ts("c", 10, 30);
398 let seg_a = segment_with_ts("a", 10, 20);
399 let seg_d = segment_with_ts("d", 5, 7);
400 let seg_b = segment_with_ts("b", 10, 20);
401
402 segments.insert(seg_c.path.clone(), seg_c);
403 segments.insert(seg_a.path.clone(), seg_a);
404 segments.insert(seg_d.path.clone(), seg_d);
405 segments.insert(seg_b.path.clone(), seg_b);
406
407 let state = TableState {
408 version: 3,
409 table_meta: sample_table_meta(),
410 segments,
411 table_coverage: None,
412 };
413
414 let ordered: Vec<(i64, i64, String)> = state
415 .segments_sorted_by_index()
416 .unwrap()
417 .iter()
418 .map(|seg| match (&seg.index_min, &seg.index_max) {
419 (IndexValue::Timestamp(min), IndexValue::Timestamp(max)) => {
420 (min.timestamp(), max.timestamp(), seg.path.clone())
421 }
422 _ => panic!("expected timestamp test bounds"),
423 })
424 .collect();
425
426 let mut expected = ordered.clone();
427 expected.sort();
428 assert_eq!(ordered, expected);
429 }
430
431 #[tokio::test]
432 async fn rebuild_table_state_happy_path() -> TestResult {
433 let (_tmp, store) = create_test_log_store();
434 let meta = sample_table_meta();
435 let seg1 = sample_segment("seg1");
436 let seg2 = sample_segment("seg2");
437
438 let v1 = store
439 .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta.clone())])
440 .await?;
441 let v2 = store
442 .commit_with_expected_version(
443 v1,
444 vec![
445 LogAction::AddSegment(seg1.clone()),
446 LogAction::AddSegment(seg2.clone()),
447 ],
448 )
449 .await?;
450 let v3 = store
451 .commit_with_expected_version(
452 v2,
453 vec![LogAction::RemoveSegment {
454 path: seg1.path.clone(),
455 }],
456 )
457 .await?;
458
459 let state = store.rebuild_table_state().await?;
460 assert_eq!(state.version, v3);
461 assert_eq!(state.table_meta, meta);
462 assert!(state.segments.contains_key(&seg2.path));
463 assert!(!state.segments.contains_key(&seg1.path));
464 Ok(())
465 }
466
467 #[tokio::test]
468 async fn replay_table_state_observes_each_action_once() -> TestResult {
469 let (_tmp, store) = create_test_log_store();
470 let segment = sample_segment("seg1");
471 store
472 .commit_with_expected_version(
473 0,
474 vec![
475 LogAction::UpdateTableMeta(sample_table_meta()),
476 LogAction::AddSegment(segment.clone()),
477 ],
478 )
479 .await?;
480 store
481 .commit_with_expected_version(1, vec![LogAction::RemoveSegment { path: segment.path }])
482 .await?;
483 let mut observed_versions = Vec::new();
484
485 let state = store
486 .replay_table_state(|version, _| observed_versions.push(version))
487 .await?;
488
489 assert_eq!(state.version, 2);
490 assert_eq!(observed_versions, [1, 1, 2]);
491 Ok(())
492 }
493
494 #[tokio::test]
495 async fn rebuild_table_state_leaves_table_kind_validation_to_callers() -> TestResult {
496 let (_tmp, store) = create_test_log_store();
497 let mut meta = sample_table_meta();
498 meta.kind = TableKind::Generic;
499 store
500 .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta.clone())])
501 .await?;
502
503 let state = store.rebuild_table_state().await?;
504
505 assert_eq!(state.table_meta, meta);
506 Ok(())
507 }
508
509 #[tokio::test]
510 async fn rebuild_table_state_errors_when_current_zero() {
511 let (_tmp, store) = create_test_log_store();
512
513 let err = store
514 .rebuild_table_state()
515 .await
516 .expect_err("expected error");
517 assert!(matches!(err, CommitError::UninitializedTableState { .. }));
518 }
519
520 #[tokio::test]
521 async fn rebuild_table_state_errors_when_no_table_meta() -> TestResult {
522 let (_tmp, store) = create_test_log_store();
523 let seg = sample_segment("seg");
524
525 store
526 .commit_with_expected_version(0, vec![LogAction::AddSegment(seg.clone())])
527 .await?;
528
529 let err = store
530 .rebuild_table_state()
531 .await
532 .expect_err("expected error");
533 assert!(matches!(err, CommitError::MissingTableMetadata { .. }));
534 Ok(())
535 }
536
537 #[tokio::test]
538 async fn rebuild_table_state_rejects_version_6_metadata() -> TestResult {
539 let (tmp, store) = create_test_log_store();
540 let log_dir = tmp.path().join(layout::LOG_DIR_NAME);
541 tokio::fs::create_dir_all(&log_dir).await?;
542 tokio::fs::write(
543 tmp.path().join(layout::commit_rel_path(1)),
544 r#"{
545 "version": 1,
546 "base_version": 0,
547 "timestamp": "2025-01-01T00:00:00Z",
548 "actions": [{
549 "UpdateTableMeta": {
550 "kind": {"TimeSeries": {
551 "column": "ts",
552 "entity_columns": ["symbol"],
553 "kind": {
554 "type": "timestamp",
555 "bucket": {"Minutes": 1}
556 }
557 }},
558 "logical_schema": null,
559 "created_at": "2025-01-01T00:00:00Z",
560 "protocol_version": 6,
561 "required_reader_features": [],
562 "required_writer_features": []
563 }
564 }]
565 }"#,
566 )
567 .await?;
568 tokio::fs::write(tmp.path().join(layout::current_rel_path()), "1\n").await?;
569
570 let err = store
571 .rebuild_table_state()
572 .await
573 .expect_err("version 6 should be rejected");
574 assert!(matches!(
575 err,
576 CommitError::Protocol {
577 source: TableProtocolError::UnsupportedVersion {
578 expected: TABLE_PROTOCOL_VERSION,
579 found: 6,
580 },
581 ..
582 }
583 ));
584 Ok(())
585 }
586
587 #[tokio::test]
588 async fn rebuild_table_state_rejects_persisted_protocol_downgrade() -> TestResult {
589 let (_tmp, store) = create_test_log_store();
590 let meta = sample_table_meta();
591 store
592 .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta.clone())])
593 .await?;
594
595 let mut downgraded = meta;
596 downgraded.protocol_version -= 1;
597 store
598 .commit_with_expected_version(1, vec![LogAction::UpdateTableMeta(downgraded)])
599 .await?;
600
601 let error = store.rebuild_table_state().await.unwrap_err();
602 assert!(matches!(
603 error,
604 CommitError::Protocol {
605 source: TableProtocolError::UnsupportedVersion {
606 expected: TABLE_PROTOCOL_VERSION,
607 found: 6,
608 },
609 ..
610 }
611 ));
612 Ok(())
613 }
614
615 #[tokio::test]
616 async fn rebuild_table_state_allows_unknown_writer_features() -> TestResult {
617 let (_tmp, store) = create_test_log_store();
618 let mut meta = sample_table_meta();
619 meta.required_writer_features
620 .insert("future_writer".to_string());
621 store
622 .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
623 .await?;
624
625 let state = store.rebuild_table_state().await?;
626 assert_eq!(
627 state.table_meta.required_writer_features(),
628 &["future_writer".to_string()].into_iter().collect()
629 );
630 Ok(())
631 }
632
633 #[tokio::test]
634 async fn rebuild_table_state_rejects_unknown_reader_features() -> TestResult {
635 let (_tmp, store) = create_test_log_store();
636 let mut meta = sample_table_meta();
637 meta.required_reader_features
638 .insert("future_reader".to_string());
639 store
640 .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
641 .await?;
642
643 let error = store.rebuild_table_state().await.unwrap_err();
644 assert!(matches!(
645 error,
646 CommitError::Protocol {
647 source: TableProtocolError::UnsupportedReaderFeatures { features },
648 ..
649 } if features == ["future_reader"]
650 ));
651 Ok(())
652 }
653
654 #[tokio::test]
655 async fn rebuild_table_state_applies_writer_feature_before_unknown_action() -> TestResult {
656 let (tmp, store) = create_test_log_store();
657 let mut meta = sample_table_meta();
658 meta.required_writer_features
659 .insert("future_action".to_string());
660 store
661 .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta.clone())])
662 .await?;
663 prepend_raw_action(
664 &tmp,
665 serde_json::json!({"FutureAction": {"payload": "ignored by readers"}}),
666 )
667 .await?;
668
669 let state = store.rebuild_table_state().await?;
670
671 assert_eq!(state.table_meta, meta);
672 Ok(())
673 }
674
675 #[tokio::test]
676 async fn rebuild_table_state_checks_reader_features_before_action_decoding() -> TestResult {
677 let (tmp, store) = create_test_log_store();
678 let mut meta = sample_table_meta();
679 meta.required_reader_features
680 .insert("future_action".to_string());
681 store
682 .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
683 .await?;
684 prepend_raw_action(&tmp, serde_json::json!({"AddSegment": {"path": false}})).await?;
685 let commit_path = tmp.path().join(layout::commit_rel_path(1));
686 let mut commit: serde_json::Value =
687 serde_json::from_slice(&tokio::fs::read(&commit_path).await?)?;
688 let metadata = commit["actions"]
689 .as_array_mut()
690 .and_then(|actions| {
691 actions
692 .iter_mut()
693 .find_map(|action| action.get_mut("UpdateTableMeta"))
694 })
695 .expect("valid committed metadata action");
696 metadata["kind"] = serde_json::json!({"FutureKind": {"payload": false}});
697 tokio::fs::write(commit_path, serde_json::to_vec(&commit)?).await?;
698
699 let error = store.rebuild_table_state().await.unwrap_err();
700
701 assert!(matches!(
702 error,
703 CommitError::Protocol {
704 source: TableProtocolError::UnsupportedReaderFeatures { features },
705 ..
706 } if features == ["future_action"]
707 ));
708 Ok(())
709 }
710
711 #[tokio::test]
712 async fn rebuild_table_state_rejects_malformed_known_action() -> TestResult {
713 let (tmp, store) = create_test_log_store();
714 store
715 .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(sample_table_meta())])
716 .await?;
717 prepend_raw_action(&tmp, serde_json::json!({"AddSegment": {"path": false}})).await?;
718
719 let error = store.rebuild_table_state().await.unwrap_err();
720
721 assert!(matches!(error, CommitError::CommitDeserialization { .. }));
722 Ok(())
723 }
724
725 #[tokio::test]
726 async fn rebuild_table_state_rejects_removed_protocol_features() -> TestResult {
727 let (_tmp, store) = create_test_log_store();
728 let mut initial = sample_table_meta();
729 initial
730 .required_writer_features
731 .insert("future_writer".to_string());
732 store
733 .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(initial.clone())])
734 .await?;
735
736 initial.required_writer_features.clear();
737 store
738 .commit_with_expected_version(1, vec![LogAction::UpdateTableMeta(initial)])
739 .await?;
740
741 let error = store.rebuild_table_state().await.unwrap_err();
742 assert!(matches!(
743 error,
744 CommitError::Protocol {
745 source: TableProtocolError::WriterFeaturesRemoved { features },
746 ..
747 } if features == ["future_writer"]
748 ));
749 Ok(())
750 }
751
752 #[tokio::test]
753 async fn rebuild_table_state_rejects_invalid_persisted_segment_bounds() -> TestResult {
754 let (_tmp, store) = create_test_log_store();
755 let mut segment = sample_segment("reversed");
756 segment.index_min =
757 IndexValue::Timestamp(chrono::Utc.timestamp_opt(2, 0).single().unwrap());
758 segment.index_max =
759 IndexValue::Timestamp(chrono::Utc.timestamp_opt(1, 0).single().unwrap());
760
761 store
762 .commit_with_expected_version(
763 0,
764 vec![
765 LogAction::UpdateTableMeta(sample_table_meta()),
766 LogAction::AddSegment(segment),
767 ],
768 )
769 .await?;
770
771 let error = store.rebuild_table_state().await.unwrap_err();
772 assert!(matches!(error, CommitError::SegmentMetadata { .. }));
773 assert!(error.to_string().contains("Invalid ordered-index bounds"));
774 Ok(())
775 }
776
777 #[tokio::test]
778 async fn rebuild_table_state_rejects_inapplicable_entity_layouts() -> TestResult {
779 let single = SegmentEntityLayout::Single(EntityIdentity::try_new(vec!["A".into()])?);
780 let cases = [
781 (
782 vec!["symbol".to_string()],
783 SegmentEntityLayout::NotApplicable,
784 "Invalid entity layout",
785 ),
786 (
787 Vec::new(),
788 SegmentEntityLayout::Mixed,
789 "Invalid entity layout",
790 ),
791 (Vec::new(), single.clone(), "Invalid entity layout"),
792 (
793 vec!["site".to_string(), "device".to_string()],
794 single,
795 "has 1 components, but the table configures 2",
796 ),
797 ];
798
799 for (entity_columns, entity_layout, expected_message) in cases {
800 let (_tmp, store) = create_test_log_store();
801 let mut table_meta = sample_table_meta();
802 let TableKind::TimeSeries(index) = &mut table_meta.kind else {
803 unreachable!("sample metadata is time-series");
804 };
805 index.entity_columns = entity_columns.clone();
806 table_meta.logical_schema = Some(schema_for_entities(&entity_columns));
807
808 let mut segment = sample_segment("invalid-layout");
809 segment.entity_layout = entity_layout;
810 store
811 .commit_with_expected_version(
812 0,
813 vec![
814 LogAction::UpdateTableMeta(table_meta),
815 LogAction::AddSegment(segment),
816 ],
817 )
818 .await?;
819
820 let error = store
821 .rebuild_table_state()
822 .await
823 .expect_err("inapplicable entity layout should be rejected");
824 if expected_message == "Invalid entity layout" {
825 assert!(matches!(
826 error,
827 CommitError::InvalidSegmentEntityLayout { .. }
828 ));
829 } else {
830 assert!(matches!(
831 error,
832 CommitError::SegmentEntityIdentitySchema { .. }
833 ));
834 }
835 assert!(error.to_string().contains(expected_message), "{error}");
836 }
837
838 Ok(())
839 }
840
841 #[tokio::test]
842 async fn rebuild_table_state_validates_persisted_entity_component_types() -> TestResult {
843 let typed_schema = LogicalSchema::new(vec![
844 LogicalField {
845 name: "ts".to_string(),
846 data_type: LogicalDataType::Timestamp {
847 unit: LogicalTimestampUnit::Millis,
848 timezone: None,
849 },
850 nullable: false,
851 },
852 LogicalField {
853 name: "symbol".to_string(),
854 data_type: LogicalDataType::Int32,
855 nullable: false,
856 },
857 ])?;
858 let mut typed_meta = sample_table_meta();
859 typed_meta.logical_schema = Some(typed_schema);
860 let mut typed_segment = sample_segment("typed-layout");
861 typed_segment.entity_layout = SegmentEntityLayout::Single(EntityIdentity::try_new(vec![
862 crate::coverage::EntityValue::Int32(-1),
863 ])?);
864
865 let (_valid_tmp, valid_store) = create_test_log_store();
866 valid_store
867 .commit_with_expected_version(
868 0,
869 vec![
870 LogAction::UpdateTableMeta(typed_meta.clone()),
871 LogAction::AddSegment(typed_segment.clone()),
872 ],
873 )
874 .await?;
875 valid_store.rebuild_table_state().await?;
876
877 let (_invalid_tmp, invalid_store) = create_test_log_store();
878 let mut string_meta = typed_meta;
879 string_meta.logical_schema = Some(schema_for_entities(&["symbol".to_string()]));
880 invalid_store
881 .commit_with_expected_version(
882 0,
883 vec![
884 LogAction::UpdateTableMeta(string_meta),
885 LogAction::AddSegment(typed_segment),
886 ],
887 )
888 .await?;
889 let error = invalid_store
890 .rebuild_table_state()
891 .await
892 .expect_err("persisted component type must match the logical schema");
893 assert!(matches!(
894 error,
895 CommitError::SegmentEntityIdentitySchema { .. }
896 ));
897 assert!(
898 error
899 .to_string()
900 .contains("column symbol has type int32; expected utf8"),
901 "{error}"
902 );
903 Ok(())
904 }
905
906 #[tokio::test]
907 async fn rebuild_table_state_validates_removed_segment_layouts() -> TestResult {
908 let (_tmp, store) = create_test_log_store();
909 let mut segment = sample_segment("removed-invalid-layout");
910 segment.entity_layout = SegmentEntityLayout::NotApplicable;
911 let path = segment.path.clone();
912
913 store
914 .commit_with_expected_version(
915 0,
916 vec![
917 LogAction::UpdateTableMeta(sample_table_meta()),
918 LogAction::AddSegment(segment),
919 LogAction::RemoveSegment { path },
920 ],
921 )
922 .await?;
923
924 let error = store
925 .rebuild_table_state()
926 .await
927 .expect_err("removed segment metadata should still be validated");
928 assert!(matches!(
929 error,
930 CommitError::InvalidSegmentEntityLayout { .. }
931 ));
932 assert!(error.to_string().contains("Invalid entity layout"));
933 Ok(())
934 }
935
936 #[tokio::test]
937 async fn rebuild_table_state_requires_valid_entity_layout_json() -> TestResult {
938 for (replacement, expected_message) in [
939 (None, "entity_layout"),
940 (
941 Some(serde_json::json!({"Single": []})),
942 "at least one component",
943 ),
944 ] {
945 let (tmp, store) = create_test_log_store();
946 store
947 .commit_with_expected_version(
948 0,
949 vec![
950 LogAction::UpdateTableMeta(sample_table_meta()),
951 LogAction::AddSegment(sample_segment("invalid-json")),
952 ],
953 )
954 .await?;
955
956 let commit_path = tmp.path().join(layout::commit_rel_path(1));
957 let mut commit: serde_json::Value =
958 serde_json::from_slice(&tokio::fs::read(&commit_path).await?)?;
959 let segment = commit["actions"][1]["AddSegment"]
960 .as_object_mut()
961 .expect("valid committed AddSegment action");
962 match replacement {
963 Some(layout) => {
964 segment.insert("entity_layout".to_string(), layout);
965 }
966 None => {
967 segment.remove("entity_layout");
968 }
969 }
970 tokio::fs::write(&commit_path, serde_json::to_vec(&commit)?).await?;
971
972 let error = store
973 .rebuild_table_state()
974 .await
975 .expect_err("missing or malformed entity layout should be rejected");
976 assert!(matches!(error, CommitError::CommitDeserialization { .. }));
977 assert!(error.to_string().contains(expected_message), "{error}");
978 }
979
980 Ok(())
981 }
982
983 #[tokio::test]
984 async fn rebuild_table_state_rejects_noncanonical_segment_action_paths() -> TestResult {
985 for path in [
986 "",
987 "/data/seg.parquet",
988 "../data/seg.parquet",
989 "data/../seg.parquet",
990 r"data\seg.parquet",
991 "data//seg.parquet",
992 r"C:\data\seg.parquet",
993 "data/C:/seg.parquet",
994 "data/C:seg.parquet",
995 ] {
996 let mut segment = sample_segment("seg");
997 segment.path = path.to_owned();
998
999 for action in [
1000 LogAction::AddSegment(segment.clone()),
1001 LogAction::RemoveSegment {
1002 path: path.to_owned(),
1003 },
1004 ] {
1005 let (_tmp, store) = create_test_log_store();
1006 store
1007 .commit_with_expected_version(
1008 0,
1009 vec![LogAction::UpdateTableMeta(sample_table_meta()), action],
1010 )
1011 .await?;
1012
1013 let err = store
1014 .rebuild_table_state()
1015 .await
1016 .expect_err("noncanonical segment action path should be rejected");
1017 assert!(matches!(err, CommitError::InvalidPersistedPath { .. }));
1018 assert!(err.to_string().contains("segment path"), "{err}");
1019 }
1020 }
1021
1022 Ok(())
1023 }
1024
1025 #[tokio::test]
1026 async fn rebuild_table_state_rejects_noncanonical_coverage_paths() -> TestResult {
1027 for path in [
1028 "",
1029 "/tmp/coverage.roar",
1030 "../coverage.roar",
1031 "_coverage/../coverage.roar",
1032 r"_coverage\segments\coverage.roar",
1033 "_coverage//segments/coverage.roar",
1034 r"C:\coverage.roar",
1035 ] {
1036 let mut segment = sample_segment("seg");
1037 segment.coverage_path = Some(path.to_owned());
1038
1039 let index_kind = match sample_table_meta().kind {
1040 TableKind::TimeSeries(index) => index.kind,
1041 TableKind::Generic => unreachable!("sample metadata is time-series"),
1042 };
1043 for (description, action) in [
1044 ("segment coverage path", LogAction::AddSegment(segment)),
1045 (
1046 "table coverage path",
1047 LogAction::UpdateTableCoverage {
1048 index_kind,
1049 coverage_path: path.to_owned(),
1050 },
1051 ),
1052 ] {
1053 let (_tmp, store) = create_test_log_store();
1054 store
1055 .commit_with_expected_version(
1056 0,
1057 vec![LogAction::UpdateTableMeta(sample_table_meta()), action],
1058 )
1059 .await?;
1060
1061 let err = store
1062 .rebuild_table_state()
1063 .await
1064 .expect_err("noncanonical coverage path should be rejected");
1065 assert!(matches!(err, CommitError::InvalidPersistedPath { .. }));
1066 assert!(err.to_string().contains(description), "{err}");
1067 }
1068 }
1069
1070 Ok(())
1071 }
1072
1073 #[tokio::test]
1074 async fn rebuild_table_state_rejects_mismatched_table_coverage_index() -> TestResult {
1075 let (_tmp, store) = create_test_log_store();
1076 store
1077 .commit_with_expected_version(
1078 0,
1079 vec![
1080 LogAction::UpdateTableMeta(sample_table_meta()),
1081 LogAction::UpdateTableCoverage {
1082 index_kind: IndexKind::Int64 {
1083 index_granularity: std::num::NonZeroU64::new(1).unwrap(),
1084 },
1085 coverage_path: "_coverage/table/1-mismatched.roar".to_string(),
1086 },
1087 ],
1088 )
1089 .await?;
1090
1091 let err = store
1092 .rebuild_table_state()
1093 .await
1094 .expect_err("mismatched coverage index should be rejected during replay");
1095 assert!(matches!(err, CommitError::CoverageIndexKindMismatch { .. }));
1096 assert!(err.to_string().contains("Table coverage index kind"));
1097 Ok(())
1098 }
1099
1100 #[tokio::test]
1101 async fn rebuild_table_state_fails_on_corrupt_commit_payload() -> TestResult {
1102 let (tmp, store) = create_test_log_store();
1103 let meta = sample_table_meta();
1104
1105 store
1106 .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
1107 .await?;
1108
1109 let commit_path = tmp.path().join(layout::commit_rel_path(1));
1110 tokio::fs::write(&commit_path, b"not-json").await?;
1111
1112 let err = store
1113 .rebuild_table_state()
1114 .await
1115 .expect_err("expected error");
1116 assert!(matches!(err, CommitError::CommitDeserialization { .. }));
1117 Ok(())
1118 }
1119
1120 #[tokio::test]
1121 async fn rebuild_table_state_fails_when_commit_missing() -> TestResult {
1122 let (tmp, store) = create_test_log_store();
1123 let meta = sample_table_meta();
1124
1125 store
1126 .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
1127 .await?;
1128
1129 let commit_path = tmp.path().join(layout::commit_rel_path(1));
1130 tokio::fs::remove_file(&commit_path).await?;
1131
1132 let err = store
1133 .rebuild_table_state()
1134 .await
1135 .expect_err("expected error");
1136 match err {
1137 CommitError::Storage { source } => match source {
1138 StorageError::NotFound { .. } => {}
1139 other => panic!("unexpected storage error: {other:?}"),
1140 },
1141 other => panic!("expected storage error, got {other:?}"),
1142 }
1143 Ok(())
1144 }
1145}