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