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