Skip to main content

timeseries_table_format/transaction_log/
table_state.rs

1//! Reconstructing the current table state by replaying log commits.
2//!
3//! `TableState` materializes the metadata stored in `_timeseries_log/` and the
4//! [`TransactionLogStore::rebuild_table_state`] helper walks all commits from version 1 up
5//! to the `CURRENT` pointer, applying their actions in order. This keeps read
6//! logic isolated from the append-only write path and documents the invariant
7//! that table readers must see a state consistent with the latest committed
8//! version.
9use 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")]
20/// Return the number of rebuilds invoked on the current thread (test-only).
21pub fn rebuild_table_state_count() -> usize {
22    REBUILD_TABLE_STATE_COUNT.with(|c| c.get())
23}
24
25#[cfg(feature = "test-counters")]
26/// Reset the rebuild counter to zero (test-only).
27pub 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/// Pointer to table coverage metadata including index descriptor, path, and version.
51#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct TableCoveragePointer {
53    /// Canonical ordered-index coverage descriptor.
54    pub index_kind: IndexKind,
55    /// Path to the coverage metadata file.
56    pub coverage_path: String,
57    /// Version number associated with this coverage pointer.
58    pub version: u64,
59}
60
61/// In-memory view of table metadata and live segments, reconstructed from the log.
62///
63/// Invariant:
64/// - `version` matches the CURRENT pointer.
65/// - `table_meta` and `segments` are the result of applying all commits from
66///   version 1 through `version` in order.
67#[derive(Debug, Clone, PartialEq, Eq)]
68pub struct TableState {
69    /// Latest committed version recorded in CURRENT.
70    pub version: u64,
71    /// Table-level metadata reconstructed from the log.
72    pub table_meta: TableMeta,
73    /// Current live segments keyed by canonical table-relative path.
74    pub segments: HashMap<String, SegmentMeta>,
75
76    /// Optional pointer to the latest table coverage metadata.
77    pub table_coverage: Option<TableCoveragePointer>,
78}
79
80impl TableState {
81    /// Return live segments sorted deterministically by ordered-index bounds.
82    ///
83    /// Ordering is by `index_min`, then `index_max`, and finally `path` as a
84    /// stable tie-breaker.
85    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    /// Rebuild the current TableState by replaying all commits up to CURRENT.
96    ///
97    /// v0.1 behavior:
98    /// - If CURRENT == 0 (no commits), this returns
99    ///   [`CommitError::UninitializedTableState`].
100    /// - The first commit must include at least one UpdateTableMeta action
101    ///   to bootstrap TableMeta; the last UpdateTableMeta wins.
102    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        // Replay all commits from 1..=current_version in order
121        for v in 1..=current_version {
122            let commit = self.load_commit(v).await?;
123
124            // Defensive: file name version should match payload
125            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                        // Full replacement of TableMeta.
164                        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}