Skip to main content

timeseries_table_format/formats/parquet/
entity_rewrite.rs

1//! Staging mixed-entity Parquet segments as single-entity replacements.
2
3use std::{collections::HashSet, io::Write, path::Path};
4
5use arrow::{array::BooleanBuilder, compute::filter_record_batch, error::ArrowError};
6use futures::StreamExt;
7use parquet::{
8    arrow::{
9        ArrowWriter,
10        arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions},
11        async_reader::ParquetRecordBatchStreamBuilder,
12    },
13    errors::ParquetError,
14    file::properties::WriterProperties,
15};
16use snafu::{Backtrace, Snafu};
17use uuid::Uuid;
18
19use crate::{
20    coverage::{
21        EntityCoverage, EntityIdentity,
22        io::{
23            CoverageSidecarError, read_entity_coverage_sidecar, write_coverage_sidecar_new_bytes,
24        },
25        layout::{
26            coverage_file_id_for_attempt, segment_coverage_key, segment_entity_coverage_id_v1,
27        },
28        serde::{CoverageCodecError, entity_coverage_to_bytes},
29    },
30    formats::parquet::{
31        INSPECTION_BATCH_SIZE, SegmentCoverageError, compute_segment_entity_coverage,
32        entity_coverage::{entity_arrays, entity_identity_at},
33        logical_schema_from_parquet,
34        segment_meta::segment_meta_from_parquet,
35    },
36    metadata::{
37        index::{IndexSpec, IndexSpecError},
38        logical_schema::LogicalSchema,
39        schema_compat::{
40            SchemaCompatibilityError, ensure_index_spec_matches_schema,
41            ensure_schema_fields_match_by_name,
42        },
43        segments::{FileFormat, SegmentEntityLayout, SegmentMeta, SegmentMetaError},
44    },
45    storage::{
46        OutputSink, StorageError, TableLocation, ensure_canonical_relative_storage_path,
47        open_new_output_sink, open_parquet_reader, remove_file_if_exists,
48    },
49    transaction_log::segments::SegmentError,
50};
51
52#[cfg(test)]
53const MAX_OPEN_WRITERS: usize = 1;
54
55/// One staged single-entity replacement and its verified exact coverage.
56#[derive(Debug, Clone, PartialEq, Eq)]
57pub struct StagedEntityReplacement {
58    /// Complete identity materialized in this replacement.
59    pub identity: EntityIdentity,
60    /// Metadata derived from the completed staged Parquet file.
61    pub meta: SegmentMeta,
62    /// Exact coverage derived from the completed staged Parquet file.
63    pub coverage: EntityCoverage,
64}
65
66/// Completed staged rewrite whose private objects are now caller-owned.
67#[derive(Debug, Clone, PartialEq, Eq)]
68pub struct StagedEntityRewrite {
69    /// Committed source segment path that was read without modification.
70    pub source_path: String,
71    /// Verified single-entity replacements in canonical identity order.
72    pub replacements: Vec<StagedEntityReplacement>,
73    /// Every staged data and sidecar path owned by the caller.
74    pub staged_object_paths: Vec<String>,
75    /// Physical rows read across all bounded source scans.
76    pub rows_read: u64,
77    /// Logical source rows written across all replacements.
78    pub rows_written: u64,
79    /// Complete identities materialized in canonical order.
80    pub materialized_identities: Vec<EntityIdentity>,
81}
82
83/// Failure while staging a mixed segment rewrite.
84#[derive(Debug, Snafu)]
85#[non_exhaustive]
86pub enum EntityRewriteError {
87    /// The committed source or table metadata violates the rewrite contract.
88    #[snafu(display("Invalid mixed-segment rewrite input: {reason}"))]
89    InvalidInput {
90        /// Inconsistent input detail.
91        reason: String,
92        /// Backtrace captured at the invalid rewrite boundary.
93        backtrace: Backtrace,
94    },
95
96    /// A completed staged output violates a rewrite invariant.
97    #[snafu(display("Invalid staged entity rewrite output: {reason}"))]
98    InvalidOutput {
99        /// Failed output invariant.
100        reason: String,
101        /// Backtrace captured at the failed rewrite invariant.
102        backtrace: Backtrace,
103    },
104
105    /// A source or sidecar path failed table-relative storage validation.
106    #[snafu(display("Invalid {description} path {path:?}: {source}"))]
107    InvalidPath {
108        /// Role of the rejected path in the rewrite.
109        description: &'static str,
110        /// Rejected path.
111        path: String,
112        /// Structured storage path validation failure.
113        #[snafu(source(from(StorageError, Box::new)), backtrace)]
114        source: Box<StorageError>,
115    },
116
117    /// The rewrite received an invalid ordered-index specification.
118    #[snafu(display("Invalid rewrite ordered-index specification: {source}"))]
119    IndexSpecValidation {
120        /// Complete ordered-index validation failure.
121        source: IndexSpecError,
122        /// Backtrace captured at the rewrite boundary.
123        backtrace: Backtrace,
124    },
125
126    /// The table schema is incompatible with its ordered-index specification.
127    #[snafu(display("Rewrite table schema validation failed: {source}"))]
128    TableSchemaValidation {
129        /// Complete schema compatibility failure.
130        #[snafu(source(from(SchemaCompatibilityError, Box::new)), backtrace)]
131        source: Box<SchemaCompatibilityError>,
132    },
133
134    /// A source or replacement segment schema is incompatible with the table schema.
135    #[snafu(display("Rewrite segment schema validation failed for {path}: {source}"))]
136    SegmentSchemaValidation {
137        /// Source or staged replacement path.
138        path: String,
139        /// Complete schema compatibility failure.
140        #[snafu(source(from(SchemaCompatibilityError, Box::new)), backtrace)]
141        source: Box<SchemaCompatibilityError>,
142    },
143
144    /// Persisted source segment metadata violates the registered ordered-index domain.
145    #[snafu(display("Invalid rewrite source metadata: {source}"))]
146    SegmentMetadataValidation {
147        /// Complete segment metadata validation failure.
148        #[snafu(source(from(SegmentMetaError, Box::new)), backtrace)]
149        source: Box<SegmentMetaError>,
150    },
151
152    /// Segment metadata or schema inspection failed.
153    #[snafu(display("Failed to inspect Parquet segment: {source}"))]
154    SegmentInspection {
155        /// Existing Parquet inspection failure.
156        #[snafu(source, backtrace)]
157        source: SegmentError,
158    },
159
160    /// Exact entity coverage inspection failed.
161    #[snafu(display("Failed to inspect exact entity coverage: {source}"))]
162    CoverageInspection {
163        /// Existing entity coverage inspection failure.
164        #[snafu(source, backtrace)]
165        source: SegmentCoverageError,
166    },
167
168    /// Coverage sidecar access failed.
169    #[snafu(display("Failed to access entity coverage sidecar: {source}"))]
170    CoverageSidecar {
171        /// Existing coverage sidecar failure.
172        #[snafu(source, backtrace)]
173        source: CoverageSidecarError,
174    },
175
176    /// Entity coverage serialization failed.
177    #[snafu(display("Failed to serialize staged entity coverage: {source}"))]
178    CoverageSerialization {
179        /// Existing coverage codec failure.
180        #[snafu(source, backtrace)]
181        source: CoverageCodecError,
182    },
183
184    /// Storage access failed.
185    #[snafu(display("Staged entity rewrite storage failure: {source}"))]
186    Storage {
187        /// Existing storage failure.
188        #[snafu(source, backtrace)]
189        source: StorageError,
190    },
191
192    /// Parquet streaming or writing failed.
193    #[snafu(display("Parquet rewrite failure at {path}: {source}"))]
194    Parquet {
195        /// Table-relative source or output path.
196        path: String,
197        /// Existing Parquet failure.
198        source: ParquetError,
199        /// Diagnostic backtrace.
200        backtrace: Backtrace,
201    },
202
203    /// Filtering a complete record batch failed.
204    #[snafu(display("Arrow row filtering failed for {path}: {source}"))]
205    Arrow {
206        /// Table-relative source path.
207        path: String,
208        /// Existing Arrow failure.
209        source: ArrowError,
210        /// Diagnostic backtrace.
211        backtrace: Backtrace,
212    },
213
214    /// Rewrite failed and one or more private objects could not be removed.
215    #[snafu(display(
216        "{source}; staged-object rollback also failed: [{}]",
217        cleanup_errors
218            .iter()
219            .map(ToString::to_string)
220            .collect::<Vec<_>>()
221            .join("; ")
222    ))]
223    Cleanup {
224        /// Primary rewrite failure.
225        #[snafu(source, backtrace)]
226        source: Box<EntityRewriteError>,
227        /// Typed cleanup failure for every private object that could not be removed.
228        cleanup_errors: Vec<StorageError>,
229    },
230}
231
232struct SinkWriter(OutputSink);
233
234impl Write for SinkWriter {
235    fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
236        self.0.writer().write(bytes)
237    }
238
239    fn flush(&mut self) -> std::io::Result<()> {
240        self.0.writer().flush()
241    }
242}
243
244fn invalid_input(reason: impl Into<String>) -> EntityRewriteError {
245    EntityRewriteError::InvalidInput {
246        reason: reason.into(),
247        backtrace: Backtrace::capture(),
248    }
249}
250
251fn invalid_output(reason: impl Into<String>) -> EntityRewriteError {
252    EntityRewriteError::InvalidOutput {
253        reason: reason.into(),
254        backtrace: Backtrace::capture(),
255    }
256}
257
258fn validate_rewrite_path(path: &str, description: &'static str) -> Result<(), EntityRewriteError> {
259    ensure_canonical_relative_storage_path(path).map_err(|source| EntityRewriteError::InvalidPath {
260        description,
261        path: path.to_string(),
262        source: Box::new(source),
263    })
264}
265
266async fn cleanup_created(location: &TableLocation, created_paths: &[String]) -> Vec<StorageError> {
267    let mut errors = Vec::new();
268    for path in created_paths.iter().rev() {
269        if let Err(error) = remove_file_if_exists(location.as_ref(), Path::new(path)).await {
270            errors.push(error);
271        }
272    }
273    errors
274}
275
276async fn stage_identity_data(
277    location: &TableLocation,
278    source_path: &str,
279    index: &IndexSpec,
280    identity: &EntityIdentity,
281    output_path: &str,
282    created_paths: &mut Vec<String>,
283) -> Result<(u64, u64), EntityRewriteError> {
284    let source_rel = Path::new(source_path);
285    let mut metadata_file = open_parquet_reader(location.as_ref(), source_rel)
286        .await
287        .map_err(|source| EntityRewriteError::Storage { source })?;
288    let metadata =
289        ArrowReaderMetadata::load_async(&mut metadata_file, ArrowReaderOptions::default())
290            .await
291            .map_err(|source| EntityRewriteError::Parquet {
292                path: source_path.to_string(),
293                source,
294                backtrace: Backtrace::capture(),
295            })?;
296    let schema = metadata.schema().clone();
297    drop(metadata_file);
298
299    let source_file = open_parquet_reader(location.as_ref(), source_rel)
300        .await
301        .map_err(|source| EntityRewriteError::Storage { source })?;
302    let mut reader = ParquetRecordBatchStreamBuilder::new_with_metadata(source_file, metadata)
303        .with_batch_size(INSPECTION_BATCH_SIZE)
304        .build()
305        .map_err(|source| EntityRewriteError::Parquet {
306            path: source_path.to_string(),
307            source,
308            backtrace: Backtrace::capture(),
309        })?;
310
311    let sink = open_new_output_sink(location.as_ref(), Path::new(output_path))
312        .await
313        .map_err(|source| EntityRewriteError::Storage { source })?;
314    created_paths.push(output_path.to_string());
315    let mut writer = ArrowWriter::try_new(
316        SinkWriter(sink),
317        schema,
318        Some(WriterProperties::builder().build()),
319    )
320    .map_err(|source| EntityRewriteError::Parquet {
321        path: output_path.to_string(),
322        source,
323        backtrace: Backtrace::capture(),
324    })?;
325
326    let mut rows_read = 0u64;
327    let mut rows_written = 0u64;
328    while let Some(batch) = reader.next().await {
329        let batch = batch.map_err(|source| EntityRewriteError::Parquet {
330            path: source_path.to_string(),
331            source,
332            backtrace: Backtrace::capture(),
333        })?;
334        let entities = entity_arrays(&batch, source_path, &index.entity_columns)
335            .map_err(|source| EntityRewriteError::CoverageInspection { source })?;
336        let mut mask = BooleanBuilder::with_capacity(batch.num_rows());
337        for row in 0..batch.num_rows() {
338            mask.append_value(
339                entity_identity_at(&entities, row, source_path)
340                    .map_err(|source| EntityRewriteError::CoverageInspection { source })?
341                    == *identity,
342            );
343        }
344        rows_read = rows_read
345            .checked_add(batch.num_rows() as u64)
346            .ok_or_else(|| invalid_output("rows-read counter overflow"))?;
347        let filtered = filter_record_batch(&batch, &mask.finish()).map_err(|source| {
348            EntityRewriteError::Arrow {
349                path: source_path.to_string(),
350                source,
351                backtrace: Backtrace::capture(),
352            }
353        })?;
354        if filtered.num_rows() == 0 {
355            continue;
356        }
357        rows_written = rows_written
358            .checked_add(filtered.num_rows() as u64)
359            .ok_or_else(|| invalid_output("rows-written counter overflow"))?;
360        writer
361            .write(&filtered)
362            .map_err(|source| EntityRewriteError::Parquet {
363                path: output_path.to_string(),
364                source,
365                backtrace: Backtrace::capture(),
366            })?;
367    }
368
369    let sink = writer
370        .into_inner()
371        .map_err(|source| EntityRewriteError::Parquet {
372            path: output_path.to_string(),
373            source,
374            backtrace: Backtrace::capture(),
375        })?
376        .0;
377    sink.finish()
378        .await
379        .map_err(|source| EntityRewriteError::Storage { source })?;
380    Ok((rows_read, rows_written))
381}
382
383async fn validate_source(
384    location: &TableLocation,
385    table_schema: &LogicalSchema,
386    index: &IndexSpec,
387    source: &SegmentMeta,
388) -> Result<EntityCoverage, EntityRewriteError> {
389    index
390        .validate()
391        .map_err(|source| EntityRewriteError::IndexSpecValidation {
392            source,
393            backtrace: Backtrace::capture(),
394        })?;
395    ensure_index_spec_matches_schema(table_schema, index).map_err(|source| {
396        EntityRewriteError::TableSchemaValidation {
397            source: Box::new(source),
398        }
399    })?;
400    if index.entity_columns.is_empty() {
401        return Err(invalid_input("table has no entity columns"));
402    }
403    if source.format != FileFormat::Parquet {
404        return Err(invalid_input("source is not Parquet"));
405    }
406    if source.entity_layout != SegmentEntityLayout::Mixed {
407        return Err(invalid_input(format!(
408            "source {} is not classified as Mixed",
409            source.path
410        )));
411    }
412    validate_rewrite_path(&source.path, "source segment")?;
413    source.validate_bounds(&index.kind).map_err(|source| {
414        EntityRewriteError::SegmentMetadataValidation {
415            source: Box::new(source),
416        }
417    })?;
418    let coverage_path = source
419        .coverage_path
420        .as_deref()
421        .ok_or_else(|| invalid_input("source has no committed entity-coverage sidecar"))?;
422    validate_rewrite_path(coverage_path, "source coverage")?;
423
424    let source_schema = logical_schema_from_parquet(location, Path::new(&source.path))
425        .await
426        .map_err(|source| EntityRewriteError::SegmentInspection { source })?;
427    ensure_schema_fields_match_by_name(table_schema, &source_schema, index).map_err(|error| {
428        EntityRewriteError::SegmentSchemaValidation {
429            path: source.path.clone(),
430            source: Box::new(error),
431        }
432    })?;
433
434    let (actual_meta, _) = segment_meta_from_parquet(location, Path::new(&source.path), index)
435        .await
436        .map_err(|source| EntityRewriteError::SegmentInspection { source })?;
437    let file_size_matches = source
438        .file_size
439        .is_none_or(|expected| actual_meta.file_size == Some(expected));
440    if actual_meta.index_min != source.index_min
441        || actual_meta.index_max != source.index_max
442        || actual_meta.row_count != source.row_count
443        || !file_size_matches
444    {
445        return Err(invalid_input(format!(
446            "source metadata does not match the committed Parquet file at {}",
447            source.path
448        )));
449    }
450
451    let committed_coverage = read_entity_coverage_sidecar(location, Path::new(coverage_path))
452        .await
453        .map_err(|source| EntityRewriteError::CoverageSidecar { source })?;
454    if committed_coverage.identity_count() < 2 {
455        return Err(invalid_input(
456            "Mixed source coverage must contain at least two identities",
457        ));
458    }
459    for (identity, coverage) in committed_coverage.iter() {
460        if identity.components().len() != index.entity_columns.len() {
461            return Err(invalid_input(format!(
462                "source identity {identity:?} has {} components, expected {}",
463                identity.components().len(),
464                index.entity_columns.len()
465            )));
466        }
467        if coverage.is_empty() {
468            return Err(invalid_input(format!(
469                "source identity {identity:?} has no covered ordered-index interval"
470            )));
471        }
472    }
473
474    let actual_coverage = compute_segment_entity_coverage(location, Path::new(&source.path), index)
475        .await
476        .map_err(|source| EntityRewriteError::CoverageInspection { source })?;
477    if actual_coverage != committed_coverage {
478        return Err(invalid_input(
479            "committed source coverage does not match the source Parquet rows",
480        ));
481    }
482    Ok(committed_coverage)
483}
484
485async fn rewrite_inner(
486    location: &TableLocation,
487    table_schema: &LogicalSchema,
488    index: &IndexSpec,
489    source: &SegmentMeta,
490    attempt_id: Uuid,
491    created_paths: &mut Vec<String>,
492) -> Result<StagedEntityRewrite, EntityRewriteError> {
493    let source_coverage = validate_source(location, table_schema, index, source).await?;
494    let mut replacements = Vec::with_capacity(source_coverage.identity_count());
495    let mut materialized_identities = Vec::with_capacity(source_coverage.identity_count());
496    let mut output_coverage = EntityCoverage::empty();
497    let mut rows_read = 0u64;
498    let mut rows_written = 0u64;
499
500    // Processing one identity at a time bounds the handle count and memory use.
501    // Batch identities only if repeated scan cost becomes material.
502    for (ordinal, (identity, expected_coverage)) in source_coverage.iter().enumerate() {
503        let data_path = format!("data/_staged/entity-rewrite/{attempt_id}/{ordinal:010}.parquet");
504        let (identity_rows_read, identity_rows_written) = stage_identity_data(
505            location,
506            &source.path,
507            index,
508            identity,
509            &data_path,
510            created_paths,
511        )
512        .await?;
513        rows_read = rows_read
514            .checked_add(identity_rows_read)
515            .ok_or_else(|| invalid_output("rows-read counter overflow"))?;
516        rows_written = rows_written
517            .checked_add(identity_rows_written)
518            .ok_or_else(|| invalid_output("rows-written counter overflow"))?;
519
520        let output_schema = logical_schema_from_parquet(location, Path::new(&data_path))
521            .await
522            .map_err(|source| EntityRewriteError::SegmentInspection { source })?;
523        ensure_schema_fields_match_by_name(table_schema, &output_schema, index).map_err(
524            |source| EntityRewriteError::SegmentSchemaValidation {
525                path: data_path.clone(),
526                source: Box::new(source),
527            },
528        )?;
529        let (mut meta, _) = segment_meta_from_parquet(location, Path::new(&data_path), index)
530            .await
531            .map_err(|source| EntityRewriteError::SegmentInspection { source })?;
532        if meta.row_count != identity_rows_written || meta.row_count == 0 {
533            return Err(invalid_output(format!(
534                "replacement {data_path} row count {} does not match written row count {identity_rows_written}",
535                meta.row_count
536            )));
537        }
538
539        let coverage = compute_segment_entity_coverage(location, Path::new(&data_path), index)
540            .await
541            .map_err(|source| EntityRewriteError::CoverageInspection { source })?;
542        let mut expected = EntityCoverage::empty();
543        expected.union_coverage(identity.clone(), expected_coverage.clone());
544        if coverage != expected {
545            return Err(invalid_output(format!(
546                "replacement {data_path} coverage does not match identity {identity:?}"
547            )));
548        }
549        if output_coverage.intersection_cardinality(&coverage) != 0 {
550            return Err(invalid_output(format!(
551                "replacement {data_path} overlaps an earlier replacement"
552            )));
553        }
554
555        let coverage_bytes = entity_coverage_to_bytes(&coverage)
556            .map_err(|source| EntityRewriteError::CoverageSerialization { source })?;
557        let coverage_id = coverage_file_id_for_attempt(
558            &segment_entity_coverage_id_v1(index, &coverage_bytes),
559            &attempt_id,
560        );
561        let coverage_path = segment_coverage_key(&coverage_id).map_err(|source| {
562            EntityRewriteError::CoverageSidecar {
563                source: CoverageSidecarError::Layout {
564                    source,
565                    backtrace: Backtrace::capture(),
566                },
567            }
568        })?;
569        let sidecar_write =
570            write_coverage_sidecar_new_bytes(location, Path::new(&coverage_path), &coverage_bytes)
571                .await;
572        if let Err(source) = sidecar_write {
573            if source.storage_cleanup_failed() {
574                created_paths.push(coverage_path);
575            }
576            return Err(EntityRewriteError::CoverageSidecar { source });
577        }
578        created_paths.push(coverage_path.clone());
579        let persisted_coverage = read_entity_coverage_sidecar(location, Path::new(&coverage_path))
580            .await
581            .map_err(|source| EntityRewriteError::CoverageSidecar { source })?;
582        if persisted_coverage != coverage {
583            return Err(invalid_output(format!(
584                "replacement sidecar {coverage_path} does not match derived coverage"
585            )));
586        }
587
588        meta.entity_layout = SegmentEntityLayout::Single(identity.clone());
589        meta.coverage_path = Some(coverage_path);
590        output_coverage.union_inplace(&coverage);
591        materialized_identities.push(identity.clone());
592        replacements.push(StagedEntityReplacement {
593            identity: identity.clone(),
594            meta,
595            coverage,
596        });
597    }
598
599    if replacements.len() != source_coverage.identity_count() {
600        return Err(invalid_output(format!(
601            "materialized {} outputs for {} source identities",
602            replacements.len(),
603            source_coverage.identity_count()
604        )));
605    }
606    if rows_written != source.row_count {
607        return Err(invalid_output(format!(
608            "wrote {rows_written} rows from a source containing {} rows",
609            source.row_count
610        )));
611    }
612    if output_coverage != source_coverage {
613        return Err(invalid_output(
614            "replacement coverage union does not equal committed source coverage",
615        ));
616    }
617    let unique_paths = created_paths.iter().collect::<HashSet<_>>();
618    if unique_paths.len() != created_paths.len() {
619        return Err(invalid_output("staged object paths are not unique"));
620    }
621
622    Ok(StagedEntityRewrite {
623        source_path: source.path.clone(),
624        replacements,
625        staged_object_paths: created_paths.clone(),
626        rows_read,
627        rows_written,
628        materialized_identities,
629    })
630}
631
632async fn rewrite_with_attempt_id(
633    location: &TableLocation,
634    table_schema: &LogicalSchema,
635    index: &IndexSpec,
636    source: &SegmentMeta,
637    attempt_id: Uuid,
638) -> Result<StagedEntityRewrite, EntityRewriteError> {
639    let mut created_paths = Vec::new();
640    match rewrite_inner(
641        location,
642        table_schema,
643        index,
644        source,
645        attempt_id,
646        &mut created_paths,
647    )
648    .await
649    {
650        Ok(rewrite) => Ok(rewrite),
651        Err(source) => {
652            let cleanup_errors = cleanup_created(location, &created_paths).await;
653            if cleanup_errors.is_empty() {
654                Err(source)
655            } else {
656                Err(EntityRewriteError::Cleanup {
657                    source: Box::new(source),
658                    cleanup_errors,
659                })
660            }
661        }
662    }
663}
664
665/// Rewrite one committed mixed Parquet segment into verified staged
666/// single-entity replacements without changing table state.
667///
668/// This implementation intentionally opens one output writer at a time. Each
669/// source scan streams complete record batches, so memory and open handles are
670/// bounded independently of identity cardinality.
671///
672/// # Errors
673///
674/// Returns [`EntityRewriteError`] when input validation, reading, writing,
675/// output verification, sidecar creation, or pre-return cleanup fails.
676pub async fn rewrite_mixed_parquet_segment(
677    location: &TableLocation,
678    table_schema: &LogicalSchema,
679    index: &IndexSpec,
680    source: &SegmentMeta,
681) -> Result<StagedEntityRewrite, EntityRewriteError> {
682    rewrite_with_attempt_id(location, table_schema, index, source, Uuid::new_v4()).await
683}
684
685#[cfg(test)]
686mod tests {
687    use std::error::Error as _;
688
689    use super::*;
690    use std::{collections::BTreeMap, fs::File, sync::Arc};
691
692    use arrow::{
693        array::{
694            ArrayRef, Float64Array, Int64Array, StringArray, StructArray, TimestampMillisecondArray,
695        },
696        datatypes::{DataType, Field, Fields, Schema, TimeUnit},
697        record_batch::RecordBatch,
698    };
699    use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
700    use parquet::file::properties::WriterProperties;
701    use snafu::ErrorCompat;
702    use tempfile::TempDir;
703
704    use crate::{
705        coverage::{
706            EntityValue, io::write_coverage_sidecar_new_bytes, serde::entity_coverage_to_bytes,
707        },
708        metadata::index::{IndexKind, TimeIndexGranularity},
709        storage::normalize_relative_storage_path,
710        table::test_util::{make_table_meta_with_unit, write_arrow_parquet_with_unit},
711        transaction_log::TableKind,
712    };
713
714    type TestResult<T = ()> = Result<T, Box<dyn std::error::Error>>;
715
716    #[test]
717    fn invalid_rewrite_path_preserves_storage_source_and_backtrace() {
718        let error = validate_rewrite_path("../outside.parquet", "source segment")
719            .expect_err("parent traversal must fail");
720        let storage = error
721            .source()
722            .and_then(|source| source.downcast_ref::<Box<StorageError>>())
723            .map(Box::as_ref)
724            .expect("storage source");
725
726        assert!(matches!(error, EntityRewriteError::InvalidPath { .. }));
727        assert!(std::ptr::eq(
728            ErrorCompat::backtrace(&error).expect("rewrite backtrace"),
729            ErrorCompat::backtrace(storage).expect("storage backtrace")
730        ));
731    }
732
733    fn read_rows(path: &Path) -> TestResult<Vec<(i64, String, f64)>> {
734        let reader = ParquetRecordBatchReaderBuilder::try_new(File::open(path)?)?.build()?;
735        let mut rows = Vec::new();
736        for batch in reader {
737            let batch = batch?;
738            let timestamps = batch
739                .column(0)
740                .as_any()
741                .downcast_ref::<TimestampMillisecondArray>()
742                .expect("timestamp column");
743            let symbols = batch
744                .column(1)
745                .as_any()
746                .downcast_ref::<StringArray>()
747                .expect("symbol column");
748            let prices = batch
749                .column(2)
750                .as_any()
751                .downcast_ref::<Float64Array>()
752                .expect("price column");
753            for row in 0..batch.num_rows() {
754                rows.push((
755                    timestamps.value(row),
756                    symbols.value(row).to_string(),
757                    prices.value(row),
758                ));
759            }
760        }
761        Ok(rows)
762    }
763
764    fn read_batch(path: &Path) -> TestResult<RecordBatch> {
765        let builder = ParquetRecordBatchReaderBuilder::try_new(File::open(path)?)?;
766        let schema = builder.schema().clone();
767        let batches = builder.build()?.collect::<Result<Vec<_>, _>>()?;
768        Ok(arrow_select::concat::concat_batches(&schema, &batches)?)
769    }
770
771    struct RewriteFixture {
772        temp: TempDir,
773        location: TableLocation,
774        table_schema: LogicalSchema,
775        index: IndexSpec,
776        source: SegmentMeta,
777        source_coverage: EntityCoverage,
778    }
779
780    async fn rewrite_fixture() -> TestResult<RewriteFixture> {
781        let temp = TempDir::new()?;
782        let location = TableLocation::local(temp.path());
783        let source_path = "data/failure-source.parquet";
784        write_arrow_parquet_with_unit(
785            &temp.path().join(source_path),
786            TimeUnit::Millisecond,
787            &[Some(1_000), Some(2_000), Some(61_000), Some(62_000)],
788            &["A", "B", "A", "B"],
789            &[10.0, 20.0, 11.0, 21.0],
790        )?;
791        let table_meta = make_table_meta_with_unit(
792            crate::metadata::logical_schema::LogicalTimestampUnit::Millis,
793        );
794        let TableKind::TimeSeries(index) = table_meta.kind else {
795            unreachable!("test metadata is time-series");
796        };
797        let table_schema = table_meta.logical_schema.expect("test table schema");
798        let source_coverage =
799            compute_segment_entity_coverage(&location, Path::new(source_path), &index).await?;
800        let source_coverage_path = "_coverage/segments/failure-source.roar";
801        write_coverage_sidecar_new_bytes(
802            &location,
803            Path::new(source_coverage_path),
804            &entity_coverage_to_bytes(&source_coverage)?,
805        )
806        .await?;
807        let (mut source, _) =
808            segment_meta_from_parquet(&location, Path::new(source_path), &index).await?;
809        source.entity_layout = SegmentEntityLayout::Mixed;
810        source.coverage_path = Some(source_coverage_path.to_string());
811        Ok(RewriteFixture {
812            temp,
813            location,
814            table_schema,
815            index,
816            source,
817            source_coverage,
818        })
819    }
820
821    fn staged_data_path(attempt_id: Uuid, ordinal: usize) -> String {
822        format!("data/_staged/entity-rewrite/{attempt_id}/{ordinal:010}.parquet")
823    }
824
825    fn staged_coverage_path(
826        fixture: &RewriteFixture,
827        attempt_id: Uuid,
828        ordinal: usize,
829    ) -> TestResult<String> {
830        let (identity, coverage) = fixture
831            .source_coverage
832            .iter()
833            .nth(ordinal)
834            .expect("fixture identity");
835        let mut output_coverage = EntityCoverage::empty();
836        output_coverage.union_coverage(identity.clone(), coverage.clone());
837        let bytes = entity_coverage_to_bytes(&output_coverage)?;
838        let coverage_id = coverage_file_id_for_attempt(
839            &segment_entity_coverage_id_v1(&fixture.index, &bytes),
840            &attempt_id,
841        );
842        Ok(segment_coverage_key(&coverage_id)?)
843    }
844
845    fn write_sentinel(root: &Path, path: &str) -> TestResult<Vec<u8>> {
846        let bytes = b"preexisting-object".to_vec();
847        let absolute = root.join(path);
848        std::fs::create_dir_all(absolute.parent().expect("object parent"))?;
849        std::fs::write(absolute, &bytes)?;
850        Ok(bytes)
851    }
852
853    fn assert_nothing_staged(fixture: &RewriteFixture) {
854        assert!(!fixture.temp.path().join("data/_staged").exists());
855    }
856
857    #[tokio::test]
858    async fn mixed_rewrite_stages_exactly_two_verified_outputs() -> TestResult {
859        let fixture = rewrite_fixture().await?;
860
861        let rewrite = rewrite_mixed_parquet_segment(
862            &fixture.location,
863            &fixture.table_schema,
864            &fixture.index,
865            &fixture.source,
866        )
867        .await?;
868
869        assert_eq!(rewrite.replacements.len(), 2);
870        assert_eq!(rewrite.rows_read, fixture.source.row_count * 2);
871        assert_eq!(rewrite.rows_written, fixture.source.row_count);
872        assert_eq!(rewrite.staged_object_paths.len(), 4);
873        let mut output_coverage = EntityCoverage::empty();
874        for replacement in &rewrite.replacements {
875            assert_eq!(
876                replacement.meta.entity_layout,
877                SegmentEntityLayout::Single(replacement.identity.clone())
878            );
879            output_coverage.union_inplace(&replacement.coverage);
880        }
881        assert_eq!(output_coverage, fixture.source_coverage);
882        assert_eq!(
883            rewrite.materialized_identities,
884            fixture
885                .source_coverage
886                .iter()
887                .map(|(identity, _)| identity.clone())
888                .collect::<Vec<_>>()
889        );
890        Ok(())
891    }
892
893    #[tokio::test]
894    async fn mixed_rewrite_stages_one_bounded_output_per_identity() -> TestResult {
895        let temp = TempDir::new()?;
896        let location = TableLocation::local(temp.path());
897        let source_path = "data/mixed.parquet";
898        let timestamps = [1_000, 2_000, 61_000, 4_000, 62_000, 64_000];
899        let symbols = [
900            "tenant-secret-b",
901            "tenant-secret-a",
902            "tenant-secret-b",
903            "tenant-secret-c",
904            "tenant-secret-a",
905            "tenant-secret-c",
906        ];
907        let prices = [10.0, 20.0, 11.0, 30.0, 21.0, 31.0];
908        write_arrow_parquet_with_unit(
909            &temp.path().join(source_path),
910            TimeUnit::Millisecond,
911            &timestamps.map(Some),
912            &symbols,
913            &prices,
914        )?;
915
916        let table_meta = make_table_meta_with_unit(
917            crate::metadata::logical_schema::LogicalTimestampUnit::Millis,
918        );
919        let TableKind::TimeSeries(index) = &table_meta.kind else {
920            unreachable!("test metadata is time-series");
921        };
922        let table_schema = table_meta
923            .logical_schema
924            .as_ref()
925            .expect("test table schema");
926        let source_coverage =
927            compute_segment_entity_coverage(&location, Path::new(source_path), index).await?;
928        assert!(source_coverage.identity_count() > MAX_OPEN_WRITERS);
929        let source_coverage_path = "_coverage/segments/source.roar";
930        write_coverage_sidecar_new_bytes(
931            &location,
932            Path::new(source_coverage_path),
933            &entity_coverage_to_bytes(&source_coverage)?,
934        )
935        .await?;
936        let (mut source, _) =
937            segment_meta_from_parquet(&location, Path::new(source_path), index).await?;
938        source.entity_layout = SegmentEntityLayout::Mixed;
939        source.coverage_path = Some(source_coverage_path.to_string());
940        let source_bytes = std::fs::read(temp.path().join(source_path))?;
941        let source_coverage_bytes = std::fs::read(temp.path().join(source_coverage_path))?;
942
943        let rewrite =
944            rewrite_mixed_parquet_segment(&location, table_schema, index, &source).await?;
945
946        assert_eq!(rewrite.source_path, source_path);
947        assert_eq!(rewrite.replacements.len(), 3);
948        assert_eq!(rewrite.rows_read, source.row_count * 3);
949        assert_eq!(rewrite.rows_written, source.row_count);
950        assert_eq!(rewrite.staged_object_paths.len(), 6);
951        assert_eq!(std::fs::read(temp.path().join(source_path))?, source_bytes);
952        assert_eq!(
953            std::fs::read(temp.path().join(source_coverage_path))?,
954            source_coverage_bytes
955        );
956        assert_eq!(
957            rewrite
958                .staged_object_paths
959                .iter()
960                .collect::<HashSet<_>>()
961                .len(),
962            rewrite.staged_object_paths.len()
963        );
964        for path in &rewrite.staged_object_paths {
965            let (canonical, _) = normalize_relative_storage_path(Path::new(path))?;
966            assert_eq!(&canonical, path);
967            for secret in ["tenant-secret-a", "tenant-secret-b", "tenant-secret-c"] {
968                assert!(!path.contains(secret));
969            }
970        }
971
972        let mut actual = BTreeMap::new();
973        for replacement in &rewrite.replacements {
974            assert_eq!(
975                replacement.meta.entity_layout,
976                SegmentEntityLayout::Single(replacement.identity.clone())
977            );
978            assert_eq!(replacement.coverage.identity_count(), 1);
979            assert_eq!(
980                read_entity_coverage_sidecar(
981                    &location,
982                    Path::new(
983                        replacement
984                            .meta
985                            .coverage_path
986                            .as_deref()
987                            .expect("replacement coverage path")
988                    )
989                )
990                .await?,
991                replacement.coverage
992            );
993            assert!(rewrite.staged_object_paths.contains(&replacement.meta.path));
994            actual.insert(
995                replacement.identity.components()[0].clone(),
996                read_rows(&temp.path().join(&replacement.meta.path))?,
997            );
998        }
999
1000        assert_eq!(
1001            actual[&EntityValue::from("tenant-secret-a")],
1002            vec![
1003                (2_000, "tenant-secret-a".to_string(), 20.0),
1004                (62_000, "tenant-secret-a".to_string(), 21.0),
1005            ]
1006        );
1007        assert_eq!(
1008            actual[&EntityValue::from("tenant-secret-b")],
1009            vec![
1010                (1_000, "tenant-secret-b".to_string(), 10.0),
1011                (61_000, "tenant-secret-b".to_string(), 11.0),
1012            ]
1013        );
1014        assert_eq!(
1015            actual[&EntityValue::from("tenant-secret-c")],
1016            vec![
1017                (4_000, "tenant-secret-c".to_string(), 30.0),
1018                (64_000, "tenant-secret-c".to_string(), 31.0),
1019            ]
1020        );
1021        Ok(())
1022    }
1023
1024    #[tokio::test]
1025    async fn rewrite_preserves_composite_rows_across_batches_and_row_groups() -> TestResult {
1026        let temp = TempDir::new()?;
1027        let location = TableLocation::local(temp.path());
1028        let source_path = "data/composite-mixed.parquet";
1029        std::fs::create_dir_all(temp.path().join("data"))?;
1030
1031        let row_count = INSPECTION_BATCH_SIZE + 17;
1032        let regions = (0..row_count)
1033            .map(|row| if row % 4 < 2 { "eu" } else { "us" })
1034            .collect::<Vec<_>>();
1035        let symbols = (0..row_count)
1036            .map(|row| if row % 2 == 0 { "A" } else { "B" })
1037            .collect::<Vec<_>>();
1038        let readings = (0..row_count)
1039            .map(|row| (row % 5 != 0).then_some(row as i64))
1040            .collect::<Vec<_>>();
1041        let notes = (0..row_count)
1042            .map(|row| (row % 7 != 0).then(|| format!("note-{row}")))
1043            .collect::<Vec<_>>();
1044        let payload_fields = Fields::from(vec![
1045            Arc::new(Field::new("reading", DataType::Int64, true)),
1046            Arc::new(Field::new("note", DataType::Utf8, true)),
1047        ]);
1048        let payload = StructArray::new(
1049            payload_fields.clone(),
1050            vec![
1051                Arc::new(Int64Array::from(readings)) as ArrayRef,
1052                Arc::new(StringArray::from(notes)) as ArrayRef,
1053            ],
1054            None,
1055        );
1056        let schema = Arc::new(Schema::new(vec![
1057            Field::new("region", DataType::Utf8, false),
1058            Field::new(
1059                "ts",
1060                DataType::Timestamp(TimeUnit::Millisecond, None),
1061                false,
1062            ),
1063            Field::new("symbol", DataType::Utf8, false),
1064            Field::new("payload", DataType::Struct(payload_fields), false),
1065            Field::new("sequence", DataType::Int64, false),
1066        ]));
1067        let source_batch = RecordBatch::try_new(
1068            Arc::clone(&schema),
1069            vec![
1070                Arc::new(StringArray::from(regions)) as ArrayRef,
1071                Arc::new(TimestampMillisecondArray::from_iter_values(
1072                    (0..row_count).map(|row| row as i64 * 60_000),
1073                )),
1074                Arc::new(StringArray::from(symbols)),
1075                Arc::new(payload),
1076                Arc::new(Int64Array::from_iter_values(
1077                    (0..row_count).map(|row| row as i64),
1078                )),
1079            ],
1080        )?;
1081        let mut writer = ArrowWriter::try_new(
1082            File::create(temp.path().join(source_path))?,
1083            schema,
1084            Some(
1085                WriterProperties::builder()
1086                    .set_max_row_group_row_count(Some(513))
1087                    .build(),
1088            ),
1089        )?;
1090        writer.write(&source_batch)?;
1091        writer.close()?;
1092        assert!(
1093            ParquetRecordBatchReaderBuilder::try_new(File::open(temp.path().join(source_path))?)?
1094                .metadata()
1095                .num_row_groups()
1096                > 1
1097        );
1098
1099        let index = IndexSpec {
1100            column: "ts".to_string(),
1101            entity_columns: vec!["region".to_string(), "symbol".to_string()],
1102            kind: IndexKind::Timestamp {
1103                index_granularity: TimeIndexGranularity::Minutes(1),
1104                timezone: None,
1105            },
1106        };
1107        let table_schema = logical_schema_from_parquet(&location, Path::new(source_path)).await?;
1108        let source_coverage =
1109            compute_segment_entity_coverage(&location, Path::new(source_path), &index).await?;
1110        let source_coverage_path = "_coverage/segments/composite-source.roar";
1111        write_coverage_sidecar_new_bytes(
1112            &location,
1113            Path::new(source_coverage_path),
1114            &entity_coverage_to_bytes(&source_coverage)?,
1115        )
1116        .await?;
1117        let (mut source, _) =
1118            segment_meta_from_parquet(&location, Path::new(source_path), &index).await?;
1119        source.entity_layout = SegmentEntityLayout::Mixed;
1120        source.coverage_path = Some(source_coverage_path.to_string());
1121
1122        let rewrite =
1123            rewrite_mixed_parquet_segment(&location, &table_schema, &index, &source).await?;
1124
1125        assert_eq!(rewrite.replacements.len(), 4);
1126        assert_eq!(rewrite.rows_read, source.row_count * 4);
1127        let source_batch = read_batch(&temp.path().join(source_path))?;
1128        let source_entities = entity_arrays(&source_batch, source_path, &index.entity_columns)?;
1129        for replacement in rewrite.replacements {
1130            let mut mask = BooleanBuilder::with_capacity(source_batch.num_rows());
1131            for row in 0..source_batch.num_rows() {
1132                mask.append_value(
1133                    entity_identity_at(&source_entities, row, source_path)? == replacement.identity,
1134                );
1135            }
1136            let expected = filter_record_batch(&source_batch, &mask.finish())?;
1137            let actual = read_batch(&temp.path().join(&replacement.meta.path))?;
1138            assert_eq!(actual, expected);
1139        }
1140        Ok(())
1141    }
1142
1143    #[tokio::test]
1144    async fn rewrite_never_overwrites_a_colliding_data_path() -> TestResult {
1145        let fixture = rewrite_fixture().await?;
1146        let attempt_id = Uuid::from_u128(1);
1147        let collision_path = staged_data_path(attempt_id, 0);
1148        let sentinel = write_sentinel(fixture.temp.path(), &collision_path)?;
1149
1150        let error = rewrite_with_attempt_id(
1151            &fixture.location,
1152            &fixture.table_schema,
1153            &fixture.index,
1154            &fixture.source,
1155            attempt_id,
1156        )
1157        .await
1158        .expect_err("data collision must fail");
1159
1160        assert!(matches!(
1161            error,
1162            EntityRewriteError::Storage {
1163                source: StorageError::AlreadyExists { .. }
1164            }
1165        ));
1166        assert_eq!(
1167            std::fs::read(fixture.temp.path().join(collision_path))?,
1168            sentinel
1169        );
1170        Ok(())
1171    }
1172
1173    #[tokio::test]
1174    async fn sidecar_collision_preserves_existing_object_and_cleans_owned_outputs() -> TestResult {
1175        let fixture = rewrite_fixture().await?;
1176        let attempt_id = Uuid::from_u128(2);
1177        let data_paths = [
1178            staged_data_path(attempt_id, 0),
1179            staged_data_path(attempt_id, 1),
1180        ];
1181        let first_coverage_path = staged_coverage_path(&fixture, attempt_id, 0)?;
1182        let collision_path = staged_coverage_path(&fixture, attempt_id, 1)?;
1183        let sentinel = write_sentinel(fixture.temp.path(), &collision_path)?;
1184
1185        let error = rewrite_with_attempt_id(
1186            &fixture.location,
1187            &fixture.table_schema,
1188            &fixture.index,
1189            &fixture.source,
1190            attempt_id,
1191        )
1192        .await
1193        .expect_err("sidecar collision must fail");
1194
1195        assert!(matches!(
1196            error,
1197            EntityRewriteError::CoverageSidecar {
1198                source: CoverageSidecarError::Storage {
1199                    source: StorageError::AlreadyExists { .. }
1200                }
1201            }
1202        ));
1203        assert_eq!(
1204            std::fs::read(fixture.temp.path().join(collision_path))?,
1205            sentinel
1206        );
1207        for path in data_paths.iter().chain([&first_coverage_path]) {
1208            assert!(!fixture.temp.path().join(path).exists(), "{path} leaked");
1209        }
1210        Ok(())
1211    }
1212
1213    #[tokio::test]
1214    async fn cleanup_reports_every_failure_without_hiding_primary_error() -> TestResult {
1215        let fixture = rewrite_fixture().await?;
1216        let attempt_id = Uuid::from_u128(3);
1217        let owned_paths = [
1218            staged_data_path(attempt_id, 0),
1219            staged_coverage_path(&fixture, attempt_id, 0)?,
1220            staged_data_path(attempt_id, 1),
1221        ];
1222        let collision_path = staged_coverage_path(&fixture, attempt_id, 1)?;
1223        let sentinel = write_sentinel(fixture.temp.path(), &collision_path)?;
1224        for path in &owned_paths {
1225            crate::storage::inject_cleanup_failure(fixture.temp.path().join(path));
1226        }
1227
1228        let error = rewrite_with_attempt_id(
1229            &fixture.location,
1230            &fixture.table_schema,
1231            &fixture.index,
1232            &fixture.source,
1233            attempt_id,
1234        )
1235        .await
1236        .expect_err("rewrite and cleanup must fail");
1237
1238        let EntityRewriteError::Cleanup {
1239            source,
1240            cleanup_errors,
1241        } = error
1242        else {
1243            panic!("expected cleanup error");
1244        };
1245        assert!(matches!(
1246            *source,
1247            EntityRewriteError::CoverageSidecar {
1248                source: CoverageSidecarError::Storage {
1249                    source: StorageError::AlreadyExists { .. }
1250                }
1251            }
1252        ));
1253        assert_eq!(cleanup_errors.len(), owned_paths.len());
1254        for path in &owned_paths {
1255            assert!(
1256                cleanup_errors
1257                    .iter()
1258                    .any(|error| error.to_string().contains(path)),
1259                "cleanup failure omitted {path}"
1260            );
1261            assert!(fixture.temp.path().join(path).exists());
1262        }
1263        assert_eq!(
1264            std::fs::read(fixture.temp.path().join(collision_path))?,
1265            sentinel
1266        );
1267        Ok(())
1268    }
1269
1270    #[tokio::test]
1271    async fn rewrite_cleans_a_sidecar_left_by_failed_write_cleanup() -> TestResult {
1272        let fixture = rewrite_fixture().await?;
1273        let attempt_id = Uuid::from_u128(4);
1274        let data_path = staged_data_path(attempt_id, 0);
1275        let coverage_path = staged_coverage_path(&fixture, attempt_id, 0)?;
1276        crate::storage::inject_write_new_failure(fixture.temp.path().join(&coverage_path), true);
1277
1278        let error = rewrite_with_attempt_id(
1279            &fixture.location,
1280            &fixture.table_schema,
1281            &fixture.index,
1282            &fixture.source,
1283            attempt_id,
1284        )
1285        .await
1286        .expect_err("sidecar write and its first cleanup must fail");
1287
1288        assert!(matches!(
1289            error,
1290            EntityRewriteError::CoverageSidecar {
1291                source: CoverageSidecarError::Storage {
1292                    source: StorageError::CleanupFailed { .. }
1293                }
1294            }
1295        ));
1296        assert!(!fixture.temp.path().join(data_path).exists());
1297        assert!(!fixture.temp.path().join(coverage_path).exists());
1298        Ok(())
1299    }
1300
1301    #[tokio::test]
1302    async fn rewrite_cleans_completed_outputs_when_a_later_finish_fails() -> TestResult {
1303        let fixture = rewrite_fixture().await?;
1304        let attempt_id = Uuid::from_u128(5);
1305        let owned_paths = [
1306            staged_data_path(attempt_id, 0),
1307            staged_coverage_path(&fixture, attempt_id, 0)?,
1308            staged_data_path(attempt_id, 1),
1309        ];
1310        crate::storage::inject_output_finish_failure(fixture.temp.path().join(&owned_paths[2]));
1311
1312        let error = rewrite_with_attempt_id(
1313            &fixture.location,
1314            &fixture.table_schema,
1315            &fixture.index,
1316            &fixture.source,
1317            attempt_id,
1318        )
1319        .await
1320        .expect_err("second output finish must fail");
1321
1322        assert!(matches!(
1323            error,
1324            EntityRewriteError::Storage {
1325                source: StorageError::OtherIo { .. }
1326            }
1327        ));
1328        for path in &owned_paths {
1329            assert!(!fixture.temp.path().join(path).exists(), "{path} leaked");
1330        }
1331        assert!(fixture.temp.path().join(&fixture.source.path).exists());
1332        assert!(
1333            fixture
1334                .temp
1335                .path()
1336                .join(
1337                    fixture
1338                        .source
1339                        .coverage_path
1340                        .as_deref()
1341                        .expect("source sidecar")
1342                )
1343                .exists()
1344        );
1345        Ok(())
1346    }
1347
1348    #[tokio::test]
1349    async fn rewrite_read_failure_never_creates_staged_objects() -> TestResult {
1350        let fixture = rewrite_fixture().await?;
1351        std::fs::remove_file(fixture.temp.path().join(&fixture.source.path))?;
1352
1353        let error = rewrite_mixed_parquet_segment(
1354            &fixture.location,
1355            &fixture.table_schema,
1356            &fixture.index,
1357            &fixture.source,
1358        )
1359        .await
1360        .expect_err("missing source must fail");
1361
1362        assert!(matches!(
1363            error,
1364            EntityRewriteError::SegmentInspection { .. }
1365        ));
1366        assert_nothing_staged(&fixture);
1367        assert!(
1368            fixture
1369                .temp
1370                .path()
1371                .join(
1372                    fixture
1373                        .source
1374                        .coverage_path
1375                        .as_deref()
1376                        .expect("source sidecar")
1377                )
1378                .exists()
1379        );
1380        Ok(())
1381    }
1382
1383    #[tokio::test]
1384    async fn rewrite_rejects_wrong_layout_and_missing_pointer_before_staging() -> TestResult {
1385        let mut fixture = rewrite_fixture().await?;
1386        let identity = fixture
1387            .source_coverage
1388            .iter()
1389            .next()
1390            .expect("fixture identity")
1391            .0
1392            .clone();
1393        fixture.source.entity_layout = SegmentEntityLayout::Single(identity);
1394        let error = rewrite_mixed_parquet_segment(
1395            &fixture.location,
1396            &fixture.table_schema,
1397            &fixture.index,
1398            &fixture.source,
1399        )
1400        .await
1401        .expect_err("single-entity source must be rejected");
1402        assert!(matches!(error, EntityRewriteError::InvalidInput { .. }));
1403
1404        fixture.source.entity_layout = SegmentEntityLayout::Mixed;
1405        fixture.source.coverage_path = None;
1406        let error = rewrite_mixed_parquet_segment(
1407            &fixture.location,
1408            &fixture.table_schema,
1409            &fixture.index,
1410            &fixture.source,
1411        )
1412        .await
1413        .expect_err("missing coverage pointer must be rejected");
1414        assert!(matches!(error, EntityRewriteError::InvalidInput { .. }));
1415        assert_nothing_staged(&fixture);
1416        Ok(())
1417    }
1418
1419    #[tokio::test]
1420    async fn rewrite_rejects_stale_metadata_and_schema_before_staging() -> TestResult {
1421        let fixture = rewrite_fixture().await?;
1422        let mut stale_source = fixture.source.clone();
1423        stale_source.row_count += 1;
1424        let error = rewrite_mixed_parquet_segment(
1425            &fixture.location,
1426            &fixture.table_schema,
1427            &fixture.index,
1428            &stale_source,
1429        )
1430        .await
1431        .expect_err("stale row count must be rejected");
1432        assert!(matches!(error, EntityRewriteError::InvalidInput { .. }));
1433
1434        let mut columns = fixture.table_schema.columns().to_vec();
1435        columns[2].nullable = true;
1436        let wrong_schema = LogicalSchema::new(columns)?;
1437        let error = rewrite_mixed_parquet_segment(
1438            &fixture.location,
1439            &wrong_schema,
1440            &fixture.index,
1441            &fixture.source,
1442        )
1443        .await
1444        .expect_err("schema mismatch must be rejected");
1445        assert!(matches!(
1446            error,
1447            EntityRewriteError::SegmentSchemaValidation { .. }
1448        ));
1449        assert_nothing_staged(&fixture);
1450        Ok(())
1451    }
1452
1453    #[tokio::test]
1454    async fn rewrite_rejects_missing_corrupt_and_stale_coverage_before_staging() -> TestResult {
1455        let fixture = rewrite_fixture().await?;
1456        let coverage_path = fixture
1457            .source
1458            .coverage_path
1459            .as_deref()
1460            .expect("fixture coverage path");
1461        let absolute_coverage_path = fixture.temp.path().join(coverage_path);
1462        std::fs::remove_file(&absolute_coverage_path)?;
1463        let error = rewrite_mixed_parquet_segment(
1464            &fixture.location,
1465            &fixture.table_schema,
1466            &fixture.index,
1467            &fixture.source,
1468        )
1469        .await
1470        .expect_err("missing coverage object must fail");
1471        assert!(matches!(
1472            error,
1473            EntityRewriteError::CoverageSidecar {
1474                source: CoverageSidecarError::Storage {
1475                    source: StorageError::NotFound { .. }
1476                }
1477            }
1478        ));
1479
1480        std::fs::write(&absolute_coverage_path, b"not entity coverage")?;
1481        let error = rewrite_mixed_parquet_segment(
1482            &fixture.location,
1483            &fixture.table_schema,
1484            &fixture.index,
1485            &fixture.source,
1486        )
1487        .await
1488        .expect_err("corrupt coverage object must fail");
1489        assert!(matches!(
1490            error,
1491            EntityRewriteError::CoverageSidecar {
1492                source: CoverageSidecarError::Codec { .. }
1493            }
1494        ));
1495
1496        let mut stale_coverage = EntityCoverage::empty();
1497        for (ordinal, (identity, coverage)) in fixture.source_coverage.iter().enumerate() {
1498            let coverage = if ordinal == 0 {
1499                coverage.union(&std::iter::once(u64::MAX).collect())
1500            } else {
1501                coverage.clone()
1502            };
1503            stale_coverage.union_coverage(identity.clone(), coverage);
1504        }
1505        std::fs::write(
1506            &absolute_coverage_path,
1507            entity_coverage_to_bytes(&stale_coverage)?,
1508        )?;
1509        let error = rewrite_mixed_parquet_segment(
1510            &fixture.location,
1511            &fixture.table_schema,
1512            &fixture.index,
1513            &fixture.source,
1514        )
1515        .await
1516        .expect_err("stale coverage object must fail");
1517        assert!(matches!(error, EntityRewriteError::InvalidInput { .. }));
1518
1519        let mut empty_identity_coverage = EntityCoverage::empty();
1520        for (ordinal, (identity, coverage)) in fixture.source_coverage.iter().enumerate() {
1521            empty_identity_coverage.union_coverage(
1522                identity.clone(),
1523                if ordinal == 0 {
1524                    crate::coverage::Coverage::empty()
1525                } else {
1526                    coverage.clone()
1527                },
1528            );
1529        }
1530        std::fs::write(
1531            &absolute_coverage_path,
1532            entity_coverage_to_bytes(&empty_identity_coverage)?,
1533        )?;
1534        let error = rewrite_mixed_parquet_segment(
1535            &fixture.location,
1536            &fixture.table_schema,
1537            &fixture.index,
1538            &fixture.source,
1539        )
1540        .await
1541        .expect_err("identity without a covered index interval must fail");
1542        assert!(matches!(error, EntityRewriteError::InvalidInput { .. }));
1543        assert_nothing_staged(&fixture);
1544        Ok(())
1545    }
1546}