Skip to main content

timeseries_table_format/table/operations/
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//! - adopting the first schema or normalizing into the registered schema,
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
11pub(super) mod error;
12
13use error::AppendError;
14
15use std::{marker::PhantomData, path::Path, sync::Arc};
16
17use arrow::{
18    array::{RecordBatch as ArrowRecordBatch, RecordBatchIterator, RecordBatchReader},
19    datatypes::SchemaRef,
20    error::ArrowError,
21};
22use parquet::arrow::ArrowWriter as ParquetArrowWriter;
23use parquet::basic::Compression;
24use parquet::file::properties::WriterProperties;
25use snafu::prelude::*;
26use uuid::Uuid;
27
28use crate::{
29    coverage::serde::{coverage_to_bytes, entity_coverage_to_bytes},
30    coverage::{
31        EntityCoverage,
32        index_interval::index_interval_for_id,
33        io::{CoverageSidecarError, write_coverage_sidecar_new_bytes},
34        layout::{
35            coverage_file_id_for_attempt, segment_coverage_id_v2, segment_coverage_key,
36            segment_entity_coverage_id_v1, table_coverage_id_v2, table_entity_coverage_id_v1,
37            table_snapshot_key,
38        },
39    },
40    formats::parquet::{
41        compute_segment_entity_coverage, coverage::compute_segment_coverage,
42        logical_schema_from_parquet, segment_meta::segment_meta_from_parquet,
43    },
44    metadata::{
45        logical_schema::LogicalSchema,
46        schema_compat::{ensure_index_spec_matches_schema, ensure_schema_fields_match_by_name},
47        segments::SegmentEntityLayout,
48    },
49    storage,
50    table::{TableError, TimeSeriesTable},
51    transaction_log::{
52        LogAction, TableState, checked_next_version, table_state::TableCoveragePointer,
53    },
54};
55
56use self::error::{
57    ArrowInputSnafu, ArrowToLogicalSchemaSnafu, GeneratedSegmentSchemaCompatibilitySnafu,
58    ParquetWriteSnafu,
59};
60use super::append_schema::AppendSchemaNormalizer;
61
62fn classify_entity_layout(
63    segment_path: &str,
64    coverage: &EntityCoverage,
65) -> Result<SegmentEntityLayout, AppendError> {
66    let first_identity = coverage
67        .iter()
68        .next()
69        .map(|(identity, _)| identity)
70        .ok_or_else(|| AppendError::EmptySegmentEntityCoverage {
71            segment_path: segment_path.to_string(),
72        })?;
73
74    if let Some((identity, _)) = coverage.iter().find(|(_, coverage)| coverage.is_empty()) {
75        return Err(AppendError::EntityWithoutIndexCoverage {
76            segment_path: segment_path.to_string(),
77            identity: identity.clone(),
78        });
79    }
80
81    Ok(if coverage.identity_count() == 1 {
82        SegmentEntityLayout::Single(first_identity.clone())
83    } else {
84        SegmentEntityLayout::Mixed
85    })
86}
87
88fn ensure_existing_segments_have_coverage(state: &TableState) -> Result<(), AppendError> {
89    for seg in state.segments.values() {
90        if seg.coverage_path.is_none() {
91            return Err(AppendError::ExistingSegmentMissingCoverageMetadata {
92                segment_path: seg.path.clone(),
93            });
94        }
95    }
96
97    Ok(())
98}
99
100fn record_append_failure<T>(result: &Result<T, TableError>) {
101    if let Err(error) = result {
102        let outcome = if matches!(
103            error,
104            TableError::Append {
105                source: AppendError::CommitAmbiguous { .. }
106            }
107        ) {
108            "ambiguous"
109        } else {
110            "failed"
111        };
112        tracing::Span::current().record("outcome", outcome);
113    }
114}
115
116/// Marker used to distinguish reader inputs from materialized batch inputs.
117///
118/// This type exists only to keep the blanket [`RecordBatchReader`]
119/// implementation disjoint from the direct batch implementations.
120#[doc(hidden)]
121pub struct RecordBatchReaderSourceKind;
122
123/// Compression used for Parquet segments written by append.
124#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
125#[non_exhaustive]
126pub enum ParquetCompression {
127    /// Write Parquet pages without compression.
128    Uncompressed,
129    /// Use Snappy compression.
130    Snappy,
131    /// Use Zstandard compression with its standard level.
132    #[default]
133    Zstd,
134}
135
136impl From<ParquetCompression> for Compression {
137    fn from(value: ParquetCompression) -> Self {
138        match value {
139            ParquetCompression::Uncompressed => Self::UNCOMPRESSED,
140            ParquetCompression::Snappy => Self::SNAPPY,
141            ParquetCompression::Zstd => Self::ZSTD(Default::default()),
142        }
143    }
144}
145
146const DEFAULT_MAX_ROWS_PER_ROW_GROUP: usize = 1024 * 1024;
147const DEFAULT_MAX_BYTES_PER_ROW_GROUP: usize = 128 * 1024 * 1024;
148
149#[derive(Debug, Clone, Copy, PartialEq, Eq)]
150struct AppendWriterSettings {
151    compression: ParquetCompression,
152    max_rows_per_row_group: usize,
153    max_bytes_per_row_group: usize,
154}
155
156impl Default for AppendWriterSettings {
157    fn default() -> Self {
158        Self {
159            compression: ParquetCompression::default(),
160            max_rows_per_row_group: DEFAULT_MAX_ROWS_PER_ROW_GROUP,
161            max_bytes_per_row_group: DEFAULT_MAX_BYTES_PER_ROW_GROUP,
162        }
163    }
164}
165
166impl AppendWriterSettings {
167    fn from_source<S, SourceKind>(source: &S) -> Self
168    where
169        S: IntoRecordBatchReader<SourceKind>,
170    {
171        Self {
172            compression: source.effective_compression().unwrap_or_default(),
173            max_rows_per_row_group: source
174                .effective_max_rows_per_row_group()
175                .unwrap_or(DEFAULT_MAX_ROWS_PER_ROW_GROUP),
176            max_bytes_per_row_group: source
177                .effective_max_bytes_per_row_group()
178                .unwrap_or(DEFAULT_MAX_BYTES_PER_ROW_GROUP),
179        }
180    }
181
182    fn validate(self) -> Result<Self, AppendError> {
183        if self.max_rows_per_row_group == 0 {
184            return Err(AppendError::InvalidMaxRowsPerRowGroup {
185                max_rows_per_row_group: 0,
186            });
187        }
188        if self.max_bytes_per_row_group == 0 {
189            return Err(AppendError::InvalidMaxBytesPerRowGroup {
190                max_bytes_per_row_group: 0,
191            });
192        }
193        Ok(self)
194    }
195
196    fn writer_properties(self) -> WriterProperties {
197        WriterProperties::builder()
198            .set_compression(self.compression.into())
199            .set_max_row_group_row_count(Some(self.max_rows_per_row_group))
200            .set_max_row_group_bytes(Some(self.max_bytes_per_row_group))
201            .build()
202    }
203}
204
205/// Convert an Arrow batch source into a schema-bearing [`RecordBatchReader`].
206///
207/// Implementations preserve non-`Send` readers. The `SourceKind` parameter is an
208/// inference-only coherence marker and should not normally be specified by
209/// callers.
210pub trait IntoRecordBatchReader<SourceKind = RecordBatchReaderSourceKind> {
211    /// Reader produced from this source.
212    type Reader: RecordBatchReader;
213
214    /// Return a per-append Parquet compression override, when configured.
215    #[doc(hidden)]
216    fn effective_compression(&self) -> Option<ParquetCompression> {
217        None
218    }
219
220    /// Return a per-append Parquet row-group limit, when configured.
221    #[doc(hidden)]
222    fn effective_max_rows_per_row_group(&self) -> Option<usize> {
223        None
224    }
225
226    /// Return a per-append Parquet row-group byte limit, when configured.
227    #[doc(hidden)]
228    fn effective_max_bytes_per_row_group(&self) -> Option<usize> {
229        None
230    }
231
232    /// Convert this source without collecting its batches.
233    fn into_record_batch_reader(self) -> Result<Self::Reader, AppendError>;
234}
235
236/// One Arrow source plus physical settings for a single append.
237///
238/// Pass the source directly to [`TimeSeriesTable::append`] to use Zstandard
239/// compression and row groups bounded by 1,048,576 rows and 128 MiB estimated
240/// encoded bytes. Wrap it in `AppendRequest` only when one append needs an
241/// explicit physical override. A nested request inherits each setting from its
242/// source unless the outer request replaces it.
243#[derive(Debug)]
244#[must_use = "an append request has no effect until passed to TimeSeriesTable::append"]
245pub struct AppendRequest<S> {
246    source: S,
247    compression: Option<ParquetCompression>,
248    max_rows_per_row_group: Option<usize>,
249    max_bytes_per_row_group: Option<usize>,
250}
251
252impl<S> AppendRequest<S> {
253    /// Create a request without adding physical writer overrides.
254    ///
255    /// If `source` already carries overrides, this request inherits each one
256    /// until the corresponding builder method replaces it.
257    pub fn new(source: S) -> Self {
258        Self {
259            source,
260            compression: None,
261            max_rows_per_row_group: None,
262            max_bytes_per_row_group: None,
263        }
264    }
265
266    /// Set Parquet compression for this append.
267    pub fn compression(mut self, compression: ParquetCompression) -> Self {
268        self.compression = Some(compression);
269        self
270    }
271
272    /// Limit output Parquet row groups to this many rows for this append.
273    ///
274    /// This controls rows, not bytes, input batches, or source-file row-group
275    /// boundaries, and replaces any limit already carried by the source. Zero
276    /// is rejected by [`TimeSeriesTable::append`] before the source is inspected
277    /// or consumed.
278    pub fn max_rows_per_row_group(mut self, max_rows_per_row_group: usize) -> Self {
279        self.max_rows_per_row_group = Some(max_rows_per_row_group);
280        self
281    }
282
283    /// Limit output Parquet row groups to this many estimated encoded bytes.
284    ///
285    /// This is not a strict process-memory ceiling, and a single oversized value
286    /// may exceed it. Zero is rejected by [`TimeSeriesTable::append`] before the
287    /// source is inspected or consumed.
288    pub fn max_bytes_per_row_group(mut self, max_bytes_per_row_group: usize) -> Self {
289        self.max_bytes_per_row_group = Some(max_bytes_per_row_group);
290        self
291    }
292}
293
294// Wrapping `SourceKind` in `PhantomData` keeps this impl disjoint from direct source
295// implementations while preserving inference without another public marker type.
296#[doc(hidden)]
297impl<S, SourceKind> IntoRecordBatchReader<PhantomData<SourceKind>> for AppendRequest<S>
298where
299    S: IntoRecordBatchReader<SourceKind>,
300{
301    type Reader = S::Reader;
302
303    fn effective_compression(&self) -> Option<ParquetCompression> {
304        self.compression
305            .or_else(|| self.source.effective_compression())
306    }
307
308    fn effective_max_rows_per_row_group(&self) -> Option<usize> {
309        self.max_rows_per_row_group
310            .or_else(|| self.source.effective_max_rows_per_row_group())
311    }
312
313    fn effective_max_bytes_per_row_group(&self) -> Option<usize> {
314        self.max_bytes_per_row_group
315            .or_else(|| self.source.effective_max_bytes_per_row_group())
316    }
317
318    fn into_record_batch_reader(self) -> Result<Self::Reader, AppendError> {
319        self.source.into_record_batch_reader()
320    }
321}
322
323impl<R> IntoRecordBatchReader<RecordBatchReaderSourceKind> for R
324where
325    R: RecordBatchReader,
326{
327    type Reader = R;
328
329    fn into_record_batch_reader(self) -> Result<Self::Reader, AppendError> {
330        Ok(self)
331    }
332}
333
334impl IntoRecordBatchReader<ArrowRecordBatch> for ArrowRecordBatch {
335    type Reader = Box<dyn RecordBatchReader + Send>;
336
337    fn into_record_batch_reader(self) -> Result<Self::Reader, AppendError> {
338        let schema = self.schema();
339        Ok(Box::new(RecordBatchIterator::new(
340            std::iter::once(Ok(self)),
341            schema,
342        )))
343    }
344}
345
346impl IntoRecordBatchReader<Vec<ArrowRecordBatch>> for Vec<ArrowRecordBatch> {
347    type Reader = Box<dyn RecordBatchReader + Send>;
348
349    fn into_record_batch_reader(self) -> Result<Self::Reader, AppendError> {
350        let schema = self
351            .first()
352            .map(ArrowRecordBatch::schema)
353            .ok_or(AppendError::EmptyInput)?;
354        Ok(Box::new(RecordBatchIterator::new(
355            self.into_iter().map(Ok::<_, ArrowError>),
356            schema,
357        )))
358    }
359}
360
361fn ensure_batch_matches_reader_schema(
362    reader_schema: &SchemaRef,
363    batch: &ArrowRecordBatch,
364) -> Result<(), AppendError> {
365    if batch.schema() == *reader_schema {
366        Ok(())
367    } else {
368        Err(ArrowError::SchemaError(
369            "record batch schema does not match its reader schema".to_string(),
370        ))
371        .context(ArrowInputSnafu)
372    }
373}
374
375impl TimeSeriesTable {
376    async fn rollback_created_artifacts(
377        &self,
378        created_paths: &[String],
379        source: AppendError,
380    ) -> AppendError {
381        let (source, mut cleanup_errors) = match source {
382            AppendError::Rollback {
383                source,
384                cleanup_errors,
385            } => (source, cleanup_errors),
386            source => (Box::new(source), Vec::new()),
387        };
388        for path in created_paths.iter().rev() {
389            if let Err(error) =
390                storage::remove_file_if_exists(self.location().as_ref(), Path::new(path)).await
391            {
392                cleanup_errors.push(error);
393            }
394        }
395
396        if cleanup_errors.is_empty() {
397            *source
398        } else {
399            AppendError::Rollback {
400                source,
401                cleanup_errors,
402            }
403        }
404    }
405
406    fn build_append_schema_normalizer(
407        &self,
408        incoming_schema: SchemaRef,
409    ) -> Result<AppendSchemaNormalizer, AppendError> {
410        ensure_existing_segments_have_coverage(&self.state)?;
411
412        match self.state.table_meta.logical_schema.as_ref() {
413            None if self.state.version == 1 => {
414                let incoming_logical_schema =
415                    LogicalSchema::try_from_arrow_schema(incoming_schema.as_ref())
416                        .context(ArrowToLogicalSchemaSnafu)?;
417                ensure_index_spec_matches_schema(&incoming_logical_schema, &self.index)
418                    .map_err(AppendError::from)?;
419                Ok(AppendSchemaNormalizer::without_conversion(incoming_schema))
420            }
421            None => Err(AppendError::MissingCanonicalTableSchema {
422                version: self.state.version,
423            }),
424            Some(table_schema) => {
425                ensure_index_spec_matches_schema(table_schema, &self.index)
426                    .map_err(AppendError::from)?;
427                let normalizer = AppendSchemaNormalizer::for_registered_schema(
428                    incoming_schema.as_ref(),
429                    table_schema,
430                )
431                .map_err(AppendError::from)?;
432                Ok(normalizer)
433            }
434        }
435    }
436
437    async fn publish_generated_parquet_segment(
438        &mut self,
439        relative_path: &str,
440        next_version: u64,
441        owned_data_guard: &mut storage::FileCleanupGuard,
442    ) -> Result<u64, AppendError> {
443        let rel_path = Path::new(relative_path);
444        let expected_version = self.state.version;
445
446        // 0) Coverage readiness checks.
447        ensure_existing_segments_have_coverage(&self.state)?;
448
449        // 1) Segment meta + schema.
450        let (mut segment_meta, _) =
451            segment_meta_from_parquet(self.location(), rel_path, &self.index)
452                .await
453                .map_err(AppendError::from)?;
454        let row_count = segment_meta.row_count;
455        let span = tracing::Span::current();
456        span.record("row_count", row_count);
457        if let Some(file_size) = segment_meta.file_size {
458            span.record("file_size_bytes", file_size);
459        }
460
461        let segment_schema = logical_schema_from_parquet(self.location(), rel_path)
462            .await
463            .map_err(AppendError::from)?;
464        ensure_index_spec_matches_schema(&segment_schema, &self.index).context(
465            GeneratedSegmentSchemaCompatibilitySnafu {
466                segment_path: relative_path.to_string(),
467            },
468        )?;
469
470        // 2) Schema behavior (return maybe_updated_meta, but do NOT build actions yet).
471        //
472        // - logical_schema == None && version == 1:
473        //     first append after create(): adopt this segment's schema.
474        // - logical_schema == None && version != 1:
475        //     table is in a bad state for v0.1: return an error.
476        // - logical_schema == Some(..): enforce no schema evolution by field name.
477        let maybe_table_schema = self.state.table_meta.logical_schema.as_ref();
478
479        let maybe_updated_meta = match maybe_table_schema {
480            None if expected_version == 1 => {
481                let mut updated_meta = self.state.table_meta.clone();
482                updated_meta.logical_schema = Some(segment_schema.clone());
483                Some(updated_meta)
484            }
485            None => {
486                return Err(AppendError::MissingCanonicalTableSchema {
487                    version: expected_version,
488                });
489            }
490            Some(table_schema) => {
491                ensure_schema_fields_match_by_name(table_schema, &segment_schema, &self.index)
492                    .context(GeneratedSegmentSchemaCompatibilitySnafu {
493                        segment_path: relative_path.to_string(),
494                    })?;
495                None
496            }
497        };
498
499        let has_entity_columns = !self.index.entity_columns.is_empty();
500
501        // 3-5) Load, compute, and compare coverage using the entity-column mode.
502        let (seg_cov_bytes, new_snap_cov_bytes, entity_layout) = if has_entity_columns {
503            let table_cov = self
504                .load_entity_coverage_with_recovery::<AppendError>()
505                .await?;
506
507            let segment_cov =
508                compute_segment_entity_coverage(self.location(), rel_path, &self.index)
509                    .await
510                    .map_err(AppendError::from)?;
511            let entity_layout = classify_entity_layout(relative_path, &segment_cov)?;
512
513            if let Some((identity, index_interval_id)) =
514                segment_cov.first_overlapping_identity_and_interval_id(&table_cov)
515            {
516                let example_index_interval =
517                    index_interval_for_id(&self.index.kind, index_interval_id)
518                        .map_err(AppendError::from)?;
519                return Err(AppendError::PersistedIndexIntervalOverlap {
520                    segment_path: relative_path.to_string(),
521                    overlap_count: segment_cov.intersection_cardinality(&table_cov),
522                    example_identity: Some(identity.clone()),
523                    example_index_interval_id: index_interval_id,
524                    example_index_interval: Box::new(example_index_interval),
525                });
526            }
527
528            let seg_bytes = entity_coverage_to_bytes(&segment_cov)
529                .map_err(CoverageSidecarError::from)
530                .map_err(AppendError::from)?;
531            let snapshot_bytes = entity_coverage_to_bytes(&table_cov.union(&segment_cov))
532                .map_err(CoverageSidecarError::from)
533                .map_err(AppendError::from)?;
534            (seg_bytes, snapshot_bytes, entity_layout)
535        } else {
536            let table_cov = self
537                .load_global_coverage_with_recovery::<AppendError>()
538                .await?;
539
540            let segment_cov = compute_segment_coverage(self.location(), rel_path, &self.index)
541                .await
542                .map_err(AppendError::from)?;
543
544            let overlap = segment_cov.intersect(&table_cov);
545            let overlap_count = overlap.cardinality();
546            if let Some(example_index_interval_id) = overlap.present().iter().next() {
547                let example_index_interval =
548                    index_interval_for_id(&self.index.kind, example_index_interval_id)
549                        .map_err(AppendError::from)?;
550                return Err(AppendError::PersistedIndexIntervalOverlap {
551                    segment_path: relative_path.to_string(),
552                    overlap_count: u128::from(overlap_count),
553                    example_identity: None,
554                    example_index_interval_id,
555                    example_index_interval: Box::new(example_index_interval),
556                });
557            }
558
559            let seg_bytes = coverage_to_bytes(&segment_cov)
560                .map_err(CoverageSidecarError::from)
561                .map_err(AppendError::from)?;
562            let snapshot_bytes = coverage_to_bytes(&table_cov.union(&segment_cov))
563                .map_err(CoverageSidecarError::from)
564                .map_err(AppendError::from)?;
565            (
566                seg_bytes,
567                snapshot_bytes,
568                SegmentEntityLayout::NotApplicable,
569            )
570        };
571        let entity_layout_name = match &entity_layout {
572            SegmentEntityLayout::NotApplicable => "not_applicable",
573            SegmentEntityLayout::Single(_) => "single",
574            SegmentEntityLayout::Mixed => "mixed",
575        };
576        span.record("entity_layout", entity_layout_name);
577
578        // 6) Give this append private sidecar paths, then write them before commit.
579        let attempt_id = Uuid::new_v4();
580        let segment_content_id = if has_entity_columns {
581            segment_entity_coverage_id_v1(&self.index, &seg_cov_bytes)
582        } else {
583            segment_coverage_id_v2(&self.index, &seg_cov_bytes)
584        };
585        let segment_file_id = coverage_file_id_for_attempt(&segment_content_id, &attempt_id);
586        let seg_cov_path = segment_coverage_key(&segment_file_id)
587            .map_err(CoverageSidecarError::from)
588            .map_err(AppendError::from)?;
589
590        let snapshot_content_id = if has_entity_columns {
591            table_entity_coverage_id_v1(&self.index, &new_snap_cov_bytes)
592        } else {
593            table_coverage_id_v2(&self.index, &new_snap_cov_bytes)
594        };
595        let snapshot_file_id = coverage_file_id_for_attempt(&snapshot_content_id, &attempt_id);
596        let snapshot_path = table_snapshot_key(next_version, &snapshot_file_id)
597            .map_err(CoverageSidecarError::from)
598            .map_err(AppendError::from)?;
599
600        let mut created_sidecars = Vec::new();
601        let mut segment_sidecar_guard = storage::FileCleanupGuard::new_disarmed(
602            self.location().as_ref(),
603            Path::new(&seg_cov_path),
604        )
605        .map_err(AppendError::from)?;
606        if let Err(source) = write_coverage_sidecar_new_bytes(
607            self.location(),
608            Path::new(&seg_cov_path),
609            &seg_cov_bytes,
610        )
611        .await
612        {
613            if source.storage_cleanup_failed() {
614                created_sidecars.push(seg_cov_path.clone());
615            }
616            let error = self
617                .rollback_created_artifacts(&created_sidecars, AppendError::from(source))
618                .await;
619            return Err(error);
620        }
621        segment_sidecar_guard.arm();
622        created_sidecars.push(seg_cov_path.clone());
623
624        let mut snapshot_sidecar_guard = storage::FileCleanupGuard::new_disarmed(
625            self.location().as_ref(),
626            Path::new(&snapshot_path),
627        )
628        .map_err(AppendError::from)?;
629        if let Err(source) = write_coverage_sidecar_new_bytes(
630            self.location(),
631            Path::new(&snapshot_path),
632            &new_snap_cov_bytes,
633        )
634        .await
635        {
636            if source.storage_cleanup_failed() {
637                created_sidecars.push(snapshot_path.clone());
638            }
639            let error = AppendError::from(source);
640            let error = self
641                .rollback_created_artifacts(&created_sidecars, error)
642                .await;
643            segment_sidecar_guard.disarm();
644            return Err(error);
645        }
646        snapshot_sidecar_guard.arm();
647        created_sidecars.push(snapshot_path.clone());
648
649        // 7) Build actions and atomically publish the commit.
650        segment_meta.coverage_path = Some(seg_cov_path);
651        segment_meta.entity_layout = entity_layout;
652
653        let mut actions = Vec::new();
654        if let Some(updated_meta) = maybe_updated_meta.clone() {
655            actions.push(LogAction::UpdateTableMeta(updated_meta));
656        }
657
658        actions.push(LogAction::AddSegment(segment_meta.clone()));
659        actions.push(LogAction::UpdateTableCoverage {
660            index_kind: self.index.kind.clone(),
661            coverage_path: snapshot_path.clone(),
662        });
663
664        let new_version = match self
665            .log
666            .commit_with_path_preservation(expected_version, actions, || {
667                owned_data_guard.disarm();
668                segment_sidecar_guard.disarm();
669                snapshot_sidecar_guard.disarm();
670            })
671            .await
672        {
673            Ok(version) => version,
674            Err(source @ crate::transaction_log::CommitError::AmbiguousOutcome { .. }) => {
675                return Err(AppendError::CommitAmbiguous {
676                    segment_path: relative_path.to_string(),
677                    source: Box::new(source),
678                });
679            }
680            Err(source) => {
681                let error = AppendError::from(source);
682                let error = self
683                    .rollback_created_artifacts(&created_sidecars, error)
684                    .await;
685                segment_sidecar_guard.disarm();
686                snapshot_sidecar_guard.disarm();
687                return Err(error);
688            }
689        };
690
691        // OCC invariant: a successful transaction commit must return
692        // the same "next" version we predicted when constructing `snapshot_path`.
693        // If this ever diverges, it indicates a severe bug between snapshot path
694        // construction and the transaction log implementation, so we panic rather
695        // than continuing with an inconsistent in-memory state.
696        assert_eq!(
697            new_version, next_version,
698            "transaction log returned unexpected version: expected {}, got {}",
699            next_version, new_version
700        );
701
702        // 8) Update in-memory state.
703        self.state.version = new_version;
704
705        if let Some(updated_meta) = maybe_updated_meta {
706            self.state.table_meta = updated_meta
707        }
708
709        self.state
710            .segments
711            .insert(segment_meta.path.clone(), segment_meta);
712
713        // Also update the snapshot pointer in state.
714        self.state.table_coverage = Some(TableCoveragePointer {
715            index_kind: self.index.kind.clone(),
716            coverage_path: snapshot_path,
717            version: new_version,
718        });
719
720        span.record("committed_version", new_version);
721        span.record("outcome", "succeeded");
722        tracing::info!(
723            name: "table.append",
724            target: "timeseries_table_format::table::append",
725            expected_version,
726            committed_version = new_version,
727            row_count,
728            entity_layout = entity_layout_name,
729            outcome = "succeeded",
730            "Appended Parquet segment"
731        );
732        Ok(new_version)
733    }
734
735    /// Check whether this client supports appending to the current table.
736    ///
737    /// Higher-level wrappers can call this before inspecting or converting an
738    /// append input. [`Self::append`] repeats the check before writing.
739    pub fn ensure_append_supported(&self) -> Result<(), TableError> {
740        self.ensure_write_compatible()
741            .map_err(AppendError::from)
742            .context(crate::table::error::AppendSnafu)
743    }
744
745    /// Append Arrow record batches into one table-managed Parquet segment.
746    ///
747    /// Rows need not be ordered by the table's ordered index. The source is
748    /// consumed incrementally and is never collected by this method.
749    /// When a registered schema exists, incoming fields are matched by name
750    /// and written in registered order. Exact types and these lossless scalar
751    /// widenings are accepted: `Int8 -> Int32/Int64`, `Int16 -> Int32/Int64`,
752    /// `Int32 -> Int64`, `UInt8/UInt16/UInt32 -> UInt64`, and
753    /// `Float32 -> Float64`.
754    /// New segments use Zstandard compression and row groups bounded by
755    /// 1,048,576 rows and 128 MiB estimated encoded bytes. Wrap the source in
756    /// [`AppendRequest`] to override those physical settings for only this
757    /// append.
758    #[tracing::instrument(
759        name = "table.append",
760        target = "timeseries_table_format::table::append",
761        level = "debug",
762        skip_all,
763        fields(
764            expected_version = self.state.version,
765            segment_path = tracing::field::Empty,
766            row_count = tracing::field::Empty,
767            file_size_bytes = tracing::field::Empty,
768            committed_version = tracing::field::Empty,
769            entity_layout = tracing::field::Empty,
770            outcome = tracing::field::Empty
771        )
772    )]
773    pub async fn append<S, SourceKind>(&mut self, source: S) -> Result<u64, TableError>
774    where
775        S: IntoRecordBatchReader<SourceKind>,
776    {
777        let append_result: Result<u64, AppendError> = async {
778            self.ensure_write_compatible().map_err(AppendError::from)?;
779
780            let writer_settings = AppendWriterSettings::from_source(&source).validate()?;
781            let next_version =
782                checked_next_version(self.state.version).map_err(AppendError::from)?;
783            let mut reader = source.into_record_batch_reader()?;
784            let incoming_schema = reader.schema();
785            let schema_normalizer =
786                self.build_append_schema_normalizer(Arc::clone(&incoming_schema))?;
787            let output_schema = Arc::clone(schema_normalizer.output_schema());
788
789            let first_batch = loop {
790                let Some(batch) = reader.next().transpose().context(ArrowInputSnafu)? else {
791                    return Err(AppendError::EmptyInput);
792                };
793                ensure_batch_matches_reader_schema(&incoming_schema, &batch)?;
794                if batch.num_rows() != 0 {
795                    break schema_normalizer
796                        .normalize_batch(&batch)
797                        .context(ArrowInputSnafu)?;
798                }
799                tokio::task::yield_now().await;
800            };
801
802            let relative_path = format!("data/{}.parquet", Uuid::new_v4());
803            tracing::Span::current().record("segment_path", relative_path.as_str());
804            let mut data_guard = storage::FileCleanupGuard::new_disarmed(
805                self.location().as_ref(),
806                Path::new(&relative_path),
807            )
808            .map_err(AppendError::from)?;
809            let sink =
810                storage::open_new_output_sink(self.location().as_ref(), Path::new(&relative_path))
811                    .await
812                    .map_err(AppendError::from)?;
813            let write_result = async {
814                let writer_properties = writer_settings.writer_properties();
815                let mut writer =
816                    ParquetArrowWriter::try_new(sink, output_schema, Some(writer_properties))
817                        .context(ParquetWriteSnafu)?;
818                writer.write(&first_batch).context(ParquetWriteSnafu)?;
819                drop(first_batch);
820                tokio::task::yield_now().await;
821
822                for batch in reader {
823                    let batch = batch.context(ArrowInputSnafu)?;
824                    ensure_batch_matches_reader_schema(&incoming_schema, &batch)?;
825                    if batch.num_rows() != 0 {
826                        let batch = schema_normalizer
827                            .normalize_batch(&batch)
828                            .context(ArrowInputSnafu)?;
829                        writer.write(&batch).context(ParquetWriteSnafu)?;
830                    }
831                    drop(batch);
832                    tokio::task::yield_now().await;
833                }
834
835                let sink = writer.into_inner().context(ParquetWriteSnafu)?;
836                sink.finish().await.map_err(AppendError::from)
837            }
838            .await;
839            if let Err(source) = write_result {
840                let error = self
841                    .rollback_created_artifacts(std::slice::from_ref(&relative_path), source)
842                    .await;
843                data_guard.disarm();
844                return Err(error);
845            }
846            data_guard.arm();
847
848            match self
849                .publish_generated_parquet_segment(&relative_path, next_version, &mut data_guard)
850                .await
851            {
852                Ok(version) => Ok(version),
853                Err(source @ AppendError::CommitAmbiguous { .. }) => Err(source),
854                Err(source) => {
855                    let error = self
856                        .rollback_created_artifacts(std::slice::from_ref(&relative_path), source)
857                        .await;
858                    data_guard.disarm();
859                    Err(error)
860                }
861            }
862        }
863        .await;
864        let result = append_result.context(crate::table::error::AppendSnafu);
865        record_append_failure(&result);
866        result
867    }
868}
869
870#[cfg(test)]
871mod tests {
872    use super::*;
873    use snafu::IntoError;
874
875    use crate::coverage::io::{
876        CoverageSidecarError, read_coverage_sidecar, read_entity_coverage_sidecar,
877    };
878    use crate::coverage::layout::{SEGMENT_COVERAGE_DIR, TABLE_SNAPSHOT_DIR};
879    use crate::coverage::serde::entity_coverage_from_bytes;
880    use crate::coverage::{EntityCoverage, EntityIdentity, EntityValue};
881    use crate::formats::parquet::SegmentCoverageError;
882    use crate::metadata::logical_schema::{
883        LogicalDataType, LogicalField, LogicalSchema, LogicalTimestampUnit,
884    };
885    use crate::metadata::segments::SegmentEntityLayout;
886    use crate::metadata::{index::IndexValue, protocol::TableProtocolError};
887    use crate::storage::layout;
888    use crate::storage::{StorageError, StorageLocation, TableLocation};
889    use crate::table::test_util::*;
890    use crate::transaction_log::{
891        CommitError, IndexKind, IndexSpec, TableMeta, TimeIndexGranularity, segments::SegmentError,
892    };
893    use arrow::{
894        array::{
895            Array, ArrayRef, BinaryArray, BooleanArray, Float32Array, Float64Array, Int16Array,
896            Int32Array, Int64Array, StringArray, TimestampMillisecondArray, UInt32Array,
897            UInt64Array, new_null_array,
898        },
899        datatypes::{DataType, Field, Fields, Schema, TimeUnit as ArrowTimeUnit},
900        record_batch::RecordBatch,
901    };
902    use futures::{FutureExt, StreamExt};
903    use parquet::arrow::{ArrowWriter, arrow_reader::ParquetRecordBatchReaderBuilder};
904    use snafu::ErrorCompat;
905    use std::cell::Cell;
906    use std::collections::{BTreeMap, HashMap};
907    use std::fs::File;
908    use std::num::NonZeroU64;
909    use std::path::PathBuf;
910    use std::rc::Rc;
911    use std::sync::{Arc, Weak};
912    use tempfile::TempDir;
913    use tracing::{Subscriber, instrument::WithSubscriber, span::Id};
914    use tracing_subscriber::{
915        Layer,
916        layer::{Context, SubscriberExt},
917        registry::LookupSpan,
918    };
919
920    #[derive(Clone, Copy)]
921    struct PanicOnCommitClose;
922
923    impl<S> Layer<S> for PanicOnCommitClose
924    where
925        S: Subscriber + for<'lookup> LookupSpan<'lookup>,
926    {
927        fn on_close(&self, id: Id, ctx: Context<'_, S>) {
928            if ctx
929                .metadata(&id)
930                .is_some_and(|metadata| metadata.name() == "transaction.commit")
931            {
932                panic!("injected transaction commit close panic");
933            }
934        }
935    }
936
937    fn panic_on_commit_close_dispatch() -> tracing::Dispatch {
938        tracing::Dispatch::new(tracing_subscriber::registry().with(PanicOnCommitClose))
939    }
940
941    #[derive(Default)]
942    struct ReaderObservations {
943        schema_calls: Cell<usize>,
944        next_calls: Cell<usize>,
945        next_before_schema: Cell<bool>,
946        previous_batch_alive: Cell<bool>,
947    }
948
949    struct InstrumentedReader {
950        schema: SchemaRef,
951        batches: std::vec::IntoIter<Result<RecordBatch, ArrowError>>,
952        observations: Rc<ReaderObservations>,
953        previous_array: Option<Weak<dyn arrow::array::Array>>,
954    }
955
956    impl InstrumentedReader {
957        fn new(
958            schema: SchemaRef,
959            batches: Vec<Result<RecordBatch, ArrowError>>,
960        ) -> (Self, Rc<ReaderObservations>) {
961            let observations = Rc::new(ReaderObservations::default());
962            (
963                Self {
964                    schema,
965                    batches: batches.into_iter(),
966                    observations: Rc::clone(&observations),
967                    previous_array: None,
968                },
969                observations,
970            )
971        }
972
973        fn one(batch: RecordBatch) -> Self {
974            let (reader, _) = Self::new(batch.schema(), vec![Ok(batch)]);
975            reader
976        }
977    }
978
979    impl Iterator for InstrumentedReader {
980        type Item = Result<RecordBatch, ArrowError>;
981
982        fn next(&mut self) -> Option<Self::Item> {
983            self.observations
984                .next_calls
985                .set(self.observations.next_calls.get() + 1);
986            if self.observations.schema_calls.get() == 0 {
987                self.observations.next_before_schema.set(true);
988            }
989            if self
990                .previous_array
991                .as_ref()
992                .is_some_and(|array| array.strong_count() != 0)
993            {
994                self.observations.previous_batch_alive.set(true);
995            }
996
997            let next = self.batches.next();
998            if let Some(Ok(batch)) = &next {
999                self.previous_array = Some(Arc::downgrade(batch.column(0)));
1000            }
1001            next
1002        }
1003    }
1004
1005    impl RecordBatchReader for InstrumentedReader {
1006        fn schema(&self) -> arrow::datatypes::SchemaRef {
1007            self.observations
1008                .schema_calls
1009                .set(self.observations.schema_calls.get() + 1);
1010            Arc::clone(&self.schema)
1011        }
1012    }
1013
1014    fn input_batch(values: Vec<i64>) -> Result<RecordBatch, ArrowError> {
1015        let schema = Arc::new(Schema::new(vec![Field::new(
1016            "value",
1017            DataType::Int64,
1018            false,
1019        )]));
1020        RecordBatch::try_new(schema, vec![Arc::new(Int64Array::from(values))])
1021    }
1022
1023    fn time_series_batch(
1024        timestamps: Vec<i64>,
1025        symbols: Vec<&str>,
1026        prices: Vec<f64>,
1027    ) -> Result<RecordBatch, ArrowError> {
1028        let schema = Arc::new(Schema::new(vec![
1029            Field::new(
1030                "ts",
1031                DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
1032                false,
1033            ),
1034            Field::new("symbol", DataType::Utf8, false),
1035            Field::new("price", DataType::Float64, false),
1036        ]));
1037        RecordBatch::try_new(
1038            schema,
1039            vec![
1040                Arc::new(TimestampMillisecondArray::from(timestamps)),
1041                Arc::new(StringArray::from(symbols)),
1042                Arc::new(Float64Array::from(prices)),
1043            ],
1044        )
1045    }
1046
1047    fn timestamp_only_index() -> IndexSpec {
1048        IndexSpec {
1049            column: "ts".to_string(),
1050            entity_columns: Vec::new(),
1051            kind: IndexKind::Timestamp {
1052                index_granularity: TimeIndexGranularity::Minutes(1),
1053                timezone: None,
1054            },
1055        }
1056    }
1057
1058    fn timestamp_only_meta() -> TableMeta {
1059        TableMeta::new_time_series_with_schema(
1060            timestamp_only_index(),
1061            LogicalSchema::new(vec![LogicalField {
1062                name: "ts".to_string(),
1063                data_type: LogicalDataType::Timestamp {
1064                    unit: LogicalTimestampUnit::Millis,
1065                    timezone: None,
1066                },
1067                nullable: false,
1068            }])
1069            .expect("valid timestamp-only schema"),
1070        )
1071    }
1072
1073    fn timestamp_only_batch(row_count: usize) -> Result<RecordBatch, ArrowError> {
1074        timestamp_only_batch_starting_at_minute_offset(0, row_count)
1075    }
1076
1077    fn timestamp_only_batch_starting_at_minute_offset(
1078        start_minute_offset: usize,
1079        row_count: usize,
1080    ) -> Result<RecordBatch, ArrowError> {
1081        timestamp_only_batch_with_millis(
1082            (start_minute_offset..start_minute_offset + row_count)
1083                .map(|minute_offset| minute_offset as i64 * 60_000),
1084        )
1085    }
1086
1087    fn timestamp_only_batch_with_millis(
1088        values: impl IntoIterator<Item = i64>,
1089    ) -> Result<RecordBatch, ArrowError> {
1090        let schema = Arc::new(Schema::new(vec![Field::new(
1091            "ts",
1092            DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
1093            false,
1094        )]));
1095        RecordBatch::try_new(
1096            schema,
1097            vec![Arc::new(TimestampMillisecondArray::from_iter_values(
1098                values,
1099            ))],
1100        )
1101    }
1102
1103    fn widening_table_meta() -> TableMeta {
1104        TableMeta::new_time_series_with_schema(
1105            IndexSpec {
1106                column: "seq".to_string(),
1107                entity_columns: vec!["device_id".to_string()],
1108                kind: IndexKind::UInt64 {
1109                    index_granularity: NonZeroU64::new(u64::from(u32::MAX) + 1).unwrap(),
1110                },
1111            },
1112            LogicalSchema::new(vec![
1113                LogicalField {
1114                    name: "seq".to_string(),
1115                    data_type: LogicalDataType::UInt64,
1116                    nullable: false,
1117                },
1118                LogicalField {
1119                    name: "device_id".to_string(),
1120                    data_type: LogicalDataType::Int32,
1121                    nullable: false,
1122                },
1123                LogicalField {
1124                    name: "reading".to_string(),
1125                    data_type: LogicalDataType::Float64,
1126                    nullable: true,
1127                },
1128                LogicalField {
1129                    name: "label".to_string(),
1130                    data_type: LogicalDataType::Utf8,
1131                    nullable: false,
1132                },
1133            ])
1134            .expect("valid widening target schema"),
1135        )
1136    }
1137
1138    fn widening_batch(
1139        seq: Vec<u32>,
1140        device_ids: Vec<i16>,
1141        readings: Vec<Option<f32>>,
1142        labels: Vec<&str>,
1143    ) -> Result<RecordBatch, ArrowError> {
1144        let schema = Arc::new(Schema::new_with_metadata(
1145            vec![
1146                Field::new("label", DataType::Utf8, false),
1147                Field::new("reading", DataType::Float32, true).with_metadata(HashMap::from([(
1148                    "source_metadata".to_string(),
1149                    "ignored".to_string(),
1150                )])),
1151                Field::new("seq", DataType::UInt32, false),
1152                Field::new("device_id", DataType::Int16, false),
1153            ],
1154            HashMap::from([("source_schema_metadata".to_string(), "ignored".to_string())]),
1155        ));
1156        RecordBatch::try_new(
1157            schema,
1158            vec![
1159                Arc::new(StringArray::from(labels)),
1160                Arc::new(Float32Array::from(readings)),
1161                Arc::new(UInt32Array::from(seq)),
1162                Arc::new(Int16Array::from(device_ids)),
1163            ],
1164        )
1165    }
1166
1167    fn declared_schema_test_meta(value_type: LogicalDataType, nullable: bool) -> TableMeta {
1168        TableMeta::new_time_series_with_schema(
1169            IndexSpec {
1170                column: "ts".to_string(),
1171                entity_columns: Vec::new(),
1172                kind: IndexKind::Timestamp {
1173                    index_granularity: TimeIndexGranularity::Minutes(1),
1174                    timezone: None,
1175                },
1176            },
1177            LogicalSchema::new(vec![
1178                LogicalField {
1179                    name: "ts".to_string(),
1180                    data_type: LogicalDataType::Timestamp {
1181                        unit: LogicalTimestampUnit::Millis,
1182                        timezone: None,
1183                    },
1184                    nullable: false,
1185                },
1186                LogicalField {
1187                    name: "value".to_string(),
1188                    data_type: value_type,
1189                    nullable,
1190                },
1191            ])
1192            .expect("valid compatibility target schema"),
1193        )
1194    }
1195
1196    async fn assert_declared_schema_rejected_before_reading(
1197        meta: TableMeta,
1198        schema: SchemaRef,
1199        batches: Vec<Result<RecordBatch, ArrowError>>,
1200        expected_column: &str,
1201    ) -> TestResult {
1202        let temp = TempDir::new()?;
1203        let mut table = TimeSeriesTable::create(TableLocation::local(temp.path()), meta).await?;
1204        let state_before = table.state().clone();
1205        let (reader, observations) = InstrumentedReader::new(schema, batches);
1206
1207        let error = table
1208            .append(reader)
1209            .await
1210            .expect_err("incompatible declared schema must fail");
1211
1212        assert!(matches!(
1213            error,
1214            TableError::Append {
1215                source: AppendError::SchemaValidation { .. }
1216            }
1217        ));
1218        assert!(error.to_string().contains(expected_column));
1219        assert!(observations.schema_calls.get() > 0);
1220        assert_eq!(observations.next_calls.get(), 0);
1221        assert_eq!(table.state(), &state_before);
1222        assert_eq!(table.log.load_current_version().await?, 1);
1223        assert!(data_files(temp.path())?.is_empty());
1224        assert!(coverage_files(temp.path())?.is_empty());
1225        assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
1226        Ok(())
1227    }
1228
1229    #[tokio::test]
1230    async fn lossy_schema_is_rejected_before_reading_or_creating_artifacts() -> TestResult {
1231        let temp = TempDir::new()?;
1232        let mut meta = timestamp_only_meta();
1233        meta.logical_schema = None;
1234        let mut table = TimeSeriesTable::create(TableLocation::local(temp.path()), meta).await?;
1235        let schema = Arc::new(Schema::new(vec![
1236            Field::new(
1237                "ts",
1238                DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
1239                false,
1240            ),
1241            Field::new("value", DataType::Int8, true),
1242        ]));
1243        let (reader, observations) = InstrumentedReader::new(schema, Vec::new());
1244
1245        assert!(matches!(
1246            table.append(reader).await,
1247            Err(TableError::Append {
1248                source: AppendError::ArrowToLogicalSchema { .. }
1249            })
1250        ));
1251        assert_eq!(observations.next_calls.get(), 0);
1252        assert_eq!(table.state().version, 1);
1253        assert!(!temp.path().join("data").exists());
1254        Ok(())
1255    }
1256
1257    #[tokio::test]
1258    async fn append_rejects_unsupported_writer_features_before_input_or_artifacts() -> TestResult {
1259        let temp = TempDir::new()?;
1260        let mut table =
1261            TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
1262                .await?;
1263        table
1264            .state
1265            .table_meta
1266            .required_writer_features
1267            .insert("future_writer".to_string());
1268        let state_before = table.state().clone();
1269        let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
1270        let (reader, observations) = InstrumentedReader::new(batch.schema(), vec![Ok(batch)]);
1271
1272        let error = table
1273            .append(reader)
1274            .await
1275            .expect_err("unsupported writer feature must reject append");
1276
1277        assert!(matches!(
1278            error,
1279            TableError::Append {
1280                source: AppendError::Protocol {
1281                    source: TableProtocolError::UnsupportedWriterFeatures { features },
1282                    ..
1283                }
1284            } if features == ["future_writer"]
1285        ));
1286        assert_eq!(observations.schema_calls.get(), 0);
1287        assert_eq!(observations.next_calls.get(), 0);
1288        assert_eq!(table.state(), &state_before);
1289        assert_eq!(table.current_version().await?, 1);
1290        assert!(data_files(temp.path())?.is_empty());
1291        assert!(coverage_files(temp.path())?.is_empty());
1292        assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
1293        Ok(())
1294    }
1295
1296    #[tokio::test]
1297    async fn append_widens_reordered_index_entity_and_data_fields_incrementally() -> TestResult {
1298        let temp = TempDir::new()?;
1299        let location = TableLocation::local(temp.path());
1300        let mut table = TimeSeriesTable::create(location.clone(), widening_table_meta()).await?;
1301        let first = widening_batch(vec![0], vec![i16::MIN], vec![Some(f32::MIN)], vec!["first"])?;
1302        let second = widening_batch(vec![u32::MAX], vec![i16::MAX], vec![None], vec!["second"])?;
1303        let (reader, observations) =
1304            InstrumentedReader::new(first.schema(), vec![Ok(first), Ok(second)]);
1305
1306        assert_eq!(table.append(reader).await?, 2);
1307        assert!(!observations.previous_batch_alive.get());
1308        assert_eq!(table.state().segments.len(), 1);
1309        let segment = table
1310            .state()
1311            .segments
1312            .values()
1313            .next()
1314            .expect("committed widened segment");
1315        assert_eq!(segment.row_count, 2);
1316        assert_eq!(segment.index_min, IndexValue::UInt64(0));
1317        assert_eq!(segment.index_max, IndexValue::UInt64(u64::from(u32::MAX)));
1318        assert_eq!(segment.entity_layout, SegmentEntityLayout::Mixed);
1319
1320        let expected_schema = table.state().table_meta.arrow_schema_ref()?;
1321        let parquet_reader =
1322            ParquetRecordBatchReaderBuilder::try_new(File::open(temp.path().join(&segment.path))?)?
1323                .build()?;
1324        assert_eq!(parquet_reader.schema(), expected_schema);
1325        for batch in parquet_reader {
1326            assert_eq!(batch?.schema(), expected_schema);
1327        }
1328
1329        let state_before_overlap = table.state().clone();
1330        let data_before_overlap = data_files(temp.path())?;
1331        let coverage_before_overlap = coverage_files(temp.path())?;
1332        let overlap = table
1333            .append(widening_batch(
1334                vec![1],
1335                vec![i16::MIN],
1336                vec![Some(1.0)],
1337                vec!["overlap"],
1338            )?)
1339            .await
1340            .expect_err("widened entity and index coverage must still reject overlap");
1341        assert!(matches!(
1342            overlap,
1343            TableError::Append {
1344                source: AppendError::PersistedIndexIntervalOverlap { .. }
1345            }
1346        ));
1347        assert_eq!(table.state(), &state_before_overlap);
1348        assert_eq!(data_files(temp.path())?, data_before_overlap);
1349        assert_eq!(coverage_files(temp.path())?, coverage_before_overlap);
1350
1351        let reopened = TimeSeriesTable::open(location).await?;
1352        let mut stream = reopened.scan_range(0_u64, u64::from(u32::MAX) + 1).await?;
1353        let mut rows = Vec::new();
1354        while let Some(batch) = stream.next().await {
1355            let batch = batch?;
1356            assert_eq!(batch.schema(), expected_schema);
1357            let seq = batch
1358                .column(0)
1359                .as_any()
1360                .downcast_ref::<UInt64Array>()
1361                .expect("seq as UInt64");
1362            let device_id = batch
1363                .column(1)
1364                .as_any()
1365                .downcast_ref::<Int32Array>()
1366                .expect("device_id as Int32");
1367            let reading = batch
1368                .column(2)
1369                .as_any()
1370                .downcast_ref::<Float64Array>()
1371                .expect("reading as Float64");
1372            let label = batch
1373                .column(3)
1374                .as_any()
1375                .downcast_ref::<StringArray>()
1376                .expect("label as Utf8");
1377            for row in 0..batch.num_rows() {
1378                rows.push((
1379                    seq.value(row),
1380                    device_id.value(row),
1381                    (!reading.is_null(row)).then(|| reading.value(row)),
1382                    label.value(row).to_string(),
1383                ));
1384            }
1385        }
1386        rows.sort_by_key(|row| row.0);
1387        assert_eq!(
1388            rows,
1389            vec![
1390                (
1391                    0,
1392                    i32::from(i16::MIN),
1393                    Some(f64::from(f32::MIN)),
1394                    "first".to_string(),
1395                ),
1396                (
1397                    u64::from(u32::MAX),
1398                    i32::from(i16::MAX),
1399                    None,
1400                    "second".to_string(),
1401                ),
1402            ]
1403        );
1404        Ok(())
1405    }
1406
1407    #[tokio::test]
1408    async fn append_rejects_incompatible_declared_schemas_without_artifacts() -> TestResult {
1409        let timestamp = || {
1410            Field::new(
1411                "ts",
1412                DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
1413                false,
1414            )
1415        };
1416        let cases = vec![
1417            (
1418                declared_schema_test_meta(LogicalDataType::Int32, false),
1419                Arc::new(Schema::new(vec![
1420                    timestamp(),
1421                    Field::new("value", DataType::Int64, false),
1422                ])),
1423                "value",
1424            ),
1425            (
1426                declared_schema_test_meta(LogicalDataType::Float64, false),
1427                Arc::new(Schema::new(vec![
1428                    timestamp(),
1429                    Field::new("value", DataType::Int32, false),
1430                ])),
1431                "value",
1432            ),
1433            (
1434                declared_schema_test_meta(LogicalDataType::Int32, false),
1435                Arc::new(Schema::new(vec![
1436                    Field::new(
1437                        "ts",
1438                        DataType::Timestamp(ArrowTimeUnit::Microsecond, None),
1439                        false,
1440                    ),
1441                    Field::new("value", DataType::Int32, false),
1442                ])),
1443                "ts",
1444            ),
1445            (
1446                declared_schema_test_meta(LogicalDataType::Int32, false),
1447                Arc::new(Schema::new(vec![
1448                    timestamp(),
1449                    Field::new("value", DataType::Int32, true),
1450                ])),
1451                "value",
1452            ),
1453            (
1454                declared_schema_test_meta(LogicalDataType::Int32, false),
1455                Arc::new(Schema::new(vec![timestamp()])),
1456                "value",
1457            ),
1458            (
1459                declared_schema_test_meta(LogicalDataType::Int32, false),
1460                Arc::new(Schema::new(vec![
1461                    timestamp(),
1462                    Field::new("value", DataType::Int32, false),
1463                    Field::new("extra", DataType::Int32, false),
1464                ])),
1465                "extra",
1466            ),
1467            (
1468                declared_schema_test_meta(LogicalDataType::Int32, false),
1469                Arc::new(Schema::new(vec![
1470                    timestamp(),
1471                    Field::new("value", DataType::Int32, false),
1472                    Field::new("value", DataType::Int32, false),
1473                ])),
1474                "value",
1475            ),
1476        ];
1477
1478        for (meta, schema, expected_column) in cases {
1479            assert_declared_schema_rejected_before_reading(
1480                meta,
1481                schema,
1482                Vec::new(),
1483                expected_column,
1484            )
1485            .await?;
1486        }
1487        Ok(())
1488    }
1489
1490    #[tokio::test]
1491    async fn append_rejects_positive_signed_index_without_inspecting_values() -> TestResult {
1492        let meta = TableMeta::new_time_series_with_schema(
1493            IndexSpec {
1494                column: "seq".to_string(),
1495                entity_columns: Vec::new(),
1496                kind: IndexKind::UInt64 {
1497                    index_granularity: NonZeroU64::new(1).unwrap(),
1498                },
1499            },
1500            LogicalSchema::new(vec![LogicalField {
1501                name: "seq".to_string(),
1502                data_type: LogicalDataType::UInt64,
1503                nullable: false,
1504            }])?,
1505        );
1506        let schema = Arc::new(Schema::new(vec![Field::new("seq", DataType::Int64, false)]));
1507        let batch = RecordBatch::try_new(
1508            Arc::clone(&schema),
1509            vec![Arc::new(Int64Array::from(vec![1]))],
1510        )?;
1511
1512        assert_declared_schema_rejected_before_reading(meta, schema, vec![Ok(batch)], "seq").await
1513    }
1514
1515    #[tokio::test]
1516    async fn widened_batch_source_failure_rolls_back_partial_output() -> TestResult {
1517        let temp = TempDir::new()?;
1518        let mut table =
1519            TimeSeriesTable::create(TableLocation::local(temp.path()), widening_table_meta())
1520                .await?;
1521        let first = widening_batch(vec![0], vec![1], vec![Some(1.0)], vec!["first"])?;
1522        let (reader, observations) = InstrumentedReader::new(
1523            first.schema(),
1524            vec![
1525                Ok(first),
1526                Err(ArrowError::ComputeError(
1527                    "injected widened source failure".to_string(),
1528                )),
1529            ],
1530        );
1531
1532        let error = table
1533            .append(reader)
1534            .await
1535            .expect_err("source failure must abort widened append");
1536
1537        assert!(matches!(
1538            error,
1539            TableError::Append {
1540                source: AppendError::ArrowInput { .. }
1541            }
1542        ));
1543        assert!(!observations.previous_batch_alive.get());
1544        assert_eq!(table.state().version, 1);
1545        assert!(table.state().segments.is_empty());
1546        assert!(data_files(temp.path())?.is_empty());
1547        assert!(coverage_files(temp.path())?.is_empty());
1548        assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
1549        Ok(())
1550    }
1551
1552    fn assert_batches<R>(reader: R, expected: &[RecordBatch]) -> TestResult
1553    where
1554        R: RecordBatchReader,
1555    {
1556        assert_eq!(reader.schema(), expected[0].schema());
1557        let actual = reader.collect::<Result<Vec<_>, _>>()?;
1558        assert_eq!(actual, expected);
1559        Ok(())
1560    }
1561
1562    #[test]
1563    fn into_record_batch_reader_accepts_all_required_input_forms() -> TestResult {
1564        let first = input_batch(vec![1, 2])?;
1565        let second = input_batch(vec![3])?;
1566
1567        assert_batches(
1568            first.clone().into_record_batch_reader()?,
1569            std::slice::from_ref(&first),
1570        )?;
1571        assert_batches(
1572            vec![first.clone(), second.clone()].into_record_batch_reader()?,
1573            &[first.clone(), second.clone()],
1574        )?;
1575
1576        let iterator = RecordBatchIterator::new(vec![Ok(first.clone())], first.schema());
1577        assert_batches(
1578            iterator.into_record_batch_reader()?,
1579            std::slice::from_ref(&first),
1580        )?;
1581        assert_batches(
1582            InstrumentedReader::one(first.clone()).into_record_batch_reader()?,
1583            std::slice::from_ref(&first),
1584        )?;
1585
1586        let boxed: Box<dyn RecordBatchReader> = Box::new(RecordBatchIterator::new(
1587            vec![Ok(first.clone())],
1588            first.schema(),
1589        ));
1590        assert_batches(
1591            boxed.into_record_batch_reader()?,
1592            std::slice::from_ref(&first),
1593        )?;
1594
1595        let sendable: Box<dyn RecordBatchReader + Send> = Box::new(RecordBatchIterator::new(
1596            vec![Ok(second.clone())],
1597            second.schema(),
1598        ));
1599        assert_batches(
1600            sendable.into_record_batch_reader()?,
1601            std::slice::from_ref(&second),
1602        )?;
1603        Ok(())
1604    }
1605
1606    #[test]
1607    fn into_record_batch_reader_rejects_schema_less_vec() {
1608        let result = Vec::<RecordBatch>::new().into_record_batch_reader();
1609        assert!(matches!(result, Err(AppendError::EmptyInput)));
1610    }
1611
1612    #[test]
1613    fn append_writer_defaults_are_compressed_and_bounded() {
1614        let settings = AppendWriterSettings::default();
1615        let properties = settings.writer_properties();
1616
1617        assert_eq!(settings.compression, ParquetCompression::Zstd);
1618        assert_eq!(
1619            properties.max_row_group_row_count(),
1620            Some(DEFAULT_MAX_ROWS_PER_ROW_GROUP)
1621        );
1622        assert_eq!(
1623            properties.max_row_group_bytes(),
1624            Some(DEFAULT_MAX_BYTES_PER_ROW_GROUP)
1625        );
1626    }
1627
1628    #[test]
1629    fn append_request_nested_writer_settings_inherit_and_override() -> TestResult {
1630        let inner = AppendRequest::new(timestamp_only_batch(1)?)
1631            .compression(ParquetCompression::Snappy)
1632            .max_rows_per_row_group(3)
1633            .max_bytes_per_row_group(30);
1634        let inherited = AppendWriterSettings::from_source(&AppendRequest::new(inner));
1635        assert_eq!(
1636            inherited,
1637            AppendWriterSettings {
1638                compression: ParquetCompression::Snappy,
1639                max_rows_per_row_group: 3,
1640                max_bytes_per_row_group: 30,
1641            }
1642        );
1643
1644        let partially_overridden = AppendWriterSettings::from_source(
1645            &AppendRequest::new(
1646                AppendRequest::new(timestamp_only_batch(1)?)
1647                    .compression(ParquetCompression::Snappy)
1648                    .max_rows_per_row_group(3)
1649                    .max_bytes_per_row_group(30),
1650            )
1651            .max_bytes_per_row_group(20),
1652        );
1653        assert_eq!(
1654            partially_overridden,
1655            AppendWriterSettings {
1656                compression: ParquetCompression::Snappy,
1657                max_rows_per_row_group: 3,
1658                max_bytes_per_row_group: 20,
1659            }
1660        );
1661
1662        let overridden = AppendWriterSettings::from_source(
1663            &AppendRequest::new(
1664                AppendRequest::new(timestamp_only_batch(1)?)
1665                    .compression(ParquetCompression::Snappy)
1666                    .max_rows_per_row_group(3)
1667                    .max_bytes_per_row_group(30),
1668            )
1669            .compression(ParquetCompression::Zstd)
1670            .max_rows_per_row_group(2)
1671            .max_bytes_per_row_group(20),
1672        );
1673        assert_eq!(
1674            overridden,
1675            AppendWriterSettings {
1676                compression: ParquetCompression::Zstd,
1677                max_rows_per_row_group: 2,
1678                max_bytes_per_row_group: 20,
1679            }
1680        );
1681        Ok(())
1682    }
1683
1684    #[tokio::test]
1685    async fn append_default_and_explicit_compression_reach_parquet_metadata() -> TestResult {
1686        for (compression, expected) in [
1687            (None, Compression::ZSTD(Default::default())),
1688            (
1689                Some(ParquetCompression::Uncompressed),
1690                Compression::UNCOMPRESSED,
1691            ),
1692            (Some(ParquetCompression::Snappy), Compression::SNAPPY),
1693            (
1694                Some(ParquetCompression::Zstd),
1695                Compression::ZSTD(Default::default()),
1696            ),
1697        ] {
1698            let temp = TempDir::new()?;
1699            let mut table =
1700                TimeSeriesTable::create(TableLocation::local(temp.path()), timestamp_only_meta())
1701                    .await?;
1702            let batch = timestamp_only_batch(1)?;
1703
1704            match compression {
1705                Some(compression) => {
1706                    table
1707                        .append(AppendRequest::new(batch).compression(compression))
1708                        .await?;
1709                }
1710                None => {
1711                    table.append(batch).await?;
1712                }
1713            }
1714
1715            let segment_path = &table
1716                .state()
1717                .segments
1718                .values()
1719                .next()
1720                .ok_or("missing appended segment")?
1721                .path;
1722            let builder = ParquetRecordBatchReaderBuilder::try_new(File::open(
1723                temp.path().join(segment_path),
1724            )?)?;
1725            assert!(builder.metadata().row_groups().iter().all(|row_group| {
1726                row_group
1727                    .columns()
1728                    .iter()
1729                    .all(|column| column.compression() == expected)
1730            }));
1731        }
1732        Ok(())
1733    }
1734
1735    #[tokio::test]
1736    async fn append_byte_limit_splits_row_groups_and_preserves_values() -> TestResult {
1737        let temp = TempDir::new()?;
1738        let location = TableLocation::local(temp.path());
1739        let mut table = TimeSeriesTable::create(
1740            location.clone(),
1741            TableMeta::new_time_series(timestamp_only_index()),
1742        )
1743        .await?;
1744        let schema = Arc::new(Schema::new(vec![
1745            Field::new(
1746                "ts",
1747                DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
1748                false,
1749            ),
1750            Field::new("payload", DataType::Binary, false),
1751        ]));
1752        let payloads = [vec![1_u8; 256], vec![2_u8; 256], vec![3_u8; 256]];
1753        let batches = payloads
1754            .iter()
1755            .enumerate()
1756            .map(|(minute, payload)| {
1757                RecordBatch::try_new(
1758                    Arc::clone(&schema),
1759                    vec![
1760                        Arc::new(TimestampMillisecondArray::from(vec![
1761                            minute as i64 * 60_000,
1762                        ])),
1763                        Arc::new(BinaryArray::from_iter_values([payload.as_slice()])),
1764                    ],
1765                )
1766            })
1767            .collect::<Result<Vec<_>, _>>()?;
1768
1769        assert_eq!(
1770            table
1771                .append(
1772                    AppendRequest::new(batches)
1773                        .compression(ParquetCompression::Uncompressed)
1774                        .max_rows_per_row_group(100)
1775                        .max_bytes_per_row_group(1),
1776                )
1777                .await?,
1778            2
1779        );
1780
1781        let segment_path = &table
1782            .state()
1783            .segments
1784            .values()
1785            .next()
1786            .ok_or("missing appended segment")?
1787            .path;
1788        let builder =
1789            ParquetRecordBatchReaderBuilder::try_new(File::open(temp.path().join(segment_path))?)?;
1790        assert_eq!(
1791            builder
1792                .metadata()
1793                .row_groups()
1794                .iter()
1795                .map(|row_group| row_group.num_rows())
1796                .collect::<Vec<_>>(),
1797            [1, 1, 1]
1798        );
1799        let batches = builder.build()?.collect::<Result<Vec<_>, _>>()?;
1800        let mut actual_payloads = Vec::new();
1801        for batch in batches {
1802            let values = batch
1803                .column_by_name("payload")
1804                .and_then(|column| column.as_any().downcast_ref::<BinaryArray>())
1805                .ok_or("missing binary payload column")?;
1806            actual_payloads.extend((0..values.len()).map(|index| values.value(index).to_vec()));
1807        }
1808        assert_eq!(actual_payloads, payloads);
1809
1810        let reopened = TimeSeriesTable::open(location).await?;
1811        assert_eq!(reopened.state().version, 2);
1812        assert_eq!(reopened.state().segments.len(), 1);
1813        Ok(())
1814    }
1815
1816    #[tokio::test]
1817    async fn append_request_nested_settings_inherit_and_override() -> TestResult {
1818        for (inner_limit, outer_limit, expected_row_groups) in
1819            [(3, None, vec![3, 3, 1]), (0, Some(2), vec![2, 2, 2, 1])]
1820        {
1821            let temp = TempDir::new()?;
1822            let location = TableLocation::local(temp.path());
1823            let mut table =
1824                TimeSeriesTable::create(location.clone(), timestamp_only_meta()).await?;
1825            let request = AppendRequest::new(
1826                AppendRequest::new(vec![
1827                    timestamp_only_batch(2)?,
1828                    timestamp_only_batch_starting_at_minute_offset(2, 5)?,
1829                ])
1830                .max_rows_per_row_group(inner_limit),
1831            );
1832            let request = match outer_limit {
1833                Some(limit) => request.max_rows_per_row_group(limit),
1834                None => request,
1835            };
1836
1837            assert_eq!(table.append(request).await?, 2);
1838
1839            let (segment_path, row_count) = table
1840                .state()
1841                .segments
1842                .values()
1843                .next()
1844                .map(|segment| (segment.path.clone(), segment.row_count))
1845                .ok_or("missing appended segment")?;
1846            assert_eq!(row_count, 7);
1847            let builder = ParquetRecordBatchReaderBuilder::try_new(File::open(
1848                temp.path().join(segment_path),
1849            )?)?;
1850            let row_group_rows = builder
1851                .metadata()
1852                .row_groups()
1853                .iter()
1854                .map(|row_group| row_group.num_rows())
1855                .collect::<Vec<_>>();
1856            assert_eq!(row_group_rows, expected_row_groups);
1857
1858            let reopened = TimeSeriesTable::open(location).await?;
1859            assert_eq!(reopened.state().version, 2);
1860            assert_eq!(reopened.state().segments.len(), 1);
1861        }
1862        Ok(())
1863    }
1864
1865    #[tokio::test]
1866    async fn append_request_nested_zero_is_rejected_before_inspecting_source() -> TestResult {
1867        let temp = TempDir::new()?;
1868        let mut table =
1869            TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
1870                .await?;
1871        let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
1872        let (reader, observations) = InstrumentedReader::new(batch.schema(), vec![Ok(batch)]);
1873        let request = AppendRequest::new(AppendRequest::new(reader).max_rows_per_row_group(0));
1874
1875        assert!(matches!(
1876            table.append(request).await,
1877            Err(TableError::Append {
1878                source: AppendError::InvalidMaxRowsPerRowGroup {
1879                    max_rows_per_row_group: 0
1880                }
1881            })
1882        ));
1883        assert_eq!(observations.schema_calls.get(), 0);
1884        assert_eq!(observations.next_calls.get(), 0);
1885        assert_eq!(table.state().version, 1);
1886        assert!(table.state().segments.is_empty());
1887        assert!(data_files(temp.path())?.is_empty());
1888        assert!(coverage_files(temp.path())?.is_empty());
1889        assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
1890        Ok(())
1891    }
1892
1893    #[tokio::test]
1894    async fn append_request_zero_byte_limit_is_rejected_before_inspecting_source() -> TestResult {
1895        let temp = TempDir::new()?;
1896        let mut table =
1897            TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
1898                .await?;
1899        let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
1900        let (reader, observations) = InstrumentedReader::new(batch.schema(), vec![Ok(batch)]);
1901
1902        assert!(matches!(
1903            table
1904                .append(AppendRequest::new(reader).max_bytes_per_row_group(0))
1905                .await,
1906            Err(TableError::Append {
1907                source: AppendError::InvalidMaxBytesPerRowGroup {
1908                    max_bytes_per_row_group: 0
1909                }
1910            })
1911        ));
1912        assert_eq!(observations.schema_calls.get(), 0);
1913        assert_eq!(observations.next_calls.get(), 0);
1914        assert_eq!(table.state().version, 1);
1915        assert!(table.state().segments.is_empty());
1916        assert!(data_files(temp.path())?.is_empty());
1917        assert!(coverage_files(temp.path())?.is_empty());
1918        assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
1919        Ok(())
1920    }
1921
1922    #[tokio::test]
1923    async fn append_request_without_overrides_matches_direct_append() -> TestResult {
1924        let direct_root = TempDir::new()?;
1925        let request_root = TempDir::new()?;
1926        let mut direct = TimeSeriesTable::create(
1927            TableLocation::local(direct_root.path()),
1928            make_basic_table_meta(),
1929        )
1930        .await?;
1931        let mut requested = TimeSeriesTable::create(
1932            TableLocation::local(request_root.path()),
1933            make_basic_table_meta(),
1934        )
1935        .await?;
1936        let source = vec![
1937            time_series_batch(vec![0], vec!["A"], vec![1.0])?,
1938            time_series_batch(vec![60_000], vec!["A"], vec![2.0])?,
1939        ];
1940
1941        assert_eq!(direct.append(source.clone()).await?, 2);
1942        assert_eq!(requested.append(AppendRequest::new(source)).await?, 2);
1943
1944        let direct_path = &direct
1945            .state()
1946            .segments
1947            .values()
1948            .next()
1949            .ok_or("missing direct segment")?
1950            .path;
1951        let request_path = &requested
1952            .state()
1953            .segments
1954            .values()
1955            .next()
1956            .ok_or("missing requested segment")?
1957            .path;
1958        assert_eq!(
1959            std::fs::read(direct_root.path().join(direct_path))?,
1960            std::fs::read(request_root.path().join(request_path))?
1961        );
1962        Ok(())
1963    }
1964
1965    #[tokio::test]
1966    async fn append_request_rejects_zero_before_inspecting_source() -> TestResult {
1967        let temp = TempDir::new()?;
1968        let mut table =
1969            TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
1970                .await?;
1971        let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
1972        let (reader, observations) = InstrumentedReader::new(batch.schema(), vec![Ok(batch)]);
1973
1974        assert!(matches!(
1975            table
1976                .append(AppendRequest::new(reader).max_rows_per_row_group(0))
1977                .await,
1978            Err(TableError::Append {
1979                source: AppendError::InvalidMaxRowsPerRowGroup {
1980                    max_rows_per_row_group: 0
1981                }
1982            })
1983        ));
1984        assert_eq!(observations.schema_calls.get(), 0);
1985        assert_eq!(observations.next_calls.get(), 0);
1986        assert_eq!(table.state().version, 1);
1987        assert!(table.state().segments.is_empty());
1988        assert!(data_files(temp.path())?.is_empty());
1989        assert!(coverage_files(temp.path())?.is_empty());
1990        assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
1991        Ok(())
1992    }
1993
1994    #[tokio::test]
1995    async fn append_rejects_version_overflow_before_inspecting_source() -> TestResult {
1996        let temp = TempDir::new()?;
1997        let mut table =
1998            TimeSeriesTable::create(TableLocation::local(temp.path()), timestamp_only_meta())
1999                .await?;
2000        table.state.version = u64::MAX;
2001        let batch = timestamp_only_batch(1)?;
2002        let (reader, observations) = InstrumentedReader::new(batch.schema(), vec![Ok(batch)]);
2003
2004        let error = table
2005            .append(reader)
2006            .await
2007            .expect_err("version overflow must fail");
2008
2009        assert!(matches!(
2010            error,
2011            TableError::Append {
2012                source: AppendError::Commit {
2013                    source: CommitError::VersionOverflow {
2014                        current_version: u64::MAX,
2015                        ..
2016                    }
2017                }
2018            }
2019        ));
2020        assert_eq!(observations.schema_calls.get(), 0);
2021        assert_eq!(observations.next_calls.get(), 0);
2022        assert_eq!(table.state().version, u64::MAX);
2023        assert!(table.state().segments.is_empty());
2024        assert!(!temp.path().join("data").exists());
2025        assert!(!temp.path().join("_coverage").exists());
2026        assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
2027        assert_eq!(table.log.load_current_version().await?, 1);
2028        Ok(())
2029    }
2030
2031    #[tokio::test]
2032    async fn append_writes_one_segment_and_returns_versions() -> TestResult {
2033        let temp = TempDir::new()?;
2034        let location = TableLocation::local(temp.path());
2035        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
2036
2037        let first = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2038        let second = time_series_batch(vec![60_000], vec!["A"], vec![2.0])?;
2039        assert_eq!(table.append(vec![first, second]).await?, 2);
2040        assert_eq!(table.state().segments.len(), 1);
2041        let Some(segment) = table.state().segments.values().next() else {
2042            return Err("missing streamed segment".into());
2043        };
2044        assert_eq!(segment.row_count, 2);
2045        assert!(temp.path().join(&segment.path).is_file());
2046
2047        let third = time_series_batch(vec![120_000], vec!["A"], vec![3.0])?;
2048        assert_eq!(table.append(third).await?, 3);
2049        assert_eq!(table.state().segments.len(), 2);
2050
2051        let reopened = TimeSeriesTable::open(location).await?;
2052        assert_eq!(reopened.state().version, 3);
2053        assert_eq!(
2054            reopened
2055                .state()
2056                .segments
2057                .values()
2058                .map(|segment| segment.row_count)
2059                .sum::<u64>(),
2060            3
2061        );
2062        Ok(())
2063    }
2064
2065    #[tokio::test]
2066    async fn append_consumes_non_send_reader_one_batch_at_a_time() -> TestResult {
2067        let temp = TempDir::new()?;
2068        let mut table =
2069            TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2070                .await?;
2071        let first = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2072        let second = time_series_batch(vec![60_000], vec!["A"], vec![2.0])?;
2073        let (reader, observations) =
2074            InstrumentedReader::new(first.schema(), vec![Ok(first), Ok(second)]);
2075
2076        assert_eq!(table.append(reader).await?, 2);
2077        assert!(observations.schema_calls.get() > 0);
2078        assert!(!observations.next_before_schema.get());
2079        assert!(!observations.previous_batch_alive.get());
2080        assert_eq!(observations.next_calls.get(), 3);
2081        Ok(())
2082    }
2083
2084    #[tokio::test(flavor = "current_thread")]
2085    async fn append_yields_while_skipping_zero_row_batches() -> TestResult {
2086        const BATCH_COUNT: usize = 64;
2087
2088        let temp = TempDir::new()?;
2089        let mut table =
2090            TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2091                .await?;
2092        let zero = time_series_batch(Vec::new(), Vec::new(), Vec::new())?;
2093        let (reader, observations) = InstrumentedReader::new(
2094            zero.schema(),
2095            (0..BATCH_COUNT).map(|_| Ok(zero.clone())).collect(),
2096        );
2097        let mut append = Box::pin(table.append(reader));
2098
2099        tokio::select! {
2100            biased;
2101            result = &mut append => panic!("append drained the reader without yielding: {result:?}"),
2102            () = async {
2103                while observations.next_calls.get() < 2 {
2104                    tokio::task::yield_now().await;
2105                }
2106            } => {}
2107        }
2108
2109        assert!(observations.next_calls.get() < BATCH_COUNT);
2110        drop(append);
2111        assert_eq!(table.state().version, 1);
2112        assert!(!temp.path().join("data").exists());
2113        Ok(())
2114    }
2115
2116    #[tokio::test(flavor = "current_thread")]
2117    async fn cancelling_append_during_batch_reads_removes_output() -> TestResult {
2118        const BATCH_COUNT: usize = 64;
2119
2120        let temp = TempDir::new()?;
2121        let mut table =
2122            TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2123                .await?;
2124        let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2125        let (reader, observations) = InstrumentedReader::new(
2126            batch.schema(),
2127            (0..BATCH_COUNT).map(|_| Ok(batch.clone())).collect(),
2128        );
2129        let mut append = Box::pin(table.append(reader));
2130
2131        tokio::select! {
2132            biased;
2133            result = &mut append => panic!("append drained the reader without yielding: {result:?}"),
2134            () = async {
2135                while observations.next_calls.get() < 2 {
2136                    tokio::task::yield_now().await;
2137                }
2138            } => {}
2139        }
2140
2141        assert!(observations.next_calls.get() < BATCH_COUNT);
2142        assert_eq!(std::fs::read_dir(temp.path().join("data"))?.count(), 1);
2143        drop(append);
2144        assert_eq!(table.state().version, 1);
2145        assert!(table.state().segments.is_empty());
2146        assert_eq!(std::fs::read_dir(temp.path().join("data"))?.count(), 0);
2147        assert!(coverage_files(temp.path())?.is_empty());
2148        assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
2149        Ok(())
2150    }
2151
2152    #[tokio::test]
2153    async fn append_rejects_empty_sources_before_creating_data() -> TestResult {
2154        let temp = TempDir::new()?;
2155        let mut table =
2156            TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2157                .await?;
2158        let zero = time_series_batch(Vec::new(), Vec::new(), Vec::new())?;
2159        let schema = zero.schema();
2160
2161        assert!(matches!(
2162            table.append(vec![zero.clone(), zero]).await,
2163            Err(TableError::Append {
2164                source: AppendError::EmptyInput
2165            })
2166        ));
2167        let empty = RecordBatchIterator::new(Vec::<Result<RecordBatch, ArrowError>>::new(), schema);
2168        assert!(matches!(
2169            table.append(empty).await,
2170            Err(TableError::Append {
2171                source: AppendError::EmptyInput
2172            })
2173        ));
2174        assert_eq!(table.state().version, 1);
2175        assert!(table.state().segments.is_empty());
2176        assert!(!temp.path().join("data").exists());
2177
2178        let data = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2179        let leading_zero = time_series_batch(Vec::new(), Vec::new(), Vec::new())?;
2180        assert_eq!(table.append(vec![leading_zero, data]).await?, 2);
2181        assert_eq!(table.state().version, 2);
2182        assert_eq!(table.state().segments.len(), 1);
2183        Ok(())
2184    }
2185
2186    #[tokio::test]
2187    async fn append_cleans_up_after_later_schema_or_source_error() -> TestResult {
2188        for (configured, source_error) in
2189            [(false, false), (false, true), (true, false), (true, true)]
2190        {
2191            let temp = TempDir::new()?;
2192            let mut table =
2193                TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2194                    .await?;
2195            let first = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2196            let later = if source_error {
2197                Err(ArrowError::ComputeError(
2198                    "injected batch source failure".to_string(),
2199                ))
2200            } else {
2201                Ok(input_batch(vec![1])?)
2202            };
2203            let (reader, _) = InstrumentedReader::new(first.schema(), vec![Ok(first), later]);
2204
2205            let result = if configured {
2206                table
2207                    .append(AppendRequest::new(reader).max_rows_per_row_group(1))
2208                    .await
2209            } else {
2210                table.append(reader).await
2211            };
2212            let error = result.expect_err("later batch must fail the append");
2213            assert!(
2214                matches!(
2215                    error,
2216                    TableError::Append {
2217                        source: AppendError::ArrowInput { .. }
2218                    }
2219                ),
2220                "configured={configured}, source_error={source_error}"
2221            );
2222            assert_eq!(
2223                table.state().version,
2224                1,
2225                "configured={configured}, source_error={source_error}"
2226            );
2227            assert!(
2228                table.state().segments.is_empty(),
2229                "configured={configured}, source_error={source_error}"
2230            );
2231            assert!(
2232                std::fs::read_dir(temp.path().join("data"))?
2233                    .next()
2234                    .is_none(),
2235                "configured={configured}, source_error={source_error}"
2236            );
2237            assert!(
2238                coverage_files(temp.path())?.is_empty(),
2239                "configured={configured}, source_error={source_error}"
2240            );
2241            assert!(
2242                !temp.path().join(layout::commit_rel_path(2)).exists(),
2243                "configured={configured}, source_error={source_error}"
2244            );
2245        }
2246        Ok(())
2247    }
2248
2249    #[tokio::test]
2250    async fn append_cleans_up_writer_lifecycle_failures() -> TestResult {
2251        #[derive(Clone, Copy, Debug)]
2252        enum Stage {
2253            Write,
2254            Close,
2255            Finish,
2256        }
2257
2258        for stage in [Stage::Write, Stage::Close, Stage::Finish] {
2259            let temp = TempDir::new()?;
2260            let mut table =
2261                TimeSeriesTable::create(TableLocation::local(temp.path()), timestamp_only_meta())
2262                    .await?;
2263            let data_dir = temp.path().join("data");
2264            match stage {
2265                Stage::Write => crate::storage::inject_output_write_failure(data_dir.clone(), 2),
2266                Stage::Close => crate::storage::inject_output_write_failure(data_dir.clone(), 1),
2267                Stage::Finish => crate::storage::inject_output_finish_failure(data_dir.clone()),
2268            }
2269            let row_count = if matches!(stage, Stage::Write) {
2270                parquet::file::properties::DEFAULT_MAX_ROW_GROUP_ROW_COUNT
2271            } else {
2272                1
2273            };
2274
2275            let result = table.append(timestamp_only_batch(row_count)?).await;
2276            let Err(error) = result else {
2277                panic!("{stage:?} writer failure was not injected");
2278            };
2279            if matches!(stage, Stage::Finish) {
2280                assert!(
2281                    matches!(
2282                        error,
2283                        TableError::Append {
2284                            source: AppendError::Storage { .. }
2285                        }
2286                    ),
2287                    "{stage:?}"
2288                );
2289            } else {
2290                assert!(
2291                    matches!(
2292                        error,
2293                        TableError::Append {
2294                            source: AppendError::ParquetWrite { .. }
2295                        }
2296                    ),
2297                    "{stage:?}: {error}"
2298                );
2299            }
2300            assert_eq!(table.state().version, 1, "{stage:?}");
2301            assert!(table.state().segments.is_empty(), "{stage:?}");
2302            assert!(std::fs::read_dir(&data_dir)?.next().is_none(), "{stage:?}");
2303            assert!(coverage_files(temp.path())?.is_empty(), "{stage:?}");
2304            assert!(
2305                !temp.path().join(layout::commit_rel_path(2)).exists(),
2306                "{stage:?}"
2307            );
2308        }
2309        Ok(())
2310    }
2311
2312    #[tokio::test]
2313    async fn cancelling_append_before_publication_removes_owned_artifacts() -> TestResult {
2314        let temp = TempDir::new()?;
2315        let location = TableLocation::local(temp.path());
2316        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
2317        let commit_path = temp.path().join(layout::commit_rel_path(2));
2318        let current_path = temp.path().join(layout::current_rel_path());
2319        let current_temp_path = current_path.with_extension("tmp");
2320        let mut pause = crate::storage::pause_atomic_write_before_rename(current_path);
2321        let mut append = Box::pin(table.append(time_series_batch(vec![0], vec!["A"], vec![1.0])?));
2322
2323        tokio::select! {
2324            () = pause.wait_until_paused() => {}
2325            result = &mut append => panic!("append completed before cancellation: {result:?}"),
2326        }
2327
2328        assert!(commit_path.is_file());
2329        assert!(current_temp_path.is_file());
2330        assert_eq!(std::fs::read_dir(temp.path().join("data"))?.count(), 1);
2331        assert_eq!(coverage_files(temp.path())?.len(), 2);
2332        let before_publication = TimeSeriesTable::open(location.clone()).await?;
2333        assert_eq!(before_publication.state().version, 1);
2334        assert!(before_publication.state().segments.is_empty());
2335
2336        drop(append);
2337        pause.release();
2338
2339        assert!(!commit_path.exists());
2340        assert!(!current_temp_path.exists());
2341        assert_eq!(std::fs::read_dir(temp.path().join("data"))?.count(), 0);
2342        assert!(coverage_files(temp.path())?.is_empty());
2343        assert_eq!(table.state().version, 1);
2344        assert!(table.state().segments.is_empty());
2345        let reopened = TimeSeriesTable::open(location).await?;
2346        assert_eq!(reopened.state().version, 1);
2347        assert!(reopened.state().segments.is_empty());
2348        Ok(())
2349    }
2350
2351    #[tokio::test]
2352    async fn post_commit_observer_panic_preserves_owned_data_files() -> TestResult {
2353        let streamed_root = TempDir::new()?;
2354        let streamed_location = TableLocation::local(streamed_root.path());
2355        let mut streamed_table =
2356            TimeSeriesTable::create(streamed_location.clone(), make_basic_table_meta()).await?;
2357        let callsite_guard = panic_on_commit_close_dispatch();
2358        let outcome = std::panic::AssertUnwindSafe(
2359            streamed_table
2360                .append(time_series_batch(vec![0], vec!["A"], vec![1.0])?)
2361                .with_subscriber(panic_on_commit_close_dispatch()),
2362        )
2363        .catch_unwind()
2364        .await;
2365        drop(callsite_guard);
2366        assert!(outcome.is_err());
2367        let reopened = TimeSeriesTable::open(streamed_location).await?;
2368        assert_eq!(reopened.state().version, 2);
2369        assert_eq!(reopened.state().segments.len(), 1);
2370        assert!(
2371            reopened
2372                .state()
2373                .segments
2374                .values()
2375                .all(|segment| streamed_root.path().join(&segment.path).is_file())
2376        );
2377
2378        Ok(())
2379    }
2380
2381    #[tokio::test]
2382    async fn simultaneous_appends_publish_one_and_clean_loser() -> TestResult {
2383        let temp = TempDir::new()?;
2384        let location = TableLocation::local(temp.path());
2385        let mut winner = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
2386        let mut loser = TimeSeriesTable::open(location.clone()).await?;
2387        let loser_state_before = loser.state().clone();
2388        let commit_path = temp.path().join(layout::commit_rel_path(2));
2389        let mut pause = crate::storage::pause_atomic_write_before_rename(
2390            temp.path().join(layout::current_rel_path()),
2391        );
2392        let mut winner_append =
2393            Box::pin(winner.append(time_series_batch(vec![0], vec!["A"], vec![1.0])?));
2394
2395        tokio::select! {
2396            () = pause.wait_until_paused() => {}
2397            result = &mut winner_append => panic!("winner completed before race: {result:?}"),
2398        }
2399
2400        let mut data_before = std::fs::read_dir(temp.path().join("data"))?
2401            .map(|entry| entry.map(|entry| entry.path()))
2402            .collect::<Result<Vec<_>, _>>()?;
2403        data_before.sort();
2404        let coverage_before = coverage_files(temp.path())?;
2405        let commit_before = std::fs::read(&commit_path)?;
2406
2407        let loser_error = loser
2408            .append(time_series_batch(vec![120_000], vec!["A"], vec![2.0])?)
2409            .await
2410            .expect_err("second simultaneous writer must lose commit 2");
2411        assert!(matches!(
2412            loser_error,
2413            TableError::Append {
2414                source: AppendError::Commit {
2415                    source: CommitError::Storage {
2416                        source: StorageError::AlreadyExists { .. },
2417                    },
2418                }
2419            }
2420        ));
2421        assert_eq!(loser.state(), &loser_state_before);
2422
2423        let mut data_after = std::fs::read_dir(temp.path().join("data"))?
2424            .map(|entry| entry.map(|entry| entry.path()))
2425            .collect::<Result<Vec<_>, _>>()?;
2426        data_after.sort();
2427        assert_eq!(data_after, data_before);
2428        assert_eq!(coverage_files(temp.path())?, coverage_before);
2429        assert_eq!(std::fs::read(&commit_path)?, commit_before);
2430
2431        pause.release();
2432        assert_eq!(winner_append.await?, 2);
2433        assert_eq!(winner.state().version, 2);
2434        assert_eq!(winner.state().segments.len(), 1);
2435        assert!(!temp.path().join(layout::commit_rel_path(3)).exists());
2436
2437        let reopened = TimeSeriesTable::open(location).await?;
2438        assert_eq!(reopened.state().version, 2);
2439        assert_eq!(reopened.state().segments.len(), 1);
2440        assert_eq!(
2441            reopened
2442                .state()
2443                .segments
2444                .values()
2445                .map(|segment| segment.row_count)
2446                .sum::<u64>(),
2447            1
2448        );
2449        Ok(())
2450    }
2451
2452    #[tokio::test]
2453    async fn guarded_output_cleans_up_parquet_writer_creation_failure() -> TestResult {
2454        let temp = TempDir::new()?;
2455        let location = StorageLocation::local(temp.path());
2456        let path = Path::new("data/create-failure.parquet");
2457        let sink = storage::open_new_output_sink(&location, path).await?;
2458        let schema = Arc::new(Schema::new(vec![Field::new(
2459            "unsupported",
2460            DataType::Struct(arrow::datatypes::Fields::empty()),
2461            true,
2462        )]));
2463
2464        assert!(ParquetArrowWriter::try_new(sink, schema, None).is_err());
2465        assert!(!temp.path().join(path).exists());
2466        Ok(())
2467    }
2468
2469    #[tokio::test]
2470    async fn append_preserves_writer_error_and_failed_cleanup_path() -> TestResult {
2471        let temp = TempDir::new()?;
2472        let mut table =
2473            TimeSeriesTable::create(TableLocation::local(temp.path()), timestamp_only_meta())
2474                .await?;
2475        let data_dir = temp.path().join("data");
2476        crate::storage::inject_output_write_failure(data_dir.clone(), 1);
2477        crate::storage::inject_cleanup_failure(data_dir.clone());
2478
2479        let error = table
2480            .append(timestamp_only_batch(1)?)
2481            .await
2482            .expect_err("writer and cleanup failures must fail append");
2483        let TableError::Append {
2484            source:
2485                AppendError::Rollback {
2486                    source,
2487                    cleanup_errors,
2488                },
2489        } = error
2490        else {
2491            panic!("unexpected error: {error}");
2492        };
2493        assert!(matches!(*source, AppendError::ParquetWrite { .. }));
2494        assert_eq!(cleanup_errors.len(), 1);
2495        let remaining = std::fs::read_dir(&data_dir)?.collect::<Result<Vec<_>, _>>()?;
2496        assert_eq!(remaining.len(), 1);
2497        let remaining_path = remaining[0].path();
2498        let remaining_name = remaining_path
2499            .file_name()
2500            .expect("remaining data filename")
2501            .to_str()
2502            .expect("UTF-8 data filename");
2503        assert!(matches!(
2504            &cleanup_errors[0],
2505            StorageError::OtherIo { path, .. } if path.ends_with(remaining_name)
2506        ));
2507        assert_eq!(table.state().version, 1);
2508        assert!(table.state().segments.is_empty());
2509        assert!(coverage_files(temp.path())?.is_empty());
2510        assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
2511        Ok(())
2512    }
2513
2514    #[tokio::test]
2515    async fn append_retries_cleanup_after_sidecar_write_cleanup_failure() -> TestResult {
2516        for sidecar_dir in [SEGMENT_COVERAGE_DIR, TABLE_SNAPSHOT_DIR] {
2517            let temp = TempDir::new()?;
2518            let mut table =
2519                TimeSeriesTable::create(TableLocation::local(temp.path()), timestamp_only_meta())
2520                    .await?;
2521            crate::storage::inject_write_new_failure(temp.path().join(sidecar_dir), true);
2522
2523            let error = table
2524                .append(timestamp_only_batch(1)?)
2525                .await
2526                .expect_err("sidecar write and its first cleanup must fail");
2527
2528            assert!(matches!(
2529                error,
2530                TableError::Append {
2531                    source: AppendError::CoverageSidecar { source }
2532                } if matches!(
2533                    source.as_ref(),
2534                    CoverageSidecarError::Storage {
2535                        source: StorageError::CleanupFailed { .. }
2536                    }
2537                )
2538            ));
2539            assert_eq!(table.state().version, 1, "{sidecar_dir}");
2540            assert!(table.state().segments.is_empty(), "{sidecar_dir}");
2541            assert!(data_files(temp.path())?.is_empty(), "{sidecar_dir}");
2542            assert!(coverage_files(temp.path())?.is_empty(), "{sidecar_dir}");
2543            assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
2544        }
2545        Ok(())
2546    }
2547
2548    #[tokio::test]
2549    async fn append_conflict_cleans_attempt_and_stays_invisible() -> TestResult {
2550        let temp = TempDir::new()?;
2551        let location = TableLocation::local(temp.path());
2552        let mut winner = TimeSeriesTable::create(location.clone(), widening_table_meta()).await?;
2553        let mut loser = TimeSeriesTable::open(location.clone()).await?;
2554        let loser_state_before = loser.state().clone();
2555
2556        winner
2557            .append(widening_batch(
2558                vec![0],
2559                vec![1],
2560                vec![Some(1.0)],
2561                vec!["winner"],
2562            )?)
2563            .await?;
2564        let data_before = data_files(temp.path())?;
2565        let coverage_before = coverage_files(temp.path())?;
2566
2567        let error = loser
2568            .append(
2569                AppendRequest::new(widening_batch(
2570                    vec![u32::MAX],
2571                    vec![2],
2572                    vec![Some(2.0)],
2573                    vec!["loser"],
2574                )?)
2575                .max_rows_per_row_group(1),
2576            )
2577            .await
2578            .expect_err("stale streaming append must conflict");
2579
2580        assert!(matches!(
2581            &error,
2582            TableError::Append {
2583                source: AppendError::Commit {
2584                    source: CommitError::Conflict {
2585                        expected: 1,
2586                        found: 2,
2587                        ..
2588                    }
2589                }
2590            }
2591        ));
2592        assert_eq!(loser.state(), &loser_state_before);
2593        assert_eq!(data_files(temp.path())?, data_before);
2594        assert_eq!(coverage_files(temp.path())?, coverage_before);
2595        assert!(!temp.path().join(layout::commit_rel_path(3)).exists());
2596
2597        let reopened = TimeSeriesTable::open(location).await?;
2598        assert_eq!(reopened.state().version, 2);
2599        assert_eq!(reopened.state().segments.len(), 1);
2600        assert_eq!(
2601            reopened
2602                .state()
2603                .segments
2604                .values()
2605                .next()
2606                .map(|segment| { temp.path().join(&segment.path) }),
2607            data_before.first().cloned()
2608        );
2609        Ok(())
2610    }
2611
2612    #[tokio::test]
2613    async fn append_ambiguous_commit_reports_and_preserves_generated_path() -> TestResult {
2614        let temp = TempDir::new()?;
2615        let location = TableLocation::local(temp.path());
2616        let mut table = TimeSeriesTable::create(location.clone(), widening_table_meta()).await?;
2617        let commit_path = temp.path().join(layout::commit_rel_path(2));
2618        crate::storage::inject_write_new_failure(commit_path.clone(), true);
2619        let capture = TraceCapture::default();
2620
2621        let error = capture
2622            .run(
2623                table.append(
2624                    AppendRequest::new(widening_batch(
2625                        vec![0],
2626                        vec![1],
2627                        vec![Some(1.0)],
2628                        vec!["ambiguous"],
2629                    )?)
2630                    .max_rows_per_row_group(1),
2631                ),
2632            )
2633            .await
2634            .expect_err("failed commit cleanup must be ambiguous");
2635        let TableError::Append {
2636            source:
2637                AppendError::CommitAmbiguous {
2638                    segment_path,
2639                    source,
2640                },
2641        } = error
2642        else {
2643            panic!("unexpected error: {error}");
2644        };
2645
2646        assert!(matches!(
2647            source.as_ref(),
2648            CommitError::AmbiguousOutcome { .. }
2649        ));
2650        assert_eq!(
2651            captured_span(&capture, "table.append").target,
2652            "timeseries_table_format::table::append"
2653        );
2654        assert!(segment_path.starts_with("data/"));
2655        assert!(temp.path().join(&segment_path).is_file());
2656        assert_eq!(
2657            data_files(temp.path())?,
2658            vec![temp.path().join(&segment_path)]
2659        );
2660        assert_eq!(table.state().version, 1);
2661        assert!(table.state().segments.is_empty());
2662        assert_eq!(table.log.load_current_version().await?, 1);
2663        assert!(commit_path.exists());
2664        assert_eq!(coverage_files(temp.path())?.len(), 2);
2665        let append_span = capture
2666            .spans()
2667            .into_iter()
2668            .find(|span| span.name == "table.append")
2669            .expect("table.append span");
2670        assert_eq!(append_span.fields.get("segment_path"), Some(&segment_path));
2671        assert_eq!(
2672            append_span.fields.get("outcome").map(String::as_str),
2673            Some("ambiguous")
2674        );
2675
2676        let reopened = TimeSeriesTable::open(location).await?;
2677        assert_eq!(reopened.state().version, 1);
2678        assert!(reopened.state().segments.is_empty());
2679        Ok(())
2680    }
2681
2682    #[tokio::test]
2683    async fn append_conflict_reports_every_failed_cleanup_path() -> TestResult {
2684        let temp = TempDir::new()?;
2685        let location = TableLocation::local(temp.path());
2686        let mut winner = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
2687        let mut loser = TimeSeriesTable::open(location.clone()).await?;
2688        winner
2689            .append(time_series_batch(vec![0], vec!["A"], vec![1.0])?)
2690            .await?;
2691        let data_before = data_files(temp.path())?;
2692        let coverage_before = coverage_files(temp.path())?;
2693
2694        for dir in [
2695            Path::new("data"),
2696            Path::new(SEGMENT_COVERAGE_DIR),
2697            Path::new(TABLE_SNAPSHOT_DIR),
2698        ] {
2699            crate::storage::inject_cleanup_failure(temp.path().join(dir));
2700        }
2701        let error = loser
2702            .append(time_series_batch(vec![120_000], vec!["A"], vec![2.0])?)
2703            .await
2704            .expect_err("conflict cleanup failures must be reported");
2705        let message = error.to_string();
2706
2707        let (source, cleanup_errors) = match &error {
2708            TableError::Append {
2709                source:
2710                    AppendError::Rollback {
2711                        source,
2712                        cleanup_errors,
2713                    },
2714            } => (source, cleanup_errors),
2715            other => panic!("unexpected error: {other:?}"),
2716        };
2717        assert!(matches!(
2718            source.as_ref(),
2719            AppendError::Commit {
2720                source: CommitError::Conflict { .. }
2721            }
2722        ));
2723        assert_eq!(cleanup_errors.len(), 3);
2724
2725        assert!(message.contains("Commit conflict"));
2726        let data_after = data_files(temp.path())?;
2727        assert_eq!(data_after.len(), data_before.len() + 1);
2728        let coverage_after = coverage_files(temp.path())?;
2729        assert_eq!(coverage_after.len(), coverage_before.len() + 2);
2730        let mut failed_paths = data_after
2731            .iter()
2732            .filter(|path| !data_before.contains(path))
2733            .cloned()
2734            .collect::<Vec<_>>();
2735        failed_paths.extend(
2736            coverage_after
2737                .keys()
2738                .filter(|path| !coverage_before.contains_key(*path))
2739                .map(|path| temp.path().join(path)),
2740        );
2741        for path in failed_paths {
2742            let filename = path
2743                .file_name()
2744                .expect("failed cleanup filename")
2745                .to_str()
2746                .expect("UTF-8 cleanup filename");
2747            assert!(
2748                cleanup_errors.iter().any(|error| matches!(
2749                    error,
2750                    StorageError::OtherIo { path, .. } if path.ends_with(filename)
2751                )),
2752                "missing typed cleanup failure for {}",
2753                path.display()
2754            );
2755            assert!(
2756                message.contains(filename),
2757                "missing cleanup diagnostic for {}",
2758                path.display()
2759            );
2760        }
2761
2762        let reopened = TimeSeriesTable::open(location).await?;
2763        assert_eq!(reopened.state().version, 2);
2764        assert_eq!(reopened.state().segments.len(), 1);
2765        Ok(())
2766    }
2767
2768    #[tokio::test]
2769    async fn append_validates_schema_before_creating_data() -> TestResult {
2770        let temp = TempDir::new()?;
2771        let mut table =
2772            TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2773                .await?;
2774
2775        assert!(matches!(
2776            table.append(input_batch(vec![1])?).await,
2777            Err(TableError::Append {
2778                source: AppendError::SchemaValidation { .. }
2779            })
2780        ));
2781        assert_eq!(table.state().version, 1);
2782        assert!(!temp.path().join("data").exists());
2783        Ok(())
2784    }
2785
2786    #[tokio::test]
2787    async fn append_removes_owned_data_after_overlap() -> TestResult {
2788        let temp = TempDir::new()?;
2789        let mut table =
2790            TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2791                .await?;
2792        let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2793        assert_eq!(table.append(batch.clone()).await?, 2);
2794
2795        assert!(matches!(
2796            table
2797                .append(AppendRequest::new(batch).max_rows_per_row_group(1))
2798                .await,
2799            Err(TableError::Append {
2800                source: AppendError::PersistedIndexIntervalOverlap { .. }
2801            })
2802        ));
2803        assert_eq!(table.state().version, 2);
2804        assert_eq!(
2805            std::fs::read_dir(temp.path().join("data"))?
2806                .collect::<Result<Vec<_>, _>>()?
2807                .len(),
2808            1
2809        );
2810        Ok(())
2811    }
2812
2813    #[tokio::test]
2814    async fn append_adopts_schema_on_first_append() -> TestResult {
2815        let temp = TempDir::new()?;
2816        let mut meta = make_basic_table_meta();
2817        meta.logical_schema = None;
2818        let mut table = TimeSeriesTable::create(TableLocation::local(temp.path()), meta).await?;
2819
2820        let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2821        assert_eq!(table.append(batch).await?, 2);
2822        assert!(table.state().table_meta.logical_schema.is_some());
2823        Ok(())
2824    }
2825
2826    #[tokio::test]
2827    async fn append_preserves_exact_schema_through_optimize() -> TestResult {
2828        let timezone = "America/Phoenix";
2829        let schema = Arc::new(Schema::new(vec![
2830            Field::new(
2831                "ts",
2832                DataType::Timestamp(ArrowTimeUnit::Millisecond, Some(timezone.into())),
2833                false,
2834            ),
2835            Field::new("symbol", DataType::Utf8, false),
2836            Field::new(
2837                "items",
2838                DataType::List(Arc::new(Field::new("item", DataType::Int64, true))),
2839                true,
2840            ),
2841            Field::new(
2842                "attrs",
2843                DataType::Map(
2844                    Arc::new(Field::new(
2845                        "entries",
2846                        DataType::Struct(Fields::from(vec![
2847                            Field::new("key", DataType::Utf8, false),
2848                            Field::new("value", DataType::Binary, true),
2849                        ])),
2850                        false,
2851                    )),
2852                    true,
2853                ),
2854                true,
2855            ),
2856        ]));
2857        let expected = LogicalSchema::try_from_arrow_schema(&schema).expect("logical schema");
2858        let batch = RecordBatch::try_new(
2859            Arc::clone(&schema),
2860            vec![
2861                Arc::new(TimestampMillisecondArray::from(vec![0, 0]).with_timezone(timezone)),
2862                Arc::new(StringArray::from(vec!["A", "B"])),
2863                new_null_array(schema.field(2).data_type(), 2),
2864                new_null_array(schema.field(3).data_type(), 2),
2865            ],
2866        )?;
2867
2868        for has_canonical_schema in [true, false] {
2869            let temp = TempDir::new()?;
2870            let location = TableLocation::local(temp.path());
2871            let mut meta = TableMeta::new_time_series_with_schema(
2872                IndexSpec {
2873                    column: "ts".to_string(),
2874                    entity_columns: vec!["symbol".to_string()],
2875                    kind: IndexKind::Timestamp {
2876                        index_granularity: TimeIndexGranularity::Minutes(1),
2877                        timezone: Some(timezone.to_string()),
2878                    },
2879                },
2880                expected.clone(),
2881            );
2882            if !has_canonical_schema {
2883                meta.logical_schema = None;
2884            }
2885            let mut table = TimeSeriesTable::create(location.clone(), meta).await?;
2886
2887            assert_eq!(table.append(batch.clone()).await?, 2);
2888            assert_eq!(
2889                table.state().table_meta.logical_schema.as_ref(),
2890                Some(&expected)
2891            );
2892
2893            let later_batch = RecordBatch::try_new(
2894                Arc::clone(&schema),
2895                vec![
2896                    Arc::new(
2897                        TimestampMillisecondArray::from(vec![120_000, 120_000])
2898                            .with_timezone(timezone),
2899                    ),
2900                    Arc::new(StringArray::from(vec!["A", "B"])),
2901                    new_null_array(schema.field(2).data_type(), 2),
2902                    new_null_array(schema.field(3).data_type(), 2),
2903                ],
2904            )?;
2905            assert_eq!(table.append(later_batch).await?, 3);
2906
2907            let report = table.optimize().await?;
2908            assert_eq!(report.committed_version, 4);
2909            assert_eq!(report.candidate_source_segments, 2);
2910            assert_eq!(report.replacement_segments_written, 4);
2911            assert_eq!(report.rows_read, 4);
2912            assert_eq!(report.rows_written, 4);
2913            assert!(
2914                table
2915                    .state()
2916                    .segments
2917                    .values()
2918                    .all(|segment| matches!(segment.entity_layout, SegmentEntityLayout::Single(_)))
2919            );
2920            let reopened = TimeSeriesTable::open(location).await?;
2921            assert_eq!(
2922                reopened.state().table_meta.logical_schema.as_ref(),
2923                Some(&expected)
2924            );
2925        }
2926        Ok(())
2927    }
2928
2929    fn registered_index(kind: IndexKind) -> IndexSpec {
2930        IndexSpec {
2931            column: "ts".to_string(),
2932            entity_columns: Vec::new(),
2933            kind,
2934        }
2935    }
2936
2937    fn write_single_index_parquet(
2938        path: &Path,
2939        data_type: DataType,
2940        values: ArrayRef,
2941    ) -> TestResult {
2942        if let Some(parent) = path.parent() {
2943            std::fs::create_dir_all(parent)?;
2944        }
2945        let schema = Arc::new(Schema::new(vec![Field::new("ts", data_type, false)]));
2946        let batch = RecordBatch::try_new(Arc::clone(&schema), vec![values])?;
2947        let mut writer = ArrowWriter::try_new(File::create(path)?, schema, None)?;
2948        writer.write(&batch)?;
2949        writer.close()?;
2950        Ok(())
2951    }
2952
2953    fn composite_entity_batch(rows: &[(i64, &str, &str, f64)]) -> Result<RecordBatch, ArrowError> {
2954        let schema = Arc::new(Schema::new(vec![
2955            Field::new(
2956                "ts",
2957                DataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
2958                false,
2959            ),
2960            Field::new("symbol", DataType::Utf8, false),
2961            Field::new("venue", DataType::Utf8, false),
2962            Field::new("price", DataType::Float64, false),
2963        ]));
2964        RecordBatch::try_new(
2965            Arc::clone(&schema),
2966            vec![
2967                Arc::new(TimestampMillisecondArray::from(
2968                    rows.iter().map(|row| row.0).collect::<Vec<_>>(),
2969                )),
2970                Arc::new(StringArray::from(
2971                    rows.iter().map(|row| row.1).collect::<Vec<_>>(),
2972                )),
2973                Arc::new(StringArray::from(
2974                    rows.iter().map(|row| row.2).collect::<Vec<_>>(),
2975                )),
2976                Arc::new(Float64Array::from(
2977                    rows.iter().map(|row| row.3).collect::<Vec<_>>(),
2978                )),
2979            ],
2980        )
2981    }
2982
2983    fn write_composite_entity_parquet(path: &Path, rows: &[(i64, &str, &str, f64)]) -> TestResult {
2984        if let Some(parent) = path.parent() {
2985            std::fs::create_dir_all(parent)?;
2986        }
2987        let batch = composite_entity_batch(rows)?;
2988        let schema = batch.schema();
2989        let mut writer = ArrowWriter::try_new(File::create(path)?, schema, None)?;
2990        writer.write(&batch)?;
2991        writer.close()?;
2992        Ok(())
2993    }
2994
2995    fn coverage_files(root: &Path) -> std::io::Result<BTreeMap<PathBuf, Vec<u8>>> {
2996        let mut files = BTreeMap::new();
2997        for rel_dir in [SEGMENT_COVERAGE_DIR, TABLE_SNAPSHOT_DIR] {
2998            let dir = root.join(rel_dir);
2999            if !dir.exists() {
3000                continue;
3001            }
3002            for entry in std::fs::read_dir(dir)? {
3003                let path = entry?.path();
3004                if path.is_file() {
3005                    files.insert(
3006                        path.strip_prefix(root)
3007                            .expect("coverage path under root")
3008                            .to_owned(),
3009                        std::fs::read(path)?,
3010                    );
3011                }
3012            }
3013        }
3014        Ok(files)
3015    }
3016
3017    fn data_files(root: &Path) -> std::io::Result<Vec<PathBuf>> {
3018        let dir = root.join("data");
3019        if !dir.exists() {
3020            return Ok(Vec::new());
3021        }
3022        let mut files = std::fs::read_dir(dir)?
3023            .map(|entry| entry.map(|entry| entry.path()))
3024            .collect::<Result<Vec<_>, _>>()?;
3025        files.sort();
3026        Ok(files)
3027    }
3028
3029    #[test]
3030    fn entity_layout_classification_rejects_empty_coverage() {
3031        assert!(matches!(
3032            classify_entity_layout("data/empty.parquet", &EntityCoverage::empty()),
3033            Err(AppendError::EmptySegmentEntityCoverage { segment_path })
3034                if segment_path == "data/empty.parquet"
3035        ));
3036    }
3037
3038    #[tokio::test]
3039    async fn duplicate_implicit_interval_rolls_back_first_append() -> TestResult {
3040        let temp = TempDir::new()?;
3041        let location = TableLocation::local(temp.path());
3042        let mut table = TimeSeriesTable::create(
3043            location.clone(),
3044            TableMeta::new_time_series(timestamp_only_index()),
3045        )
3046        .await?;
3047        let state_before = table.state().clone();
3048
3049        let error = table
3050            .append(timestamp_only_batch_with_millis([0, 30_000])?)
3051            .await
3052            .expect_err("duplicate implicit interval must fail");
3053
3054        let message = error.to_string();
3055        match error {
3056            TableError::Append {
3057                source: AppendError::GeneratedSegmentCoverage { source, .. },
3058            } => match *source {
3059                SegmentCoverageError::DuplicateIndexInterval {
3060                    path: segment_path,
3061                    example_identity: None,
3062                    example_index_interval,
3063                } => {
3064                    assert!(segment_path.starts_with("data/"));
3065                    assert!(segment_path.ends_with(".parquet"));
3066                    assert!(message.contains(&segment_path));
3067                    assert!(message.contains("Duplicate ordered-index interval"));
3068                    assert!(message.contains("ordered-index interval"));
3069                    assert_eq!(
3070                        example_index_interval.to_string(),
3071                        "[1970-01-01T00:00:00Z, 1970-01-01T00:01:00Z)"
3072                    );
3073                }
3074                other => panic!("expected duplicate interval source, got {other:?}"),
3075            },
3076            other => panic!("expected duplicate interval error, got {other:?}"),
3077        }
3078        assert_eq!(table.state(), &state_before);
3079        assert_eq!(table.log.load_current_version().await?, 1);
3080        assert!(data_files(temp.path())?.is_empty());
3081        assert!(coverage_files(temp.path())?.is_empty());
3082        assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
3083        assert_eq!(
3084            TimeSeriesTable::open(location).await?.state(),
3085            &state_before
3086        );
3087
3088        assert_eq!(table.append(timestamp_only_batch(1)?).await?, 2);
3089        assert!(table.state().table_meta.logical_schema.is_some());
3090        Ok(())
3091    }
3092
3093    #[tokio::test]
3094    async fn duplicate_interval_across_input_batches_is_rejected() -> TestResult {
3095        let temp = TempDir::new()?;
3096        let location = TableLocation::local(temp.path());
3097        let mut table = TimeSeriesTable::create(location.clone(), timestamp_only_meta()).await?;
3098        let state_before = table.state().clone();
3099
3100        let error = table
3101            .append(
3102                AppendRequest::new(vec![
3103                    timestamp_only_batch_with_millis([0])?,
3104                    timestamp_only_batch_with_millis([30_000])?,
3105                ])
3106                .max_rows_per_row_group(1),
3107            )
3108            .await
3109            .expect_err("duplicate split across input batches and row groups must fail");
3110
3111        assert!(matches!(
3112            error,
3113            TableError::Append {
3114                source: AppendError::GeneratedSegmentCoverage { source, .. },
3115            } if matches!(source.as_ref(),
3116                SegmentCoverageError::DuplicateIndexInterval {
3117                    example_identity: None,
3118                    example_index_interval,
3119                    ..
3120                } if example_index_interval.to_string()
3121                    == "[1970-01-01T00:00:00Z, 1970-01-01T00:01:00Z)"
3122            )
3123        ));
3124        assert_eq!(table.state(), &state_before);
3125        assert!(data_files(temp.path())?.is_empty());
3126        assert!(coverage_files(temp.path())?.is_empty());
3127        assert_eq!(
3128            TimeSeriesTable::open(location).await?.state(),
3129            &state_before
3130        );
3131        Ok(())
3132    }
3133
3134    #[tokio::test]
3135    async fn composite_identity_uses_every_component_for_duplicates() -> TestResult {
3136        let temp = TempDir::new()?;
3137        let index = IndexSpec {
3138            column: "ts".to_string(),
3139            entity_columns: vec!["symbol".to_string(), "venue".to_string()],
3140            kind: IndexKind::Timestamp {
3141                index_granularity: TimeIndexGranularity::Minutes(1),
3142                timezone: None,
3143            },
3144        };
3145        let mut table = TimeSeriesTable::create(
3146            TableLocation::local(temp.path()),
3147            TableMeta::new_time_series(index),
3148        )
3149        .await?;
3150
3151        assert_eq!(
3152            table
3153                .append(composite_entity_batch(&[
3154                    (0, "A", "X", 1.0),
3155                    (0, "A", "Y", 2.0),
3156                ])?)
3157                .await?,
3158            2
3159        );
3160        let state_before = table.state().clone();
3161
3162        let error = table
3163            .append(composite_entity_batch(&[
3164                (60_000, "A", "Z", 3.0),
3165                (90_000, "A", "Z", 4.0),
3166            ])?)
3167            .await
3168            .expect_err("matching composite identity must be rejected");
3169
3170        assert!(matches!(
3171            error,
3172            TableError::Append {
3173                source: AppendError::GeneratedSegmentCoverage { source, .. },
3174            } if matches!(source.as_ref(),
3175                SegmentCoverageError::DuplicateIndexInterval {
3176                    example_identity: Some(example_identity),
3177                    ..
3178                } if example_identity.components()
3179                    == [EntityValue::from("A"), EntityValue::from("Z")]
3180            )
3181        ));
3182        assert_eq!(table.state(), &state_before);
3183        Ok(())
3184    }
3185
3186    #[tokio::test]
3187    async fn duplicate_entity_interval_leaves_nonempty_table_unchanged() -> TestResult {
3188        let temp = TempDir::new()?;
3189        let location = TableLocation::local(temp.path());
3190        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
3191        assert_eq!(
3192            table
3193                .append(time_series_batch(vec![0], vec!["A"], vec![1.0])?)
3194                .await?,
3195            2
3196        );
3197        let state_before = table.state().clone();
3198        let data_before = data_files(temp.path())?;
3199        let coverage_before = coverage_files(temp.path())?;
3200
3201        let error = table
3202            .append(
3203                AppendRequest::new(time_series_batch(
3204                    vec![60_000, 90_000],
3205                    vec!["A", "A"],
3206                    vec![2.0, 3.0],
3207                )?)
3208                .max_rows_per_row_group(1),
3209            )
3210            .await
3211            .expect_err("duplicate entity interval must fail");
3212        let expected_identity = EntityIdentity::try_new(vec!["A".into()])?;
3213
3214        assert!(matches!(
3215            error,
3216            TableError::Append {
3217                source: AppendError::GeneratedSegmentCoverage { source, .. },
3218            } if matches!(source.as_ref(),
3219                SegmentCoverageError::DuplicateIndexInterval {
3220                    example_identity: Some(example_identity),
3221                    example_index_interval,
3222                    ..
3223                } if example_identity == &expected_identity
3224                    && example_index_interval.to_string()
3225                        == "[1970-01-01T00:01:00Z, 1970-01-01T00:02:00Z)"
3226            )
3227        ));
3228        assert_eq!(table.state(), &state_before);
3229        assert_eq!(table.log.load_current_version().await?, 2);
3230        assert_eq!(data_files(temp.path())?, data_before);
3231        assert_eq!(coverage_files(temp.path())?, coverage_before);
3232        assert!(!temp.path().join(layout::commit_rel_path(3)).exists());
3233        assert_eq!(
3234            TimeSeriesTable::open(location).await?.state(),
3235            &state_before
3236        );
3237        Ok(())
3238    }
3239
3240    #[tokio::test]
3241    async fn append_updates_state_and_log() -> TestResult {
3242        let tmp = TempDir::new()?;
3243        let location = TableLocation::local(tmp.path());
3244        let meta = make_basic_table_meta();
3245
3246        let mut table = TimeSeriesTable::create(location.clone(), meta).await?;
3247
3248        let rel_path = "data/seg1.parquet";
3249        let abs_path = tmp.path().join(rel_path);
3250        write_test_parquet(
3251            &abs_path,
3252            true,
3253            false,
3254            &[TestRow {
3255                ts_millis: 1_000,
3256                symbol: "A",
3257                price: 10.0,
3258            }],
3259        )?;
3260
3261        let new_version = append_parquet_fixture(&mut table, rel_path).await?;
3262
3263        assert_eq!(new_version, 2);
3264        assert_eq!(table.state.version, 2);
3265        let seg = table
3266            .state
3267            .segments
3268            .values()
3269            .next()
3270            .expect("segment present");
3271        assert_eq!(seg.row_count, 1);
3272        assert_eq!(
3273            seg.entity_layout,
3274            SegmentEntityLayout::Single(EntityIdentity::try_new(vec!["A".into()])?)
3275        );
3276        assert!(matches!(
3277            &seg.index_min,
3278            IndexValue::Timestamp(value) if value.timestamp_millis() == 1_000
3279        ));
3280        assert!(matches!(
3281            &seg.index_max,
3282            IndexValue::Timestamp(value) if value.timestamp_millis() == 1_000
3283        ));
3284
3285        let commit_path = tmp.path().join(layout::commit_rel_path(2));
3286        assert!(commit_path.is_file());
3287        let current =
3288            tokio::fs::read_to_string(tmp.path().join(layout::current_rel_path())).await?;
3289        assert_eq!(current.trim(), "2");
3290
3291        let reopened = TimeSeriesTable::open(location).await?;
3292        assert_eq!(reopened.state, table.state);
3293        Ok(())
3294    }
3295
3296    #[tokio::test]
3297    async fn append_without_entities_uses_global_coverage() -> TestResult {
3298        let tmp = TempDir::new()?;
3299        let location = TableLocation::local(tmp.path());
3300        let index = registered_index(IndexKind::Int64 {
3301            index_granularity: NonZeroU64::new(10).unwrap(),
3302        });
3303        let mut table =
3304            TimeSeriesTable::create(location.clone(), TableMeta::new_time_series(index.clone()))
3305                .await?;
3306        let rel_path = "data/int64.parquet";
3307        write_arrow_parquet_int_time(
3308            &tmp.path().join(rel_path),
3309            &[i64::MIN, -1, 0, i64::MAX],
3310            &["A", "A", "A", "A"],
3311            &[1.0, 2.0, 3.0, 4.0],
3312        )?;
3313
3314        assert_eq!(append_parquet_fixture(&mut table, rel_path).await?, 2);
3315
3316        let segment = table
3317            .state
3318            .segments
3319            .values()
3320            .next()
3321            .expect("segment present");
3322        assert_eq!(segment.entity_layout, SegmentEntityLayout::NotApplicable);
3323        assert_eq!(segment.index_min, IndexValue::Int64(i64::MIN));
3324        assert_eq!(segment.index_max, IndexValue::Int64(i64::MAX));
3325        let pointer = table.state.table_coverage.as_ref().expect("table coverage");
3326        assert_eq!(pointer.index_kind, index.kind);
3327        let persisted = read_coverage_sidecar(&location, Path::new(&pointer.coverage_path)).await?;
3328        let expected =
3329            compute_segment_coverage(&location, Path::new(&segment.path), table.index_spec())
3330                .await?;
3331        assert_eq!(persisted, expected);
3332        let reopened = TimeSeriesTable::open(location).await?;
3333        assert_eq!(reopened.state, table.state);
3334        Ok(())
3335    }
3336
3337    #[tokio::test]
3338    async fn int64_appends_enforce_coverage_and_exact_later_schema() -> TestResult {
3339        let tmp = TempDir::new()?;
3340        let location = TableLocation::local(tmp.path());
3341        let index = registered_index(IndexKind::Int64 {
3342            index_granularity: NonZeroU64::new(10).unwrap(),
3343        });
3344        let mut table =
3345            TimeSeriesTable::create(location, TableMeta::new_time_series(index)).await?;
3346
3347        for (path, values) in [
3348            ("data/negative.parquet", &[-25, -15][..]),
3349            ("data/positive.parquet", &[5, 15][..]),
3350        ] {
3351            write_arrow_parquet_int_time(&tmp.path().join(path), values, &["A", "A"], &[1.0, 2.0])?;
3352            append_parquet_fixture(&mut table, path).await?;
3353        }
3354        assert_eq!(table.state.version, 3);
3355
3356        let state_before = table.state.clone();
3357        let coverage_before = coverage_files(tmp.path())?;
3358        let overlap_path = "data/negative-overlap.parquet";
3359        write_arrow_parquet_int_time(&tmp.path().join(overlap_path), &[-19], &["A"], &[3.0])?;
3360        let overlap_error = append_parquet_fixture(&mut table, overlap_path)
3361            .await
3362            .expect_err("negative index interval overlap must fail");
3363        assert!(matches!(
3364            &overlap_error,
3365            TableError::Append {
3366                source: AppendError::PersistedIndexIntervalOverlap {
3367                    example_identity: None,
3368                    example_index_interval,
3369                    ..
3370                }
3371            } if example_index_interval.to_string() == "[-20, -10)"
3372        ));
3373        assert!(
3374            overlap_error
3375                .to_string()
3376                .contains("example_index_interval=[-20, -10)")
3377        );
3378
3379        let mismatch_path = "data/schema-mismatch.parquet";
3380        write_single_index_parquet(
3381            &tmp.path().join(mismatch_path),
3382            DataType::Int64,
3383            Arc::new(Int64Array::from(vec![100])),
3384        )?;
3385        assert!(matches!(
3386            append_parquet_fixture(&mut table, mismatch_path)
3387                .await
3388                .expect_err("later schema mismatch must fail"),
3389            TableError::Append {
3390                source: AppendError::SchemaValidation { .. }
3391            }
3392        ));
3393        assert_eq!(table.state, state_before);
3394        assert_eq!(coverage_files(tmp.path())?, coverage_before);
3395        Ok(())
3396    }
3397
3398    #[tokio::test]
3399    async fn append_supports_registered_uint64_index() -> TestResult {
3400        let tmp = TempDir::new()?;
3401        let location = TableLocation::local(tmp.path());
3402        let index = registered_index(IndexKind::UInt64 {
3403            index_granularity: NonZeroU64::new(10).unwrap(),
3404        });
3405        let mut table =
3406            TimeSeriesTable::create(location.clone(), TableMeta::new_time_series(index.clone()))
3407                .await?;
3408        let rel_path = "data/uint64.parquet";
3409        write_single_index_parquet(
3410            &tmp.path().join(rel_path),
3411            DataType::UInt64,
3412            Arc::new(UInt64Array::from(vec![0, i64::MAX as u64 + 1, u64::MAX])),
3413        )?;
3414
3415        assert_eq!(append_parquet_fixture(&mut table, rel_path).await?, 2);
3416
3417        let segment = table
3418            .state
3419            .segments
3420            .values()
3421            .next()
3422            .expect("segment present");
3423        assert_eq!(segment.index_min, IndexValue::UInt64(0));
3424        assert_eq!(segment.index_max, IndexValue::UInt64(u64::MAX));
3425        assert_eq!(
3426            table
3427                .state
3428                .table_meta
3429                .logical_schema
3430                .as_ref()
3431                .expect("schema adopted")
3432                .columns()[0]
3433                .data_type,
3434            LogicalDataType::UInt64
3435        );
3436        assert_eq!(
3437            table
3438                .state
3439                .table_coverage
3440                .as_ref()
3441                .expect("table coverage")
3442                .index_kind,
3443            index.kind
3444        );
3445
3446        let non_overlap_path = "data/uint64-non-overlap.parquet";
3447        write_single_index_parquet(
3448            &tmp.path().join(non_overlap_path),
3449            DataType::UInt64,
3450            Arc::new(UInt64Array::from(vec![u64::MAX - 20])),
3451        )?;
3452        assert_eq!(
3453            append_parquet_fixture(&mut table, non_overlap_path).await?,
3454            3
3455        );
3456
3457        let state_before = table.state.clone();
3458        let coverage_before = coverage_files(tmp.path())?;
3459        let overlap_path = "data/uint64-overlap.parquet";
3460        write_single_index_parquet(
3461            &tmp.path().join(overlap_path),
3462            DataType::UInt64,
3463            Arc::new(UInt64Array::from(vec![u64::MAX - 1])),
3464        )?;
3465        let overlap_error = append_parquet_fixture(&mut table, overlap_path)
3466            .await
3467            .expect_err("large uint64 index interval overlap must fail");
3468        assert!(matches!(
3469            overlap_error,
3470            TableError::Append {
3471                source: AppendError::PersistedIndexIntervalOverlap {
3472                    example_identity: None,
3473                    example_index_interval,
3474                    ..
3475                }
3476            } if example_index_interval.to_string()
3477                == "[18446744073709551610, 18446744073709551615]"
3478        ));
3479        assert_eq!(table.state, state_before);
3480        assert_eq!(coverage_files(tmp.path())?, coverage_before);
3481        let reopened = TimeSeriesTable::open(location).await?;
3482        assert_eq!(reopened.state, table.state);
3483        Ok(())
3484    }
3485
3486    #[tokio::test]
3487    async fn append_rejects_signed_data_for_uint64_index_without_mutation() -> TestResult {
3488        let tmp = TempDir::new()?;
3489        let location = TableLocation::local(tmp.path());
3490        let index = registered_index(IndexKind::UInt64 {
3491            index_granularity: NonZeroU64::new(1).unwrap(),
3492        });
3493        let mut table =
3494            TimeSeriesTable::create(location.clone(), TableMeta::new_time_series(index)).await?;
3495        let state_before = table.state.clone();
3496        let coverage_before = coverage_files(tmp.path())?;
3497        let rel_path = "data/signed.parquet";
3498        write_arrow_parquet_int_time(&tmp.path().join(rel_path), &[1], &["A"], &[1.0])?;
3499
3500        let error = append_parquet_fixture(&mut table, rel_path)
3501            .await
3502            .expect_err("signed data must not append to a uint64 index");
3503
3504        assert!(
3505            matches!(
3506                error,
3507                TableError::Append {
3508                    source: AppendError::SchemaValidation {
3509                        ref source,
3510                        ..
3511                    }
3512                } if matches!(
3513                    source.as_ref(),
3514                    crate::metadata::schema_compat::SchemaCompatibilityError::IndexKindMismatch {
3515                        expected: "uint64",
3516                        actual: LogicalDataType::Int64,
3517                        ..
3518                    }
3519                )
3520            ),
3521            "unexpected error: {error:?}"
3522        );
3523        assert_eq!(table.state, state_before);
3524        assert_eq!(table.log.load_current_version().await?, 1);
3525        assert_eq!(coverage_files(tmp.path())?, coverage_before);
3526        Ok(())
3527    }
3528
3529    #[tokio::test]
3530    async fn append_records_single_layout_for_each_entity_segment() -> TestResult {
3531        let tmp = TempDir::new()?;
3532        let location = TableLocation::local(tmp.path());
3533        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
3534
3535        for (path, symbol) in [
3536            ("data/entity-a.parquet", "A"),
3537            ("data/entity-b.parquet", "B"),
3538        ] {
3539            write_test_parquet(
3540                &tmp.path().join(path),
3541                true,
3542                false,
3543                &[TestRow {
3544                    ts_millis: 1_000,
3545                    symbol,
3546                    price: 10.0,
3547                }],
3548            )?;
3549            append_parquet_fixture(&mut table, path).await?;
3550            let expected_layout =
3551                SegmentEntityLayout::Single(EntityIdentity::try_new(vec![symbol.into()])?);
3552            assert!(
3553                table
3554                    .state
3555                    .segments
3556                    .values()
3557                    .any(|segment| segment.entity_layout == expected_layout)
3558            );
3559        }
3560
3561        assert_eq!(table.state.version, 3);
3562        let pointer = table
3563            .state
3564            .table_coverage
3565            .as_ref()
3566            .expect("table coverage pointer");
3567        let coverage =
3568            read_entity_coverage_sidecar(&location, Path::new(&pointer.coverage_path)).await?;
3569        assert_eq!(coverage.identity_count(), 2);
3570        assert_eq!(coverage.cardinality(), 2);
3571        Ok(())
3572    }
3573
3574    #[tokio::test]
3575    async fn append_records_mixed_layout_for_multiple_identities() -> TestResult {
3576        let tmp = TempDir::new()?;
3577        let location = TableLocation::local(tmp.path());
3578        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
3579        let path = "data/multiple-identities.parquet";
3580        write_test_parquet(
3581            &tmp.path().join(path),
3582            true,
3583            false,
3584            &[
3585                TestRow {
3586                    ts_millis: 1_000,
3587                    symbol: "A",
3588                    price: 10.0,
3589                },
3590                TestRow {
3591                    ts_millis: 1_000,
3592                    symbol: "B",
3593                    price: 20.0,
3594                },
3595            ],
3596        )?;
3597
3598        append_parquet_fixture(&mut table, path).await?;
3599
3600        let segment = table
3601            .state
3602            .segments
3603            .values()
3604            .next()
3605            .expect("segment present");
3606        assert_eq!(segment.entity_layout, SegmentEntityLayout::Mixed);
3607        let coverage = read_entity_coverage_sidecar(
3608            &location,
3609            Path::new(segment.coverage_path.as_ref().expect("coverage path")),
3610        )
3611        .await?;
3612        assert_eq!(coverage.identity_count(), 2);
3613        assert_eq!(coverage.cardinality(), 2);
3614        Ok(())
3615    }
3616
3617    #[tokio::test]
3618    async fn numeric_entities_append_overlap_and_recover_with_exact_types() -> TestResult {
3619        let tmp = TempDir::new()?;
3620        let location = TableLocation::local(tmp.path());
3621        let mut table =
3622            TimeSeriesTable::create(location.clone(), make_int32_entity_table_meta()).await?;
3623
3624        let negative_path = "data/negative-device.parquet";
3625        write_int32_entity_parquet(
3626            &tmp.path().join(negative_path),
3627            &[1_000, 61_000],
3628            &[-1, -1],
3629            &[10.0, 11.0],
3630        )?;
3631        append_parquet_fixture(&mut table, negative_path).await?;
3632        let negative_identity = EntityIdentity::try_new(vec![EntityValue::Int32(-1)])?;
3633        let negative_layout = SegmentEntityLayout::Single(negative_identity.clone());
3634        assert!(
3635            table
3636                .state
3637                .segments
3638                .values()
3639                .any(|segment| segment.entity_layout == negative_layout)
3640        );
3641
3642        let maximum_path = "data/maximum-device.parquet";
3643        write_int32_entity_parquet(
3644            &tmp.path().join(maximum_path),
3645            &[1_000],
3646            &[i32::MAX],
3647            &[20.0],
3648        )?;
3649        append_parquet_fixture(&mut table, maximum_path).await?;
3650        let maximum_layout =
3651            SegmentEntityLayout::Single(EntityIdentity::try_new(vec![EntityValue::Int32(
3652                i32::MAX,
3653            )])?);
3654        assert!(
3655            table
3656                .state
3657                .segments
3658                .values()
3659                .any(|segment| segment.entity_layout == maximum_layout)
3660        );
3661
3662        let overlap_path = "data/negative-overlap.parquet";
3663        write_int32_entity_parquet(&tmp.path().join(overlap_path), &[1_500], &[-1], &[12.0])?;
3664        let error = append_parquet_fixture(&mut table, overlap_path)
3665            .await
3666            .expect_err("same typed identity and index interval must overlap");
3667        assert!(matches!(
3668            error,
3669            TableError::Append {
3670                source: AppendError::PersistedIndexIntervalOverlap {
3671                    overlap_count: 1,
3672                    example_identity: Some(example_identity),
3673                    ..
3674                }
3675            } if example_identity == negative_identity
3676        ));
3677
3678        let snapshot = table
3679            .load_entity_coverage_with_recovery::<AppendError>()
3680            .await?;
3681        let reopened = TimeSeriesTable::open(location).await?;
3682        assert_eq!(reopened.state(), table.state());
3683        assert_eq!(
3684            reopened
3685                .recover_entity_coverage_from_segments::<AppendError>()
3686                .await?,
3687            snapshot
3688        );
3689        Ok(())
3690    }
3691
3692    #[tokio::test]
3693    async fn append_preserves_composite_identity_order_in_layout() -> TestResult {
3694        let tmp = TempDir::new()?;
3695        let location = TableLocation::local(tmp.path());
3696        let index = IndexSpec {
3697            column: "ts".to_string(),
3698            entity_columns: vec!["symbol".to_string(), "venue".to_string()],
3699            kind: IndexKind::Timestamp {
3700                index_granularity: TimeIndexGranularity::Minutes(1),
3701                timezone: None,
3702            },
3703        };
3704        let schema = LogicalSchema::new(vec![
3705            LogicalField {
3706                name: "ts".to_string(),
3707                data_type: LogicalDataType::Timestamp {
3708                    unit: LogicalTimestampUnit::Millis,
3709                    timezone: None,
3710                },
3711                nullable: false,
3712            },
3713            LogicalField {
3714                name: "symbol".to_string(),
3715                data_type: LogicalDataType::Utf8,
3716                nullable: false,
3717            },
3718            LogicalField {
3719                name: "venue".to_string(),
3720                data_type: LogicalDataType::Utf8,
3721                nullable: false,
3722            },
3723            LogicalField {
3724                name: "price".to_string(),
3725                data_type: LogicalDataType::Float64,
3726                nullable: false,
3727            },
3728        ])?;
3729        let mut table = TimeSeriesTable::create(
3730            location,
3731            TableMeta::new_time_series_with_schema(index, schema),
3732        )
3733        .await?;
3734
3735        for (path, venue) in [
3736            ("data/composite-x.parquet", "X"),
3737            ("data/composite-y.parquet", "Y"),
3738        ] {
3739            write_composite_entity_parquet(&tmp.path().join(path), &[(1_000, "A", venue, 10.0)])?;
3740            append_parquet_fixture(&mut table, path).await?;
3741            let expected_layout = SegmentEntityLayout::Single(EntityIdentity::try_new(vec![
3742                "A".into(),
3743                venue.into(),
3744            ])?);
3745            assert!(
3746                table
3747                    .state
3748                    .segments
3749                    .values()
3750                    .any(|segment| segment.entity_layout == expected_layout)
3751            );
3752        }
3753
3754        let overlap_path = "data/composite-x-overlap.parquet";
3755        write_composite_entity_parquet(&tmp.path().join(overlap_path), &[(1_500, "A", "X", 20.0)])?;
3756        let error = append_parquet_fixture(&mut table, overlap_path)
3757            .await
3758            .expect_err("matching composite identity and index interval must overlap");
3759        assert!(matches!(
3760            error,
3761            TableError::Append {
3762                source: AppendError::PersistedIndexIntervalOverlap {
3763                    overlap_count: 1,
3764                    example_identity: Some(example_identity),
3765                    ..
3766                }
3767            } if example_identity.components()
3768                == [EntityValue::from("A"), EntityValue::from("X")]
3769        ));
3770        Ok(())
3771    }
3772
3773    #[tokio::test]
3774    async fn entity_with_only_null_index_values_is_rejected() -> TestResult {
3775        let tmp = TempDir::new()?;
3776        let location = TableLocation::local(tmp.path());
3777        let mut table = TimeSeriesTable::create(
3778            location,
3779            make_table_meta_with_unit(LogicalTimestampUnit::Millis),
3780        )
3781        .await?;
3782        let state_before = table.state.clone();
3783        let path = "data/entity-without-index-coverage.parquet";
3784        write_arrow_parquet_with_unit(
3785            &tmp.path().join(path),
3786            ArrowTimeUnit::Millisecond,
3787            &[Some(1_000), None],
3788            &["A", "B"],
3789            &[10.0, 20.0],
3790        )?;
3791
3792        let error = append_parquet_fixture(&mut table, path)
3793            .await
3794            .expect_err("identity without index coverage must be rejected");
3795
3796        match error {
3797            TableError::Append {
3798                source: AppendError::EntityWithoutIndexCoverage { identity, .. },
3799            } => {
3800                assert_eq!(identity.components(), [EntityValue::from("B")]);
3801            }
3802            other => panic!("unexpected error: {other:?}"),
3803        }
3804        assert_eq!(table.state, state_before);
3805        assert!(coverage_files(tmp.path())?.is_empty());
3806        Ok(())
3807    }
3808
3809    #[tokio::test]
3810    async fn append_rejects_unsupported_entity_type_without_publication() -> TestResult {
3811        let tmp = TempDir::new()?;
3812        let location = TableLocation::local(tmp.path());
3813        let index = IndexSpec {
3814            column: "ts".to_string(),
3815            entity_columns: vec!["device_id".to_string()],
3816            kind: IndexKind::Timestamp {
3817                index_granularity: TimeIndexGranularity::Minutes(1),
3818                timezone: None,
3819            },
3820        };
3821        let mut table =
3822            TimeSeriesTable::create(location, TableMeta::new_time_series(index)).await?;
3823        let schema = Arc::new(Schema::new(vec![
3824            Field::new(
3825                "ts",
3826                DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
3827                false,
3828            ),
3829            Field::new("device_id", DataType::Boolean, false),
3830        ]));
3831        let batch = RecordBatch::try_new(
3832            Arc::clone(&schema),
3833            vec![
3834                Arc::new(TimestampMillisecondArray::from(vec![1_000])),
3835                Arc::new(BooleanArray::from(vec![true])),
3836            ],
3837        )?;
3838        let state_before = table.state.clone();
3839
3840        let error = table
3841            .append(batch)
3842            .await
3843            .expect_err("Boolean entity columns must be rejected");
3844
3845        assert!(matches!(
3846            error,
3847            TableError::Append {
3848                source: AppendError::SchemaValidation {
3849                    ref source,
3850                    ..
3851                }
3852            } if matches!(
3853                source.as_ref(),
3854                crate::metadata::schema_compat::SchemaCompatibilityError::UnsupportedEntityColumnType {
3855                    column,
3856                    actual: LogicalDataType::Bool,
3857                } if column == "device_id"
3858            )
3859        ));
3860        assert_eq!(table.state, state_before);
3861        assert!(coverage_files(tmp.path())?.is_empty());
3862        assert_eq!(table.log.load_current_version().await?, 1);
3863        Ok(())
3864    }
3865
3866    #[tokio::test]
3867    async fn append_updates_snapshot() -> TestResult {
3868        let tmp = TempDir::new()?;
3869        let location = TableLocation::local(tmp.path());
3870        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
3871
3872        let rel1 = "data/seg-auto-1.parquet";
3873        let rel2 = "data/seg-auto-2.parquet";
3874        let path1 = tmp.path().join(rel1);
3875        let path2 = tmp.path().join(rel2);
3876
3877        write_test_parquet(
3878            &path1,
3879            true,
3880            false,
3881            &[
3882                TestRow {
3883                    ts_millis: 1_000,
3884                    symbol: "A",
3885                    price: 10.0,
3886                },
3887                TestRow {
3888                    ts_millis: 61_000,
3889                    symbol: "A",
3890                    price: 20.0,
3891                },
3892            ],
3893        )?;
3894        write_test_parquet(
3895            &path2,
3896            true,
3897            false,
3898            &[
3899                TestRow {
3900                    ts_millis: 120_000,
3901                    symbol: "A",
3902                    price: 30.0,
3903                },
3904                TestRow {
3905                    ts_millis: 180_000,
3906                    symbol: "A",
3907                    price: 40.0,
3908                },
3909            ],
3910        )?;
3911
3912        let v2 = append_parquet_fixture(&mut table, rel1).await?;
3913        let v3 = append_parquet_fixture(&mut table, rel2).await?;
3914        assert_eq!(v2, 2);
3915        assert_eq!(v3, 3);
3916
3917        assert_eq!(table.state.segments.len(), 2);
3918        assert!(
3919            table
3920                .state
3921                .segments
3922                .values()
3923                .all(|segment| segment.coverage_path.is_some())
3924        );
3925        let expected_snapshot = table
3926            .recover_entity_coverage_from_segments::<AppendError>()
3927            .await?;
3928
3929        let ptr = table
3930            .state
3931            .table_coverage
3932            .as_ref()
3933            .expect("table snapshot pointer present after append");
3934        assert_eq!(ptr.version, v3);
3935        assert_eq!(ptr.index_kind, table.index_spec().kind);
3936
3937        let snapshot_cov =
3938            read_entity_coverage_sidecar(&location, Path::new(&ptr.coverage_path)).await?;
3939
3940        assert_eq!(snapshot_cov, expected_snapshot);
3941        Ok(())
3942    }
3943
3944    #[tokio::test]
3945    async fn append_rejects_overlap() -> TestResult {
3946        let tmp = TempDir::new()?;
3947        let location = TableLocation::local(tmp.path());
3948        let mut table = TimeSeriesTable::create(location, make_basic_table_meta()).await?;
3949
3950        let rel1 = "data/seg-overlap-a.parquet";
3951        let rel2 = "data/seg-overlap-b.parquet";
3952        let path1 = tmp.path().join(rel1);
3953        let path2 = tmp.path().join(rel2);
3954
3955        write_test_parquet(
3956            &path1,
3957            true,
3958            false,
3959            &[
3960                TestRow {
3961                    ts_millis: 1_000,
3962                    symbol: "B",
3963                    price: 10.0,
3964                },
3965                TestRow {
3966                    ts_millis: 1_000,
3967                    symbol: "A",
3968                    price: 20.0,
3969                },
3970                TestRow {
3971                    ts_millis: 61_000,
3972                    symbol: "A",
3973                    price: 30.0,
3974                },
3975            ],
3976        )?;
3977        write_test_parquet(
3978            &path2,
3979            true,
3980            false,
3981            &[
3982                TestRow {
3983                    ts_millis: 1_500,
3984                    symbol: "B",
3985                    price: 40.0,
3986                },
3987                TestRow {
3988                    ts_millis: 1_500,
3989                    symbol: "A",
3990                    price: 50.0,
3991                },
3992                TestRow {
3993                    ts_millis: 61_500,
3994                    symbol: "A",
3995                    price: 60.0,
3996                },
3997            ],
3998        )?;
3999
4000        append_parquet_fixture(&mut table, rel1).await?;
4001
4002        let err = append_parquet_fixture(&mut table, rel2)
4003            .await
4004            .expect_err("overlapping append should fail");
4005
4006        assert!(matches!(
4007            err,
4008            TableError::Append {
4009                source: AppendError::PersistedIndexIntervalOverlap {
4010                    overlap_count: 3,
4011                    example_identity: Some(example_identity),
4012                    example_index_interval_id: 0x8000_0000_0000_0000,
4013                    example_index_interval,
4014                    ..
4015                }
4016            } if example_identity.components() == [EntityValue::from("A")]
4017                && example_index_interval.to_string()
4018                    == "[1970-01-01T00:00:00Z, 1970-01-01T00:01:00Z)"
4019        ));
4020        Ok(())
4021    }
4022
4023    #[tokio::test]
4024    async fn append_snapshot_survives_reopen() -> TestResult {
4025        let tmp = TempDir::new()?;
4026        let location = TableLocation::local(tmp.path());
4027        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4028
4029        let rel1 = "data/seg-reopen-a.parquet";
4030        let rel2 = "data/seg-reopen-b.parquet";
4031        let path1 = tmp.path().join(rel1);
4032        let path2 = tmp.path().join(rel2);
4033
4034        write_test_parquet(
4035            &path1,
4036            true,
4037            false,
4038            &[TestRow {
4039                ts_millis: 1_000,
4040                symbol: "A",
4041                price: 10.0,
4042            }],
4043        )?;
4044        write_test_parquet(
4045            &path2,
4046            true,
4047            false,
4048            &[TestRow {
4049                ts_millis: 120_000,
4050                symbol: "A",
4051                price: 20.0,
4052            }],
4053        )?;
4054
4055        append_parquet_fixture(&mut table, rel1).await?;
4056        append_parquet_fixture(&mut table, rel2).await?;
4057
4058        let reopened = TimeSeriesTable::open(location.clone()).await?;
4059        let ptr = reopened
4060            .state()
4061            .table_coverage
4062            .as_ref()
4063            .expect("table snapshot pointer present after reopen");
4064
4065        assert_eq!(ptr.index_kind, reopened.index_spec().kind);
4066
4067        let expected = reopened
4068            .recover_entity_coverage_from_segments::<AppendError>()
4069            .await?;
4070
4071        let snapshot_cov =
4072            read_entity_coverage_sidecar(&location, Path::new(&ptr.coverage_path)).await?;
4073        assert_eq!(snapshot_cov, expected);
4074        Ok(())
4075    }
4076
4077    #[tokio::test]
4078    async fn load_snapshot_recovers_when_missing_file() -> TestResult {
4079        let tmp = TempDir::new()?;
4080        let location = TableLocation::local(tmp.path());
4081        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4082
4083        // Append two segments so we have segment sidecars plus a table snapshot pointer.
4084        let rel1 = "data/seg-missing-snap-a.parquet";
4085        let rel2 = "data/seg-missing-snap-b.parquet";
4086        let path1 = tmp.path().join(rel1);
4087        let path2 = tmp.path().join(rel2);
4088        write_test_parquet(
4089            &path1,
4090            true,
4091            false,
4092            &[TestRow {
4093                ts_millis: 1_000,
4094                symbol: "A",
4095                price: 10.0,
4096            }],
4097        )?;
4098        write_test_parquet(
4099            &path2,
4100            true,
4101            false,
4102            &[TestRow {
4103                ts_millis: 120_000,
4104                symbol: "A",
4105                price: 20.0,
4106            }],
4107        )?;
4108
4109        append_parquet_fixture(&mut table, rel1).await?;
4110        append_parquet_fixture(&mut table, rel2).await?;
4111
4112        let state = table.state.clone();
4113        let ptr = state
4114            .table_coverage
4115            .as_ref()
4116            .expect("snapshot pointer present");
4117        let snapshot_abs = match &location.as_ref() {
4118            StorageLocation::Local(root) => root.join(&ptr.coverage_path),
4119        };
4120
4121        tokio::fs::remove_file(&snapshot_abs).await?;
4122
4123        let recovered = table
4124            .load_entity_coverage_with_recovery::<AppendError>()
4125            .await?;
4126
4127        let mut expected = EntityCoverage::empty();
4128        for seg in state.segments.values() {
4129            let cov_path = seg.coverage_path.as_ref().expect("coverage path");
4130            let cov = read_entity_coverage_sidecar(&location, Path::new(cov_path)).await?;
4131            expected.union_inplace(&cov);
4132        }
4133
4134        assert_eq!(recovered, expected);
4135        Ok(())
4136    }
4137
4138    #[tokio::test]
4139    async fn load_snapshot_recovers_when_corrupt_file() -> TestResult {
4140        let tmp = TempDir::new()?;
4141        let location = TableLocation::local(tmp.path());
4142        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4143
4144        let rel1 = "data/seg-corrupt-snap-a.parquet";
4145        let rel2 = "data/seg-corrupt-snap-b.parquet";
4146        let path1 = tmp.path().join(rel1);
4147        let path2 = tmp.path().join(rel2);
4148        write_test_parquet(
4149            &path1,
4150            true,
4151            false,
4152            &[TestRow {
4153                ts_millis: 1_000,
4154                symbol: "A",
4155                price: 10.0,
4156            }],
4157        )?;
4158        write_test_parquet(
4159            &path2,
4160            true,
4161            false,
4162            &[TestRow {
4163                ts_millis: 120_000,
4164                symbol: "A",
4165                price: 20.0,
4166            }],
4167        )?;
4168
4169        append_parquet_fixture(&mut table, rel1).await?;
4170        append_parquet_fixture(&mut table, rel2).await?;
4171
4172        let state = table.state.clone();
4173        let ptr = state
4174            .table_coverage
4175            .as_ref()
4176            .expect("snapshot pointer present");
4177        let snapshot_abs = match &location.as_ref() {
4178            StorageLocation::Local(root) => root.join(&ptr.coverage_path),
4179        };
4180
4181        tokio::fs::write(&snapshot_abs, b"garbage").await?;
4182
4183        let recovered = table
4184            .load_entity_coverage_with_recovery::<AppendError>()
4185            .await?;
4186
4187        let mut expected = EntityCoverage::empty();
4188        for seg in state.segments.values() {
4189            let cov_path = seg.coverage_path.as_ref().expect("coverage path");
4190            let cov = read_entity_coverage_sidecar(&location, Path::new(cov_path)).await?;
4191            expected.union_inplace(&cov);
4192        }
4193
4194        assert_eq!(recovered, expected);
4195        Ok(())
4196    }
4197
4198    #[tokio::test]
4199    async fn rejected_append_does_not_heal_corrupt_snapshot() -> TestResult {
4200        let tmp = TempDir::new()?;
4201        let location = TableLocation::local(tmp.path());
4202        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4203
4204        let existing = "data/existing.parquet";
4205        write_test_parquet(
4206            &tmp.path().join(existing),
4207            true,
4208            false,
4209            &[TestRow {
4210                ts_millis: 1_000,
4211                symbol: "A",
4212                price: 10.0,
4213            }],
4214        )?;
4215        append_parquet_fixture(&mut table, existing).await?;
4216
4217        let snapshot_path = table
4218            .state
4219            .table_coverage
4220            .as_ref()
4221            .expect("snapshot pointer present")
4222            .coverage_path
4223            .clone();
4224        let snapshot_abs = tmp.path().join(snapshot_path);
4225        tokio::fs::write(&snapshot_abs, b"garbage").await?;
4226
4227        let overlapping = "data/overlapping.parquet";
4228        write_test_parquet(
4229            &tmp.path().join(overlapping),
4230            true,
4231            false,
4232            &[TestRow {
4233                ts_millis: 1_000,
4234                symbol: "A",
4235                price: 20.0,
4236            }],
4237        )?;
4238
4239        let err = append_parquet_fixture(&mut table, overlapping)
4240            .await
4241            .expect_err("overlap must be rejected");
4242        assert!(matches!(
4243            err,
4244            TableError::Append {
4245                source: AppendError::PersistedIndexIntervalOverlap { .. }
4246            }
4247        ));
4248        assert_eq!(tokio::fs::read(snapshot_abs).await?, b"garbage");
4249        Ok(())
4250    }
4251
4252    #[tokio::test]
4253    async fn load_snapshot_errors_when_segment_missing_coverage_path() -> TestResult {
4254        let tmp = TempDir::new()?;
4255        let location = TableLocation::local(tmp.path());
4256        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4257
4258        let rel1 = "data/seg-missing-cov-path.parquet";
4259        let path1 = tmp.path().join(rel1);
4260        write_test_parquet(
4261            &path1,
4262            true,
4263            false,
4264            &[TestRow {
4265                ts_millis: 1_000,
4266                symbol: "A",
4267                price: 10.0,
4268            }],
4269        )?;
4270
4271        append_parquet_fixture(&mut table, rel1).await?;
4272
4273        let mut state = table.state.clone();
4274        state.table_coverage = None;
4275
4276        let segment_path = state
4277            .segments
4278            .keys()
4279            .next()
4280            .expect("segment present")
4281            .clone();
4282        state
4283            .segments
4284            .get_mut(&segment_path)
4285            .expect("segment present")
4286            .coverage_path = None;
4287
4288        // Overwrite table state with the modified snapshot missing coverage_path.
4289        table.state = state;
4290
4291        let err = table
4292            .load_entity_coverage_with_recovery::<AppendError>()
4293            .await
4294            .expect_err("missing coverage_path should error");
4295
4296        assert!(matches!(
4297            err,
4298            AppendError::ExistingSegmentMissingCoverageMetadata { .. }
4299        ));
4300        Ok(())
4301    }
4302
4303    #[tokio::test]
4304    async fn load_snapshot_errors_when_segment_sidecar_corrupt() -> TestResult {
4305        let tmp = TempDir::new()?;
4306        let location = TableLocation::local(tmp.path());
4307        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4308
4309        let rel1 = "data/seg-corrupt-sidecar.parquet";
4310        let rel2 = "data/seg-corrupt-sidecar-ok.parquet";
4311        let path1 = tmp.path().join(rel1);
4312        let path2 = tmp.path().join(rel2);
4313        write_test_parquet(
4314            &path1,
4315            true,
4316            false,
4317            &[TestRow {
4318                ts_millis: 1_000,
4319                symbol: "A",
4320                price: 10.0,
4321            }],
4322        )?;
4323        write_test_parquet(
4324            &path2,
4325            true,
4326            false,
4327            &[TestRow {
4328                ts_millis: 120_000,
4329                symbol: "A",
4330                price: 20.0,
4331            }],
4332        )?;
4333
4334        append_parquet_fixture(&mut table, rel1).await?;
4335        append_parquet_fixture(&mut table, rel2).await?;
4336
4337        let mut state = table.state.clone();
4338        state.table_coverage = None;
4339        let (corrupt_segment_path, corrupt_cov_path) = state
4340            .segments
4341            .values()
4342            .next()
4343            .map(|meta| {
4344                (
4345                    meta.path.clone(),
4346                    meta.coverage_path.as_ref().expect("coverage path").clone(),
4347                )
4348            })
4349            .expect("at least one segment");
4350        table.state = state;
4351
4352        let corrupt_abs = match &location.as_ref() {
4353            StorageLocation::Local(root) => root.join(&corrupt_cov_path),
4354        };
4355        tokio::fs::write(&corrupt_abs, b"not a coverage bitmap").await?;
4356
4357        let err = table
4358            .load_entity_coverage_with_recovery::<AppendError>()
4359            .await
4360            .expect_err("corrupt sidecar should error");
4361
4362        match err {
4363            AppendError::ExistingSegmentCoverageSidecarRead {
4364                segment_path,
4365                coverage_path,
4366                ..
4367            } => {
4368                assert_eq!(segment_path, corrupt_segment_path);
4369                assert_eq!(coverage_path, corrupt_cov_path);
4370            }
4371            other => panic!("unexpected error: {other:?}"),
4372        }
4373
4374        Ok(())
4375    }
4376
4377    #[tokio::test]
4378    async fn entity_aware_stale_append_cleans_sidecars_without_state_mutation() -> TestResult {
4379        let tmp = TempDir::new()?;
4380        let location = TableLocation::local(tmp.path());
4381        let meta = make_basic_table_meta();
4382        let mut winner = TimeSeriesTable::create(location.clone(), meta).await?;
4383        let mut loser = TimeSeriesTable::open(location.clone()).await?;
4384        let loser_state_before = loser.state.clone();
4385
4386        let winner_path = "data/winner.parquet";
4387        let loser_path = "data/loser.parquet";
4388        write_test_parquet(
4389            &tmp.path().join(winner_path),
4390            true,
4391            false,
4392            &[TestRow {
4393                ts_millis: 10_000,
4394                symbol: "X",
4395                price: 100.0,
4396            }],
4397        )?;
4398        write_test_parquet(
4399            &tmp.path().join(loser_path),
4400            true,
4401            false,
4402            &[TestRow {
4403                ts_millis: 120_000,
4404                symbol: "X",
4405                price: 200.0,
4406            }],
4407        )?;
4408
4409        assert_eq!(append_parquet_fixture(&mut winner, winner_path).await?, 2);
4410        let coverage_before = coverage_files(tmp.path())?;
4411
4412        let err = append_parquet_fixture(&mut loser, loser_path)
4413            .await
4414            .expect_err("expected conflict due to stale version");
4415
4416        match err {
4417            TableError::Append {
4418                source: AppendError::Commit { source },
4419            } => {
4420                assert!(matches!(
4421                    source,
4422                    CommitError::Conflict {
4423                        expected: 1,
4424                        found: 2,
4425                        ..
4426                    }
4427                ));
4428            }
4429            other => panic!("unexpected error: {other:?}"),
4430        }
4431
4432        assert_eq!(loser.state, loser_state_before);
4433        assert_eq!(loser.log.load_current_version().await?, 2);
4434        let committed = loser.load_latest_state().await?;
4435        assert_eq!(committed, winner.state);
4436        assert_eq!(coverage_files(tmp.path())?, coverage_before);
4437        for bytes in coverage_before.values() {
4438            entity_coverage_from_bytes(bytes)?;
4439        }
4440        assert!(!tmp.path().join(layout::commit_rel_path(3)).exists());
4441        Ok(())
4442    }
4443
4444    #[tokio::test]
4445    async fn stale_int64_append_cleans_only_its_writer_owned_sidecars() -> TestResult {
4446        let tmp = TempDir::new()?;
4447        let location = TableLocation::local(tmp.path());
4448        let index = registered_index(IndexKind::Int64 {
4449            index_granularity: NonZeroU64::new(10).unwrap(),
4450        });
4451        let mut winner =
4452            TimeSeriesTable::create(location.clone(), TableMeta::new_time_series(index)).await?;
4453        let mut loser = TimeSeriesTable::open(location).await?;
4454        let winner_path = "data/writer-owned-winner.parquet";
4455        let loser_path = "data/writer-owned-loser.parquet";
4456
4457        write_arrow_parquet_int_time(&tmp.path().join(winner_path), &[0], &["X"], &[100.0])?;
4458        write_arrow_parquet_int_time(&tmp.path().join(loser_path), &[100], &["X"], &[200.0])?;
4459
4460        append_parquet_fixture(&mut winner, winner_path).await?;
4461        let coverage_before = coverage_files(tmp.path())?;
4462
4463        let err = append_parquet_fixture(&mut loser, loser_path)
4464            .await
4465            .expect_err("stale append should conflict");
4466
4467        assert!(matches!(
4468            err,
4469            TableError::Append {
4470                source: AppendError::Commit {
4471                    source: CommitError::Conflict { .. }
4472                }
4473            }
4474        ));
4475        assert_eq!(coverage_files(tmp.path())?, coverage_before);
4476        Ok(())
4477    }
4478
4479    #[tokio::test]
4480    async fn ambiguous_int64_commit_retains_writer_owned_sidecars() -> TestResult {
4481        let tmp = TempDir::new()?;
4482        let location = TableLocation::local(tmp.path());
4483        let index = registered_index(IndexKind::Int64 {
4484            index_granularity: NonZeroU64::new(10).unwrap(),
4485        });
4486        let mut table =
4487            TimeSeriesTable::create(location, TableMeta::new_time_series(index)).await?;
4488        let state_before = table.state.clone();
4489        let coverage_before = coverage_files(tmp.path())?;
4490        let segment_path = "data/ambiguous.parquet";
4491
4492        write_arrow_parquet_int_time(&tmp.path().join(segment_path), &[10], &["X"], &[100.0])?;
4493
4494        let commit_path = tmp.path().join(layout::commit_rel_path(2));
4495        crate::storage::inject_write_new_failure(commit_path.clone(), true);
4496
4497        let err = append_parquet_fixture(&mut table, segment_path)
4498            .await
4499            .expect_err("failed commit cleanup should make the outcome ambiguous");
4500
4501        assert!(
4502            matches!(
4503                &err,
4504                TableError::Append {
4505                    source: AppendError::CommitAmbiguous {
4506                        source,
4507                        ..
4508                    }
4509                } if matches!(source.as_ref(), CommitError::AmbiguousOutcome { .. })
4510            ),
4511            "unexpected error: {err:?}"
4512        );
4513        assert_eq!(table.state, state_before);
4514        assert_eq!(table.log.load_current_version().await?, 1);
4515        assert!(commit_path.exists());
4516        assert_eq!(coverage_files(tmp.path())?.len(), coverage_before.len() + 2);
4517        Ok(())
4518    }
4519
4520    #[tokio::test]
4521    async fn entity_sidecar_cleanup_failures_preserve_error_and_reverse_order() -> TestResult {
4522        let tmp = TempDir::new()?;
4523        let location = TableLocation::local(tmp.path());
4524        let table = TimeSeriesTable::create(location, make_basic_table_meta()).await?;
4525        let sidecars = [
4526            format!("{SEGMENT_COVERAGE_DIR}/first-stuck.roar"),
4527            format!("{TABLE_SNAPSHOT_DIR}/second-stuck.roar"),
4528        ];
4529        for sidecar in &sidecars {
4530            tokio::fs::create_dir_all(tmp.path().join(sidecar)).await?;
4531        }
4532        let example_index_interval_id =
4533            crate::coverage::index_interval::index_interval_id_for_value(
4534                &table.index.kind,
4535                &IndexValue::Timestamp(utc_datetime(1970, 1, 1, 0, 0, 0)),
4536            )?;
4537        let source = AppendError::PersistedIndexIntervalOverlap {
4538            segment_path: "data/failed.parquet".to_string(),
4539            overlap_count: 1,
4540            example_identity: Some(EntityIdentity::try_new(vec!["A".into()])?),
4541            example_index_interval_id,
4542            example_index_interval: Box::new(index_interval_for_id(
4543                &table.index.kind,
4544                example_index_interval_id,
4545            )?),
4546        };
4547        let err = table.rollback_created_artifacts(&sidecars, source).await;
4548        let message = err.to_string();
4549
4550        assert!(matches!(
4551            err,
4552            AppendError::Rollback {
4553                source,
4554                cleanup_errors,
4555            } if matches!(
4556                *source,
4557                AppendError::PersistedIndexIntervalOverlap { .. }
4558            )
4559                && cleanup_errors.len() == 2
4560                && matches!(
4561                    &cleanup_errors[0],
4562                    StorageError::OtherIo { path, .. } if path.ends_with("second-stuck.roar")
4563                )
4564                && matches!(
4565                    &cleanup_errors[1],
4566                    StorageError::OtherIo { path, .. } if path.ends_with("first-stuck.roar")
4567                )
4568        ));
4569        assert!(message.contains("data/failed.parquet"));
4570        assert!(message.contains("first-stuck.roar"));
4571        assert!(message.contains("second-stuck.roar"));
4572        Ok(())
4573    }
4574
4575    #[tokio::test]
4576    async fn append_fails_when_existing_segment_missing_coverage_path() -> TestResult {
4577        let tmp = TempDir::new()?;
4578        let location = TableLocation::local(tmp.path());
4579        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4580
4581        let rel1 = "data/seg-missing-cov.parquet";
4582        let rel2 = "data/seg-next.parquet";
4583        let path1 = tmp.path().join(rel1);
4584        let path2 = tmp.path().join(rel2);
4585
4586        write_test_parquet(
4587            &path1,
4588            true,
4589            false,
4590            &[TestRow {
4591                ts_millis: 1_000,
4592                symbol: "A",
4593                price: 10.0,
4594            }],
4595        )?;
4596        write_test_parquet(
4597            &path2,
4598            true,
4599            false,
4600            &[TestRow {
4601                ts_millis: 120_000,
4602                symbol: "A",
4603                price: 20.0,
4604            }],
4605        )?;
4606
4607        append_parquet_fixture(&mut table, rel1).await?;
4608
4609        // Simulate legacy/bad state: drop coverage_path on the existing segment.
4610        let seg = table
4611            .state
4612            .segments
4613            .values_mut()
4614            .next()
4615            .expect("segment present");
4616        seg.coverage_path = None;
4617
4618        let err = append_parquet_fixture(&mut table, rel2)
4619            .await
4620            .expect_err("append should fail when existing segment lacks coverage");
4621
4622        assert!(matches!(
4623            err,
4624            TableError::Append {
4625                source: AppendError::ExistingSegmentMissingCoverageMetadata { .. }
4626            }
4627        ));
4628        Ok(())
4629    }
4630
4631    #[tokio::test]
4632    // Unlike load_snapshot_recovers_when_missing_file (which exercises recovery when
4633    // the pointer exists but the snapshot file is gone), this covers the case where
4634    // the in-memory pointer itself is missing while segments exist, and append
4635    // must rebuild + rewrite the pointer as part of the append flow.
4636    async fn append_recovers_when_table_snapshot_pointer_missing() -> TestResult {
4637        let tmp = TempDir::new()?;
4638        let location = TableLocation::local(tmp.path());
4639        let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4640
4641        let rel1 = "data/seg-no-pointer-a.parquet";
4642        let rel2 = "data/seg-no-pointer-b.parquet";
4643        let path1 = tmp.path().join(rel1);
4644        let path2 = tmp.path().join(rel2);
4645
4646        write_test_parquet(
4647            &path1,
4648            true,
4649            false,
4650            &[TestRow {
4651                ts_millis: 1_000,
4652                symbol: "A",
4653                price: 10.0,
4654            }],
4655        )?;
4656        write_test_parquet(
4657            &path2,
4658            true,
4659            false,
4660            &[TestRow {
4661                ts_millis: 120_000,
4662                symbol: "A",
4663                price: 20.0,
4664            }],
4665        )?;
4666
4667        append_parquet_fixture(&mut table, rel1).await?;
4668
4669        // Simulate missing snapshot pointer while segments exist.
4670        table.state.table_coverage = None;
4671
4672        append_parquet_fixture(&mut table, rel2).await?;
4673
4674        // Snapshot pointer should be restored after a successful append.
4675        let ptr = table
4676            .state
4677            .table_coverage
4678            .as_ref()
4679            .expect("snapshot pointer restored");
4680
4681        let cov = read_entity_coverage_sidecar(&location, Path::new(&ptr.coverage_path)).await?;
4682
4683        let mut expected = EntityCoverage::empty();
4684        for seg in table.state.segments.values() {
4685            let path = seg.coverage_path.as_ref().expect("coverage path");
4686            let seg_cov = read_entity_coverage_sidecar(&location, Path::new(path)).await?;
4687            expected.union_inplace(&seg_cov);
4688        }
4689
4690        assert_eq!(cov, expected);
4691        Ok(())
4692    }
4693
4694    #[tokio::test]
4695    async fn generated_segment_metadata_failure_preserves_source_cleans_data_and_allows_retry()
4696    -> TestResult {
4697        let tmp = TempDir::new()?;
4698        let location = TableLocation::local(tmp.path());
4699        let mut table = TimeSeriesTable::create(location, timestamp_only_meta()).await?;
4700        let state_before = table.state().clone();
4701        let segment_path = "data/corrupt-generated.parquet";
4702        tokio::fs::create_dir_all(tmp.path().join("data")).await?;
4703        tokio::fs::write(tmp.path().join(segment_path), b"not parquet data").await?;
4704
4705        let mut data_guard = storage::FileCleanupGuard::new_disarmed(
4706            table.location().as_ref(),
4707            Path::new(segment_path),
4708        )?;
4709        data_guard.arm();
4710        let source = table
4711            .publish_generated_parquet_segment(segment_path, 2, &mut data_guard)
4712            .await
4713            .expect_err("corrupt generated Parquet must fail metadata loading");
4714        let source = table
4715            .rollback_created_artifacts(&[segment_path.to_string()], source)
4716            .await;
4717        data_guard.disarm();
4718        let error = crate::table::error::AppendSnafu.into_error(source);
4719
4720        assert!(matches!(
4721            &error,
4722            TableError::Append {
4723                source: AppendError::SegmentMetadata { source, .. }
4724            } if matches!(source.as_ref(), SegmentError::Metadata { .. })
4725        ));
4726        assert!(error.to_string().contains(segment_path));
4727        assert!(ErrorCompat::backtrace(&error).is_some());
4728        assert_eq!(table.state(), &state_before);
4729        assert_eq!(table.log.load_current_version().await?, 1);
4730        assert!(!tmp.path().join(segment_path).exists());
4731        assert!(data_files(tmp.path())?.is_empty());
4732        assert!(coverage_files(tmp.path())?.is_empty());
4733
4734        assert_eq!(table.append(timestamp_only_batch(1)?).await?, 2);
4735        Ok(())
4736    }
4737
4738    #[tokio::test]
4739    async fn append_missing_or_corrupt_segment_sidecar_preserves_source_and_rolls_back()
4740    -> TestResult {
4741        #[derive(Debug, Clone, Copy)]
4742        enum Damage {
4743            Missing,
4744            Corrupt,
4745        }
4746
4747        for damage in [Damage::Missing, Damage::Corrupt] {
4748            let tmp = TempDir::new()?;
4749            let location = TableLocation::local(tmp.path());
4750            let mut table =
4751                TimeSeriesTable::create(location.clone(), timestamp_only_meta()).await?;
4752            assert_eq!(table.append(timestamp_only_batch(1)?).await?, 2);
4753
4754            let (segment_path, coverage_path) = table
4755                .state()
4756                .segments
4757                .values()
4758                .next()
4759                .map(|segment| {
4760                    (
4761                        segment.path.clone(),
4762                        segment.coverage_path.clone().expect("coverage path"),
4763                    )
4764                })
4765                .expect("committed segment");
4766            let coverage_abs = tmp.path().join(&coverage_path);
4767            let original_coverage = tokio::fs::read(&coverage_abs).await?;
4768            match damage {
4769                Damage::Missing => tokio::fs::remove_file(&coverage_abs).await?,
4770                Damage::Corrupt => {
4771                    tokio::fs::write(&coverage_abs, b"not a coverage sidecar").await?
4772                }
4773            }
4774            table.state.table_coverage = None;
4775
4776            let state_before = table.state().clone();
4777            let data_before = data_files(tmp.path())?;
4778            let coverage_before = coverage_files(tmp.path())?;
4779            let error = table
4780                .append(timestamp_only_batch_starting_at_minute_offset(2, 1)?)
4781                .await
4782                .expect_err("damaged segment sidecar must fail append recovery");
4783
4784            let sidecar_source = match &error {
4785                TableError::Append {
4786                    source:
4787                        AppendError::ExistingSegmentCoverageSidecarRead {
4788                            segment_path: actual_segment_path,
4789                            coverage_path: actual_coverage_path,
4790                            source,
4791                        },
4792                } => {
4793                    assert_eq!(actual_segment_path, &segment_path);
4794                    assert_eq!(actual_coverage_path, &coverage_path);
4795                    source.as_ref()
4796                }
4797                other => panic!("unexpected {damage:?} error: {other:?}"),
4798            };
4799            match damage {
4800                Damage::Missing => assert!(matches!(
4801                    sidecar_source,
4802                    CoverageSidecarError::Storage {
4803                        source: StorageError::NotFound { .. }
4804                    }
4805                )),
4806                Damage::Corrupt => {
4807                    assert!(matches!(sidecar_source, CoverageSidecarError::Codec { .. }))
4808                }
4809            }
4810            assert!(ErrorCompat::backtrace(&error).is_some());
4811            assert_eq!(table.state(), &state_before, "{damage:?}");
4812            assert_eq!(data_files(tmp.path())?, data_before, "{damage:?}");
4813            assert_eq!(coverage_files(tmp.path())?, coverage_before, "{damage:?}");
4814
4815            tokio::fs::write(&coverage_abs, &original_coverage).await?;
4816            assert_eq!(
4817                table
4818                    .append(timestamp_only_batch_starting_at_minute_offset(2, 1)?)
4819                    .await?,
4820                3,
4821                "{damage:?}"
4822            );
4823        }
4824        Ok(())
4825    }
4826}