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        self.replay_table_state(|_, _| {}).await
104    }
105
106    /// Rebuild table state while observing each decoded action once.
107    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        // Replay all commits from 1..=current_version in order
132        for v in 1..=current_version {
133            let commit = self.load_commit(v).await?;
134
135            // Defensive: file name version should match payload
136            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                        // Full replacement of TableMeta.
176                        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}