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, path::Path};
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::{segments::cmp_segment_meta_by_time, table_metadata::TABLE_FORMAT_VERSION},
33    storage::normalize_relative_segment_path,
34    transaction_log::*,
35};
36
37fn validate_persisted_segment_path(path: &str) -> Result<(), CommitError> {
38    let (canonical, _) = match normalize_relative_segment_path(Path::new(path)) {
39        Ok(path) => path,
40        Err(source) => {
41            return CorruptStateSnafu {
42                msg: format!("Invalid persisted segment path {path:?}: {source}"),
43            }
44            .fail();
45        }
46    };
47
48    if canonical != path {
49        return CorruptStateSnafu {
50            msg: format!(
51                "Non-canonical persisted segment path {path:?}; canonical form is {canonical:?}"
52            ),
53        }
54        .fail();
55    }
56
57    Ok(())
58}
59
60/// Pointer to table coverage metadata including bucket specification, path, and version.
61#[derive(Debug, Clone, PartialEq, Eq)]
62pub struct TableCoveragePointer {
63    /// Time bucket specification for the coverage metadata.
64    pub bucket_spec: TimeBucket,
65    /// Path to the coverage metadata file.
66    pub coverage_path: String,
67    /// Version number associated with this coverage pointer.
68    pub version: u64,
69}
70
71/// In-memory view of table metadata and live segments, reconstructed from the log.
72///
73/// Invariant:
74/// - `version` matches the CURRENT pointer.
75/// - `table_meta` and `segments` are the result of applying all commits from
76///   version 1 through `version` in order.
77#[derive(Debug, Clone, PartialEq, Eq)]
78pub struct TableState {
79    /// Latest committed version recorded in CURRENT.
80    pub version: u64,
81    /// Table-level metadata reconstructed from the log.
82    pub table_meta: TableMeta,
83    /// Current live segments keyed by canonical table-relative path.
84    pub segments: HashMap<String, SegmentMeta>,
85
86    /// Optional pointer to the latest table coverage metadata.
87    pub table_coverage: Option<TableCoveragePointer>,
88}
89
90impl TableState {
91    /// Return live segments sorted deterministically by time.
92    ///
93    /// Ordering is by `ts_min`, then `ts_max`, and finally `path` as a
94    /// stable tie-breaker.
95    pub fn segments_sorted_by_time(&self) -> Vec<&SegmentMeta> {
96        let mut v: Vec<&SegmentMeta> = self.segments.values().collect();
97        v.sort_unstable_by(|a, b| cmp_segment_meta_by_time(a, b));
98        v
99    }
100}
101
102impl TransactionLogStore {
103    /// Rebuild the current TableState by replaying all commits up to CURRENT.
104    ///
105    /// v0.1 behavior:
106    /// - If CURRENT == 0 (no commits), this returns CommitError::CorruptState.
107    /// - The first commit must include at least one UpdateTableMeta action
108    ///   to bootstrap TableMeta; the last UpdateTableMeta wins.
109    pub async fn rebuild_table_state(&self) -> Result<TableState, CommitError> {
110        #[cfg(feature = "test-counters")]
111        REBUILD_TABLE_STATE_COUNT.with(|c| c.set(c.get() + 1));
112
113        let current_version = self.load_current_version().await?;
114
115        if current_version == 0 {
116            // v0.1: treat "no commits" as an uninitialized / corrupt table.
117            return CorruptStateSnafu {
118                msg: "Cannot rebuild TableState: CURRENT is 0 (no commits)".to_string(),
119            }
120            .fail();
121        }
122
123        let mut table_meta: Option<TableMeta> = None;
124        let mut segments: HashMap<String, SegmentMeta> = HashMap::new();
125
126        let mut table_coverage: Option<TableCoveragePointer> = None;
127
128        // Replay all commits from 1..=current_version in order
129        for v in 1..=current_version {
130            let commit = self.load_commit(v).await?;
131
132            // Defensive: file name version should match payload
133            if commit.version != v {
134                return CorruptStateSnafu {
135                    msg: format!(
136                        "Commit version mismatch: expected {v}, found {} in payload",
137                        commit.version
138                    ),
139                }
140                .fail();
141            }
142
143            for action in commit.actions {
144                match action {
145                    LogAction::AddSegment(meta) => {
146                        validate_persisted_segment_path(&meta.path)?;
147                        if segments.contains_key(&meta.path) {
148                            return CorruptStateSnafu {
149                                msg: format!("Duplicate live segment path: {}", meta.path),
150                            }
151                            .fail();
152                        }
153                        segments.insert(meta.path.clone(), meta);
154                    }
155                    LogAction::RemoveSegment { path } => {
156                        validate_persisted_segment_path(&path)?;
157                        segments.remove(&path);
158                    }
159                    LogAction::UpdateTableMeta(delta) => {
160                        if delta.format_version() != TABLE_FORMAT_VERSION {
161                            return CorruptStateSnafu {
162                                msg: format!(
163                                    "Unsupported table format version: expected {TABLE_FORMAT_VERSION}, found {}",
164                                    delta.format_version()
165                                ),
166                            }
167                            .fail();
168                        }
169                        // v0.1: full replacement of TableMeta
170                        table_meta = Some(delta);
171                    }
172                    LogAction::UpdateTableCoverage {
173                        bucket_spec,
174                        coverage_path,
175                    } => {
176                        table_coverage = Some(TableCoveragePointer {
177                            bucket_spec,
178                            coverage_path,
179                            version: v,
180                        })
181                    }
182                }
183            }
184        }
185
186        let table_meta = table_meta.context(CorruptStateSnafu {
187            msg: format!("No TableMeta found in commits up to version {current_version}",),
188        })?;
189
190        Ok(TableState {
191            version: current_version,
192            table_meta,
193            segments,
194            table_coverage,
195        })
196    }
197}
198
199#[cfg(test)]
200mod tests {
201    use super::*;
202    use crate::storage::layout;
203    use crate::storage::{StorageError, TableLocation};
204    use crate::transaction_log::{
205        FileFormat, LogAction, SegmentMeta, TableKind, TableMeta, TimeBucket, TimeIndexSpec,
206        TransactionLogStore,
207    };
208    use chrono::TimeZone;
209    use tempfile::TempDir;
210
211    type TestResult = Result<(), Box<dyn std::error::Error>>;
212
213    fn create_test_log_store() -> (TempDir, TransactionLogStore) {
214        let tmp = TempDir::new().expect("create temp dir");
215        let location = TableLocation::local(tmp.path());
216        let store = TransactionLogStore::new(location);
217        (tmp, store)
218    }
219
220    fn sample_table_meta() -> TableMeta {
221        TableMeta {
222            kind: TableKind::TimeSeries(TimeIndexSpec {
223                timestamp_column: "ts".to_string(),
224                entity_columns: vec!["symbol".to_string()],
225                bucket: TimeBucket::Minutes(1),
226                timezone: None,
227            }),
228            logical_schema: None,
229            created_at: chrono::Utc
230                .with_ymd_and_hms(2025, 1, 1, 0, 0, 0)
231                .single()
232                .expect("valid sample table metadata timestamp"),
233            format_version: TABLE_FORMAT_VERSION,
234            entity_identity: None,
235        }
236    }
237
238    fn sample_segment(id: &str) -> SegmentMeta {
239        SegmentMeta {
240            path: format!("data/{id}.parquet"),
241            format: FileFormat::Parquet,
242            ts_min: chrono::Utc
243                .with_ymd_and_hms(2025, 1, 1, 0, 0, 0)
244                .single()
245                .expect("valid sample segment ts_min"),
246            ts_max: chrono::Utc
247                .with_ymd_and_hms(2025, 1, 1, 1, 0, 0)
248                .single()
249                .expect("valid sample segment ts_max"),
250            row_count: 42,
251            file_size: None,
252            coverage_path: None,
253        }
254    }
255
256    fn segment_with_ts(id: &str, ts_min: i64, ts_max: i64) -> SegmentMeta {
257        SegmentMeta {
258            path: format!("data/{id}.parquet"),
259            format: FileFormat::Parquet,
260            ts_min: chrono::Utc.timestamp_opt(ts_min, 0).single().unwrap(),
261            ts_max: chrono::Utc.timestamp_opt(ts_max, 0).single().unwrap(),
262            row_count: 1,
263            file_size: None,
264            coverage_path: None,
265        }
266    }
267
268    #[test]
269    fn segments_sorted_by_time_orders_hashmap_deterministically() {
270        let mut segments = HashMap::new();
271        let seg_c = segment_with_ts("c", 10, 30);
272        let seg_a = segment_with_ts("a", 10, 20);
273        let seg_d = segment_with_ts("d", 5, 7);
274        let seg_b = segment_with_ts("b", 10, 20);
275
276        segments.insert(seg_c.path.clone(), seg_c);
277        segments.insert(seg_a.path.clone(), seg_a);
278        segments.insert(seg_d.path.clone(), seg_d);
279        segments.insert(seg_b.path.clone(), seg_b);
280
281        let state = TableState {
282            version: 3,
283            table_meta: sample_table_meta(),
284            segments,
285            table_coverage: None,
286        };
287
288        let ordered: Vec<(i64, i64, String)> = state
289            .segments_sorted_by_time()
290            .iter()
291            .map(|seg| {
292                (
293                    seg.ts_min.timestamp(),
294                    seg.ts_max.timestamp(),
295                    seg.path.clone(),
296                )
297            })
298            .collect();
299
300        let mut expected = ordered.clone();
301        expected.sort();
302        assert_eq!(ordered, expected);
303    }
304
305    #[tokio::test]
306    async fn rebuild_table_state_happy_path() -> TestResult {
307        let (_tmp, store) = create_test_log_store();
308        let meta = sample_table_meta();
309        let seg1 = sample_segment("seg1");
310        let seg2 = sample_segment("seg2");
311
312        let v1 = store
313            .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta.clone())])
314            .await?;
315        let v2 = store
316            .commit_with_expected_version(
317                v1,
318                vec![
319                    LogAction::AddSegment(seg1.clone()),
320                    LogAction::AddSegment(seg2.clone()),
321                ],
322            )
323            .await?;
324        let v3 = store
325            .commit_with_expected_version(
326                v2,
327                vec![LogAction::RemoveSegment {
328                    path: seg1.path.clone(),
329                }],
330            )
331            .await?;
332
333        let state = store.rebuild_table_state().await?;
334        assert_eq!(state.version, v3);
335        assert_eq!(state.table_meta, meta);
336        assert!(state.segments.contains_key(&seg2.path));
337        assert!(!state.segments.contains_key(&seg1.path));
338        Ok(())
339    }
340
341    #[tokio::test]
342    async fn rebuild_table_state_errors_when_current_zero() {
343        let (_tmp, store) = create_test_log_store();
344
345        let err = store
346            .rebuild_table_state()
347            .await
348            .expect_err("expected error");
349        assert!(matches!(err, CommitError::CorruptState { .. }));
350    }
351
352    #[tokio::test]
353    async fn rebuild_table_state_errors_when_no_table_meta() -> TestResult {
354        let (_tmp, store) = create_test_log_store();
355        let seg = sample_segment("seg");
356
357        store
358            .commit_with_expected_version(0, vec![LogAction::AddSegment(seg.clone())])
359            .await?;
360
361        let err = store
362            .rebuild_table_state()
363            .await
364            .expect_err("expected error");
365        assert!(matches!(err, CommitError::CorruptState { .. }));
366        Ok(())
367    }
368
369    #[tokio::test]
370    async fn rebuild_table_state_rejects_old_format_version() -> TestResult {
371        let (_tmp, store) = create_test_log_store();
372        let mut meta = sample_table_meta();
373        meta.format_version = TABLE_FORMAT_VERSION - 1;
374
375        store
376            .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
377            .await?;
378
379        let err = store
380            .rebuild_table_state()
381            .await
382            .expect_err("old format version should be rejected");
383        assert!(matches!(err, CommitError::CorruptState { .. }));
384        assert!(err.to_string().contains(&format!(
385            "expected {TABLE_FORMAT_VERSION}, found {}",
386            TABLE_FORMAT_VERSION - 1
387        )));
388        Ok(())
389    }
390
391    #[tokio::test]
392    async fn rebuild_table_state_rejects_noncanonical_segment_action_paths() -> TestResult {
393        for path in [
394            "",
395            "/data/seg.parquet",
396            "../data/seg.parquet",
397            "data/../seg.parquet",
398            r"data\seg.parquet",
399            "data//seg.parquet",
400            r"C:\data\seg.parquet",
401            "data/C:/seg.parquet",
402            "data/C:seg.parquet",
403        ] {
404            let mut segment = sample_segment("seg");
405            segment.path = path.to_owned();
406
407            for action in [
408                LogAction::AddSegment(segment.clone()),
409                LogAction::RemoveSegment {
410                    path: path.to_owned(),
411                },
412            ] {
413                let (_tmp, store) = create_test_log_store();
414                store
415                    .commit_with_expected_version(
416                        0,
417                        vec![LogAction::UpdateTableMeta(sample_table_meta()), action],
418                    )
419                    .await?;
420
421                let err = store
422                    .rebuild_table_state()
423                    .await
424                    .expect_err("noncanonical segment action path should be rejected");
425                assert!(matches!(err, CommitError::CorruptState { .. }));
426                assert!(err.to_string().contains("segment path"), "{err}");
427            }
428        }
429
430        Ok(())
431    }
432
433    #[tokio::test]
434    async fn rebuild_table_state_fails_on_corrupt_commit_payload() -> TestResult {
435        let (tmp, store) = create_test_log_store();
436        let meta = sample_table_meta();
437
438        store
439            .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
440            .await?;
441
442        let commit_path = tmp.path().join(layout::commit_rel_path(1));
443        tokio::fs::write(&commit_path, b"not-json").await?;
444
445        let err = store
446            .rebuild_table_state()
447            .await
448            .expect_err("expected error");
449        assert!(matches!(err, CommitError::CorruptState { .. }));
450        Ok(())
451    }
452
453    #[tokio::test]
454    async fn rebuild_table_state_fails_when_commit_missing() -> TestResult {
455        let (tmp, store) = create_test_log_store();
456        let meta = sample_table_meta();
457
458        store
459            .commit_with_expected_version(0, vec![LogAction::UpdateTableMeta(meta)])
460            .await?;
461
462        let commit_path = tmp.path().join(layout::commit_rel_path(1));
463        tokio::fs::remove_file(&commit_path).await?;
464
465        let err = store
466            .rebuild_table_state()
467            .await
468            .expect_err("expected error");
469        match err {
470            CommitError::Storage { source } => match source {
471                StorageError::NotFound { .. } => {}
472                other => panic!("unexpected storage error: {other:?}"),
473            },
474            other => panic!("expected storage error, got {other:?}"),
475        }
476        Ok(())
477    }
478}