1use std::path::Path;
12use std::time::Instant;
13
14use snafu::prelude::*;
15use uuid::Uuid;
16
17use crate::{
18 coverage::serde::{coverage_to_bytes, entity_coverage_to_bytes},
19 coverage::{
20 EntityCoverage,
21 bucket::logical_bucket_range,
22 io::{CoverageError, write_coverage_sidecar_new_bytes},
23 layout::{
24 coverage_file_id_for_attempt, segment_coverage_id_v2, segment_coverage_key,
25 segment_entity_coverage_id_v1, table_coverage_id_v2, table_entity_coverage_id_v1,
26 table_snapshot_key,
27 },
28 },
29 formats::parquet::{
30 compute_segment_entity_coverage, coverage::compute_segment_coverage,
31 logical_schema_from_parquet, segment_meta::segment_meta_from_parquet,
32 },
33 metadata::{
34 schema_compat::{ensure_index_spec_matches_schema, ensure_schema_exact_match},
35 segments::SegmentEntityLayout,
36 },
37 storage,
38 transaction_log::{CommitError, LogAction, TableState, table_state::TableCoveragePointer},
39};
40
41use super::{
42 TimeSeriesTable,
43 append_report::{AppendReport, AppendReportBuilder},
44 error::{
45 CoverageBucketSnafu, CoverageOverlapSnafu, DuplicateSegmentPathSnafu,
46 EmptySegmentEntityCoverageSnafu, EntityCoverageOverlapSnafu,
47 EntityWithoutIndexCoverageSnafu, ExistingSegmentMissingCoverageSnafu,
48 MissingCanonicalSchemaSnafu, SegmentCoverageSnafu, SegmentMetaSnafu,
49 SegmentSchemaCompatibilitySnafu, StorageSnafu, TableError,
50 },
51};
52
53fn classify_entity_layout(
54 segment_path: &str,
55 coverage: &EntityCoverage,
56) -> Result<SegmentEntityLayout, TableError> {
57 let first_identity = coverage
58 .iter()
59 .next()
60 .map(|(identity, _)| identity)
61 .context(EmptySegmentEntityCoverageSnafu {
62 segment_path: segment_path.to_string(),
63 })?;
64
65 if let Some((identity, _)) = coverage.iter().find(|(_, coverage)| coverage.is_empty()) {
66 return EntityWithoutIndexCoverageSnafu {
67 segment_path: segment_path.to_string(),
68 identity: identity.clone(),
69 }
70 .fail();
71 }
72
73 Ok(if coverage.identity_count() == 1 {
74 SegmentEntityLayout::Single(first_identity.clone())
75 } else {
76 SegmentEntityLayout::Mixed
77 })
78}
79
80fn ensure_existing_segments_have_coverage(state: &TableState) -> Result<(), TableError> {
81 for seg in state.segments.values() {
82 if seg.coverage_path.is_none() {
83 return ExistingSegmentMissingCoverageSnafu {
84 path: seg.path.clone(),
85 }
86 .fail();
87 }
88 }
89
90 Ok(())
91}
92
93impl TimeSeriesTable {
94 async fn rollback_created_sidecars(
95 &self,
96 created_sidecars: &[String],
97 source: TableError,
98 ) -> TableError {
99 let mut cleanup_errors = Vec::new();
100 for path in created_sidecars.iter().rev() {
101 if let Err(error) =
102 storage::remove_file(self.location().as_ref(), Path::new(path)).await
103 {
104 cleanup_errors.push(format!("{path}: {error}"));
105 }
106 }
107
108 if cleanup_errors.is_empty() {
109 source
110 } else {
111 TableError::AppendRollback {
112 source: Box::new(source),
113 cleanup_errors,
114 }
115 }
116 }
117
118 async fn normalize_new_segment_path(&self, relative_path: &str) -> Result<String, TableError> {
119 let supplied_path = Path::new(relative_path);
120 let (normalized, native_path) =
121 storage::normalize_relative_storage_path(supplied_path).context(StorageSnafu)?;
122
123 if self.state.segments.contains_key(&normalized) {
124 return DuplicateSegmentPathSnafu { path: normalized }.fail();
125 }
126
127 self.location()
128 .validate_segment_file(supplied_path, &native_path)
129 .await
130 .context(StorageSnafu)?;
131
132 Ok(normalized)
133 }
134
135 async fn append_parquet_path_file(
136 &mut self,
137 parquet_path: &Path,
138 mut report: Option<&mut AppendReportBuilder>,
139 ) -> Result<(u64, String), TableError> {
140 let prepared = self
141 .location()
142 .prepare_parquet_under_root(parquet_path)
143 .await
144 .context(StorageSnafu)?;
145 let prepared_path = prepared.relative_path.to_string_lossy().into_owned();
146
147 let append_result = async {
148 let relative_path = self.normalize_new_segment_path(&prepared_path).await?;
149 if let Some(r) = report.as_mut() {
150 r.set_context("relative_path", &relative_path);
151 }
152 let version = self
153 .append_parquet_segment_file(&relative_path, report)
154 .await?;
155 Ok((version, relative_path))
156 }
157 .await;
158
159 match append_result {
160 Ok(result) => Ok(result),
161 Err(
162 error @ TableError::TransactionLog {
163 source: CommitError::AmbiguousOutcome { .. },
164 },
165 ) => Err(error),
166 Err(source) if prepared.created => {
167 match storage::remove_file(self.location().as_ref(), &prepared.relative_path).await
168 {
169 Ok(()) => Err(source),
170 Err(cleanup_error) => Err(TableError::ExternalParquetRollback {
171 path: prepared.relative_path.display().to_string(),
172 source: Box::new(source),
173 cleanup_error,
174 }),
175 }
176 }
177 Err(source) => Err(source),
178 }
179 }
180
181 async fn append_parquet_segment_file(
182 &mut self,
183 relative_path: &str,
184 mut report: Option<&mut AppendReportBuilder>,
185 ) -> Result<u64, TableError> {
186 let rel_path = Path::new(relative_path);
187 let expected_version = self.state.version;
188 if let Some(r) = report.as_mut() {
189 r.set_context("index_column", self.index.column.as_str());
190 }
191
192 ensure_existing_segments_have_coverage(&self.state)?;
194
195 let step_start = Instant::now();
197 let (mut segment_meta, meta_report) =
198 segment_meta_from_parquet(self.location(), rel_path, &self.index)
199 .await
200 .context(SegmentMetaSnafu)?;
201 if let Some(r) = report.as_mut() {
202 if let Some(file_size) = segment_meta.file_size {
203 r.set_context("file_size_bytes", file_size.to_string());
204 }
205 let fields = vec![
206 ("row_groups".to_string(), meta_report.row_groups.to_string()),
207 ("row_count".to_string(), meta_report.row_count.to_string()),
208 ("used_stats".to_string(), meta_report.used_stats.to_string()),
209 (
210 "scanned_rows".to_string(),
211 meta_report.scanned_rows.to_string(),
212 ),
213 ];
214 r.push_step("segment_meta", step_start.elapsed(), fields);
215 }
216
217 let step_start = Instant::now();
218 let segment_schema = logical_schema_from_parquet(self.location(), rel_path)
219 .await
220 .context(SegmentMetaSnafu)?;
221 ensure_index_spec_matches_schema(&segment_schema, &self.index).context(
222 SegmentSchemaCompatibilitySnafu {
223 path: relative_path.to_string(),
224 },
225 )?;
226 if let Some(r) = report.as_mut() {
227 r.push_step("logical_schema", step_start.elapsed(), Vec::new());
228 }
229
230 let maybe_table_schema = self.state.table_meta.logical_schema.as_ref();
239
240 let maybe_updated_meta = match maybe_table_schema {
241 None if expected_version == 1 => {
242 let mut updated_meta = self.state.table_meta.clone();
243 updated_meta.logical_schema = Some(segment_schema.clone());
244 Some(updated_meta)
245 }
246 None => {
247 return MissingCanonicalSchemaSnafu {
248 version: expected_version,
249 }
250 .fail();
251 }
252 Some(table_schema) => {
253 ensure_schema_exact_match(table_schema, &segment_schema, &self.index).context(
254 SegmentSchemaCompatibilitySnafu {
255 path: relative_path.to_string(),
256 },
257 )?;
258 None
259 }
260 };
261
262 let has_entity_columns = !self.index.entity_columns.is_empty();
263
264 let (seg_cov_bytes, new_snap_cov_bytes, entity_layout) = if has_entity_columns {
266 let step_start = Instant::now();
267 let table_cov = self.load_table_entity_snapshot_coverage_readonly().await?;
268 if let Some(r) = report.as_mut() {
269 r.push_step("load_table_snapshot", step_start.elapsed(), Vec::new());
270 }
271
272 let step_start = Instant::now();
273 let segment_cov =
274 compute_segment_entity_coverage(self.location(), rel_path, &self.index)
275 .await
276 .context(SegmentCoverageSnafu)?;
277 let entity_layout = classify_entity_layout(relative_path, &segment_cov)?;
278 if let Some(r) = report.as_mut() {
279 r.push_step("segment_coverage", step_start.elapsed(), Vec::new());
280 }
281
282 let step_start = Instant::now();
283 if let Some((identity, bucket)) = segment_cov.overlap_example(&table_cov) {
284 let example_bucket_range =
285 logical_bucket_range(&self.index.kind, bucket).context(CoverageBucketSnafu)?;
286 return EntityCoverageOverlapSnafu {
287 segment_path: relative_path.to_string(),
288 overlap_count: segment_cov.intersection_cardinality(&table_cov),
289 example_identity: identity.clone(),
290 example_bucket: bucket,
291 example_bucket_range,
292 }
293 .fail();
294 }
295 if let Some(r) = report.as_mut() {
296 r.push_step("overlap_check", step_start.elapsed(), Vec::new());
297 }
298
299 let seg_bytes = entity_coverage_to_bytes(&segment_cov).map_err(|source| {
300 TableError::CoverageSidecar {
301 source: CoverageError::EntitySerde { source },
302 }
303 })?;
304 let snapshot_bytes =
305 entity_coverage_to_bytes(&table_cov.union(&segment_cov)).map_err(|source| {
306 TableError::CoverageSidecar {
307 source: CoverageError::EntitySerde { source },
308 }
309 })?;
310 (seg_bytes, snapshot_bytes, entity_layout)
311 } else {
312 let step_start = Instant::now();
313 let table_cov = self.load_table_snapshot_coverage_readonly().await?;
314 if let Some(r) = report.as_mut() {
315 r.push_step("load_table_snapshot", step_start.elapsed(), Vec::new());
316 }
317
318 let step_start = Instant::now();
319 let segment_cov = compute_segment_coverage(self.location(), rel_path, &self.index)
320 .await
321 .context(SegmentCoverageSnafu)?;
322 if let Some(r) = report.as_mut() {
323 r.push_step("segment_coverage", step_start.elapsed(), Vec::new());
324 }
325
326 let step_start = Instant::now();
327 let overlap = segment_cov.intersect(&table_cov);
328 let overlap_count = overlap.cardinality();
329 if let Some(example_bucket) = overlap.present().iter().next() {
330 let example_bucket_range = logical_bucket_range(&self.index.kind, example_bucket)
331 .context(CoverageBucketSnafu)?;
332 return CoverageOverlapSnafu {
333 segment_path: relative_path.to_string(),
334 overlap_count,
335 example_bucket: Some(example_bucket),
336 example_bucket_range,
337 }
338 .fail();
339 }
340 if let Some(r) = report.as_mut() {
341 r.push_step("overlap_check", step_start.elapsed(), Vec::new());
342 }
343
344 let seg_bytes =
345 coverage_to_bytes(&segment_cov).map_err(|source| TableError::CoverageSidecar {
346 source: CoverageError::Serde { source },
347 })?;
348 let snapshot_bytes =
349 coverage_to_bytes(&table_cov.union(&segment_cov)).map_err(|source| {
350 TableError::CoverageSidecar {
351 source: CoverageError::Serde { source },
352 }
353 })?;
354 (
355 seg_bytes,
356 snapshot_bytes,
357 SegmentEntityLayout::NotApplicable,
358 )
359 };
360
361 let attempt_id = Uuid::new_v4();
363 let segment_content_id = if has_entity_columns {
364 segment_entity_coverage_id_v1(&self.index, &seg_cov_bytes)
365 } else {
366 segment_coverage_id_v2(&self.index, &seg_cov_bytes)
367 };
368 let segment_file_id = coverage_file_id_for_attempt(&segment_content_id, &attempt_id);
369 let seg_cov_path = segment_coverage_key(&segment_file_id).map_err(|source| {
370 TableError::CoverageSidecar {
371 source: CoverageError::Layout { source },
372 }
373 })?;
374
375 let new_version_guess = expected_version + 1;
376 let snapshot_content_id = if has_entity_columns {
377 table_entity_coverage_id_v1(&self.index, &new_snap_cov_bytes)
378 } else {
379 table_coverage_id_v2(&self.index, &new_snap_cov_bytes)
380 };
381 let snapshot_file_id = coverage_file_id_for_attempt(&snapshot_content_id, &attempt_id);
382 let snapshot_path =
383 table_snapshot_key(new_version_guess, &snapshot_file_id).map_err(|source| {
384 TableError::CoverageSidecar {
385 source: CoverageError::Layout { source },
386 }
387 })?;
388
389 let step_start = Instant::now();
390 let mut created_sidecars = Vec::new();
391 write_coverage_sidecar_new_bytes(self.location(), Path::new(&seg_cov_path), &seg_cov_bytes)
392 .await
393 .map_err(|source| TableError::CoverageSidecar { source })?;
394 created_sidecars.push(seg_cov_path.clone());
395 if let Some(r) = report.as_mut() {
396 r.push_step("write_segment_sidecar", step_start.elapsed(), Vec::new());
397 }
398
399 let step_start = Instant::now();
400 if let Err(source) = write_coverage_sidecar_new_bytes(
401 self.location(),
402 Path::new(&snapshot_path),
403 &new_snap_cov_bytes,
404 )
405 .await
406 {
407 let error = TableError::CoverageSidecar { source };
408 return Err(self
409 .rollback_created_sidecars(&created_sidecars, error)
410 .await);
411 }
412 created_sidecars.push(snapshot_path.clone());
413 if let Some(r) = report.as_mut() {
414 r.push_step("write_snapshot_sidecar", step_start.elapsed(), Vec::new());
415 }
416
417 segment_meta.coverage_path = Some(seg_cov_path);
419 segment_meta.entity_layout = entity_layout;
420
421 let mut actions = Vec::new();
422 if let Some(updated_meta) = maybe_updated_meta.clone() {
423 actions.push(LogAction::UpdateTableMeta(updated_meta));
424 }
425
426 actions.push(LogAction::AddSegment(segment_meta.clone()));
427 actions.push(LogAction::UpdateTableCoverage {
428 index_kind: self.index.kind.clone(),
429 coverage_path: snapshot_path.clone(),
430 });
431
432 let step_start = Instant::now();
433 let new_version = match self
434 .log
435 .commit_with_expected_version(expected_version, actions)
436 .await
437 {
438 Ok(version) => version,
439 Err(source @ crate::transaction_log::CommitError::AmbiguousOutcome { .. }) => {
440 return Err(TableError::TransactionLog { source });
441 }
442 Err(source) => {
443 let error = TableError::TransactionLog { source };
444 return Err(self
445 .rollback_created_sidecars(&created_sidecars, error)
446 .await);
447 }
448 };
449 if let Some(r) = report.as_mut() {
450 r.push_step("commit_log", step_start.elapsed(), Vec::new());
451 }
452
453 assert_eq!(
459 new_version, new_version_guess,
460 "transaction log returned unexpected version: expected {}, got {}",
461 new_version_guess, new_version
462 );
463
464 let step_start = Instant::now();
466 self.state.version = new_version;
467
468 if let Some(updated_meta) = maybe_updated_meta {
469 self.state.table_meta = updated_meta
470 }
471
472 self.state
473 .segments
474 .insert(segment_meta.path.clone(), segment_meta);
475
476 self.state.table_coverage = Some(TableCoveragePointer {
478 index_kind: self.index.kind.clone(),
479 coverage_path: snapshot_path,
480 version: new_version,
481 });
482 if let Some(r) = report.as_mut() {
483 r.push_step("state_update", step_start.elapsed(), Vec::new());
484 }
485
486 Ok(new_version)
487 }
488
489 pub async fn append_parquet_segment(&mut self, relative_path: &str) -> Result<u64, TableError> {
493 let relative_path = self.normalize_new_segment_path(relative_path).await?;
494 self.append_parquet_segment_file(&relative_path, None).await
495 }
496
497 pub async fn append_parquet_from_path(
506 &mut self,
507 parquet_path: &Path,
508 ) -> Result<(u64, String), TableError> {
509 self.append_parquet_path_file(parquet_path, None).await
510 }
511
512 pub async fn append_parquet_from_path_with_report(
515 &mut self,
516 parquet_path: &Path,
517 ) -> Result<(u64, String, AppendReport), TableError> {
518 let mut report = AppendReportBuilder::new();
519 let (version, relative_path) = self
520 .append_parquet_path_file(parquet_path, Some(&mut report))
521 .await?;
522 Ok((version, relative_path, report.finish()))
523 }
524
525 pub async fn append_parquet_segment_with_report(
527 &mut self,
528 relative_path: &str,
529 ) -> Result<(u64, AppendReport), TableError> {
530 let relative_path = self.normalize_new_segment_path(relative_path).await?;
531 let mut report = AppendReportBuilder::new();
532 report.set_context("relative_path", &relative_path);
533
534 let version = self
535 .append_parquet_segment_file(&relative_path, Some(&mut report))
536 .await?;
537
538 Ok((version, report.finish()))
539 }
540}
541
542#[cfg(test)]
543mod tests {
544 use super::super::test_util::*;
545 use super::*;
546 use crate::coverage::io::{read_coverage_sidecar, read_entity_coverage_sidecar};
547 use crate::coverage::serde::entity_coverage_from_bytes;
548 use crate::coverage::{EntityCoverage, EntityIdentity, EntityValue};
549 use crate::metadata::logical_schema::{
550 LogicalDataType, LogicalField, LogicalSchema, LogicalTimestampUnit,
551 };
552 use crate::metadata::segments::{ParquetIndexColumnError, SegmentEntityLayout};
553 use crate::metadata::table_metadata::{IndexValue, TABLE_FORMAT_VERSION};
554 use crate::storage::layout;
555 use crate::storage::{StorageError, StorageLocation, TableLocation};
556 use crate::transaction_log::segments::{SegmentError, SegmentMetaError};
557 use crate::transaction_log::{
558 CommitError, IndexKind, IndexSpec, TableKind, TableMeta, TimeBucket,
559 };
560 use arrow::{
561 array::{
562 ArrayRef, BooleanArray, Float64Array, Int64Array, StringArray,
563 TimestampMillisecondArray, UInt64Array,
564 },
565 datatypes::{DataType, Field, Schema, TimeUnit as ArrowTimeUnit},
566 record_batch::RecordBatch,
567 };
568 use parquet::arrow::ArrowWriter;
569 use parquet::file::reader::{FileReader, SerializedFileReader};
570 use std::collections::BTreeMap;
571 use std::fs::{File, OpenOptions};
572 use std::io::{Seek, SeekFrom, Write};
573 use std::num::NonZeroU64;
574 use std::path::PathBuf;
575 use std::sync::Arc;
576 use tempfile::TempDir;
577
578 fn registered_index(kind: IndexKind) -> IndexSpec {
579 IndexSpec {
580 column: "ts".to_string(),
581 entity_columns: Vec::new(),
582 kind,
583 }
584 }
585
586 fn write_single_index_parquet(
587 path: &Path,
588 data_type: DataType,
589 values: ArrayRef,
590 ) -> TestResult {
591 if let Some(parent) = path.parent() {
592 std::fs::create_dir_all(parent)?;
593 }
594 let schema = Arc::new(Schema::new(vec![Field::new("ts", data_type, false)]));
595 let batch = RecordBatch::try_new(Arc::clone(&schema), vec![values])?;
596 let mut writer = ArrowWriter::try_new(File::create(path)?, schema, None)?;
597 writer.write(&batch)?;
598 writer.close()?;
599 Ok(())
600 }
601
602 fn write_composite_entity_parquet(path: &Path, rows: &[(i64, &str, &str, f64)]) -> TestResult {
603 if let Some(parent) = path.parent() {
604 std::fs::create_dir_all(parent)?;
605 }
606 let schema = Arc::new(Schema::new(vec![
607 Field::new(
608 "ts",
609 DataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
610 false,
611 ),
612 Field::new("symbol", DataType::Utf8, false),
613 Field::new("venue", DataType::Utf8, false),
614 Field::new("price", DataType::Float64, false),
615 ]));
616 let batch = RecordBatch::try_new(
617 Arc::clone(&schema),
618 vec![
619 Arc::new(TimestampMillisecondArray::from(
620 rows.iter().map(|row| row.0).collect::<Vec<_>>(),
621 )),
622 Arc::new(StringArray::from(
623 rows.iter().map(|row| row.1).collect::<Vec<_>>(),
624 )),
625 Arc::new(StringArray::from(
626 rows.iter().map(|row| row.2).collect::<Vec<_>>(),
627 )),
628 Arc::new(Float64Array::from(
629 rows.iter().map(|row| row.3).collect::<Vec<_>>(),
630 )),
631 ],
632 )?;
633 let mut writer = ArrowWriter::try_new(File::create(path)?, schema, None)?;
634 writer.write(&batch)?;
635 writer.close()?;
636 Ok(())
637 }
638
639 fn coverage_files(root: &Path) -> std::io::Result<BTreeMap<PathBuf, Vec<u8>>> {
640 let mut files = BTreeMap::new();
641 for rel_dir in [layout::SEGMENT_COVERAGE_DIR, layout::TABLE_SNAPSHOT_DIR] {
642 let dir = root.join(rel_dir);
643 if !dir.exists() {
644 continue;
645 }
646 for entry in std::fs::read_dir(dir)? {
647 let path = entry?.path();
648 if path.is_file() {
649 files.insert(
650 path.strip_prefix(root)
651 .expect("coverage path under root")
652 .to_owned(),
653 std::fs::read(path)?,
654 );
655 }
656 }
657 }
658 Ok(files)
659 }
660
661 #[test]
662 fn entity_layout_classification_rejects_empty_coverage() {
663 assert!(matches!(
664 classify_entity_layout("data/empty.parquet", &EntityCoverage::empty()),
665 Err(TableError::EmptySegmentEntityCoverage { segment_path })
666 if segment_path == "data/empty.parquet"
667 ));
668 }
669
670 #[tokio::test]
671 async fn append_parquet_segment_missing_time_column_errors() -> TestResult {
672 let tmp = TempDir::new()?;
673 let location = TableLocation::local(tmp.path());
674 let meta = make_basic_table_meta();
675 let mut table = TimeSeriesTable::create(location.clone(), meta).await?;
676
677 let rel = "data/seg-no-ts.parquet";
678 let path = tmp.path().join(rel);
679 write_parquet_without_time_column(&path, &["A"], &[1.0])?;
680
681 let err = table
682 .append_parquet_segment(rel)
683 .await
684 .expect_err("expected missing time column");
685
686 match err {
687 TableError::SegmentMeta { source } => {
688 assert!(matches!(
689 source,
690 SegmentError::Meta {
691 source: SegmentMetaError::OrderedIndexColumn {
692 source: ParquetIndexColumnError {
693 expected_domain: "timestamp",
694 observed_type,
695 ..
696 }
697 }
698 } if observed_type == "missing",
699 ));
700 }
701 other => panic!("unexpected error: {other:?}"),
702 }
703
704 Ok(())
705 }
706
707 #[tokio::test]
708 async fn entity_aware_validation_failure_leaves_no_state_or_sidecars() -> TestResult {
709 let tmp = TempDir::new()?;
710 let table_root = tmp.path().join("table");
711 let location = TableLocation::local(&table_root);
712 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
713 let state_before = table.state.clone();
714 let coverage_before = coverage_files(&table_root)?;
715 let source = tmp.path().join("wrong-time-column.parquet");
716 write_parquet_without_time_column(&source, &["A"], &[1.0])?;
717 let source_bytes = std::fs::read(&source)?;
718
719 let err = table
720 .append_parquet_from_path(&source)
721 .await
722 .expect_err("missing time column should fail");
723
724 assert!(matches!(err, TableError::SegmentMeta { .. }));
725 assert!(!table_root.join("data/wrong-time-column.parquet").exists());
726 assert_eq!(std::fs::read(source)?, source_bytes);
727 assert_eq!(table.state, state_before);
728 assert_eq!(table.log.load_current_version().await?, 1);
729 assert_eq!(coverage_files(&table_root)?, coverage_before);
730 assert!(!table_root.join(layout::commit_rel_path(2)).exists());
731 Ok(())
732 }
733
734 #[tokio::test]
735 async fn append_parquet_from_path_preserves_failed_in_root_file() -> TestResult {
736 let tmp = TempDir::new()?;
737 let location = TableLocation::local(tmp.path());
738 let mut table = TimeSeriesTable::create(location, make_basic_table_meta()).await?;
739 let source = tmp.path().join("data/in-root-invalid.parquet");
740 write_parquet_without_time_column(&source, &["A"], &[1.0])?;
741 let source_bytes = std::fs::read(&source)?;
742
743 let err = table
744 .append_parquet_from_path(&source)
745 .await
746 .expect_err("missing time column should fail");
747
748 assert!(matches!(err, TableError::SegmentMeta { .. }));
749 assert_eq!(std::fs::read(source)?, source_bytes);
750 Ok(())
751 }
752
753 #[tokio::test]
754 async fn append_parquet_from_path_retains_successful_external_copy() -> TestResult {
755 let tmp = TempDir::new()?;
756 let table_root = tmp.path().join("table");
757 let location = TableLocation::local(&table_root);
758 let mut table = TimeSeriesTable::create(location, make_basic_table_meta()).await?;
759 let source = tmp.path().join("external-success.parquet");
760 write_test_parquet(
761 &source,
762 true,
763 false,
764 &[TestRow {
765 ts_millis: 10_000,
766 symbol: "X",
767 price: 100.0,
768 }],
769 )?;
770 let source_bytes = std::fs::read(&source)?;
771
772 let (version, relative_path) = table.append_parquet_from_path(&source).await?;
773
774 assert_eq!(version, 2);
775 assert_eq!(relative_path, "data/external-success.parquet");
776 assert_eq!(
777 std::fs::read(table_root.join(&relative_path))?,
778 source_bytes
779 );
780 assert_eq!(std::fs::read(source)?, source_bytes);
781 assert!(table.state.segments.contains_key(&relative_path));
782 Ok(())
783 }
784
785 #[tokio::test]
786 async fn entity_aware_ambiguous_commit_retains_copy_and_sidecars() -> TestResult {
787 let tmp = TempDir::new()?;
788 let table_root = tmp.path().join("table");
789 let location = TableLocation::local(&table_root);
790 let mut table = TimeSeriesTable::create(location, make_basic_table_meta()).await?;
791 let state_before = table.state.clone();
792 let coverage_before = coverage_files(&table_root)?;
793 let source = tmp.path().join("ambiguous-external.parquet");
794 write_test_parquet(
795 &source,
796 true,
797 false,
798 &[TestRow {
799 ts_millis: 10_000,
800 symbol: "X",
801 price: 100.0,
802 }],
803 )?;
804 let source_bytes = std::fs::read(&source)?;
805 let commit_path = table_root.join(layout::commit_rel_path(2));
806 crate::storage::inject_write_new_failure(commit_path.clone(), true);
807
808 let err = table
809 .append_parquet_from_path(&source)
810 .await
811 .expect_err("commit outcome should be ambiguous");
812
813 assert!(matches!(
814 err,
815 TableError::TransactionLog {
816 source: CommitError::AmbiguousOutcome { .. }
817 }
818 ));
819 assert_eq!(
820 std::fs::read(table_root.join("data/ambiguous-external.parquet"))?,
821 source_bytes
822 );
823 assert_eq!(std::fs::read(source)?, source_bytes);
824 assert_eq!(table.state, state_before);
825 assert_eq!(table.log.load_current_version().await?, 1);
826 assert!(commit_path.exists());
827 let coverage_after = coverage_files(&table_root)?;
828 assert_eq!(coverage_after.len(), coverage_before.len() + 2);
829 for bytes in coverage_after.values() {
830 entity_coverage_from_bytes(bytes)?;
831 }
832 Ok(())
833 }
834
835 #[tokio::test]
836 async fn append_parquet_from_path_reports_copy_rollback_failure() -> TestResult {
837 let tmp = TempDir::new()?;
838 let table_root = tmp.path().join("table");
839 let location = TableLocation::local(&table_root);
840 let mut table = TimeSeriesTable::create(location, make_basic_table_meta()).await?;
841 let source = tmp.path().join("rollback-cleanup.parquet");
842 write_parquet_without_time_column(&source, &["A"], &[1.0])?;
843 let destination = table_root.join("data/rollback-cleanup.parquet");
844 crate::storage::inject_cleanup_failure(destination.clone());
845
846 let err = table
847 .append_parquet_from_path(&source)
848 .await
849 .expect_err("copy rollback should fail");
850 let message = err.to_string();
851
852 assert!(matches!(
853 err,
854 TableError::ExternalParquetRollback {
855 path,
856 source,
857 cleanup_error: StorageError::OtherIo { .. },
858 } if path.contains("rollback-cleanup.parquet")
859 && matches!(*source, TableError::SegmentMeta { .. })
860 ));
861 assert!(message.contains("rollback-cleanup.parquet"));
862 assert!(message.contains("injected cleanup failure"));
863 assert!(destination.exists());
864 tokio::fs::remove_file(destination).await?;
865 Ok(())
866 }
867
868 #[tokio::test]
869 async fn append_parquet_segment_updates_state_and_log() -> TestResult {
870 let tmp = TempDir::new()?;
871 let location = TableLocation::local(tmp.path());
872 let meta = make_basic_table_meta();
873
874 let mut table = TimeSeriesTable::create(location.clone(), meta).await?;
875
876 let rel_path = "data/seg1.parquet";
877 let abs_path = tmp.path().join(rel_path);
878 write_test_parquet(
879 &abs_path,
880 true,
881 false,
882 &[TestRow {
883 ts_millis: 1_000,
884 symbol: "A",
885 price: 10.0,
886 }],
887 )?;
888
889 let new_version = table.append_parquet_segment(rel_path).await?;
890
891 assert_eq!(new_version, 2);
892 assert_eq!(table.state.version, 2);
893 let seg = table.state.segments.get(rel_path).expect("segment present");
894 assert_eq!(seg.path, rel_path);
895 assert_eq!(seg.row_count, 1);
896 assert_eq!(
897 seg.entity_layout,
898 SegmentEntityLayout::Single(EntityIdentity::try_new(vec!["A".into()])?)
899 );
900 assert!(matches!(
901 &seg.index_min,
902 IndexValue::Timestamp(value) if value.timestamp_millis() == 1_000
903 ));
904 assert!(matches!(
905 &seg.index_max,
906 IndexValue::Timestamp(value) if value.timestamp_millis() == 1_000
907 ));
908
909 let commit_path = tmp.path().join(layout::commit_rel_path(2));
910 assert!(commit_path.is_file());
911 let current =
912 tokio::fs::read_to_string(tmp.path().join(layout::current_rel_path())).await?;
913 assert_eq!(current.trim(), "2");
914
915 let reopened = TimeSeriesTable::open(location).await?;
916 assert_eq!(reopened.state.segments.get(rel_path), Some(seg));
917 Ok(())
918 }
919
920 #[tokio::test]
921 async fn version_six_no_entity_int64_append_uses_global_coverage() -> TestResult {
922 let tmp = TempDir::new()?;
923 let location = TableLocation::local(tmp.path());
924 let index = registered_index(IndexKind::Int64 {
925 bucket_width: NonZeroU64::new(10).unwrap(),
926 });
927 let mut table =
928 TimeSeriesTable::create(location.clone(), TableMeta::new_time_series(index.clone()))
929 .await?;
930 let rel_path = "data/int64.parquet";
931 write_arrow_parquet_int_time(
932 &tmp.path().join(rel_path),
933 &[i64::MIN, -1, 0, i64::MAX],
934 &["A", "A", "A", "A"],
935 &[1.0, 2.0, 3.0, 4.0],
936 )?;
937
938 assert_eq!(table.append_parquet_segment(rel_path).await?, 2);
939
940 let segment = table.state.segments.get(rel_path).expect("segment present");
941 assert_eq!(segment.entity_layout, SegmentEntityLayout::NotApplicable);
942 assert_eq!(segment.index_min, IndexValue::Int64(i64::MIN));
943 assert_eq!(segment.index_max, IndexValue::Int64(i64::MAX));
944 let pointer = table.state.table_coverage.as_ref().expect("table coverage");
945 assert_eq!(pointer.index_kind, index.kind);
946 let persisted = read_coverage_sidecar(&location, Path::new(&pointer.coverage_path)).await?;
947 let expected =
948 compute_segment_coverage(&location, Path::new(rel_path), table.index_spec()).await?;
949 assert_eq!(persisted, expected);
950 let reopened = TimeSeriesTable::open(location).await?;
951 assert_eq!(reopened.state, table.state);
952 Ok(())
953 }
954
955 #[tokio::test]
956 async fn int64_appends_enforce_coverage_and_exact_later_schema() -> TestResult {
957 let tmp = TempDir::new()?;
958 let location = TableLocation::local(tmp.path());
959 let index = registered_index(IndexKind::Int64 {
960 bucket_width: NonZeroU64::new(10).unwrap(),
961 });
962 let mut table =
963 TimeSeriesTable::create(location, TableMeta::new_time_series(index)).await?;
964
965 for (path, values) in [
966 ("data/negative.parquet", &[-25, -15][..]),
967 ("data/positive.parquet", &[5, 15][..]),
968 ] {
969 write_arrow_parquet_int_time(&tmp.path().join(path), values, &["A", "A"], &[1.0, 2.0])?;
970 table.append_parquet_segment(path).await?;
971 }
972 assert_eq!(table.state.version, 3);
973
974 let state_before = table.state.clone();
975 let coverage_before = coverage_files(tmp.path())?;
976 let overlap_path = "data/negative-overlap.parquet";
977 write_arrow_parquet_int_time(&tmp.path().join(overlap_path), &[-19], &["A"], &[3.0])?;
978 let overlap_error = table
979 .append_parquet_segment(overlap_path)
980 .await
981 .expect_err("negative bucket overlap must fail");
982 assert!(matches!(
983 &overlap_error,
984 TableError::CoverageOverlap {
985 example_bucket_range,
986 ..
987 } if example_bucket_range.to_string() == "[-20, -10)"
988 ));
989 assert!(
990 overlap_error
991 .to_string()
992 .contains("example_bucket_range=[-20, -10)")
993 );
994
995 let mismatch_path = "data/schema-mismatch.parquet";
996 write_single_index_parquet(
997 &tmp.path().join(mismatch_path),
998 DataType::Int64,
999 Arc::new(Int64Array::from(vec![100])),
1000 )?;
1001 assert!(matches!(
1002 table
1003 .append_parquet_segment(mismatch_path)
1004 .await
1005 .expect_err("later schema mismatch must fail"),
1006 TableError::SegmentSchemaCompatibility { .. }
1007 ));
1008 assert_eq!(table.state, state_before);
1009 assert_eq!(coverage_files(tmp.path())?, coverage_before);
1010 Ok(())
1011 }
1012
1013 #[tokio::test]
1014 async fn append_parquet_segment_supports_registered_uint64_index() -> TestResult {
1015 let tmp = TempDir::new()?;
1016 let location = TableLocation::local(tmp.path());
1017 let index = registered_index(IndexKind::UInt64 {
1018 bucket_width: NonZeroU64::new(10).unwrap(),
1019 });
1020 let mut table =
1021 TimeSeriesTable::create(location.clone(), TableMeta::new_time_series(index.clone()))
1022 .await?;
1023 let rel_path = "data/uint64.parquet";
1024 write_single_index_parquet(
1025 &tmp.path().join(rel_path),
1026 DataType::UInt64,
1027 Arc::new(UInt64Array::from(vec![0, i64::MAX as u64 + 1, u64::MAX])),
1028 )?;
1029
1030 assert_eq!(table.append_parquet_segment(rel_path).await?, 2);
1031
1032 let segment = table.state.segments.get(rel_path).expect("segment present");
1033 assert_eq!(segment.index_min, IndexValue::UInt64(0));
1034 assert_eq!(segment.index_max, IndexValue::UInt64(u64::MAX));
1035 assert_eq!(
1036 table
1037 .state
1038 .table_meta
1039 .logical_schema
1040 .as_ref()
1041 .expect("schema adopted")
1042 .columns()[0]
1043 .data_type,
1044 LogicalDataType::UInt64
1045 );
1046 assert_eq!(
1047 table
1048 .state
1049 .table_coverage
1050 .as_ref()
1051 .expect("table coverage")
1052 .index_kind,
1053 index.kind
1054 );
1055
1056 let non_overlap_path = "data/uint64-non-overlap.parquet";
1057 write_single_index_parquet(
1058 &tmp.path().join(non_overlap_path),
1059 DataType::UInt64,
1060 Arc::new(UInt64Array::from(vec![u64::MAX - 20])),
1061 )?;
1062 assert_eq!(table.append_parquet_segment(non_overlap_path).await?, 3);
1063
1064 let state_before = table.state.clone();
1065 let coverage_before = coverage_files(tmp.path())?;
1066 let overlap_path = "data/uint64-overlap.parquet";
1067 write_single_index_parquet(
1068 &tmp.path().join(overlap_path),
1069 DataType::UInt64,
1070 Arc::new(UInt64Array::from(vec![u64::MAX - 1])),
1071 )?;
1072 let overlap_error = table
1073 .append_parquet_segment(overlap_path)
1074 .await
1075 .expect_err("large uint64 bucket overlap must fail");
1076 assert!(matches!(
1077 overlap_error,
1078 TableError::CoverageOverlap {
1079 example_bucket_range,
1080 ..
1081 } if example_bucket_range.to_string()
1082 == "[18446744073709551610, 18446744073709551615]"
1083 ));
1084 assert_eq!(table.state, state_before);
1085 assert_eq!(coverage_files(tmp.path())?, coverage_before);
1086 let reopened = TimeSeriesTable::open(location).await?;
1087 assert_eq!(reopened.state, table.state);
1088 Ok(())
1089 }
1090
1091 #[tokio::test]
1092 async fn append_rejects_signed_data_for_uint64_index_without_mutation() -> TestResult {
1093 let tmp = TempDir::new()?;
1094 let location = TableLocation::local(tmp.path());
1095 let index = registered_index(IndexKind::UInt64 {
1096 bucket_width: NonZeroU64::new(1).unwrap(),
1097 });
1098 let mut table =
1099 TimeSeriesTable::create(location.clone(), TableMeta::new_time_series(index)).await?;
1100 let state_before = table.state.clone();
1101 let coverage_before = coverage_files(tmp.path())?;
1102 let rel_path = "data/signed.parquet";
1103 write_arrow_parquet_int_time(&tmp.path().join(rel_path), &[1], &["A"], &[1.0])?;
1104
1105 let error = table
1106 .append_parquet_segment(rel_path)
1107 .await
1108 .expect_err("signed data must not append to a uint64 index");
1109
1110 assert!(matches!(
1111 error,
1112 TableError::SegmentMeta {
1113 source: SegmentError::Meta {
1114 source: SegmentMetaError::OrderedIndexColumn {
1115 source: ParquetIndexColumnError {
1116 expected_domain: "uint64",
1117 observed_type,
1118 ..
1119 }
1120 }
1121 }
1122 } if observed_type.contains("logical=None")
1123 ));
1124 assert_eq!(table.state, state_before);
1125 assert_eq!(table.log.load_current_version().await?, 1);
1126 assert_eq!(coverage_files(tmp.path())?, coverage_before);
1127 Ok(())
1128 }
1129
1130 #[tokio::test]
1131 async fn append_inspects_file_without_reading_unrelated_column_data() -> TestResult {
1132 let tmp = TempDir::new()?;
1133 let location = TableLocation::local(tmp.path());
1134 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
1135 let rel_path = "data/corrupt-price.parquet";
1136 let abs_path = tmp.path().join(rel_path);
1137 write_test_parquet(
1138 &abs_path,
1139 true,
1140 false,
1141 &[TestRow {
1142 ts_millis: 1_000,
1143 symbol: "A",
1144 price: 10.0,
1145 }],
1146 )?;
1147
1148 let reader = SerializedFileReader::new(File::open(&abs_path)?)?;
1149 let price_page = reader.metadata().row_group(0).column(2).data_page_offset() as u64;
1150 drop(reader);
1151 let mut file = OpenOptions::new().read(true).write(true).open(&abs_path)?;
1152 file.seek(SeekFrom::Start(price_page))?;
1153 file.write_all(&[0xFF; 16])?;
1154 file.flush()?;
1155 drop(file);
1156
1157 let file_size = std::fs::metadata(&abs_path)?.len().to_string();
1158 let (version, report) = table.append_parquet_segment_with_report(rel_path).await?;
1159
1160 assert_eq!(version, 2);
1161 assert_eq!(
1162 report.context,
1163 vec![
1164 ("relative_path".to_string(), rel_path.to_string()),
1165 ("index_column".to_string(), "ts".to_string()),
1166 ("file_size_bytes".to_string(), file_size),
1167 ]
1168 );
1169 assert_eq!(
1170 report
1171 .steps
1172 .iter()
1173 .map(|step| step.name.as_str())
1174 .collect::<Vec<_>>(),
1175 vec![
1176 "segment_meta",
1177 "logical_schema",
1178 "load_table_snapshot",
1179 "segment_coverage",
1180 "overlap_check",
1181 "write_segment_sidecar",
1182 "write_snapshot_sidecar",
1183 "commit_log",
1184 "state_update",
1185 ]
1186 );
1187 assert_eq!(
1188 report.steps[0]
1189 .fields
1190 .iter()
1191 .map(|(key, _)| key.as_str())
1192 .collect::<Vec<_>>(),
1193 vec!["row_groups", "row_count", "used_stats", "scanned_rows"]
1194 );
1195 assert!(report.steps[1..].iter().all(|step| step.fields.is_empty()));
1196 Ok(())
1197 }
1198
1199 #[tokio::test]
1200 async fn version_six_records_single_layout_for_each_entity_segment() -> TestResult {
1201 let tmp = TempDir::new()?;
1202 let location = TableLocation::local(tmp.path());
1203 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
1204
1205 for (path, symbol) in [
1206 ("data/entity-a.parquet", "A"),
1207 ("data/entity-b.parquet", "B"),
1208 ] {
1209 write_test_parquet(
1210 &tmp.path().join(path),
1211 true,
1212 false,
1213 &[TestRow {
1214 ts_millis: 1_000,
1215 symbol,
1216 price: 10.0,
1217 }],
1218 )?;
1219 table.append_parquet_segment(path).await?;
1220 assert_eq!(
1221 table
1222 .state
1223 .segments
1224 .get(path)
1225 .expect("segment present")
1226 .entity_layout,
1227 SegmentEntityLayout::Single(EntityIdentity::try_new(vec![symbol.into()])?)
1228 );
1229 }
1230
1231 assert_eq!(table.state.version, 3);
1232 let pointer = table
1233 .state
1234 .table_coverage
1235 .as_ref()
1236 .expect("table coverage pointer");
1237 let coverage =
1238 read_entity_coverage_sidecar(&location, Path::new(&pointer.coverage_path)).await?;
1239 assert_eq!(coverage.identity_count(), 2);
1240 assert_eq!(coverage.cardinality(), 2);
1241 Ok(())
1242 }
1243
1244 #[tokio::test]
1245 async fn version_six_records_mixed_layout_for_multiple_identities() -> TestResult {
1246 let tmp = TempDir::new()?;
1247 let location = TableLocation::local(tmp.path());
1248 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
1249 let path = "data/multiple-identities.parquet";
1250 write_test_parquet(
1251 &tmp.path().join(path),
1252 true,
1253 false,
1254 &[
1255 TestRow {
1256 ts_millis: 1_000,
1257 symbol: "A",
1258 price: 10.0,
1259 },
1260 TestRow {
1261 ts_millis: 1_000,
1262 symbol: "B",
1263 price: 20.0,
1264 },
1265 ],
1266 )?;
1267
1268 table.append_parquet_segment(path).await?;
1269
1270 let segment = table.state.segments.get(path).expect("segment present");
1271 assert_eq!(segment.entity_layout, SegmentEntityLayout::Mixed);
1272 let coverage = read_entity_coverage_sidecar(
1273 &location,
1274 Path::new(segment.coverage_path.as_ref().expect("coverage path")),
1275 )
1276 .await?;
1277 assert_eq!(coverage.identity_count(), 2);
1278 assert_eq!(coverage.cardinality(), 2);
1279 Ok(())
1280 }
1281
1282 #[tokio::test]
1283 async fn numeric_entities_append_overlap_and_recover_with_exact_types() -> TestResult {
1284 let tmp = TempDir::new()?;
1285 let location = TableLocation::local(tmp.path());
1286 let mut table =
1287 TimeSeriesTable::create(location.clone(), make_int32_entity_table_meta()).await?;
1288
1289 let negative_path = "data/negative-device.parquet";
1290 write_int32_entity_parquet(
1291 &tmp.path().join(negative_path),
1292 &[1_000, 61_000],
1293 &[-1, -1],
1294 &[10.0, 11.0],
1295 )?;
1296 table.append_parquet_segment(negative_path).await?;
1297 let negative_identity = EntityIdentity::try_new(vec![EntityValue::Int32(-1)])?;
1298 assert_eq!(
1299 table.state.segments[negative_path].entity_layout,
1300 SegmentEntityLayout::Single(negative_identity.clone())
1301 );
1302
1303 let maximum_path = "data/maximum-device.parquet";
1304 write_int32_entity_parquet(
1305 &tmp.path().join(maximum_path),
1306 &[1_000],
1307 &[i32::MAX],
1308 &[20.0],
1309 )?;
1310 table.append_parquet_segment(maximum_path).await?;
1311 assert_eq!(
1312 table.state.segments[maximum_path].entity_layout,
1313 SegmentEntityLayout::Single(EntityIdentity::try_new(vec![EntityValue::Int32(
1314 i32::MAX,
1315 )])?)
1316 );
1317
1318 let overlap_path = "data/negative-overlap.parquet";
1319 write_int32_entity_parquet(&tmp.path().join(overlap_path), &[1_500], &[-1], &[12.0])?;
1320 let error = table
1321 .append_parquet_segment(overlap_path)
1322 .await
1323 .expect_err("same typed identity and bucket must overlap");
1324 assert!(matches!(
1325 error,
1326 TableError::EntityCoverageOverlap {
1327 overlap_count: 1,
1328 example_identity,
1329 ..
1330 } if example_identity == negative_identity
1331 ));
1332
1333 let snapshot = table.load_table_entity_snapshot_coverage_readonly().await?;
1334 let reopened = TimeSeriesTable::open(location).await?;
1335 assert_eq!(reopened.state(), table.state());
1336 assert_eq!(
1337 reopened
1338 .recover_table_entity_coverage_from_segments()
1339 .await?,
1340 snapshot
1341 );
1342 Ok(())
1343 }
1344
1345 #[tokio::test]
1346 async fn version_six_preserves_composite_identity_order_in_layout() -> TestResult {
1347 let tmp = TempDir::new()?;
1348 let location = TableLocation::local(tmp.path());
1349 let index = IndexSpec {
1350 column: "ts".to_string(),
1351 entity_columns: vec!["symbol".to_string(), "venue".to_string()],
1352 kind: IndexKind::Timestamp {
1353 bucket: TimeBucket::Minutes(1),
1354 timezone: None,
1355 },
1356 };
1357 let schema = LogicalSchema::new(vec![
1358 LogicalField {
1359 name: "ts".to_string(),
1360 data_type: LogicalDataType::Timestamp {
1361 unit: LogicalTimestampUnit::Millis,
1362 timezone: None,
1363 },
1364 nullable: false,
1365 },
1366 LogicalField {
1367 name: "symbol".to_string(),
1368 data_type: LogicalDataType::Utf8,
1369 nullable: false,
1370 },
1371 LogicalField {
1372 name: "venue".to_string(),
1373 data_type: LogicalDataType::Utf8,
1374 nullable: false,
1375 },
1376 LogicalField {
1377 name: "price".to_string(),
1378 data_type: LogicalDataType::Float64,
1379 nullable: false,
1380 },
1381 ])?;
1382 let mut table = TimeSeriesTable::create(
1383 location,
1384 TableMeta::new_time_series_with_schema(index, schema),
1385 )
1386 .await?;
1387
1388 for (path, venue) in [
1389 ("data/composite-x.parquet", "X"),
1390 ("data/composite-y.parquet", "Y"),
1391 ] {
1392 write_composite_entity_parquet(&tmp.path().join(path), &[(1_000, "A", venue, 10.0)])?;
1393 table.append_parquet_segment(path).await?;
1394 assert_eq!(
1395 table
1396 .state
1397 .segments
1398 .get(path)
1399 .expect("segment present")
1400 .entity_layout,
1401 SegmentEntityLayout::Single(EntityIdentity::try_new(vec![
1402 "A".into(),
1403 venue.into(),
1404 ])?)
1405 );
1406 }
1407
1408 let overlap_path = "data/composite-x-overlap.parquet";
1409 write_composite_entity_parquet(&tmp.path().join(overlap_path), &[(1_500, "A", "X", 20.0)])?;
1410 let error = table
1411 .append_parquet_segment(overlap_path)
1412 .await
1413 .expect_err("matching composite identity and bucket must overlap");
1414 assert!(matches!(
1415 error,
1416 TableError::EntityCoverageOverlap {
1417 overlap_count: 1,
1418 example_identity,
1419 ..
1420 } if example_identity.components()
1421 == [EntityValue::from("A"), EntityValue::from("X")]
1422 ));
1423 Ok(())
1424 }
1425
1426 #[tokio::test]
1427 async fn entity_with_only_null_index_values_is_rejected() -> TestResult {
1428 let tmp = TempDir::new()?;
1429 let location = TableLocation::local(tmp.path());
1430 let mut table = TimeSeriesTable::create(
1431 location,
1432 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
1433 )
1434 .await?;
1435 let state_before = table.state.clone();
1436 let path = "data/entity-without-index-coverage.parquet";
1437 write_arrow_parquet_with_unit(
1438 &tmp.path().join(path),
1439 ArrowTimeUnit::Millisecond,
1440 &[Some(1_000), None],
1441 &["A", "B"],
1442 &[10.0, 20.0],
1443 )?;
1444
1445 let error = table
1446 .append_parquet_segment(path)
1447 .await
1448 .expect_err("identity without index coverage must be rejected");
1449
1450 match error {
1451 TableError::EntityWithoutIndexCoverage {
1452 segment_path,
1453 identity,
1454 } => {
1455 assert_eq!(segment_path, path);
1456 assert_eq!(identity.components(), [EntityValue::from("B")]);
1457 }
1458 other => panic!("unexpected error: {other:?}"),
1459 }
1460 assert_eq!(table.state, state_before);
1461 assert!(coverage_files(tmp.path())?.is_empty());
1462 Ok(())
1463 }
1464
1465 #[tokio::test]
1466 async fn append_parquet_segment_adopts_schema_when_missing() -> TestResult {
1467 let tmp = TempDir::new()?;
1468 let location = TableLocation::local(tmp.path());
1469
1470 let index = IndexSpec {
1471 column: "ts".to_string(),
1472 entity_columns: vec![],
1473 kind: IndexKind::Timestamp {
1474 bucket: TimeBucket::Minutes(1),
1475 timezone: None,
1476 },
1477 };
1478 let meta = TableMeta {
1479 kind: TableKind::TimeSeries(index),
1480 logical_schema: None,
1481 created_at: utc_datetime(2025, 1, 1, 0, 0, 0),
1482 format_version: TABLE_FORMAT_VERSION,
1483 };
1484
1485 let mut table = TimeSeriesTable::create(location, meta).await?;
1486
1487 let rel_path = "data/seg-adopt.parquet";
1488 let abs_path = tmp.path().join(rel_path);
1489 write_test_parquet(
1490 &abs_path,
1491 true,
1492 false,
1493 &[TestRow {
1494 ts_millis: 5_000,
1495 symbol: "B",
1496 price: 20.0,
1497 }],
1498 )?;
1499
1500 let new_version = table.append_parquet_segment(rel_path).await?;
1501
1502 assert_eq!(new_version, 2);
1503 let schema = table
1504 .state
1505 .table_meta
1506 .logical_schema
1507 .as_ref()
1508 .expect("schema adopted");
1509 let names: Vec<_> = schema.columns().iter().map(|c| c.name.as_str()).collect();
1510 assert_eq!(names, vec!["ts", "symbol", "price"]);
1511 let ts_col = &schema.columns()[0];
1512 assert_eq!(
1513 ts_col.data_type,
1514 LogicalDataType::Timestamp {
1515 unit: LogicalTimestampUnit::Millis,
1516 timezone: None,
1517 }
1518 );
1519 Ok(())
1520 }
1521
1522 #[tokio::test]
1523 async fn append_parquet_segment_rejects_schema_mismatch() -> TestResult {
1524 let tmp = TempDir::new()?;
1525 let location = TableLocation::local(tmp.path());
1526 let meta = make_basic_table_meta();
1527 let mut table = TimeSeriesTable::create(location, meta).await?;
1528
1529 let rel_path = "data/seg-missing-symbol.parquet";
1530 let abs_path = tmp.path().join(rel_path);
1531 write_test_parquet(
1532 &abs_path,
1533 false,
1534 false,
1535 &[TestRow {
1536 ts_millis: 10_000,
1537 symbol: "C",
1538 price: 30.0,
1539 }],
1540 )?;
1541
1542 let err = table
1543 .append_parquet_segment(rel_path)
1544 .await
1545 .expect_err("expected schema mismatch");
1546
1547 match err {
1548 TableError::SegmentSchemaCompatibility { path, source } => {
1549 assert_eq!(path, rel_path);
1550 assert!(matches!(
1551 source,
1552 crate::metadata::schema_compat::SchemaCompatibilityError::MissingEntityColumn { .. }
1553 ));
1554 }
1555 other => panic!("unexpected error: {other:?}"),
1556 }
1557 Ok(())
1558 }
1559
1560 #[tokio::test]
1561 async fn first_append_rejects_unsupported_entity_type_without_publication() -> TestResult {
1562 let tmp = TempDir::new()?;
1563 let location = TableLocation::local(tmp.path());
1564 let index = IndexSpec {
1565 column: "ts".to_string(),
1566 entity_columns: vec!["device_id".to_string()],
1567 kind: IndexKind::Timestamp {
1568 bucket: TimeBucket::Minutes(1),
1569 timezone: None,
1570 },
1571 };
1572 let mut table =
1573 TimeSeriesTable::create(location, TableMeta::new_time_series(index)).await?;
1574 let rel_path = "data/unsupported-entity.parquet";
1575 let abs_path = tmp.path().join(rel_path);
1576 std::fs::create_dir_all(abs_path.parent().expect("data parent"))?;
1577 let schema = Arc::new(Schema::new(vec![
1578 Field::new(
1579 "ts",
1580 DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
1581 false,
1582 ),
1583 Field::new("device_id", DataType::Boolean, false),
1584 ]));
1585 let batch = RecordBatch::try_new(
1586 Arc::clone(&schema),
1587 vec![
1588 Arc::new(TimestampMillisecondArray::from(vec![1_000])),
1589 Arc::new(BooleanArray::from(vec![true])),
1590 ],
1591 )?;
1592 let mut writer = ArrowWriter::try_new(File::create(abs_path)?, schema, None)?;
1593 writer.write(&batch)?;
1594 writer.close()?;
1595 let state_before = table.state.clone();
1596
1597 let error = table
1598 .append_parquet_segment(rel_path)
1599 .await
1600 .expect_err("Boolean entity columns must be rejected");
1601
1602 assert!(matches!(
1603 error,
1604 TableError::SegmentSchemaCompatibility { path, source }
1605 if path == rel_path
1606 && matches!(
1607 &source,
1608 crate::metadata::schema_compat::SchemaCompatibilityError::UnsupportedEntityColumnType {
1609 column,
1610 actual: LogicalDataType::Bool,
1611 } if column == "device_id"
1612 )
1613 ));
1614 assert_eq!(table.state, state_before);
1615 assert!(coverage_files(tmp.path())?.is_empty());
1616 assert_eq!(table.log.load_current_version().await?, 1);
1617 Ok(())
1618 }
1619
1620 #[tokio::test]
1621 async fn append_rejects_duplicate_path_before_parquet_read_without_mutation() -> TestResult {
1622 let tmp = TempDir::new()?;
1623 let location = TableLocation::local(tmp.path());
1624 let meta = make_basic_table_meta();
1625 let mut table = TimeSeriesTable::create(location.clone(), meta).await?;
1626
1627 let rel_path = "data/dup.parquet";
1628 let abs_path = tmp.path().join(rel_path);
1629
1630 write_test_parquet(
1631 &abs_path,
1632 true,
1633 false,
1634 &[TestRow {
1635 ts_millis: 1_000,
1636 symbol: "A",
1637 price: 10.0,
1638 }],
1639 )?;
1640 table.append_parquet_segment(rel_path).await?;
1641 let state_before = table.state.clone();
1642 let sidecar_counts_before = [
1643 std::fs::read_dir(tmp.path().join(layout::SEGMENT_COVERAGE_DIR))?.count(),
1644 std::fs::read_dir(tmp.path().join(layout::TABLE_SNAPSHOT_DIR))?.count(),
1645 ];
1646
1647 tokio::fs::remove_file(&abs_path).await?;
1650
1651 let err = table
1652 .append_parquet_segment(rel_path)
1653 .await
1654 .expect_err("live path must be rejected");
1655 assert!(matches!(
1656 err,
1657 TableError::DuplicateSegmentPath { ref path } if path == rel_path
1658 ));
1659
1660 let err = table
1661 .append_parquet_segment_with_report(r"data\dup.parquet")
1662 .await
1663 .expect_err("normalized live path must be rejected");
1664 assert!(matches!(
1665 err,
1666 TableError::DuplicateSegmentPath { ref path } if path == rel_path
1667 ));
1668
1669 assert_eq!(table.state, state_before);
1670 assert_eq!(table.log.load_current_version().await?, 2);
1671 assert!(!tmp.path().join(layout::commit_rel_path(3)).exists());
1672 assert_eq!(
1673 [
1674 std::fs::read_dir(tmp.path().join(layout::SEGMENT_COVERAGE_DIR))?.count(),
1675 std::fs::read_dir(tmp.path().join(layout::TABLE_SNAPSHOT_DIR))?.count(),
1676 ],
1677 sidecar_counts_before
1678 );
1679 Ok(())
1680 }
1681
1682 #[tokio::test]
1683 async fn append_parquet_segment_keys_paths_and_updates_snapshot() -> TestResult {
1684 let tmp = TempDir::new()?;
1685 let location = TableLocation::local(tmp.path());
1686 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
1687
1688 let rel1 = "data/seg-auto-1.parquet";
1689 let rel2 = "data/seg-auto-2.parquet";
1690 let path1 = tmp.path().join(rel1);
1691 let path2 = tmp.path().join(rel2);
1692
1693 write_test_parquet(
1694 &path1,
1695 true,
1696 false,
1697 &[
1698 TestRow {
1699 ts_millis: 1_000,
1700 symbol: "A",
1701 price: 10.0,
1702 },
1703 TestRow {
1704 ts_millis: 2_000,
1705 symbol: "A",
1706 price: 20.0,
1707 },
1708 ],
1709 )?;
1710 write_test_parquet(
1711 &path2,
1712 true,
1713 false,
1714 &[
1715 TestRow {
1716 ts_millis: 120_000,
1717 symbol: "A",
1718 price: 30.0,
1719 },
1720 TestRow {
1721 ts_millis: 121_000,
1722 symbol: "A",
1723 price: 40.0,
1724 },
1725 ],
1726 )?;
1727
1728 let v2 = table.append_parquet_segment(rel1).await?;
1729 let v3 = table.append_parquet_segment(rel2).await?;
1730 assert_eq!(v2, 2);
1731 assert_eq!(v3, 3);
1732
1733 let seg1 = table.state.segments.get(rel1).expect("segment 1 present");
1734 let seg2 = table.state.segments.get(rel2).expect("segment 2 present");
1735 assert_eq!(seg1.path, rel1);
1736 assert_eq!(seg2.path, rel2);
1737 assert!(seg1.coverage_path.is_some());
1738 assert!(seg2.coverage_path.is_some());
1739
1740 let cov1 =
1741 compute_segment_entity_coverage(&location, Path::new(rel1), table.index_spec()).await?;
1742 let cov2 =
1743 compute_segment_entity_coverage(&location, Path::new(rel2), table.index_spec()).await?;
1744 let expected_snapshot = cov1.union(&cov2);
1745
1746 let ptr = table
1747 .state
1748 .table_coverage
1749 .as_ref()
1750 .expect("table snapshot pointer present after append");
1751 assert_eq!(ptr.version, v3);
1752 assert_eq!(ptr.index_kind, table.index_spec().kind);
1753
1754 let snapshot_cov =
1755 read_entity_coverage_sidecar(&location, Path::new(&ptr.coverage_path)).await?;
1756
1757 assert_eq!(snapshot_cov, expected_snapshot);
1758 Ok(())
1759 }
1760
1761 #[tokio::test]
1762 async fn append_parquet_segment_rejects_overlap() -> TestResult {
1763 let tmp = TempDir::new()?;
1764 let location = TableLocation::local(tmp.path());
1765 let mut table = TimeSeriesTable::create(location, make_basic_table_meta()).await?;
1766
1767 let rel1 = "data/seg-overlap-a.parquet";
1768 let rel2 = "data/seg-overlap-b.parquet";
1769 let path1 = tmp.path().join(rel1);
1770 let path2 = tmp.path().join(rel2);
1771
1772 write_test_parquet(
1773 &path1,
1774 true,
1775 false,
1776 &[
1777 TestRow {
1778 ts_millis: 1_000,
1779 symbol: "B",
1780 price: 10.0,
1781 },
1782 TestRow {
1783 ts_millis: 1_000,
1784 symbol: "A",
1785 price: 20.0,
1786 },
1787 TestRow {
1788 ts_millis: 61_000,
1789 symbol: "A",
1790 price: 30.0,
1791 },
1792 ],
1793 )?;
1794 write_test_parquet(
1795 &path2,
1796 true,
1797 false,
1798 &[
1799 TestRow {
1800 ts_millis: 1_500,
1801 symbol: "B",
1802 price: 40.0,
1803 },
1804 TestRow {
1805 ts_millis: 1_500,
1806 symbol: "A",
1807 price: 50.0,
1808 },
1809 TestRow {
1810 ts_millis: 61_500,
1811 symbol: "A",
1812 price: 60.0,
1813 },
1814 ],
1815 )?;
1816
1817 table.append_parquet_segment(rel1).await?;
1818
1819 let err = table
1820 .append_parquet_segment(rel2)
1821 .await
1822 .expect_err("overlapping append should fail");
1823
1824 assert!(matches!(
1825 err,
1826 TableError::EntityCoverageOverlap {
1827 segment_path,
1828 overlap_count: 3,
1829 example_identity,
1830 example_bucket: 0x8000_0000_0000_0000,
1831 example_bucket_range,
1832 } if segment_path == rel2
1833 && example_identity.components() == [EntityValue::from("A")]
1834 && example_bucket_range.to_string()
1835 == "[1970-01-01T00:00:00Z, 1970-01-01T00:01:00Z)"
1836 ));
1837 Ok(())
1838 }
1839
1840 #[tokio::test]
1841 async fn append_parquet_segment_snapshot_survives_reopen() -> TestResult {
1842 let tmp = TempDir::new()?;
1843 let location = TableLocation::local(tmp.path());
1844 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
1845
1846 let rel1 = "data/seg-reopen-a.parquet";
1847 let rel2 = "data/seg-reopen-b.parquet";
1848 let path1 = tmp.path().join(rel1);
1849 let path2 = tmp.path().join(rel2);
1850
1851 write_test_parquet(
1852 &path1,
1853 true,
1854 false,
1855 &[TestRow {
1856 ts_millis: 1_000,
1857 symbol: "A",
1858 price: 10.0,
1859 }],
1860 )?;
1861 write_test_parquet(
1862 &path2,
1863 true,
1864 false,
1865 &[TestRow {
1866 ts_millis: 120_000,
1867 symbol: "A",
1868 price: 20.0,
1869 }],
1870 )?;
1871
1872 table.append_parquet_segment(rel1).await?;
1873 table.append_parquet_segment(rel2).await?;
1874
1875 let reopened = TimeSeriesTable::open(location.clone()).await?;
1876 let ptr = reopened
1877 .state()
1878 .table_coverage
1879 .as_ref()
1880 .expect("table snapshot pointer present after reopen");
1881
1882 assert_eq!(ptr.index_kind, reopened.index_spec().kind);
1883
1884 let cov1 =
1885 compute_segment_entity_coverage(&location, Path::new(rel1), reopened.index_spec())
1886 .await?;
1887 let cov2 =
1888 compute_segment_entity_coverage(&location, Path::new(rel2), reopened.index_spec())
1889 .await?;
1890 let expected = cov1.union(&cov2);
1891
1892 let snapshot_cov =
1893 read_entity_coverage_sidecar(&location, Path::new(&ptr.coverage_path)).await?;
1894 assert_eq!(snapshot_cov, expected);
1895 Ok(())
1896 }
1897
1898 #[tokio::test]
1899 async fn load_snapshot_recovers_when_missing_file() -> TestResult {
1900 let tmp = TempDir::new()?;
1901 let location = TableLocation::local(tmp.path());
1902 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
1903
1904 let rel1 = "data/seg-missing-snap-a.parquet";
1906 let rel2 = "data/seg-missing-snap-b.parquet";
1907 let path1 = tmp.path().join(rel1);
1908 let path2 = tmp.path().join(rel2);
1909 write_test_parquet(
1910 &path1,
1911 true,
1912 false,
1913 &[TestRow {
1914 ts_millis: 1_000,
1915 symbol: "A",
1916 price: 10.0,
1917 }],
1918 )?;
1919 write_test_parquet(
1920 &path2,
1921 true,
1922 false,
1923 &[TestRow {
1924 ts_millis: 120_000,
1925 symbol: "A",
1926 price: 20.0,
1927 }],
1928 )?;
1929
1930 table.append_parquet_segment(rel1).await?;
1931 table.append_parquet_segment(rel2).await?;
1932
1933 let state = table.state.clone();
1934 let ptr = state
1935 .table_coverage
1936 .as_ref()
1937 .expect("snapshot pointer present");
1938 let snapshot_abs = match &location.as_ref() {
1939 StorageLocation::Local(root) => root.join(&ptr.coverage_path),
1940 };
1941
1942 tokio::fs::remove_file(&snapshot_abs).await?;
1943
1944 let recovered = table.load_table_entity_snapshot_coverage_readonly().await?;
1945
1946 let mut expected = EntityCoverage::empty();
1947 for seg in state.segments.values() {
1948 let cov_path = seg.coverage_path.as_ref().expect("coverage path");
1949 let cov = read_entity_coverage_sidecar(&location, Path::new(cov_path)).await?;
1950 expected.union_inplace(&cov);
1951 }
1952
1953 assert_eq!(recovered, expected);
1954 Ok(())
1955 }
1956
1957 #[tokio::test]
1958 async fn load_snapshot_recovers_when_corrupt_file() -> TestResult {
1959 let tmp = TempDir::new()?;
1960 let location = TableLocation::local(tmp.path());
1961 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
1962
1963 let rel1 = "data/seg-corrupt-snap-a.parquet";
1964 let rel2 = "data/seg-corrupt-snap-b.parquet";
1965 let path1 = tmp.path().join(rel1);
1966 let path2 = tmp.path().join(rel2);
1967 write_test_parquet(
1968 &path1,
1969 true,
1970 false,
1971 &[TestRow {
1972 ts_millis: 1_000,
1973 symbol: "A",
1974 price: 10.0,
1975 }],
1976 )?;
1977 write_test_parquet(
1978 &path2,
1979 true,
1980 false,
1981 &[TestRow {
1982 ts_millis: 120_000,
1983 symbol: "A",
1984 price: 20.0,
1985 }],
1986 )?;
1987
1988 table.append_parquet_segment(rel1).await?;
1989 table.append_parquet_segment(rel2).await?;
1990
1991 let state = table.state.clone();
1992 let ptr = state
1993 .table_coverage
1994 .as_ref()
1995 .expect("snapshot pointer present");
1996 let snapshot_abs = match &location.as_ref() {
1997 StorageLocation::Local(root) => root.join(&ptr.coverage_path),
1998 };
1999
2000 tokio::fs::write(&snapshot_abs, b"garbage").await?;
2001
2002 let recovered = table.load_table_entity_snapshot_coverage_readonly().await?;
2003
2004 let mut expected = EntityCoverage::empty();
2005 for seg in state.segments.values() {
2006 let cov_path = seg.coverage_path.as_ref().expect("coverage path");
2007 let cov = read_entity_coverage_sidecar(&location, Path::new(cov_path)).await?;
2008 expected.union_inplace(&cov);
2009 }
2010
2011 assert_eq!(recovered, expected);
2012 Ok(())
2013 }
2014
2015 #[tokio::test]
2016 async fn rejected_append_does_not_heal_corrupt_snapshot() -> TestResult {
2017 let tmp = TempDir::new()?;
2018 let location = TableLocation::local(tmp.path());
2019 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
2020
2021 let existing = "data/existing.parquet";
2022 write_test_parquet(
2023 &tmp.path().join(existing),
2024 true,
2025 false,
2026 &[TestRow {
2027 ts_millis: 1_000,
2028 symbol: "A",
2029 price: 10.0,
2030 }],
2031 )?;
2032 table.append_parquet_segment(existing).await?;
2033
2034 let snapshot_path = table
2035 .state
2036 .table_coverage
2037 .as_ref()
2038 .expect("snapshot pointer present")
2039 .coverage_path
2040 .clone();
2041 let snapshot_abs = tmp.path().join(snapshot_path);
2042 tokio::fs::write(&snapshot_abs, b"garbage").await?;
2043
2044 let overlapping = "data/overlapping.parquet";
2045 write_test_parquet(
2046 &tmp.path().join(overlapping),
2047 true,
2048 false,
2049 &[TestRow {
2050 ts_millis: 1_000,
2051 symbol: "A",
2052 price: 20.0,
2053 }],
2054 )?;
2055
2056 let err = table
2057 .append_parquet_segment(overlapping)
2058 .await
2059 .expect_err("overlap must be rejected");
2060 assert!(matches!(err, TableError::EntityCoverageOverlap { .. }));
2061 assert_eq!(tokio::fs::read(snapshot_abs).await?, b"garbage");
2062 Ok(())
2063 }
2064
2065 #[tokio::test]
2066 async fn load_snapshot_errors_when_segment_missing_coverage_path() -> TestResult {
2067 let tmp = TempDir::new()?;
2068 let location = TableLocation::local(tmp.path());
2069 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
2070
2071 let rel1 = "data/seg-missing-cov-path.parquet";
2072 let path1 = tmp.path().join(rel1);
2073 write_test_parquet(
2074 &path1,
2075 true,
2076 false,
2077 &[TestRow {
2078 ts_millis: 1_000,
2079 symbol: "A",
2080 price: 10.0,
2081 }],
2082 )?;
2083
2084 table.append_parquet_segment(rel1).await?;
2085
2086 let mut state = table.state.clone();
2087 state.table_coverage = None;
2088
2089 let segment_path = state
2090 .segments
2091 .keys()
2092 .next()
2093 .expect("segment present")
2094 .clone();
2095 state
2096 .segments
2097 .get_mut(&segment_path)
2098 .expect("segment present")
2099 .coverage_path = None;
2100
2101 table.state = state;
2103
2104 let err = table
2105 .load_table_entity_snapshot_coverage_readonly()
2106 .await
2107 .expect_err("missing coverage_path should error");
2108
2109 assert!(matches!(
2110 err,
2111 TableError::ExistingSegmentMissingCoverage { .. }
2112 ));
2113 Ok(())
2114 }
2115
2116 #[tokio::test]
2117 async fn load_snapshot_errors_when_segment_sidecar_corrupt() -> TestResult {
2118 let tmp = TempDir::new()?;
2119 let location = TableLocation::local(tmp.path());
2120 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
2121
2122 let rel1 = "data/seg-corrupt-sidecar.parquet";
2123 let rel2 = "data/seg-corrupt-sidecar-ok.parquet";
2124 let path1 = tmp.path().join(rel1);
2125 let path2 = tmp.path().join(rel2);
2126 write_test_parquet(
2127 &path1,
2128 true,
2129 false,
2130 &[TestRow {
2131 ts_millis: 1_000,
2132 symbol: "A",
2133 price: 10.0,
2134 }],
2135 )?;
2136 write_test_parquet(
2137 &path2,
2138 true,
2139 false,
2140 &[TestRow {
2141 ts_millis: 120_000,
2142 symbol: "A",
2143 price: 20.0,
2144 }],
2145 )?;
2146
2147 table.append_parquet_segment(rel1).await?;
2148 table.append_parquet_segment(rel2).await?;
2149
2150 let mut state = table.state.clone();
2151 state.table_coverage = None;
2152 let (corrupt_segment_path, corrupt_cov_path) = state
2153 .segments
2154 .values()
2155 .next()
2156 .map(|meta| {
2157 (
2158 meta.path.clone(),
2159 meta.coverage_path.as_ref().expect("coverage path").clone(),
2160 )
2161 })
2162 .expect("at least one segment");
2163 table.state = state;
2164
2165 let corrupt_abs = match &location.as_ref() {
2166 StorageLocation::Local(root) => root.join(&corrupt_cov_path),
2167 };
2168 tokio::fs::write(&corrupt_abs, b"not a coverage bitmap").await?;
2169
2170 let err = table
2171 .load_table_entity_snapshot_coverage_readonly()
2172 .await
2173 .expect_err("corrupt sidecar should error");
2174
2175 match err {
2176 TableError::SegmentCoverageSidecarRead {
2177 path,
2178 coverage_path,
2179 ..
2180 } => {
2181 assert_eq!(path, corrupt_segment_path);
2182 assert_eq!(coverage_path, corrupt_cov_path);
2183 }
2184 other => panic!("unexpected error: {other:?}"),
2185 }
2186
2187 Ok(())
2188 }
2189
2190 #[tokio::test]
2191 async fn entity_aware_stale_append_cleans_sidecars_without_state_mutation() -> TestResult {
2192 let tmp = TempDir::new()?;
2193 let location = TableLocation::local(tmp.path());
2194 let meta = make_basic_table_meta();
2195 let mut winner = TimeSeriesTable::create(location.clone(), meta).await?;
2196 let mut loser = TimeSeriesTable::open(location.clone()).await?;
2197 let loser_state_before = loser.state.clone();
2198
2199 let winner_path = "data/winner.parquet";
2200 let loser_path = "data/loser.parquet";
2201 write_test_parquet(
2202 &tmp.path().join(winner_path),
2203 true,
2204 false,
2205 &[TestRow {
2206 ts_millis: 10_000,
2207 symbol: "X",
2208 price: 100.0,
2209 }],
2210 )?;
2211 write_test_parquet(
2212 &tmp.path().join(loser_path),
2213 true,
2214 false,
2215 &[TestRow {
2216 ts_millis: 120_000,
2217 symbol: "X",
2218 price: 200.0,
2219 }],
2220 )?;
2221
2222 assert_eq!(winner.append_parquet_segment(winner_path).await?, 2);
2223 let coverage_before = coverage_files(tmp.path())?;
2224
2225 let err = loser
2226 .append_parquet_segment(loser_path)
2227 .await
2228 .expect_err("expected conflict due to stale version");
2229
2230 match err {
2231 TableError::TransactionLog { source } => {
2232 assert!(matches!(
2233 source,
2234 CommitError::Conflict {
2235 expected: 1,
2236 found: 2,
2237 ..
2238 }
2239 ));
2240 }
2241 other => panic!("unexpected error: {other:?}"),
2242 }
2243
2244 assert_eq!(loser.state, loser_state_before);
2245 assert_eq!(loser.log.load_current_version().await?, 2);
2246 let committed = loser.load_latest_state().await?;
2247 assert!(committed.segments.contains_key(winner_path));
2248 assert!(!committed.segments.contains_key(loser_path));
2249 assert_eq!(coverage_files(tmp.path())?, coverage_before);
2250 for bytes in coverage_before.values() {
2251 entity_coverage_from_bytes(bytes)?;
2252 }
2253 assert!(!tmp.path().join(layout::commit_rel_path(3)).exists());
2254 Ok(())
2255 }
2256
2257 #[tokio::test]
2258 async fn stale_int64_append_cleans_only_its_writer_owned_sidecars() -> TestResult {
2259 let tmp = TempDir::new()?;
2260 let location = TableLocation::local(tmp.path());
2261 let index = registered_index(IndexKind::Int64 {
2262 bucket_width: NonZeroU64::new(10).unwrap(),
2263 });
2264 let mut winner =
2265 TimeSeriesTable::create(location.clone(), TableMeta::new_time_series(index)).await?;
2266 let mut loser = TimeSeriesTable::open(location).await?;
2267 let winner_path = "data/writer-owned-winner.parquet";
2268 let loser_path = "data/writer-owned-loser.parquet";
2269
2270 write_arrow_parquet_int_time(&tmp.path().join(winner_path), &[0], &["X"], &[100.0])?;
2271 write_arrow_parquet_int_time(&tmp.path().join(loser_path), &[100], &["X"], &[200.0])?;
2272
2273 winner.append_parquet_segment(winner_path).await?;
2274 let coverage_before = coverage_files(tmp.path())?;
2275
2276 let err = loser
2277 .append_parquet_segment(loser_path)
2278 .await
2279 .expect_err("stale append should conflict");
2280
2281 assert!(matches!(
2282 err,
2283 TableError::TransactionLog {
2284 source: CommitError::Conflict { .. }
2285 }
2286 ));
2287 assert_eq!(coverage_files(tmp.path())?, coverage_before);
2288 Ok(())
2289 }
2290
2291 #[tokio::test]
2292 async fn ambiguous_int64_commit_retains_writer_owned_sidecars() -> TestResult {
2293 let tmp = TempDir::new()?;
2294 let location = TableLocation::local(tmp.path());
2295 let index = registered_index(IndexKind::Int64 {
2296 bucket_width: NonZeroU64::new(10).unwrap(),
2297 });
2298 let mut table =
2299 TimeSeriesTable::create(location, TableMeta::new_time_series(index)).await?;
2300 let state_before = table.state.clone();
2301 let coverage_before = coverage_files(tmp.path())?;
2302 let segment_path = "data/ambiguous.parquet";
2303
2304 write_arrow_parquet_int_time(&tmp.path().join(segment_path), &[10], &["X"], &[100.0])?;
2305
2306 let commit_path = tmp.path().join(layout::commit_rel_path(2));
2307 crate::storage::inject_write_new_failure(commit_path.clone(), true);
2308
2309 let err = table
2310 .append_parquet_segment(segment_path)
2311 .await
2312 .expect_err("failed commit cleanup should make the outcome ambiguous");
2313
2314 assert!(matches!(
2315 err,
2316 TableError::TransactionLog {
2317 source: CommitError::AmbiguousOutcome { .. }
2318 }
2319 ));
2320 assert_eq!(table.state, state_before);
2321 assert_eq!(table.log.load_current_version().await?, 1);
2322 assert!(commit_path.exists());
2323 assert_eq!(coverage_files(tmp.path())?.len(), coverage_before.len() + 2);
2324 Ok(())
2325 }
2326
2327 #[tokio::test]
2328 async fn entity_sidecar_cleanup_failures_preserve_error_and_reverse_order() -> TestResult {
2329 let tmp = TempDir::new()?;
2330 let location = TableLocation::local(tmp.path());
2331 let table = TimeSeriesTable::create(location, make_basic_table_meta()).await?;
2332 let sidecars = [
2333 format!("{}/first-stuck.roar", layout::SEGMENT_COVERAGE_DIR),
2334 format!("{}/second-stuck.roar", layout::TABLE_SNAPSHOT_DIR),
2335 ];
2336 for sidecar in &sidecars {
2337 tokio::fs::create_dir_all(tmp.path().join(sidecar)).await?;
2338 }
2339 let example_bucket = crate::coverage::bucket::bucket_id(
2340 &table.index.kind,
2341 &IndexValue::Timestamp(utc_datetime(1970, 1, 1, 0, 0, 0)),
2342 )?;
2343 let source = TableError::EntityCoverageOverlap {
2344 segment_path: "data/failed.parquet".to_string(),
2345 overlap_count: 1,
2346 example_identity: EntityIdentity::try_new(vec!["A".into()])?,
2347 example_bucket,
2348 example_bucket_range: logical_bucket_range(&table.index.kind, example_bucket)?,
2349 };
2350 let err = table.rollback_created_sidecars(&sidecars, source).await;
2351 let message = err.to_string();
2352
2353 assert!(matches!(
2354 err,
2355 TableError::AppendRollback {
2356 source,
2357 cleanup_errors,
2358 } if matches!(*source, TableError::EntityCoverageOverlap { .. })
2359 && cleanup_errors.len() == 2
2360 && cleanup_errors[0].contains("second-stuck.roar")
2361 && cleanup_errors[1].contains("first-stuck.roar")
2362 ));
2363 assert!(message.contains("data/failed.parquet"));
2364 assert!(message.contains("first-stuck.roar"));
2365 assert!(message.contains("second-stuck.roar"));
2366 Ok(())
2367 }
2368
2369 #[tokio::test]
2370 async fn append_fails_when_existing_segment_missing_coverage_path() -> TestResult {
2371 let tmp = TempDir::new()?;
2372 let location = TableLocation::local(tmp.path());
2373 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
2374
2375 let rel1 = "data/seg-missing-cov.parquet";
2376 let rel2 = "data/seg-next.parquet";
2377 let path1 = tmp.path().join(rel1);
2378 let path2 = tmp.path().join(rel2);
2379
2380 write_test_parquet(
2381 &path1,
2382 true,
2383 false,
2384 &[TestRow {
2385 ts_millis: 1_000,
2386 symbol: "A",
2387 price: 10.0,
2388 }],
2389 )?;
2390 write_test_parquet(
2391 &path2,
2392 true,
2393 false,
2394 &[TestRow {
2395 ts_millis: 120_000,
2396 symbol: "A",
2397 price: 20.0,
2398 }],
2399 )?;
2400
2401 table.append_parquet_segment(rel1).await?;
2402
2403 let seg = table.state.segments.get_mut(rel1).expect("segment present");
2405 seg.coverage_path = None;
2406
2407 let err = table
2408 .append_parquet_segment(rel2)
2409 .await
2410 .expect_err("append should fail when existing segment lacks coverage");
2411
2412 assert!(matches!(
2413 err,
2414 TableError::ExistingSegmentMissingCoverage { .. }
2415 ));
2416 Ok(())
2417 }
2418
2419 #[tokio::test]
2420 async fn append_recovers_when_table_snapshot_pointer_missing() -> TestResult {
2425 let tmp = TempDir::new()?;
2426 let location = TableLocation::local(tmp.path());
2427 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
2428
2429 let rel1 = "data/seg-no-pointer-a.parquet";
2430 let rel2 = "data/seg-no-pointer-b.parquet";
2431 let path1 = tmp.path().join(rel1);
2432 let path2 = tmp.path().join(rel2);
2433
2434 write_test_parquet(
2435 &path1,
2436 true,
2437 false,
2438 &[TestRow {
2439 ts_millis: 1_000,
2440 symbol: "A",
2441 price: 10.0,
2442 }],
2443 )?;
2444 write_test_parquet(
2445 &path2,
2446 true,
2447 false,
2448 &[TestRow {
2449 ts_millis: 120_000,
2450 symbol: "A",
2451 price: 20.0,
2452 }],
2453 )?;
2454
2455 table.append_parquet_segment(rel1).await?;
2456
2457 table.state.table_coverage = None;
2459
2460 table.append_parquet_segment(rel2).await?;
2461
2462 let ptr = table
2464 .state
2465 .table_coverage
2466 .as_ref()
2467 .expect("snapshot pointer restored");
2468
2469 let cov = read_entity_coverage_sidecar(&location, Path::new(&ptr.coverage_path)).await?;
2470
2471 let mut expected = EntityCoverage::empty();
2472 for seg in table.state.segments.values() {
2473 let path = seg.coverage_path.as_ref().expect("coverage path");
2474 let seg_cov = read_entity_coverage_sidecar(&location, Path::new(path)).await?;
2475 expected.union_inplace(&seg_cov);
2476 }
2477
2478 assert_eq!(cov, expected);
2479 Ok(())
2480 }
2481
2482 #[tokio::test]
2483 async fn append_fails_when_table_snapshot_bucket_mismatches_index() -> TestResult {
2484 let tmp = TempDir::new()?;
2485 let location = TableLocation::local(tmp.path());
2486 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
2487
2488 let rel1 = "data/seg-bucket-a.parquet";
2489 let rel2 = "data/seg-bucket-b.parquet";
2490 let path1 = tmp.path().join(rel1);
2491 let path2 = tmp.path().join(rel2);
2492
2493 write_test_parquet(
2494 &path1,
2495 true,
2496 false,
2497 &[TestRow {
2498 ts_millis: 1_000,
2499 symbol: "A",
2500 price: 10.0,
2501 }],
2502 )?;
2503 write_test_parquet(
2504 &path2,
2505 true,
2506 false,
2507 &[TestRow {
2508 ts_millis: 120_000,
2509 symbol: "A",
2510 price: 20.0,
2511 }],
2512 )?;
2513
2514 table.append_parquet_segment(rel1).await?;
2515
2516 let bad_bucket = TimeBucket::Hours(1);
2518 let ptr = table
2519 .state
2520 .table_coverage
2521 .as_ref()
2522 .expect("pointer present")
2523 .clone();
2524 table.state.table_coverage = Some(TableCoveragePointer {
2525 index_kind: IndexKind::Timestamp {
2526 bucket: bad_bucket.clone(),
2527 timezone: None,
2528 },
2529 coverage_path: ptr.coverage_path.clone(),
2530 version: ptr.version,
2531 });
2532
2533 let err = table
2534 .append_parquet_segment(rel2)
2535 .await
2536 .expect_err("append should fail when snapshot bucket mismatches index");
2537
2538 assert!(matches!(
2539 err,
2540 TableError::TableCoverageIndexKindMismatch { .. }
2541 ));
2542 Ok(())
2543 }
2544}