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