1use 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#[derive(Debug, Clone, PartialEq, Eq)]
58pub struct StagedEntityReplacement {
59 pub identity: EntityIdentity,
61 pub meta: SegmentMeta,
63 pub coverage: EntityCoverage,
65}
66
67#[derive(Debug, Clone, PartialEq, Eq)]
69pub struct StagedEntityRewrite {
70 pub source_path: String,
72 pub replacements: Vec<StagedEntityReplacement>,
74 pub staged_object_paths: Vec<String>,
76 pub rows_read: u64,
78 pub rows_written: u64,
80 pub materialized_identities: Vec<EntityIdentity>,
82}
83
84#[derive(Debug, Snafu)]
86#[non_exhaustive]
87pub enum EntityRewriteError {
88 #[snafu(display("Invalid mixed-segment rewrite input: {reason}"))]
90 InvalidInput {
91 reason: String,
93 backtrace: Backtrace,
95 },
96
97 #[snafu(display("Invalid staged entity rewrite output: {reason}"))]
99 InvalidOutput {
100 reason: String,
102 backtrace: Backtrace,
104 },
105
106 #[snafu(display("Invalid {description} path {path:?}: {source}"))]
108 InvalidPath {
109 description: &'static str,
111 path: String,
113 #[snafu(source(from(StorageError, Box::new)), backtrace)]
115 source: Box<StorageError>,
116 },
117
118 #[snafu(display("Invalid rewrite ordered-index specification: {source}"))]
120 IndexSpecValidation {
121 source: IndexSpecError,
123 backtrace: Backtrace,
125 },
126
127 #[snafu(display("Rewrite table schema validation failed: {source}"))]
129 TableSchemaValidation {
130 #[snafu(source(from(SchemaCompatibilityError, Box::new)), backtrace)]
132 source: Box<SchemaCompatibilityError>,
133 },
134
135 #[snafu(display("Rewrite segment schema validation failed for {path}: {source}"))]
137 SegmentSchemaValidation {
138 path: String,
140 #[snafu(source(from(SchemaCompatibilityError, Box::new)), backtrace)]
142 source: Box<SchemaCompatibilityError>,
143 },
144
145 #[snafu(display("Invalid rewrite source metadata: {source}"))]
147 SegmentMetadataValidation {
148 #[snafu(source(from(SegmentMetaError, Box::new)), backtrace)]
150 source: Box<SegmentMetaError>,
151 },
152
153 #[snafu(display("Failed to inspect Parquet segment: {source}"))]
155 SegmentInspection {
156 #[snafu(source, backtrace)]
158 source: SegmentError,
159 },
160
161 #[snafu(display("Failed to inspect exact entity coverage: {source}"))]
163 CoverageInspection {
164 #[snafu(source, backtrace)]
166 source: SegmentCoverageError,
167 },
168
169 #[snafu(display("Failed to access entity coverage sidecar: {source}"))]
171 CoverageSidecar {
172 #[snafu(source, backtrace)]
174 source: CoverageSidecarError,
175 },
176
177 #[snafu(display("Failed to serialize staged entity coverage: {source}"))]
179 CoverageSerialization {
180 #[snafu(source, backtrace)]
182 source: CoverageCodecError,
183 },
184
185 #[snafu(display("Staged entity rewrite storage failure: {source}"))]
187 Storage {
188 #[snafu(source, backtrace)]
190 source: StorageError,
191 },
192
193 #[snafu(display("Parquet rewrite failure at {path}: {source}"))]
195 Parquet {
196 path: String,
198 source: ParquetError,
200 backtrace: Backtrace,
202 },
203
204 #[snafu(display("Arrow row filtering failed for {path}: {source}"))]
206 Arrow {
207 path: String,
209 source: ArrowError,
211 backtrace: Backtrace,
213 },
214
215 #[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 #[snafu(source, backtrace)]
227 source: Box<EntityRewriteError>,
228 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 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
666pub 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 ×tamps.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}