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