1use 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")]
20pub fn rebuild_table_state_count() -> usize {
22 REBUILD_TABLE_STATE_COUNT.with(|c| c.get())
23}
24
25#[cfg(feature = "test-counters")]
26pub fn reset_rebuild_table_state_count() {
28 REBUILD_TABLE_STATE_COUNT.with(|c| c.set(0));
29}
30
31use crate::{
32 metadata::{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#[derive(Debug, Clone, PartialEq, Eq)]
62pub struct TableCoveragePointer {
63 pub bucket_spec: TimeBucket,
65 pub coverage_path: String,
67 pub version: u64,
69}
70
71#[derive(Debug, Clone, PartialEq, Eq)]
78pub struct TableState {
79 pub version: u64,
81 pub table_meta: TableMeta,
83 pub segments: HashMap<String, SegmentMeta>,
85
86 pub table_coverage: Option<TableCoveragePointer>,
88}
89
90impl TableState {
91 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 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 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 for v in 1..=current_version {
130 let commit = self.load_commit(v).await?;
131
132 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 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}