Skip to main content

timeseries_table_format/table/
append.rs

1//! Append pipeline for `TimeSeriesTable`.
2//!
3//! This module contains the core append implementation plus the public
4//! wrappers. It is responsible for:
5//! - loading/deriving segment metadata and logical schema,
6//! - enforcing v0.1 schema rules (adopt on first append, otherwise exact match),
7//! - computing segment coverage, detecting overlaps, and writing coverage sidecars,
8//! - optimistic commit to the transaction log and in-memory state update.
9//!   Keep new append-time invariants here so the flow remains centralized.
10
11use 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        // 0) Coverage readiness checks.
193        ensure_existing_segments_have_coverage(&self.state)?;
194
195        // 1) Segment meta + schema.
196        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        // 2) Schema behavior (return maybe_updated_meta, but do NOT build actions yet).
231        //
232        // - logical_schema == None && version == 1:
233        //     first append after create() — adopt this segment’s schema.
234        // - logical_schema == None && version != 1:
235        //     table is in a bad state for v0.1 → error.
236        // - logical_schema == Some(..):
237        //     enforce “no schema evolution” via ensure_schema_exact_match.
238        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        // 3-5) Load, compute, and compare coverage using the entity-column mode.
265        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        // 6) Give this append private sidecar paths, then write them before commit.
362        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        // 7) Build actions and atomically publish the commit.
418        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        // OCC invariant: a successful commit_with_expected_version must return
454        // the same "next" version we predicted when constructing `snapshot_path`.
455        // If this ever diverges, it indicates a severe bug between snapshot path
456        // construction and the transaction log implementation, so we panic rather
457        // than continuing with an inconsistent in-memory state.
458        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        // 8) Update in-memory state.
465        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        // Also update the snapshot pointer in state.
477        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    /// Append a Parquet segment using its canonical relative path as identity.
490    ///
491    /// Rows need not be ordered by the table's ordered index.
492    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    /// Copy an external Parquet file into the table when needed and append it.
498    ///
499    /// Rows need not be ordered by the table's ordered index.
500    ///
501    /// A copy created by this operation is removed when append fails before
502    /// publication. Files already under the table root and copies involved in
503    /// an ambiguous commit outcome are preserved.
504    /// Returns the committed version and normalized table-relative segment path.
505    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    /// Copy and append a Parquet file while collecting a profiling report.
513    /// Returns the committed version, normalized table-relative path, and report.
514    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    /// Append a Parquet segment and return a profiling report.
526    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        // Removing the file proves duplicate detection depends only on the
1648        // normalized live identity, not filesystem or Parquet inspection.
1649        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        // Append two segments so we have segment sidecars plus a table snapshot pointer.
1905        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        // Overwrite table state with the modified snapshot missing coverage_path.
2102        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        // Simulate legacy/bad state: drop coverage_path on the existing segment.
2404        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    // Unlike load_snapshot_recovers_when_missing_file (which exercises recovery when
2421    // the pointer exists but the snapshot file is gone), this covers the case where
2422    // the in-memory pointer itself is missing while segments exist, and append
2423    // must rebuild + rewrite the pointer as part of the append flow.
2424    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        // Simulate missing snapshot pointer while segments exist.
2458        table.state.table_coverage = None;
2459
2460        table.append_parquet_segment(rel2).await?;
2461
2462        // Snapshot pointer should be restored after a successful append.
2463        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        // Tamper snapshot pointer to a mismatching bucket spec.
2517        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}