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