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