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