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