1pub(super) mod error;
12
13use error::AppendError;
14
15use std::{marker::PhantomData, path::Path, sync::Arc, time::Instant};
16
17use arrow::{
18 array::{RecordBatch as ArrowRecordBatch, RecordBatchIterator, RecordBatchReader},
19 datatypes::SchemaRef,
20 error::ArrowError,
21};
22use parquet::arrow::ArrowWriter as ParquetArrowWriter;
23use parquet::basic::Compression;
24use parquet::file::properties::WriterProperties;
25use snafu::prelude::*;
26use uuid::Uuid;
27
28use crate::{
29 coverage::serde::{coverage_to_bytes, entity_coverage_to_bytes},
30 coverage::{
31 EntityCoverage,
32 index_interval::index_interval_for_id,
33 io::{CoverageSidecarError, write_coverage_sidecar_new_bytes},
34 layout::{
35 coverage_file_id_for_attempt, segment_coverage_id_v2, segment_coverage_key,
36 segment_entity_coverage_id_v1, table_coverage_id_v2, table_entity_coverage_id_v1,
37 table_snapshot_key,
38 },
39 },
40 formats::parquet::{
41 compute_segment_entity_coverage, coverage::compute_segment_coverage,
42 logical_schema_from_parquet, segment_meta::segment_meta_from_parquet,
43 },
44 metadata::{
45 logical_schema::LogicalSchema,
46 schema_compat::{ensure_index_spec_matches_schema, ensure_schema_fields_match_by_name},
47 segments::SegmentEntityLayout,
48 },
49 storage,
50 table::{TableError, TimeSeriesTable},
51 transaction_log::{
52 LogAction, TableState, checked_next_version, table_state::TableCoveragePointer,
53 },
54};
55
56use self::error::{
57 ArrowInputSnafu, ArrowToLogicalSchemaSnafu, GeneratedSegmentSchemaCompatibilitySnafu,
58 ParquetWriteSnafu,
59};
60use super::append_schema::AppendSchemaNormalizer;
61
62fn classify_entity_layout(
63 segment_path: &str,
64 coverage: &EntityCoverage,
65) -> Result<SegmentEntityLayout, AppendError> {
66 let first_identity = coverage
67 .iter()
68 .next()
69 .map(|(identity, _)| identity)
70 .ok_or_else(|| AppendError::EmptySegmentEntityCoverage {
71 segment_path: segment_path.to_string(),
72 })?;
73
74 if let Some((identity, _)) = coverage.iter().find(|(_, coverage)| coverage.is_empty()) {
75 return Err(AppendError::EntityWithoutIndexCoverage {
76 segment_path: segment_path.to_string(),
77 identity: identity.clone(),
78 });
79 }
80
81 Ok(if coverage.identity_count() == 1 {
82 SegmentEntityLayout::Single(first_identity.clone())
83 } else {
84 SegmentEntityLayout::Mixed
85 })
86}
87
88fn ensure_existing_segments_have_coverage(state: &TableState) -> Result<(), AppendError> {
89 for seg in state.segments.values() {
90 if seg.coverage_path.is_none() {
91 return Err(AppendError::ExistingSegmentMissingCoverageMetadata {
92 segment_path: seg.path.clone(),
93 });
94 }
95 }
96
97 Ok(())
98}
99
100fn elapsed_ms(started: Instant) -> u64 {
101 u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX)
102}
103
104struct AppendTelemetry {
105 track_input_counts: bool,
106 append_started: Instant,
107 phase_started: Instant,
108 phase: &'static str,
109 last_completed_phase: Option<&'static str>,
110 batches_consumed: u64,
111 rows_consumed: u64,
112 failed_phase_duration_ms: Option<u64>,
113 cleanup_duration_ms: u64,
114 cleanup_outcome: Option<&'static str>,
115}
116
117impl AppendTelemetry {
118 fn new(track_input_counts: bool) -> Self {
119 let now = Instant::now();
120 Self {
121 track_input_counts,
122 append_started: now,
123 phase_started: now,
124 phase: "source_schema_preparation",
125 last_completed_phase: None,
126 batches_consumed: 0,
127 rows_consumed: 0,
128 failed_phase_duration_ms: None,
129 cleanup_duration_ms: 0,
130 cleanup_outcome: None,
131 }
132 }
133
134 fn start_phase(&mut self, phase: &'static str) {
135 self.phase = phase;
136 self.phase_started = Instant::now();
137 }
138
139 fn complete_phase(&mut self) {
140 tracing::debug!(
141 name: "table.append",
142 target: "timeseries_table_format::table::append",
143 phase = self.phase,
144 duration_ms = elapsed_ms(self.phase_started),
145 outcome = "succeeded",
146 "Completed append phase"
147 );
148 self.last_completed_phase = Some(self.phase);
149 }
150
151 fn record_batch(&mut self, batch: &ArrowRecordBatch) {
152 if !self.track_input_counts {
153 return;
154 }
155 self.batches_consumed = self.batches_consumed.saturating_add(1);
156 self.rows_consumed = self
157 .rows_consumed
158 .saturating_add(u64::try_from(batch.num_rows()).unwrap_or(u64::MAX));
159 }
160
161 fn record_cleanup(&mut self, started: Instant, error: &AppendError) {
162 let phase_duration_ms = u64::try_from(
163 started
164 .saturating_duration_since(self.phase_started)
165 .as_millis(),
166 )
167 .unwrap_or(u64::MAX);
168 self.failed_phase_duration_ms
169 .get_or_insert(phase_duration_ms);
170 self.cleanup_duration_ms = self.cleanup_duration_ms.saturating_add(elapsed_ms(started));
171 if self.cleanup_outcome != Some("failed") {
172 self.cleanup_outcome = Some(if matches!(error, AppendError::Rollback { .. }) {
173 "failed"
174 } else {
175 "succeeded"
176 });
177 }
178 }
179
180 fn record_failure(&self, error: &AppendError) {
181 let outcome = if matches!(error, AppendError::CommitAmbiguous { .. }) {
182 "ambiguous"
183 } else {
184 "failed"
185 };
186 let last_completed_phase = self.last_completed_phase.unwrap_or("none");
187 let phase_duration_ms = self
188 .failed_phase_duration_ms
189 .unwrap_or_else(|| elapsed_ms(self.phase_started));
190 let span = tracing::Span::current();
191 span.record("outcome", outcome);
192 span.record("failed_phase", self.phase);
193 span.record("last_completed_phase", last_completed_phase);
194 span.record("batches_consumed", self.batches_consumed);
195 span.record("rows_consumed", self.rows_consumed);
196
197 let (cleanup_duration_ms, cleanup_outcome, retained_path) =
198 if let Some(cleanup_outcome) = self.cleanup_outcome {
199 (Some(self.cleanup_duration_ms), Some(cleanup_outcome), None)
200 } else if let AppendError::CommitAmbiguous { segment_path, .. } = error {
201 (None, Some("not_attempted"), Some(segment_path.as_str()))
202 } else {
203 (None, None, None)
204 };
205 if let Some(cleanup_outcome) = cleanup_outcome {
206 span.record("cleanup_outcome", cleanup_outcome);
207 }
208 tracing::debug!(
209 name: "table.append",
210 target: "timeseries_table_format::table::append",
211 phase = self.phase,
212 last_completed_phase,
213 duration_ms = phase_duration_ms,
214 total_duration_ms = elapsed_ms(self.append_started),
215 batches_consumed = self.batches_consumed,
216 rows_consumed = self.rows_consumed,
217 cleanup_duration_ms,
218 cleanup_outcome,
219 retained_path,
220 outcome,
221 "Append phase failed"
222 );
223 }
224}
225
226#[doc(hidden)]
231pub struct RecordBatchReaderSourceKind;
232
233#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
235#[non_exhaustive]
236pub enum ParquetCompression {
237 Uncompressed,
239 Snappy,
241 #[default]
243 Zstd,
244}
245
246impl From<ParquetCompression> for Compression {
247 fn from(value: ParquetCompression) -> Self {
248 match value {
249 ParquetCompression::Uncompressed => Self::UNCOMPRESSED,
250 ParquetCompression::Snappy => Self::SNAPPY,
251 ParquetCompression::Zstd => Self::ZSTD(Default::default()),
252 }
253 }
254}
255
256impl ParquetCompression {
257 pub fn as_str(self) -> &'static str {
259 match self {
260 Self::Uncompressed => "uncompressed",
261 Self::Snappy => "snappy",
262 Self::Zstd => "zstd",
263 }
264 }
265}
266
267const DEFAULT_MAX_ROWS_PER_ROW_GROUP: usize = 1024 * 1024;
268const DEFAULT_MAX_BYTES_PER_ROW_GROUP: usize = 128 * 1024 * 1024;
269
270#[derive(Debug, Clone, Copy, PartialEq, Eq)]
271struct AppendWriterSettings {
272 compression: ParquetCompression,
273 max_rows_per_row_group: usize,
274 max_bytes_per_row_group: usize,
275}
276
277impl Default for AppendWriterSettings {
278 fn default() -> Self {
279 Self {
280 compression: ParquetCompression::default(),
281 max_rows_per_row_group: DEFAULT_MAX_ROWS_PER_ROW_GROUP,
282 max_bytes_per_row_group: DEFAULT_MAX_BYTES_PER_ROW_GROUP,
283 }
284 }
285}
286
287impl AppendWriterSettings {
288 fn from_source<S, SourceKind>(source: &S) -> Self
289 where
290 S: IntoRecordBatchReader<SourceKind>,
291 {
292 Self {
293 compression: source.effective_compression().unwrap_or_default(),
294 max_rows_per_row_group: source
295 .effective_max_rows_per_row_group()
296 .unwrap_or(DEFAULT_MAX_ROWS_PER_ROW_GROUP),
297 max_bytes_per_row_group: source
298 .effective_max_bytes_per_row_group()
299 .unwrap_or(DEFAULT_MAX_BYTES_PER_ROW_GROUP),
300 }
301 }
302
303 fn validate(self) -> Result<Self, AppendError> {
304 if self.max_rows_per_row_group == 0 {
305 return Err(AppendError::InvalidMaxRowsPerRowGroup {
306 max_rows_per_row_group: 0,
307 });
308 }
309 if self.max_bytes_per_row_group == 0 {
310 return Err(AppendError::InvalidMaxBytesPerRowGroup {
311 max_bytes_per_row_group: 0,
312 });
313 }
314 Ok(self)
315 }
316
317 fn writer_properties(self) -> WriterProperties {
318 WriterProperties::builder()
319 .set_compression(self.compression.into())
320 .set_max_row_group_row_count(Some(self.max_rows_per_row_group))
321 .set_max_row_group_bytes(Some(self.max_bytes_per_row_group))
322 .build()
323 }
324}
325
326#[derive(Debug, Clone, PartialEq, Eq)]
328#[non_exhaustive]
329pub struct AppendReport {
330 pub starting_version: u64,
332 pub committed_version: u64,
334 pub segment_path: String,
336 pub row_count: u64,
338 pub row_group_count: usize,
340 pub file_size_bytes: u64,
342 pub compression: ParquetCompression,
344 pub max_rows_per_row_group: usize,
346 pub max_bytes_per_row_group: usize,
348}
349
350pub trait IntoRecordBatchReader<SourceKind = RecordBatchReaderSourceKind> {
356 type Reader: RecordBatchReader;
358
359 #[doc(hidden)]
361 fn effective_compression(&self) -> Option<ParquetCompression> {
362 None
363 }
364
365 #[doc(hidden)]
367 fn effective_max_rows_per_row_group(&self) -> Option<usize> {
368 None
369 }
370
371 #[doc(hidden)]
373 fn effective_max_bytes_per_row_group(&self) -> Option<usize> {
374 None
375 }
376
377 fn into_record_batch_reader(self) -> Result<Self::Reader, AppendError>;
379}
380
381#[derive(Debug)]
389#[must_use = "an append request has no effect until passed to TimeSeriesTable::append"]
390pub struct AppendRequest<S> {
391 source: S,
392 compression: Option<ParquetCompression>,
393 max_rows_per_row_group: Option<usize>,
394 max_bytes_per_row_group: Option<usize>,
395}
396
397impl<S> AppendRequest<S> {
398 pub fn new(source: S) -> Self {
403 Self {
404 source,
405 compression: None,
406 max_rows_per_row_group: None,
407 max_bytes_per_row_group: None,
408 }
409 }
410
411 pub fn compression(mut self, compression: ParquetCompression) -> Self {
413 self.compression = Some(compression);
414 self
415 }
416
417 pub fn max_rows_per_row_group(mut self, max_rows_per_row_group: usize) -> Self {
424 self.max_rows_per_row_group = Some(max_rows_per_row_group);
425 self
426 }
427
428 pub fn max_bytes_per_row_group(mut self, max_bytes_per_row_group: usize) -> Self {
434 self.max_bytes_per_row_group = Some(max_bytes_per_row_group);
435 self
436 }
437}
438
439#[doc(hidden)]
442impl<S, SourceKind> IntoRecordBatchReader<PhantomData<SourceKind>> for AppendRequest<S>
443where
444 S: IntoRecordBatchReader<SourceKind>,
445{
446 type Reader = S::Reader;
447
448 fn effective_compression(&self) -> Option<ParquetCompression> {
449 self.compression
450 .or_else(|| self.source.effective_compression())
451 }
452
453 fn effective_max_rows_per_row_group(&self) -> Option<usize> {
454 self.max_rows_per_row_group
455 .or_else(|| self.source.effective_max_rows_per_row_group())
456 }
457
458 fn effective_max_bytes_per_row_group(&self) -> Option<usize> {
459 self.max_bytes_per_row_group
460 .or_else(|| self.source.effective_max_bytes_per_row_group())
461 }
462
463 fn into_record_batch_reader(self) -> Result<Self::Reader, AppendError> {
464 self.source.into_record_batch_reader()
465 }
466}
467
468impl<R> IntoRecordBatchReader<RecordBatchReaderSourceKind> for R
469where
470 R: RecordBatchReader,
471{
472 type Reader = R;
473
474 fn into_record_batch_reader(self) -> Result<Self::Reader, AppendError> {
475 Ok(self)
476 }
477}
478
479impl IntoRecordBatchReader<ArrowRecordBatch> for ArrowRecordBatch {
480 type Reader = Box<dyn RecordBatchReader + Send>;
481
482 fn into_record_batch_reader(self) -> Result<Self::Reader, AppendError> {
483 let schema = self.schema();
484 Ok(Box::new(RecordBatchIterator::new(
485 std::iter::once(Ok(self)),
486 schema,
487 )))
488 }
489}
490
491impl IntoRecordBatchReader<Vec<ArrowRecordBatch>> for Vec<ArrowRecordBatch> {
492 type Reader = Box<dyn RecordBatchReader + Send>;
493
494 fn into_record_batch_reader(self) -> Result<Self::Reader, AppendError> {
495 let schema = self
496 .first()
497 .map(ArrowRecordBatch::schema)
498 .ok_or(AppendError::EmptyInput)?;
499 Ok(Box::new(RecordBatchIterator::new(
500 self.into_iter().map(Ok::<_, ArrowError>),
501 schema,
502 )))
503 }
504}
505
506fn ensure_batch_matches_reader_schema(
507 reader_schema: &SchemaRef,
508 batch: &ArrowRecordBatch,
509) -> Result<(), AppendError> {
510 if batch.schema() == *reader_schema {
511 Ok(())
512 } else {
513 Err(ArrowError::SchemaError(
514 "record batch schema does not match its reader schema".to_string(),
515 ))
516 .context(ArrowInputSnafu)
517 }
518}
519
520impl TimeSeriesTable {
521 async fn rollback_created_artifacts(
522 &self,
523 created_paths: &[String],
524 source: AppendError,
525 ) -> AppendError {
526 let (source, mut cleanup_errors) = match source {
527 AppendError::Rollback {
528 source,
529 cleanup_errors,
530 } => (source, cleanup_errors),
531 source => (Box::new(source), Vec::new()),
532 };
533 for path in created_paths.iter().rev() {
534 if let Err(error) =
535 storage::remove_file_if_exists(self.location().as_ref(), Path::new(path)).await
536 {
537 cleanup_errors.push(error);
538 }
539 }
540
541 if cleanup_errors.is_empty() {
542 *source
543 } else {
544 AppendError::Rollback {
545 source,
546 cleanup_errors,
547 }
548 }
549 }
550
551 fn build_append_schema_normalizer(
552 &self,
553 incoming_schema: SchemaRef,
554 ) -> Result<AppendSchemaNormalizer, AppendError> {
555 ensure_existing_segments_have_coverage(&self.state)?;
556
557 match self.state.table_meta.logical_schema.as_ref() {
558 None if self.state.version == 1 => {
559 let incoming_logical_schema =
560 LogicalSchema::try_from_arrow_schema(incoming_schema.as_ref())
561 .context(ArrowToLogicalSchemaSnafu)?;
562 ensure_index_spec_matches_schema(&incoming_logical_schema, &self.index)
563 .map_err(AppendError::from)?;
564 Ok(AppendSchemaNormalizer::without_conversion(incoming_schema))
565 }
566 None => Err(AppendError::MissingCanonicalTableSchema {
567 version: self.state.version,
568 }),
569 Some(table_schema) => {
570 ensure_index_spec_matches_schema(table_schema, &self.index)
571 .map_err(AppendError::from)?;
572 let normalizer = AppendSchemaNormalizer::for_registered_schema(
573 incoming_schema.as_ref(),
574 table_schema,
575 )
576 .map_err(AppendError::from)?;
577 Ok(normalizer)
578 }
579 }
580 }
581
582 async fn publish_generated_parquet_segment(
583 &mut self,
584 relative_path: &str,
585 next_version: u64,
586 owned_data_guard: &mut storage::FileCleanupGuard,
587 telemetry: &mut AppendTelemetry,
588 writer_settings: AppendWriterSettings,
589 ) -> Result<AppendReport, AppendError> {
590 let rel_path = Path::new(relative_path);
591 let expected_version = self.state.version;
592
593 telemetry.start_phase("segment_validation");
595 ensure_existing_segments_have_coverage(&self.state)?;
596
597 let (mut segment_meta, meta_report) =
599 segment_meta_from_parquet(self.location(), rel_path, &self.index)
600 .await
601 .map_err(AppendError::from)?;
602 let row_count = meta_report.row_count;
603 let file_size_bytes = meta_report.file_size_bytes;
604 let span = tracing::Span::current();
605 span.record("row_count", row_count);
606 span.record("file_size_bytes", file_size_bytes);
607
608 let segment_schema = logical_schema_from_parquet(self.location(), rel_path)
609 .await
610 .map_err(AppendError::from)?;
611 ensure_index_spec_matches_schema(&segment_schema, &self.index).context(
612 GeneratedSegmentSchemaCompatibilitySnafu {
613 segment_path: relative_path.to_string(),
614 },
615 )?;
616
617 let maybe_table_schema = self.state.table_meta.logical_schema.as_ref();
625
626 let maybe_updated_meta = match maybe_table_schema {
627 None if expected_version == 1 => {
628 let mut updated_meta = self.state.table_meta.clone();
629 updated_meta.logical_schema = Some(segment_schema.clone());
630 Some(updated_meta)
631 }
632 None => {
633 return Err(AppendError::MissingCanonicalTableSchema {
634 version: expected_version,
635 });
636 }
637 Some(table_schema) => {
638 ensure_schema_fields_match_by_name(table_schema, &segment_schema, &self.index)
639 .context(GeneratedSegmentSchemaCompatibilitySnafu {
640 segment_path: relative_path.to_string(),
641 })?;
642 None
643 }
644 };
645 telemetry.complete_phase();
646
647 let has_entity_columns = !self.index.entity_columns.is_empty();
648
649 telemetry.start_phase("coverage_and_overlap");
651 let (seg_cov_bytes, new_snap_cov_bytes, entity_layout) = if has_entity_columns {
652 let table_cov = self
653 .load_entity_coverage_with_recovery::<AppendError>()
654 .await?;
655
656 let segment_cov =
657 compute_segment_entity_coverage(self.location(), rel_path, &self.index)
658 .await
659 .map_err(AppendError::from)?;
660 let entity_layout = classify_entity_layout(relative_path, &segment_cov)?;
661
662 if let Some((identity, index_interval_id)) =
663 segment_cov.first_overlapping_identity_and_interval_id(&table_cov)
664 {
665 let example_index_interval =
666 index_interval_for_id(&self.index.kind, index_interval_id)
667 .map_err(AppendError::from)?;
668 return Err(AppendError::PersistedIndexIntervalOverlap {
669 segment_path: relative_path.to_string(),
670 overlap_count: segment_cov.intersection_cardinality(&table_cov),
671 example_identity: Some(identity.clone()),
672 example_index_interval_id: index_interval_id,
673 example_index_interval: Box::new(example_index_interval),
674 });
675 }
676
677 let seg_bytes = entity_coverage_to_bytes(&segment_cov)
678 .map_err(CoverageSidecarError::from)
679 .map_err(AppendError::from)?;
680 let snapshot_bytes = entity_coverage_to_bytes(&table_cov.union(&segment_cov))
681 .map_err(CoverageSidecarError::from)
682 .map_err(AppendError::from)?;
683 (seg_bytes, snapshot_bytes, entity_layout)
684 } else {
685 let table_cov = self
686 .load_global_coverage_with_recovery::<AppendError>()
687 .await?;
688
689 let segment_cov = compute_segment_coverage(self.location(), rel_path, &self.index)
690 .await
691 .map_err(AppendError::from)?;
692
693 let overlap = segment_cov.intersect(&table_cov);
694 let overlap_count = overlap.cardinality();
695 if let Some(example_index_interval_id) = overlap.present().iter().next() {
696 let example_index_interval =
697 index_interval_for_id(&self.index.kind, example_index_interval_id)
698 .map_err(AppendError::from)?;
699 return Err(AppendError::PersistedIndexIntervalOverlap {
700 segment_path: relative_path.to_string(),
701 overlap_count: u128::from(overlap_count),
702 example_identity: None,
703 example_index_interval_id,
704 example_index_interval: Box::new(example_index_interval),
705 });
706 }
707
708 let seg_bytes = coverage_to_bytes(&segment_cov)
709 .map_err(CoverageSidecarError::from)
710 .map_err(AppendError::from)?;
711 let snapshot_bytes = coverage_to_bytes(&table_cov.union(&segment_cov))
712 .map_err(CoverageSidecarError::from)
713 .map_err(AppendError::from)?;
714 (
715 seg_bytes,
716 snapshot_bytes,
717 SegmentEntityLayout::NotApplicable,
718 )
719 };
720 let entity_layout_name = match &entity_layout {
721 SegmentEntityLayout::NotApplicable => "not_applicable",
722 SegmentEntityLayout::Single(_) => "single",
723 SegmentEntityLayout::Mixed => "mixed",
724 };
725 span.record("entity_layout", entity_layout_name);
726 span.record("row_group_count", meta_report.row_group_count);
727 telemetry.complete_phase();
728
729 telemetry.start_phase("sidecar_write");
731 let attempt_id = Uuid::new_v4();
732 let segment_content_id = if has_entity_columns {
733 segment_entity_coverage_id_v1(&self.index, &seg_cov_bytes)
734 } else {
735 segment_coverage_id_v2(&self.index, &seg_cov_bytes)
736 };
737 let segment_file_id = coverage_file_id_for_attempt(&segment_content_id, &attempt_id);
738 let seg_cov_path = segment_coverage_key(&segment_file_id)
739 .map_err(CoverageSidecarError::from)
740 .map_err(AppendError::from)?;
741
742 let snapshot_content_id = if has_entity_columns {
743 table_entity_coverage_id_v1(&self.index, &new_snap_cov_bytes)
744 } else {
745 table_coverage_id_v2(&self.index, &new_snap_cov_bytes)
746 };
747 let snapshot_file_id = coverage_file_id_for_attempt(&snapshot_content_id, &attempt_id);
748 let snapshot_path = table_snapshot_key(next_version, &snapshot_file_id)
749 .map_err(CoverageSidecarError::from)
750 .map_err(AppendError::from)?;
751
752 let mut created_sidecars = Vec::new();
753 let mut segment_sidecar_guard = storage::FileCleanupGuard::new_disarmed(
754 self.location().as_ref(),
755 Path::new(&seg_cov_path),
756 )
757 .map_err(AppendError::from)?;
758 if let Err(source) = write_coverage_sidecar_new_bytes(
759 self.location(),
760 Path::new(&seg_cov_path),
761 &seg_cov_bytes,
762 )
763 .await
764 {
765 if source.storage_cleanup_failed() {
766 created_sidecars.push(seg_cov_path.clone());
767 }
768 let cleanup_started = Instant::now();
769 let error = self
770 .rollback_created_artifacts(&created_sidecars, AppendError::from(source))
771 .await;
772 telemetry.record_cleanup(cleanup_started, &error);
773 return Err(error);
774 }
775 segment_sidecar_guard.arm();
776 created_sidecars.push(seg_cov_path.clone());
777
778 let mut snapshot_sidecar_guard = storage::FileCleanupGuard::new_disarmed(
779 self.location().as_ref(),
780 Path::new(&snapshot_path),
781 )
782 .map_err(AppendError::from)?;
783 if let Err(source) = write_coverage_sidecar_new_bytes(
784 self.location(),
785 Path::new(&snapshot_path),
786 &new_snap_cov_bytes,
787 )
788 .await
789 {
790 if source.storage_cleanup_failed() {
791 created_sidecars.push(snapshot_path.clone());
792 }
793 let error = AppendError::from(source);
794 let cleanup_started = Instant::now();
795 let error = self
796 .rollback_created_artifacts(&created_sidecars, error)
797 .await;
798 telemetry.record_cleanup(cleanup_started, &error);
799 segment_sidecar_guard.disarm();
800 return Err(error);
801 }
802 snapshot_sidecar_guard.arm();
803 created_sidecars.push(snapshot_path.clone());
804 telemetry.complete_phase();
805
806 segment_meta.coverage_path = Some(seg_cov_path);
808 segment_meta.entity_layout = entity_layout;
809
810 let mut actions = Vec::new();
811 if let Some(updated_meta) = maybe_updated_meta.clone() {
812 actions.push(LogAction::UpdateTableMeta(updated_meta));
813 }
814
815 actions.push(LogAction::AddSegment(segment_meta.clone()));
816 actions.push(LogAction::UpdateTableCoverage {
817 index_kind: self.index.kind.clone(),
818 coverage_path: snapshot_path.clone(),
819 });
820
821 telemetry.start_phase("transaction_commit");
822 let new_version = match self
823 .log
824 .commit_with_path_preservation(expected_version, actions, || {
825 owned_data_guard.disarm();
826 segment_sidecar_guard.disarm();
827 snapshot_sidecar_guard.disarm();
828 })
829 .await
830 {
831 Ok(version) => version,
832 Err(source @ crate::transaction_log::CommitError::AmbiguousOutcome { .. }) => {
833 return Err(AppendError::CommitAmbiguous {
834 segment_path: relative_path.to_string(),
835 source: Box::new(source),
836 });
837 }
838 Err(source) => {
839 let error = AppendError::from(source);
840 let cleanup_started = Instant::now();
841 let error = self
842 .rollback_created_artifacts(&created_sidecars, error)
843 .await;
844 telemetry.record_cleanup(cleanup_started, &error);
845 segment_sidecar_guard.disarm();
846 snapshot_sidecar_guard.disarm();
847 return Err(error);
848 }
849 };
850 telemetry.complete_phase();
851
852 assert_eq!(
858 new_version, next_version,
859 "transaction log returned unexpected version: expected {}, got {}",
860 next_version, new_version
861 );
862
863 self.state.version = new_version;
865
866 if let Some(updated_meta) = maybe_updated_meta {
867 self.state.table_meta = updated_meta
868 }
869
870 self.state
871 .segments
872 .insert(segment_meta.path.clone(), segment_meta);
873
874 self.state.table_coverage = Some(TableCoveragePointer {
876 index_kind: self.index.kind.clone(),
877 coverage_path: snapshot_path,
878 version: new_version,
879 });
880
881 span.record("committed_version", new_version);
882 span.record("outcome", "succeeded");
883 let row_group_statistics = meta_report.row_group_statistics;
884 tracing::info!(
885 name: "table.append",
886 target: "timeseries_table_format::table::append",
887 expected_version,
888 committed_version = new_version,
889 row_count,
890 batches_consumed = telemetry.batches_consumed,
891 row_group_count = meta_report.row_group_count,
892 row_group_rows_min = row_group_statistics.rows.min,
893 row_group_rows_median = row_group_statistics.rows.median,
894 row_group_rows_max = row_group_statistics.rows.max,
895 row_group_compressed_bytes_min = row_group_statistics.compressed_bytes.min,
896 row_group_compressed_bytes_median = row_group_statistics.compressed_bytes.median,
897 row_group_compressed_bytes_max = row_group_statistics.compressed_bytes.max,
898 row_group_uncompressed_bytes_min = row_group_statistics.uncompressed_bytes.min,
899 row_group_uncompressed_bytes_median = row_group_statistics.uncompressed_bytes.median,
900 row_group_uncompressed_bytes_max = row_group_statistics.uncompressed_bytes.max,
901 file_size_bytes,
902 entity_layout = entity_layout_name,
903 compression = writer_settings.compression.as_str(),
904 max_rows_per_row_group = writer_settings.max_rows_per_row_group,
905 max_bytes_per_row_group = writer_settings.max_bytes_per_row_group,
906 total_duration_ms = elapsed_ms(telemetry.append_started),
907 outcome = "succeeded",
908 "Appended Parquet segment"
909 );
910 Ok(AppendReport {
911 starting_version: expected_version,
912 committed_version: new_version,
913 segment_path: relative_path.to_string(),
914 row_count,
915 row_group_count: meta_report.row_group_count,
916 file_size_bytes,
917 compression: writer_settings.compression,
918 max_rows_per_row_group: writer_settings.max_rows_per_row_group,
919 max_bytes_per_row_group: writer_settings.max_bytes_per_row_group,
920 })
921 }
922
923 pub fn ensure_append_supported(&self) -> Result<(), TableError> {
928 self.ensure_write_compatible()
929 .map_err(AppendError::from)
930 .context(crate::table::error::AppendSnafu)
931 }
932
933 #[tracing::instrument(
949 name = "table.append",
950 target = "timeseries_table_format::table::append",
951 level = "debug",
952 skip_all,
953 fields(
954 expected_version = self.state.version,
955 segment_path = tracing::field::Empty,
956 row_count = tracing::field::Empty,
957 file_size_bytes = tracing::field::Empty,
958 row_group_count = tracing::field::Empty,
959 batches_consumed = tracing::field::Empty,
960 rows_consumed = tracing::field::Empty,
961 compression = tracing::field::Empty,
962 max_rows_per_row_group = tracing::field::Empty,
963 max_bytes_per_row_group = tracing::field::Empty,
964 committed_version = tracing::field::Empty,
965 entity_layout = tracing::field::Empty,
966 failed_phase = tracing::field::Empty,
967 last_completed_phase = tracing::field::Empty,
968 cleanup_outcome = tracing::field::Empty,
969 outcome = tracing::field::Empty
970 )
971 )]
972 pub async fn append<S, SourceKind>(&mut self, source: S) -> Result<AppendReport, TableError>
973 where
974 S: IntoRecordBatchReader<SourceKind>,
975 {
976 let tracing_enabled = tracing::enabled!(
977 target: "timeseries_table_format::table::append",
978 tracing::Level::INFO
979 ) || tracing::enabled!(
980 target: "timeseries_table_format::table::append",
981 tracing::Level::DEBUG
982 );
983 let log_enabled = log::log_enabled!(
985 target: "timeseries_table_format::table::append",
986 log::Level::Info
987 ) || log::log_enabled!(
988 target: "timeseries_table_format::table::append",
989 log::Level::Debug
990 );
991 let mut telemetry = AppendTelemetry::new(tracing_enabled || log_enabled);
992 let append_result: Result<AppendReport, AppendError> = async {
993 self.ensure_write_compatible().map_err(AppendError::from)?;
994
995 let writer_settings = AppendWriterSettings::from_source(&source).validate()?;
996 let span = tracing::Span::current();
997 span.record("compression", writer_settings.compression.as_str());
998 span.record(
999 "max_rows_per_row_group",
1000 writer_settings.max_rows_per_row_group,
1001 );
1002 span.record(
1003 "max_bytes_per_row_group",
1004 writer_settings.max_bytes_per_row_group,
1005 );
1006 let next_version =
1007 checked_next_version(self.state.version).map_err(AppendError::from)?;
1008 let mut reader = source.into_record_batch_reader()?;
1009 let incoming_schema = reader.schema();
1010 let schema_normalizer =
1011 self.build_append_schema_normalizer(Arc::clone(&incoming_schema))?;
1012 let output_schema = Arc::clone(schema_normalizer.output_schema());
1013 telemetry.complete_phase();
1014
1015 telemetry.start_phase("source_consumption_and_parquet_write");
1016 let first_batch = loop {
1017 let Some(batch) = reader.next().transpose().context(ArrowInputSnafu)? else {
1018 return Err(AppendError::EmptyInput);
1019 };
1020 telemetry.record_batch(&batch);
1021 ensure_batch_matches_reader_schema(&incoming_schema, &batch)?;
1022 if batch.num_rows() != 0 {
1023 break schema_normalizer
1024 .normalize_batch(&batch)
1025 .context(ArrowInputSnafu)?;
1026 }
1027 tokio::task::yield_now().await;
1028 };
1029
1030 let relative_path = format!(
1031 "{}/{}.parquet",
1032 storage::layout::APPEND_DATA_DIR,
1033 Uuid::new_v4()
1034 );
1035 tracing::Span::current().record("segment_path", relative_path.as_str());
1036 let mut data_guard = storage::FileCleanupGuard::new_disarmed(
1037 self.location().as_ref(),
1038 Path::new(&relative_path),
1039 )
1040 .map_err(AppendError::from)?;
1041 let sink =
1042 storage::open_new_output_sink(self.location().as_ref(), Path::new(&relative_path))
1043 .await
1044 .map_err(AppendError::from)?;
1045 let write_result = async {
1046 let writer_properties = writer_settings.writer_properties();
1047 let mut writer =
1048 ParquetArrowWriter::try_new(sink, output_schema, Some(writer_properties))
1049 .context(ParquetWriteSnafu)?;
1050 writer.write(&first_batch).context(ParquetWriteSnafu)?;
1051 drop(first_batch);
1052 tokio::task::yield_now().await;
1053
1054 for batch in reader {
1055 let batch = batch.context(ArrowInputSnafu)?;
1056 telemetry.record_batch(&batch);
1057 ensure_batch_matches_reader_schema(&incoming_schema, &batch)?;
1058 if batch.num_rows() != 0 {
1059 let batch = schema_normalizer
1060 .normalize_batch(&batch)
1061 .context(ArrowInputSnafu)?;
1062 writer.write(&batch).context(ParquetWriteSnafu)?;
1063 }
1064 drop(batch);
1065 tokio::task::yield_now().await;
1066 }
1067
1068 telemetry.complete_phase();
1069 telemetry.start_phase("writer_finalize_and_sync");
1070 let sink = writer.into_inner().context(ParquetWriteSnafu)?;
1071 sink.finish().await.map_err(AppendError::from)?;
1072 telemetry.complete_phase();
1073 Ok(())
1074 }
1075 .await;
1076 if let Err(source) = write_result {
1077 let cleanup_started = Instant::now();
1078 let error = self
1079 .rollback_created_artifacts(std::slice::from_ref(&relative_path), source)
1080 .await;
1081 telemetry.record_cleanup(cleanup_started, &error);
1082 data_guard.disarm();
1083 return Err(error);
1084 }
1085 data_guard.arm();
1086
1087 match self
1088 .publish_generated_parquet_segment(
1089 &relative_path,
1090 next_version,
1091 &mut data_guard,
1092 &mut telemetry,
1093 writer_settings,
1094 )
1095 .await
1096 {
1097 Ok(report) => Ok(report),
1098 Err(source @ AppendError::CommitAmbiguous { .. }) => Err(source),
1099 Err(source) => {
1100 let cleanup_started = Instant::now();
1101 let error = self
1102 .rollback_created_artifacts(std::slice::from_ref(&relative_path), source)
1103 .await;
1104 telemetry.record_cleanup(cleanup_started, &error);
1105 data_guard.disarm();
1106 Err(error)
1107 }
1108 }
1109 }
1110 .await;
1111 let span = tracing::Span::current();
1112 span.record("batches_consumed", telemetry.batches_consumed);
1113 span.record("rows_consumed", telemetry.rows_consumed);
1114 if let Err(error) = &append_result {
1115 telemetry.record_failure(error);
1116 }
1117 append_result.context(crate::table::error::AppendSnafu)
1118 }
1119}
1120
1121#[cfg(test)]
1122mod tests {
1123 use super::*;
1124 use snafu::IntoError;
1125
1126 use crate::coverage::io::{
1127 CoverageSidecarError, read_coverage_sidecar, read_entity_coverage_sidecar,
1128 };
1129 use crate::coverage::layout::{SEGMENT_COVERAGE_DIR, TABLE_SNAPSHOT_DIR};
1130 use crate::coverage::serde::entity_coverage_from_bytes;
1131 use crate::coverage::{EntityCoverage, EntityIdentity, EntityValue};
1132 use crate::formats::parquet::SegmentCoverageError;
1133 use crate::metadata::logical_schema::{
1134 LogicalDataType, LogicalField, LogicalSchema, LogicalTimestampUnit,
1135 };
1136 use crate::metadata::segments::SegmentEntityLayout;
1137 use crate::metadata::{index::IndexValue, protocol::TableProtocolError};
1138 use crate::storage::layout;
1139 use crate::storage::{StorageError, StorageLocation, TableLocation};
1140 use crate::table::test_util::*;
1141 use crate::transaction_log::{
1142 CommitError, IndexKind, IndexSpec, TableMeta, TimeIndexGranularity, segments::SegmentError,
1143 };
1144 use arrow::{
1145 array::{
1146 Array, ArrayRef, BinaryArray, BooleanArray, Float32Array, Float64Array, Int16Array,
1147 Int32Array, Int64Array, StringArray, TimestampMillisecondArray, UInt32Array,
1148 UInt64Array, new_null_array,
1149 },
1150 datatypes::{DataType, Field, Fields, Schema, TimeUnit as ArrowTimeUnit},
1151 record_batch::RecordBatch,
1152 };
1153 use futures::{FutureExt, StreamExt};
1154 use parquet::arrow::{ArrowWriter, arrow_reader::ParquetRecordBatchReaderBuilder};
1155 use snafu::ErrorCompat;
1156 use std::cell::Cell;
1157 use std::collections::{BTreeMap, HashMap};
1158 use std::fs::File;
1159 use std::num::NonZeroU64;
1160 use std::path::PathBuf;
1161 use std::rc::Rc;
1162 use std::sync::{Arc, Weak};
1163 use tempfile::TempDir;
1164 use tracing::{Subscriber, instrument::WithSubscriber, span::Id};
1165 use tracing_subscriber::{
1166 Layer,
1167 layer::{Context, SubscriberExt},
1168 registry::LookupSpan,
1169 };
1170
1171 const ROW_GROUP_STATISTIC_FIELDS: [&str; 9] = [
1172 "row_group_rows_min",
1173 "row_group_rows_median",
1174 "row_group_rows_max",
1175 "row_group_compressed_bytes_min",
1176 "row_group_compressed_bytes_median",
1177 "row_group_compressed_bytes_max",
1178 "row_group_uncompressed_bytes_min",
1179 "row_group_uncompressed_bytes_median",
1180 "row_group_uncompressed_bytes_max",
1181 ];
1182
1183 fn assert_no_row_group_statistics(fields: &BTreeMap<&'static str, String>) {
1184 for field in ROW_GROUP_STATISTIC_FIELDS {
1185 assert!(!fields.contains_key(field), "unexpected {field}");
1186 }
1187 }
1188
1189 #[derive(Clone, Copy)]
1190 struct PanicOnCommitClose;
1191
1192 impl<S> Layer<S> for PanicOnCommitClose
1193 where
1194 S: Subscriber + for<'lookup> LookupSpan<'lookup>,
1195 {
1196 fn on_close(&self, id: Id, ctx: Context<'_, S>) {
1197 if ctx
1198 .metadata(&id)
1199 .is_some_and(|metadata| metadata.name() == "transaction.commit")
1200 {
1201 panic!("injected transaction commit close panic");
1202 }
1203 }
1204 }
1205
1206 fn panic_on_commit_close_dispatch() -> tracing::Dispatch {
1207 tracing::Dispatch::new(tracing_subscriber::registry().with(PanicOnCommitClose))
1208 }
1209
1210 #[derive(Default)]
1211 struct ReaderObservations {
1212 schema_calls: Cell<usize>,
1213 next_calls: Cell<usize>,
1214 next_before_schema: Cell<bool>,
1215 previous_batch_alive: Cell<bool>,
1216 }
1217
1218 struct InstrumentedReader {
1219 schema: SchemaRef,
1220 batches: std::vec::IntoIter<Result<RecordBatch, ArrowError>>,
1221 observations: Rc<ReaderObservations>,
1222 previous_array: Option<Weak<dyn arrow::array::Array>>,
1223 }
1224
1225 impl InstrumentedReader {
1226 fn new(
1227 schema: SchemaRef,
1228 batches: Vec<Result<RecordBatch, ArrowError>>,
1229 ) -> (Self, Rc<ReaderObservations>) {
1230 let observations = Rc::new(ReaderObservations::default());
1231 (
1232 Self {
1233 schema,
1234 batches: batches.into_iter(),
1235 observations: Rc::clone(&observations),
1236 previous_array: None,
1237 },
1238 observations,
1239 )
1240 }
1241
1242 fn one(batch: RecordBatch) -> Self {
1243 let (reader, _) = Self::new(batch.schema(), vec![Ok(batch)]);
1244 reader
1245 }
1246 }
1247
1248 impl Iterator for InstrumentedReader {
1249 type Item = Result<RecordBatch, ArrowError>;
1250
1251 fn next(&mut self) -> Option<Self::Item> {
1252 self.observations
1253 .next_calls
1254 .set(self.observations.next_calls.get() + 1);
1255 if self.observations.schema_calls.get() == 0 {
1256 self.observations.next_before_schema.set(true);
1257 }
1258 if self
1259 .previous_array
1260 .as_ref()
1261 .is_some_and(|array| array.strong_count() != 0)
1262 {
1263 self.observations.previous_batch_alive.set(true);
1264 }
1265
1266 let next = self.batches.next();
1267 if let Some(Ok(batch)) = &next {
1268 self.previous_array = Some(Arc::downgrade(batch.column(0)));
1269 }
1270 next
1271 }
1272 }
1273
1274 impl RecordBatchReader for InstrumentedReader {
1275 fn schema(&self) -> arrow::datatypes::SchemaRef {
1276 self.observations
1277 .schema_calls
1278 .set(self.observations.schema_calls.get() + 1);
1279 Arc::clone(&self.schema)
1280 }
1281 }
1282
1283 fn input_batch(values: Vec<i64>) -> Result<RecordBatch, ArrowError> {
1284 let schema = Arc::new(Schema::new(vec![Field::new(
1285 "value",
1286 DataType::Int64,
1287 false,
1288 )]));
1289 RecordBatch::try_new(schema, vec![Arc::new(Int64Array::from(values))])
1290 }
1291
1292 fn time_series_batch(
1293 timestamps: Vec<i64>,
1294 symbols: Vec<&str>,
1295 prices: Vec<f64>,
1296 ) -> Result<RecordBatch, ArrowError> {
1297 let schema = Arc::new(Schema::new(vec![
1298 Field::new(
1299 "ts",
1300 DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
1301 false,
1302 ),
1303 Field::new("symbol", DataType::Utf8, false),
1304 Field::new("price", DataType::Float64, false),
1305 ]));
1306 RecordBatch::try_new(
1307 schema,
1308 vec![
1309 Arc::new(TimestampMillisecondArray::from(timestamps)),
1310 Arc::new(StringArray::from(symbols)),
1311 Arc::new(Float64Array::from(prices)),
1312 ],
1313 )
1314 }
1315
1316 fn binary_heavy_time_series_batch() -> Result<RecordBatch, ArrowError> {
1317 let payloads = (0_u64..3)
1318 .map(|row| {
1319 let mut payload = vec![0; 64 * 1_024];
1320 payload[..8].copy_from_slice(&row.to_le_bytes());
1321 payload
1322 })
1323 .collect::<Vec<_>>();
1324 let schema = Arc::new(Schema::new(vec![
1325 Field::new(
1326 "ts",
1327 DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
1328 false,
1329 ),
1330 Field::new("symbol", DataType::Utf8, false),
1331 Field::new("payload", DataType::Binary, false),
1332 ]));
1333 RecordBatch::try_new(
1334 schema,
1335 vec![
1336 Arc::new(TimestampMillisecondArray::from(vec![0, 60_000, 120_000])),
1337 Arc::new(StringArray::from(vec!["A"; 3])),
1338 Arc::new(BinaryArray::from_iter_values(
1339 payloads.iter().map(Vec::as_slice),
1340 )),
1341 ],
1342 )
1343 }
1344
1345 fn binary_heavy_table_meta() -> TableMeta {
1346 TableMeta::new_time_series(IndexSpec {
1347 column: "ts".to_string(),
1348 entity_columns: vec!["symbol".to_string()],
1349 kind: IndexKind::Timestamp {
1350 index_granularity: TimeIndexGranularity::Minutes(1),
1351 timezone: None,
1352 },
1353 })
1354 }
1355
1356 fn timestamp_only_index() -> IndexSpec {
1357 IndexSpec {
1358 column: "ts".to_string(),
1359 entity_columns: Vec::new(),
1360 kind: IndexKind::Timestamp {
1361 index_granularity: TimeIndexGranularity::Minutes(1),
1362 timezone: None,
1363 },
1364 }
1365 }
1366
1367 fn timestamp_only_meta() -> TableMeta {
1368 TableMeta::new_time_series_with_schema(
1369 timestamp_only_index(),
1370 LogicalSchema::new(vec![LogicalField {
1371 name: "ts".to_string(),
1372 data_type: LogicalDataType::Timestamp {
1373 unit: LogicalTimestampUnit::Millis,
1374 timezone: None,
1375 },
1376 nullable: false,
1377 }])
1378 .expect("valid timestamp-only schema"),
1379 )
1380 }
1381
1382 fn timestamp_only_batch(row_count: usize) -> Result<RecordBatch, ArrowError> {
1383 timestamp_only_batch_starting_at_minute_offset(0, row_count)
1384 }
1385
1386 fn timestamp_only_batch_starting_at_minute_offset(
1387 start_minute_offset: usize,
1388 row_count: usize,
1389 ) -> Result<RecordBatch, ArrowError> {
1390 timestamp_only_batch_with_millis(
1391 (start_minute_offset..start_minute_offset + row_count)
1392 .map(|minute_offset| minute_offset as i64 * 60_000),
1393 )
1394 }
1395
1396 fn timestamp_only_batch_with_millis(
1397 values: impl IntoIterator<Item = i64>,
1398 ) -> Result<RecordBatch, ArrowError> {
1399 let schema = Arc::new(Schema::new(vec![Field::new(
1400 "ts",
1401 DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
1402 false,
1403 )]));
1404 RecordBatch::try_new(
1405 schema,
1406 vec![Arc::new(TimestampMillisecondArray::from_iter_values(
1407 values,
1408 ))],
1409 )
1410 }
1411
1412 fn widening_table_meta() -> TableMeta {
1413 TableMeta::new_time_series_with_schema(
1414 IndexSpec {
1415 column: "seq".to_string(),
1416 entity_columns: vec!["device_id".to_string()],
1417 kind: IndexKind::UInt64 {
1418 index_granularity: NonZeroU64::new(u64::from(u32::MAX) + 1).unwrap(),
1419 },
1420 },
1421 LogicalSchema::new(vec![
1422 LogicalField {
1423 name: "seq".to_string(),
1424 data_type: LogicalDataType::UInt64,
1425 nullable: false,
1426 },
1427 LogicalField {
1428 name: "device_id".to_string(),
1429 data_type: LogicalDataType::Int32,
1430 nullable: false,
1431 },
1432 LogicalField {
1433 name: "reading".to_string(),
1434 data_type: LogicalDataType::Float64,
1435 nullable: true,
1436 },
1437 LogicalField {
1438 name: "label".to_string(),
1439 data_type: LogicalDataType::Utf8,
1440 nullable: false,
1441 },
1442 ])
1443 .expect("valid widening target schema"),
1444 )
1445 }
1446
1447 fn widening_batch(
1448 seq: Vec<u32>,
1449 device_ids: Vec<i16>,
1450 readings: Vec<Option<f32>>,
1451 labels: Vec<&str>,
1452 ) -> Result<RecordBatch, ArrowError> {
1453 let schema = Arc::new(Schema::new_with_metadata(
1454 vec![
1455 Field::new("label", DataType::Utf8, false),
1456 Field::new("reading", DataType::Float32, true).with_metadata(HashMap::from([(
1457 "source_metadata".to_string(),
1458 "ignored".to_string(),
1459 )])),
1460 Field::new("seq", DataType::UInt32, false),
1461 Field::new("device_id", DataType::Int16, false),
1462 ],
1463 HashMap::from([("source_schema_metadata".to_string(), "ignored".to_string())]),
1464 ));
1465 RecordBatch::try_new(
1466 schema,
1467 vec![
1468 Arc::new(StringArray::from(labels)),
1469 Arc::new(Float32Array::from(readings)),
1470 Arc::new(UInt32Array::from(seq)),
1471 Arc::new(Int16Array::from(device_ids)),
1472 ],
1473 )
1474 }
1475
1476 fn declared_schema_test_meta(value_type: LogicalDataType, nullable: bool) -> TableMeta {
1477 TableMeta::new_time_series_with_schema(
1478 IndexSpec {
1479 column: "ts".to_string(),
1480 entity_columns: Vec::new(),
1481 kind: IndexKind::Timestamp {
1482 index_granularity: TimeIndexGranularity::Minutes(1),
1483 timezone: None,
1484 },
1485 },
1486 LogicalSchema::new(vec![
1487 LogicalField {
1488 name: "ts".to_string(),
1489 data_type: LogicalDataType::Timestamp {
1490 unit: LogicalTimestampUnit::Millis,
1491 timezone: None,
1492 },
1493 nullable: false,
1494 },
1495 LogicalField {
1496 name: "value".to_string(),
1497 data_type: value_type,
1498 nullable,
1499 },
1500 ])
1501 .expect("valid compatibility target schema"),
1502 )
1503 }
1504
1505 async fn assert_declared_schema_rejected_before_reading(
1506 meta: TableMeta,
1507 schema: SchemaRef,
1508 batches: Vec<Result<RecordBatch, ArrowError>>,
1509 expected_column: &str,
1510 ) -> TestResult {
1511 let temp = TempDir::new()?;
1512 let mut table = TimeSeriesTable::create(TableLocation::local(temp.path()), meta).await?;
1513 let state_before = table.state().clone();
1514 let (reader, observations) = InstrumentedReader::new(schema, batches);
1515
1516 let error = table
1517 .append(reader)
1518 .await
1519 .expect_err("incompatible declared schema must fail");
1520
1521 assert!(matches!(
1522 error,
1523 TableError::Append {
1524 source: AppendError::SchemaValidation { .. }
1525 }
1526 ));
1527 assert!(error.to_string().contains(expected_column));
1528 assert!(observations.schema_calls.get() > 0);
1529 assert_eq!(observations.next_calls.get(), 0);
1530 assert_eq!(table.state(), &state_before);
1531 assert_eq!(table.log.load_current_version().await?, 1);
1532 assert!(data_files(temp.path())?.is_empty());
1533 assert!(coverage_files(temp.path())?.is_empty());
1534 assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
1535 Ok(())
1536 }
1537
1538 #[tokio::test]
1539 async fn lossy_schema_is_rejected_before_reading_or_creating_artifacts() -> TestResult {
1540 let temp = TempDir::new()?;
1541 let mut meta = timestamp_only_meta();
1542 meta.logical_schema = None;
1543 let mut table = TimeSeriesTable::create(TableLocation::local(temp.path()), meta).await?;
1544 let schema = Arc::new(Schema::new(vec![
1545 Field::new(
1546 "ts",
1547 DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
1548 false,
1549 ),
1550 Field::new("value", DataType::Int8, true),
1551 ]));
1552 let (reader, observations) = InstrumentedReader::new(schema, Vec::new());
1553
1554 assert!(matches!(
1555 table.append(reader).await,
1556 Err(TableError::Append {
1557 source: AppendError::ArrowToLogicalSchema { .. }
1558 })
1559 ));
1560 assert_eq!(observations.next_calls.get(), 0);
1561 assert_eq!(table.state().version, 1);
1562 assert!(!temp.path().join("data").exists());
1563 Ok(())
1564 }
1565
1566 #[tokio::test]
1567 async fn append_rejects_unsupported_writer_features_before_input_or_artifacts() -> TestResult {
1568 let temp = TempDir::new()?;
1569 let mut table =
1570 TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
1571 .await?;
1572 table
1573 .state
1574 .table_meta
1575 .required_writer_features
1576 .insert("future_writer".to_string());
1577 let state_before = table.state().clone();
1578 let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
1579 let (reader, observations) = InstrumentedReader::new(batch.schema(), vec![Ok(batch)]);
1580
1581 let error = table
1582 .append(reader)
1583 .await
1584 .expect_err("unsupported writer feature must reject append");
1585
1586 assert!(matches!(
1587 error,
1588 TableError::Append {
1589 source: AppendError::Protocol {
1590 source: TableProtocolError::UnsupportedWriterFeatures { features },
1591 ..
1592 }
1593 } if features == ["future_writer"]
1594 ));
1595 assert_eq!(observations.schema_calls.get(), 0);
1596 assert_eq!(observations.next_calls.get(), 0);
1597 assert_eq!(table.state(), &state_before);
1598 assert_eq!(table.current_version().await?, 1);
1599 assert!(data_files(temp.path())?.is_empty());
1600 assert!(coverage_files(temp.path())?.is_empty());
1601 assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
1602 Ok(())
1603 }
1604
1605 #[tokio::test]
1606 async fn append_widens_reordered_index_entity_and_data_fields_incrementally() -> TestResult {
1607 let temp = TempDir::new()?;
1608 let location = TableLocation::local(temp.path());
1609 let mut table = TimeSeriesTable::create(location.clone(), widening_table_meta()).await?;
1610 let first = widening_batch(vec![0], vec![i16::MIN], vec![Some(f32::MIN)], vec!["first"])?;
1611 let second = widening_batch(vec![u32::MAX], vec![i16::MAX], vec![None], vec!["second"])?;
1612 let (reader, observations) =
1613 InstrumentedReader::new(first.schema(), vec![Ok(first), Ok(second)]);
1614
1615 assert_eq!(table.append(reader).await?.committed_version, 2);
1616 assert!(!observations.previous_batch_alive.get());
1617 assert_eq!(table.state().segments.len(), 1);
1618 let segment = table
1619 .state()
1620 .segments
1621 .values()
1622 .next()
1623 .expect("committed widened segment");
1624 assert_eq!(segment.row_count, 2);
1625 assert_eq!(segment.index_min, IndexValue::UInt64(0));
1626 assert_eq!(segment.index_max, IndexValue::UInt64(u64::from(u32::MAX)));
1627 assert_eq!(segment.entity_layout, SegmentEntityLayout::Mixed);
1628
1629 let expected_schema = table.state().table_meta.arrow_schema_ref()?;
1630 let parquet_reader =
1631 ParquetRecordBatchReaderBuilder::try_new(File::open(temp.path().join(&segment.path))?)?
1632 .build()?;
1633 assert_eq!(parquet_reader.schema(), expected_schema);
1634 for batch in parquet_reader {
1635 assert_eq!(batch?.schema(), expected_schema);
1636 }
1637
1638 let state_before_overlap = table.state().clone();
1639 let data_before_overlap = data_files(temp.path())?;
1640 let coverage_before_overlap = coverage_files(temp.path())?;
1641 let overlap = table
1642 .append(widening_batch(
1643 vec![1],
1644 vec![i16::MIN],
1645 vec![Some(1.0)],
1646 vec!["overlap"],
1647 )?)
1648 .await
1649 .expect_err("widened entity and index coverage must still reject overlap");
1650 assert!(matches!(
1651 overlap,
1652 TableError::Append {
1653 source: AppendError::PersistedIndexIntervalOverlap { .. }
1654 }
1655 ));
1656 assert_eq!(table.state(), &state_before_overlap);
1657 assert_eq!(data_files(temp.path())?, data_before_overlap);
1658 assert_eq!(coverage_files(temp.path())?, coverage_before_overlap);
1659
1660 let reopened = TimeSeriesTable::open(location).await?;
1661 let mut stream = reopened.scan_range(0_u64, u64::from(u32::MAX) + 1).await?;
1662 let mut rows = Vec::new();
1663 while let Some(batch) = stream.next().await {
1664 let batch = batch?;
1665 assert_eq!(batch.schema(), expected_schema);
1666 let seq = batch
1667 .column(0)
1668 .as_any()
1669 .downcast_ref::<UInt64Array>()
1670 .expect("seq as UInt64");
1671 let device_id = batch
1672 .column(1)
1673 .as_any()
1674 .downcast_ref::<Int32Array>()
1675 .expect("device_id as Int32");
1676 let reading = batch
1677 .column(2)
1678 .as_any()
1679 .downcast_ref::<Float64Array>()
1680 .expect("reading as Float64");
1681 let label = batch
1682 .column(3)
1683 .as_any()
1684 .downcast_ref::<StringArray>()
1685 .expect("label as Utf8");
1686 for row in 0..batch.num_rows() {
1687 rows.push((
1688 seq.value(row),
1689 device_id.value(row),
1690 (!reading.is_null(row)).then(|| reading.value(row)),
1691 label.value(row).to_string(),
1692 ));
1693 }
1694 }
1695 rows.sort_by_key(|row| row.0);
1696 assert_eq!(
1697 rows,
1698 vec![
1699 (
1700 0,
1701 i32::from(i16::MIN),
1702 Some(f64::from(f32::MIN)),
1703 "first".to_string(),
1704 ),
1705 (
1706 u64::from(u32::MAX),
1707 i32::from(i16::MAX),
1708 None,
1709 "second".to_string(),
1710 ),
1711 ]
1712 );
1713 Ok(())
1714 }
1715
1716 #[tokio::test]
1717 async fn append_rejects_incompatible_declared_schemas_without_artifacts() -> TestResult {
1718 let timestamp = || {
1719 Field::new(
1720 "ts",
1721 DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
1722 false,
1723 )
1724 };
1725 let cases = vec![
1726 (
1727 declared_schema_test_meta(LogicalDataType::Int32, false),
1728 Arc::new(Schema::new(vec![
1729 timestamp(),
1730 Field::new("value", DataType::Int64, false),
1731 ])),
1732 "value",
1733 ),
1734 (
1735 declared_schema_test_meta(LogicalDataType::Float64, false),
1736 Arc::new(Schema::new(vec![
1737 timestamp(),
1738 Field::new("value", DataType::Int32, false),
1739 ])),
1740 "value",
1741 ),
1742 (
1743 declared_schema_test_meta(LogicalDataType::Int32, false),
1744 Arc::new(Schema::new(vec![
1745 Field::new(
1746 "ts",
1747 DataType::Timestamp(ArrowTimeUnit::Microsecond, None),
1748 false,
1749 ),
1750 Field::new("value", DataType::Int32, false),
1751 ])),
1752 "ts",
1753 ),
1754 (
1755 declared_schema_test_meta(LogicalDataType::Int32, false),
1756 Arc::new(Schema::new(vec![
1757 timestamp(),
1758 Field::new("value", DataType::Int32, true),
1759 ])),
1760 "value",
1761 ),
1762 (
1763 declared_schema_test_meta(LogicalDataType::Int32, false),
1764 Arc::new(Schema::new(vec![timestamp()])),
1765 "value",
1766 ),
1767 (
1768 declared_schema_test_meta(LogicalDataType::Int32, false),
1769 Arc::new(Schema::new(vec![
1770 timestamp(),
1771 Field::new("value", DataType::Int32, false),
1772 Field::new("extra", DataType::Int32, false),
1773 ])),
1774 "extra",
1775 ),
1776 (
1777 declared_schema_test_meta(LogicalDataType::Int32, false),
1778 Arc::new(Schema::new(vec![
1779 timestamp(),
1780 Field::new("value", DataType::Int32, false),
1781 Field::new("value", DataType::Int32, false),
1782 ])),
1783 "value",
1784 ),
1785 ];
1786
1787 for (meta, schema, expected_column) in cases {
1788 assert_declared_schema_rejected_before_reading(
1789 meta,
1790 schema,
1791 Vec::new(),
1792 expected_column,
1793 )
1794 .await?;
1795 }
1796 Ok(())
1797 }
1798
1799 #[tokio::test]
1800 async fn append_rejects_positive_signed_index_without_inspecting_values() -> TestResult {
1801 let meta = TableMeta::new_time_series_with_schema(
1802 IndexSpec {
1803 column: "seq".to_string(),
1804 entity_columns: Vec::new(),
1805 kind: IndexKind::UInt64 {
1806 index_granularity: NonZeroU64::new(1).unwrap(),
1807 },
1808 },
1809 LogicalSchema::new(vec![LogicalField {
1810 name: "seq".to_string(),
1811 data_type: LogicalDataType::UInt64,
1812 nullable: false,
1813 }])?,
1814 );
1815 let schema = Arc::new(Schema::new(vec![Field::new("seq", DataType::Int64, false)]));
1816 let batch = RecordBatch::try_new(
1817 Arc::clone(&schema),
1818 vec![Arc::new(Int64Array::from(vec![1]))],
1819 )?;
1820
1821 assert_declared_schema_rejected_before_reading(meta, schema, vec![Ok(batch)], "seq").await
1822 }
1823
1824 #[tokio::test]
1825 async fn widened_batch_source_failure_rolls_back_partial_output() -> TestResult {
1826 let temp = TempDir::new()?;
1827 let mut table =
1828 TimeSeriesTable::create(TableLocation::local(temp.path()), widening_table_meta())
1829 .await?;
1830 let first = widening_batch(vec![0], vec![1], vec![Some(1.0)], vec!["first"])?;
1831 let (reader, observations) = InstrumentedReader::new(
1832 first.schema(),
1833 vec![
1834 Ok(first),
1835 Err(ArrowError::ComputeError(
1836 "injected widened source failure".to_string(),
1837 )),
1838 ],
1839 );
1840
1841 let error = table
1842 .append(reader)
1843 .await
1844 .expect_err("source failure must abort widened append");
1845
1846 assert!(matches!(
1847 error,
1848 TableError::Append {
1849 source: AppendError::ArrowInput { .. }
1850 }
1851 ));
1852 assert!(!observations.previous_batch_alive.get());
1853 assert_eq!(table.state().version, 1);
1854 assert!(table.state().segments.is_empty());
1855 assert!(data_files(temp.path())?.is_empty());
1856 assert!(coverage_files(temp.path())?.is_empty());
1857 assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
1858 Ok(())
1859 }
1860
1861 fn assert_batches<R>(reader: R, expected: &[RecordBatch]) -> TestResult
1862 where
1863 R: RecordBatchReader,
1864 {
1865 assert_eq!(reader.schema(), expected[0].schema());
1866 let actual = reader.collect::<Result<Vec<_>, _>>()?;
1867 assert_eq!(actual, expected);
1868 Ok(())
1869 }
1870
1871 #[test]
1872 fn into_record_batch_reader_accepts_all_required_input_forms() -> TestResult {
1873 let first = input_batch(vec![1, 2])?;
1874 let second = input_batch(vec![3])?;
1875
1876 assert_batches(
1877 first.clone().into_record_batch_reader()?,
1878 std::slice::from_ref(&first),
1879 )?;
1880 assert_batches(
1881 vec![first.clone(), second.clone()].into_record_batch_reader()?,
1882 &[first.clone(), second.clone()],
1883 )?;
1884
1885 let iterator = RecordBatchIterator::new(vec![Ok(first.clone())], first.schema());
1886 assert_batches(
1887 iterator.into_record_batch_reader()?,
1888 std::slice::from_ref(&first),
1889 )?;
1890 assert_batches(
1891 InstrumentedReader::one(first.clone()).into_record_batch_reader()?,
1892 std::slice::from_ref(&first),
1893 )?;
1894
1895 let boxed: Box<dyn RecordBatchReader> = Box::new(RecordBatchIterator::new(
1896 vec![Ok(first.clone())],
1897 first.schema(),
1898 ));
1899 assert_batches(
1900 boxed.into_record_batch_reader()?,
1901 std::slice::from_ref(&first),
1902 )?;
1903
1904 let sendable: Box<dyn RecordBatchReader + Send> = Box::new(RecordBatchIterator::new(
1905 vec![Ok(second.clone())],
1906 second.schema(),
1907 ));
1908 assert_batches(
1909 sendable.into_record_batch_reader()?,
1910 std::slice::from_ref(&second),
1911 )?;
1912 Ok(())
1913 }
1914
1915 #[test]
1916 fn into_record_batch_reader_rejects_schema_less_vec() {
1917 let result = Vec::<RecordBatch>::new().into_record_batch_reader();
1918 assert!(matches!(result, Err(AppendError::EmptyInput)));
1919 }
1920
1921 #[test]
1922 fn append_writer_defaults_are_compressed_and_bounded() {
1923 let settings = AppendWriterSettings::default();
1924 let properties = settings.writer_properties();
1925
1926 assert_eq!(settings.compression, ParquetCompression::Zstd);
1927 assert_eq!(
1928 properties.max_row_group_row_count(),
1929 Some(DEFAULT_MAX_ROWS_PER_ROW_GROUP)
1930 );
1931 assert_eq!(
1932 properties.max_row_group_bytes(),
1933 Some(DEFAULT_MAX_BYTES_PER_ROW_GROUP)
1934 );
1935 }
1936
1937 #[test]
1938 fn append_request_nested_writer_settings_inherit_and_override() -> TestResult {
1939 let inner = AppendRequest::new(timestamp_only_batch(1)?)
1940 .compression(ParquetCompression::Snappy)
1941 .max_rows_per_row_group(3)
1942 .max_bytes_per_row_group(30);
1943 let inherited = AppendWriterSettings::from_source(&AppendRequest::new(inner));
1944 assert_eq!(
1945 inherited,
1946 AppendWriterSettings {
1947 compression: ParquetCompression::Snappy,
1948 max_rows_per_row_group: 3,
1949 max_bytes_per_row_group: 30,
1950 }
1951 );
1952
1953 let partially_overridden = AppendWriterSettings::from_source(
1954 &AppendRequest::new(
1955 AppendRequest::new(timestamp_only_batch(1)?)
1956 .compression(ParquetCompression::Snappy)
1957 .max_rows_per_row_group(3)
1958 .max_bytes_per_row_group(30),
1959 )
1960 .max_bytes_per_row_group(20),
1961 );
1962 assert_eq!(
1963 partially_overridden,
1964 AppendWriterSettings {
1965 compression: ParquetCompression::Snappy,
1966 max_rows_per_row_group: 3,
1967 max_bytes_per_row_group: 20,
1968 }
1969 );
1970
1971 let overridden = AppendWriterSettings::from_source(
1972 &AppendRequest::new(
1973 AppendRequest::new(timestamp_only_batch(1)?)
1974 .compression(ParquetCompression::Snappy)
1975 .max_rows_per_row_group(3)
1976 .max_bytes_per_row_group(30),
1977 )
1978 .compression(ParquetCompression::Zstd)
1979 .max_rows_per_row_group(2)
1980 .max_bytes_per_row_group(20),
1981 );
1982 assert_eq!(
1983 overridden,
1984 AppendWriterSettings {
1985 compression: ParquetCompression::Zstd,
1986 max_rows_per_row_group: 2,
1987 max_bytes_per_row_group: 20,
1988 }
1989 );
1990 Ok(())
1991 }
1992
1993 #[tokio::test]
1994 async fn append_default_and_explicit_compression_reach_parquet_metadata() -> TestResult {
1995 for (compression, expected) in [
1996 (None, Compression::ZSTD(Default::default())),
1997 (
1998 Some(ParquetCompression::Uncompressed),
1999 Compression::UNCOMPRESSED,
2000 ),
2001 (Some(ParquetCompression::Snappy), Compression::SNAPPY),
2002 (
2003 Some(ParquetCompression::Zstd),
2004 Compression::ZSTD(Default::default()),
2005 ),
2006 ] {
2007 let temp = TempDir::new()?;
2008 let mut table =
2009 TimeSeriesTable::create(TableLocation::local(temp.path()), timestamp_only_meta())
2010 .await?;
2011 let batch = timestamp_only_batch(1)?;
2012
2013 match compression {
2014 Some(compression) => {
2015 table
2016 .append(AppendRequest::new(batch).compression(compression))
2017 .await?;
2018 }
2019 None => {
2020 table.append(batch).await?;
2021 }
2022 }
2023
2024 let segment_path = &table
2025 .state()
2026 .segments
2027 .values()
2028 .next()
2029 .ok_or("missing appended segment")?
2030 .path;
2031 let builder = ParquetRecordBatchReaderBuilder::try_new(File::open(
2032 temp.path().join(segment_path),
2033 )?)?;
2034 assert!(builder.metadata().row_groups().iter().all(|row_group| {
2035 row_group
2036 .columns()
2037 .iter()
2038 .all(|column| column.compression() == expected)
2039 }));
2040 }
2041 Ok(())
2042 }
2043
2044 #[tokio::test]
2045 async fn append_byte_limit_splits_row_groups_and_preserves_values() -> TestResult {
2046 let temp = TempDir::new()?;
2047 let location = TableLocation::local(temp.path());
2048 let mut table = TimeSeriesTable::create(
2049 location.clone(),
2050 TableMeta::new_time_series(timestamp_only_index()),
2051 )
2052 .await?;
2053 let schema = Arc::new(Schema::new(vec![
2054 Field::new(
2055 "ts",
2056 DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
2057 false,
2058 ),
2059 Field::new("payload", DataType::Binary, false),
2060 ]));
2061 let payloads = [vec![1_u8; 256], vec![2_u8; 256], vec![3_u8; 256]];
2062 let batches = payloads
2063 .iter()
2064 .enumerate()
2065 .map(|(minute, payload)| {
2066 RecordBatch::try_new(
2067 Arc::clone(&schema),
2068 vec![
2069 Arc::new(TimestampMillisecondArray::from(vec![
2070 minute as i64 * 60_000,
2071 ])),
2072 Arc::new(BinaryArray::from_iter_values([payload.as_slice()])),
2073 ],
2074 )
2075 })
2076 .collect::<Result<Vec<_>, _>>()?;
2077
2078 assert_eq!(
2079 table
2080 .append(
2081 AppendRequest::new(batches)
2082 .compression(ParquetCompression::Uncompressed)
2083 .max_rows_per_row_group(100)
2084 .max_bytes_per_row_group(1),
2085 )
2086 .await?
2087 .committed_version,
2088 2
2089 );
2090
2091 let segment_path = &table
2092 .state()
2093 .segments
2094 .values()
2095 .next()
2096 .ok_or("missing appended segment")?
2097 .path;
2098 let builder =
2099 ParquetRecordBatchReaderBuilder::try_new(File::open(temp.path().join(segment_path))?)?;
2100 assert_eq!(
2101 builder
2102 .metadata()
2103 .row_groups()
2104 .iter()
2105 .map(|row_group| row_group.num_rows())
2106 .collect::<Vec<_>>(),
2107 [1, 1, 1]
2108 );
2109 let batches = builder.build()?.collect::<Result<Vec<_>, _>>()?;
2110 let mut actual_payloads = Vec::new();
2111 for batch in batches {
2112 let values = batch
2113 .column_by_name("payload")
2114 .and_then(|column| column.as_any().downcast_ref::<BinaryArray>())
2115 .ok_or("missing binary payload column")?;
2116 actual_payloads.extend((0..values.len()).map(|index| values.value(index).to_vec()));
2117 }
2118 assert_eq!(actual_payloads, payloads);
2119
2120 let reopened = TimeSeriesTable::open(location).await?;
2121 assert_eq!(reopened.state().version, 2);
2122 assert_eq!(reopened.state().segments.len(), 1);
2123 Ok(())
2124 }
2125
2126 #[tokio::test]
2127 async fn append_request_nested_settings_inherit_and_override() -> TestResult {
2128 for (inner_limit, outer_limit, expected_row_groups) in
2129 [(3, None, vec![3, 3, 1]), (0, Some(2), vec![2, 2, 2, 1])]
2130 {
2131 let temp = TempDir::new()?;
2132 let location = TableLocation::local(temp.path());
2133 let mut table =
2134 TimeSeriesTable::create(location.clone(), timestamp_only_meta()).await?;
2135 let request = AppendRequest::new(
2136 AppendRequest::new(vec![
2137 timestamp_only_batch(2)?,
2138 timestamp_only_batch_starting_at_minute_offset(2, 5)?,
2139 ])
2140 .max_rows_per_row_group(inner_limit),
2141 );
2142 let request = match outer_limit {
2143 Some(limit) => request.max_rows_per_row_group(limit),
2144 None => request,
2145 };
2146
2147 assert_eq!(table.append(request).await?.committed_version, 2);
2148
2149 let (segment_path, row_count) = table
2150 .state()
2151 .segments
2152 .values()
2153 .next()
2154 .map(|segment| (segment.path.clone(), segment.row_count))
2155 .ok_or("missing appended segment")?;
2156 assert_eq!(row_count, 7);
2157 let builder = ParquetRecordBatchReaderBuilder::try_new(File::open(
2158 temp.path().join(segment_path),
2159 )?)?;
2160 let row_group_rows = builder
2161 .metadata()
2162 .row_groups()
2163 .iter()
2164 .map(|row_group| row_group.num_rows())
2165 .collect::<Vec<_>>();
2166 assert_eq!(row_group_rows, expected_row_groups);
2167
2168 let reopened = TimeSeriesTable::open(location).await?;
2169 assert_eq!(reopened.state().version, 2);
2170 assert_eq!(reopened.state().segments.len(), 1);
2171 }
2172 Ok(())
2173 }
2174
2175 #[tokio::test]
2176 async fn append_request_nested_zero_is_rejected_before_inspecting_source() -> TestResult {
2177 let temp = TempDir::new()?;
2178 let mut table =
2179 TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2180 .await?;
2181 let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2182 let (reader, observations) = InstrumentedReader::new(batch.schema(), vec![Ok(batch)]);
2183 let request = AppendRequest::new(AppendRequest::new(reader).max_rows_per_row_group(0));
2184
2185 assert!(matches!(
2186 table.append(request).await,
2187 Err(TableError::Append {
2188 source: AppendError::InvalidMaxRowsPerRowGroup {
2189 max_rows_per_row_group: 0
2190 }
2191 })
2192 ));
2193 assert_eq!(observations.schema_calls.get(), 0);
2194 assert_eq!(observations.next_calls.get(), 0);
2195 assert_eq!(table.state().version, 1);
2196 assert!(table.state().segments.is_empty());
2197 assert!(data_files(temp.path())?.is_empty());
2198 assert!(coverage_files(temp.path())?.is_empty());
2199 assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
2200 Ok(())
2201 }
2202
2203 #[tokio::test]
2204 async fn append_request_zero_byte_limit_is_rejected_before_inspecting_source() -> TestResult {
2205 let temp = TempDir::new()?;
2206 let mut table =
2207 TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2208 .await?;
2209 let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2210 let (reader, observations) = InstrumentedReader::new(batch.schema(), vec![Ok(batch)]);
2211
2212 assert!(matches!(
2213 table
2214 .append(AppendRequest::new(reader).max_bytes_per_row_group(0))
2215 .await,
2216 Err(TableError::Append {
2217 source: AppendError::InvalidMaxBytesPerRowGroup {
2218 max_bytes_per_row_group: 0
2219 }
2220 })
2221 ));
2222 assert_eq!(observations.schema_calls.get(), 0);
2223 assert_eq!(observations.next_calls.get(), 0);
2224 assert_eq!(table.state().version, 1);
2225 assert!(table.state().segments.is_empty());
2226 assert!(data_files(temp.path())?.is_empty());
2227 assert!(coverage_files(temp.path())?.is_empty());
2228 assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
2229 Ok(())
2230 }
2231
2232 #[tokio::test]
2233 async fn append_request_without_overrides_matches_direct_append() -> TestResult {
2234 let direct_root = TempDir::new()?;
2235 let request_root = TempDir::new()?;
2236 let mut direct = TimeSeriesTable::create(
2237 TableLocation::local(direct_root.path()),
2238 make_basic_table_meta(),
2239 )
2240 .await?;
2241 let mut requested = TimeSeriesTable::create(
2242 TableLocation::local(request_root.path()),
2243 make_basic_table_meta(),
2244 )
2245 .await?;
2246 let source = vec![
2247 time_series_batch(vec![0], vec!["A"], vec![1.0])?,
2248 time_series_batch(vec![60_000], vec!["A"], vec![2.0])?,
2249 ];
2250
2251 assert_eq!(direct.append(source.clone()).await?.committed_version, 2);
2252 assert_eq!(
2253 requested
2254 .append(AppendRequest::new(source))
2255 .await?
2256 .committed_version,
2257 2
2258 );
2259
2260 let direct_path = &direct
2261 .state()
2262 .segments
2263 .values()
2264 .next()
2265 .ok_or("missing direct segment")?
2266 .path;
2267 let request_path = &requested
2268 .state()
2269 .segments
2270 .values()
2271 .next()
2272 .ok_or("missing requested segment")?
2273 .path;
2274 assert_eq!(
2275 std::fs::read(direct_root.path().join(direct_path))?,
2276 std::fs::read(request_root.path().join(request_path))?
2277 );
2278 Ok(())
2279 }
2280
2281 #[tokio::test]
2282 async fn append_request_rejects_zero_before_inspecting_source() -> TestResult {
2283 let temp = TempDir::new()?;
2284 let mut table =
2285 TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2286 .await?;
2287 let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2288 let (reader, observations) = InstrumentedReader::new(batch.schema(), vec![Ok(batch)]);
2289
2290 assert!(matches!(
2291 table
2292 .append(AppendRequest::new(reader).max_rows_per_row_group(0))
2293 .await,
2294 Err(TableError::Append {
2295 source: AppendError::InvalidMaxRowsPerRowGroup {
2296 max_rows_per_row_group: 0
2297 }
2298 })
2299 ));
2300 assert_eq!(observations.schema_calls.get(), 0);
2301 assert_eq!(observations.next_calls.get(), 0);
2302 assert_eq!(table.state().version, 1);
2303 assert!(table.state().segments.is_empty());
2304 assert!(data_files(temp.path())?.is_empty());
2305 assert!(coverage_files(temp.path())?.is_empty());
2306 assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
2307 Ok(())
2308 }
2309
2310 #[tokio::test]
2311 async fn append_rejects_version_overflow_before_inspecting_source() -> TestResult {
2312 let temp = TempDir::new()?;
2313 let mut table =
2314 TimeSeriesTable::create(TableLocation::local(temp.path()), timestamp_only_meta())
2315 .await?;
2316 table.state.version = u64::MAX;
2317 let batch = timestamp_only_batch(1)?;
2318 let (reader, observations) = InstrumentedReader::new(batch.schema(), vec![Ok(batch)]);
2319
2320 let error = table
2321 .append(reader)
2322 .await
2323 .expect_err("version overflow must fail");
2324
2325 assert!(matches!(
2326 error,
2327 TableError::Append {
2328 source: AppendError::Commit {
2329 source: CommitError::VersionOverflow {
2330 current_version: u64::MAX,
2331 ..
2332 }
2333 }
2334 }
2335 ));
2336 assert_eq!(observations.schema_calls.get(), 0);
2337 assert_eq!(observations.next_calls.get(), 0);
2338 assert_eq!(table.state().version, u64::MAX);
2339 assert!(table.state().segments.is_empty());
2340 assert!(!temp.path().join("data").exists());
2341 assert!(!temp.path().join("_coverage").exists());
2342 assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
2343 assert_eq!(table.log.load_current_version().await?, 1);
2344 Ok(())
2345 }
2346
2347 #[tokio::test]
2348 async fn append_report_matches_committed_segment_and_parquet_footer() -> TestResult {
2349 let temp = TempDir::new()?;
2350 let location = TableLocation::local(temp.path());
2351 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
2352
2353 let first = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2354 let second = time_series_batch(vec![60_000], vec!["A"], vec![2.0])?;
2355 let (reader, observations) =
2356 InstrumentedReader::new(first.schema(), vec![Ok(first), Ok(second)]);
2357 let report = table
2358 .append(
2359 AppendRequest::new(reader)
2360 .compression(ParquetCompression::Snappy)
2361 .max_rows_per_row_group(1)
2362 .max_bytes_per_row_group(1_024),
2363 )
2364 .await?;
2365
2366 assert_eq!(observations.next_calls.get(), 3);
2367 assert!(!observations.next_before_schema.get());
2368 assert_eq!(report.starting_version, 1);
2369 assert_eq!(report.committed_version, 2);
2370 assert_eq!(report.row_count, 2);
2371 assert_eq!(report.row_group_count, 2);
2372 assert_eq!(report.compression, ParquetCompression::Snappy);
2373 assert_eq!(report.max_rows_per_row_group, 1);
2374 assert_eq!(report.max_bytes_per_row_group, 1_024);
2375 assert_eq!(table.state().segments.len(), 1);
2376 let segment = table
2377 .state()
2378 .segments
2379 .get(&report.segment_path)
2380 .ok_or("reported segment is not committed")?;
2381 assert_eq!(report.row_count, segment.row_count);
2382 assert_eq!(
2383 report.file_size_bytes,
2384 segment.file_size.unwrap_or_default()
2385 );
2386
2387 let segment_path = temp.path().join(&report.segment_path);
2388 assert_eq!(
2389 report.file_size_bytes,
2390 std::fs::metadata(&segment_path)?.len()
2391 );
2392 let parquet = ParquetRecordBatchReaderBuilder::try_new(File::open(segment_path)?)?;
2393 assert_eq!(report.row_group_count, parquet.metadata().num_row_groups());
2394 assert!(parquet.metadata().row_groups().iter().all(|row_group| {
2395 row_group
2396 .columns()
2397 .iter()
2398 .all(|column| column.compression() == Compression::SNAPPY)
2399 }));
2400
2401 let third = time_series_batch(vec![120_000], vec!["A"], vec![3.0])?;
2402 let direct_report = table.append(third).await?;
2403 assert_eq!(direct_report.starting_version, 2);
2404 assert_eq!(direct_report.committed_version, 3);
2405 assert_eq!(direct_report.compression, ParquetCompression::Zstd);
2406 assert_eq!(
2407 direct_report.max_rows_per_row_group,
2408 DEFAULT_MAX_ROWS_PER_ROW_GROUP
2409 );
2410 assert_eq!(
2411 direct_report.max_bytes_per_row_group,
2412 DEFAULT_MAX_BYTES_PER_ROW_GROUP
2413 );
2414 assert_eq!(table.state().segments.len(), 2);
2415
2416 let reopened = TimeSeriesTable::open(location).await?;
2417 assert_eq!(reopened.state().version, 3);
2418 assert_eq!(
2419 reopened
2420 .state()
2421 .segments
2422 .values()
2423 .map(|segment| segment.row_count)
2424 .sum::<u64>(),
2425 3
2426 );
2427 Ok(())
2428 }
2429
2430 #[tokio::test]
2431 async fn successful_append_reports_phases_and_footer_statistics() -> TestResult {
2432 let temp = TempDir::new()?;
2433 let mut table =
2434 TimeSeriesTable::create(TableLocation::local(temp.path()), binary_heavy_table_meta())
2435 .await?;
2436 let reader = InstrumentedReader::one(binary_heavy_time_series_batch()?);
2437 let batches_consumed = 1;
2438 let expected_version = table.state().version;
2439 let capture = TraceCapture::default();
2440
2441 let report = capture
2442 .run(
2443 table.append(
2444 AppendRequest::new(reader)
2445 .compression(ParquetCompression::Snappy)
2446 .max_rows_per_row_group(1)
2447 .max_bytes_per_row_group(1_024),
2448 ),
2449 )
2450 .await?;
2451 let committed_version = report.committed_version;
2452 assert_eq!(committed_version, expected_version + 1);
2453 assert_eq!(table.state().version, committed_version);
2454
2455 let segment = table
2456 .state()
2457 .segments
2458 .values()
2459 .next()
2460 .ok_or("missing committed segment")?;
2461 let segment_path = temp.path().join(&segment.path);
2462 let file_size_bytes = std::fs::metadata(&segment_path)?.len();
2463 assert_eq!(segment.file_size, Some(file_size_bytes));
2464 let parquet = ParquetRecordBatchReaderBuilder::try_new(File::open(&segment_path)?)?;
2465 let row_group_count = parquet.metadata().row_groups().len();
2466 let parquet_row_count = u64::try_from(parquet.metadata().file_metadata().num_rows())?;
2467 let summarize = |mut values: Vec<u64>| {
2468 values.sort_unstable();
2469 assert_eq!(values.len(), 3);
2470 [values[0], values[1], values[2]]
2471 };
2472 let [
2473 row_group_rows_min,
2474 row_group_rows_median,
2475 row_group_rows_max,
2476 ] = summarize(
2477 parquet
2478 .metadata()
2479 .row_groups()
2480 .iter()
2481 .map(|row_group| u64::try_from(row_group.num_rows()))
2482 .collect::<Result<Vec<_>, _>>()?,
2483 );
2484 let [
2485 row_group_compressed_bytes_min,
2486 row_group_compressed_bytes_median,
2487 row_group_compressed_bytes_max,
2488 ] = summarize(
2489 parquet
2490 .metadata()
2491 .row_groups()
2492 .iter()
2493 .map(|row_group| u64::try_from(row_group.compressed_size()))
2494 .collect::<Result<Vec<_>, _>>()?,
2495 );
2496 let [
2497 row_group_uncompressed_bytes_min,
2498 row_group_uncompressed_bytes_median,
2499 row_group_uncompressed_bytes_max,
2500 ] = summarize(
2501 parquet
2502 .metadata()
2503 .row_groups()
2504 .iter()
2505 .map(|row_group| u64::try_from(row_group.total_byte_size()))
2506 .collect::<Result<Vec<_>, _>>()?,
2507 );
2508 assert_eq!(segment.row_count, parquet_row_count);
2509 assert!(
2510 row_group_uncompressed_bytes_min >= row_group_compressed_bytes_max.saturating_mul(2)
2511 );
2512
2513 let events = capture.events();
2514 let phase_events = events
2515 .iter()
2516 .filter(|event| {
2517 event.name == "table.append"
2518 && event.level == tracing::Level::DEBUG
2519 && event.fields.contains_key("phase")
2520 })
2521 .collect::<Vec<_>>();
2522 assert_eq!(
2523 phase_events
2524 .iter()
2525 .map(|event| event.fields["phase"].as_str())
2526 .collect::<Vec<_>>(),
2527 [
2528 "source_schema_preparation",
2529 "source_consumption_and_parquet_write",
2530 "writer_finalize_and_sync",
2531 "segment_validation",
2532 "coverage_and_overlap",
2533 "sidecar_write",
2534 "transaction_commit",
2535 ]
2536 );
2537 for event in phase_events {
2538 assert_no_row_group_statistics(&event.fields);
2539 assert_eq!(event.target, "timeseries_table_format::table::append");
2540 assert_eq!(
2541 event.fields.get("outcome").map(String::as_str),
2542 Some("succeeded")
2543 );
2544 event.fields["duration_ms"].parse::<u64>()?;
2545 }
2546
2547 let completion_events = events
2548 .iter()
2549 .filter(|event| event.name == "table.append" && event.level == tracing::Level::INFO)
2550 .collect::<Vec<_>>();
2551 assert_eq!(completion_events.len(), 1);
2552 let completion = completion_events[0];
2553 let span = captured_span(&capture, "table.append");
2554 assert_no_row_group_statistics(&span.fields);
2555 for (field, expected) in [
2556 ("expected_version", expected_version.to_string()),
2557 ("committed_version", committed_version.to_string()),
2558 ("row_count", parquet_row_count.to_string()),
2559 ("batches_consumed", batches_consumed.to_string()),
2560 ("row_group_count", row_group_count.to_string()),
2561 ("file_size_bytes", file_size_bytes.to_string()),
2562 ("entity_layout", "single".to_string()),
2563 ("compression", "snappy".to_string()),
2564 ("max_rows_per_row_group", "1".to_string()),
2565 ("max_bytes_per_row_group", "1024".to_string()),
2566 ("outcome", "succeeded".to_string()),
2567 ] {
2568 assert_eq!(
2569 completion.fields.get(field),
2570 Some(&expected),
2571 "event.{field}"
2572 );
2573 assert_eq!(span.fields.get(field), Some(&expected), "span.{field}");
2574 }
2575 for (field, expected) in [
2576 ("row_group_rows_min", row_group_rows_min),
2577 ("row_group_rows_median", row_group_rows_median),
2578 ("row_group_rows_max", row_group_rows_max),
2579 (
2580 "row_group_compressed_bytes_min",
2581 row_group_compressed_bytes_min,
2582 ),
2583 (
2584 "row_group_compressed_bytes_median",
2585 row_group_compressed_bytes_median,
2586 ),
2587 (
2588 "row_group_compressed_bytes_max",
2589 row_group_compressed_bytes_max,
2590 ),
2591 (
2592 "row_group_uncompressed_bytes_min",
2593 row_group_uncompressed_bytes_min,
2594 ),
2595 (
2596 "row_group_uncompressed_bytes_median",
2597 row_group_uncompressed_bytes_median,
2598 ),
2599 (
2600 "row_group_uncompressed_bytes_max",
2601 row_group_uncompressed_bytes_max,
2602 ),
2603 ] {
2604 assert_eq!(
2605 completion.fields.get(field),
2606 Some(&expected.to_string()),
2607 "event.{field}"
2608 );
2609 }
2610 completion.fields["total_duration_ms"].parse::<u64>()?;
2611 assert_eq!(
2612 span.fields.get("rows_consumed"),
2613 Some(&parquet_row_count.to_string())
2614 );
2615 Ok(())
2616 }
2617
2618 #[test]
2619 fn disabled_append_telemetry_skips_input_counting() -> TestResult {
2620 let mut telemetry = AppendTelemetry::new(false);
2621
2622 telemetry.record_batch(×tamp_only_batch(1)?);
2623
2624 assert_eq!(telemetry.batches_consumed, 0);
2625 assert_eq!(telemetry.rows_consumed, 0);
2626 Ok(())
2627 }
2628
2629 #[tokio::test]
2630 async fn midstream_source_failure_reports_phase_and_cleanup() -> TestResult {
2631 let temp = TempDir::new()?;
2632 let mut table =
2633 TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2634 .await?;
2635 let first = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2636 let reader = InstrumentedReader::new(
2637 first.schema(),
2638 vec![
2639 Ok(first),
2640 Err(ArrowError::ComputeError(
2641 "injected batch source failure".to_string(),
2642 )),
2643 ],
2644 )
2645 .0;
2646 let capture = TraceCapture::default();
2647
2648 let error = capture
2649 .run(table.append(reader))
2650 .await
2651 .expect_err("source failure must fail append");
2652 assert!(matches!(
2653 error,
2654 TableError::Append {
2655 source: AppendError::ArrowInput { .. }
2656 }
2657 ));
2658
2659 let events = capture.events();
2660 for event in &events {
2661 assert_no_row_group_statistics(&event.fields);
2662 }
2663 let failures = events
2664 .iter()
2665 .filter(|event| {
2666 event.name == "table.append"
2667 && event.fields.get("outcome").map(String::as_str) == Some("failed")
2668 })
2669 .collect::<Vec<_>>();
2670 assert_eq!(failures.len(), 1);
2671 let failure = failures[0];
2672 assert_eq!(
2673 failure.fields.get("phase").map(String::as_str),
2674 Some("source_consumption_and_parquet_write")
2675 );
2676 assert_eq!(
2677 failure
2678 .fields
2679 .get("last_completed_phase")
2680 .map(String::as_str),
2681 Some("source_schema_preparation")
2682 );
2683 assert_eq!(
2684 failure.fields.get("cleanup_outcome").map(String::as_str),
2685 Some("succeeded")
2686 );
2687 assert_eq!(
2688 failure.fields.get("batches_consumed"),
2689 Some(&"1".to_string())
2690 );
2691 assert_eq!(failure.fields.get("rows_consumed"), Some(&"1".to_string()));
2692 failure.fields["duration_ms"].parse::<u64>()?;
2693 failure.fields["cleanup_duration_ms"].parse::<u64>()?;
2694 failure.fields["total_duration_ms"].parse::<u64>()?;
2695 assert!(!events.iter().any(|event| {
2696 event.name == "table.append"
2697 && event.fields.get("outcome").map(String::as_str) == Some("succeeded")
2698 && event.fields.get("phase").map(String::as_str)
2699 == Some("source_consumption_and_parquet_write")
2700 }));
2701 assert_capture_excludes(&capture, &["injected batch source failure"]);
2702 Ok(())
2703 }
2704
2705 #[tokio::test]
2706 async fn writer_finalization_failure_reports_last_completed_phase() -> TestResult {
2707 let temp = TempDir::new()?;
2708 let mut table =
2709 TimeSeriesTable::create(TableLocation::local(temp.path()), timestamp_only_meta())
2710 .await?;
2711 crate::storage::inject_output_finish_failure(temp.path().join("data"));
2712 let capture = TraceCapture::default();
2713
2714 let error = capture
2715 .run(table.append(timestamp_only_batch(1)?))
2716 .await
2717 .expect_err("writer finalization must fail");
2718 assert!(matches!(
2719 error,
2720 TableError::Append {
2721 source: AppendError::Storage { .. }
2722 }
2723 ));
2724
2725 let failures = capture
2726 .events()
2727 .into_iter()
2728 .filter(|event| {
2729 event.name == "table.append"
2730 && event.fields.get("outcome").map(String::as_str) == Some("failed")
2731 })
2732 .collect::<Vec<_>>();
2733 assert_eq!(failures.len(), 1);
2734 assert_eq!(
2735 failures[0].fields.get("phase").map(String::as_str),
2736 Some("writer_finalize_and_sync")
2737 );
2738 assert_eq!(
2739 failures[0]
2740 .fields
2741 .get("last_completed_phase")
2742 .map(String::as_str),
2743 Some("source_consumption_and_parquet_write")
2744 );
2745 assert_eq!(
2746 failures[0]
2747 .fields
2748 .get("cleanup_outcome")
2749 .map(String::as_str),
2750 Some("succeeded")
2751 );
2752 Ok(())
2753 }
2754
2755 #[tokio::test]
2756 async fn append_consumes_non_send_reader_one_batch_at_a_time() -> TestResult {
2757 let temp = TempDir::new()?;
2758 let mut table =
2759 TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2760 .await?;
2761 let first = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2762 let second = time_series_batch(vec![60_000], vec!["A"], vec![2.0])?;
2763 let (reader, observations) =
2764 InstrumentedReader::new(first.schema(), vec![Ok(first), Ok(second)]);
2765
2766 assert_eq!(table.append(reader).await?.committed_version, 2);
2767 assert!(observations.schema_calls.get() > 0);
2768 assert!(!observations.next_before_schema.get());
2769 assert!(!observations.previous_batch_alive.get());
2770 assert_eq!(observations.next_calls.get(), 3);
2771 Ok(())
2772 }
2773
2774 #[tokio::test(flavor = "current_thread")]
2775 async fn append_yields_while_skipping_zero_row_batches() -> TestResult {
2776 const BATCH_COUNT: usize = 64;
2777
2778 let temp = TempDir::new()?;
2779 let mut table =
2780 TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2781 .await?;
2782 let zero = time_series_batch(Vec::new(), Vec::new(), Vec::new())?;
2783 let (reader, observations) = InstrumentedReader::new(
2784 zero.schema(),
2785 (0..BATCH_COUNT).map(|_| Ok(zero.clone())).collect(),
2786 );
2787 let mut append = Box::pin(table.append(reader));
2788
2789 tokio::select! {
2790 biased;
2791 result = &mut append => panic!("append drained the reader without yielding: {result:?}"),
2792 () = async {
2793 while observations.next_calls.get() < 2 {
2794 tokio::task::yield_now().await;
2795 }
2796 } => {}
2797 }
2798
2799 assert!(observations.next_calls.get() < BATCH_COUNT);
2800 drop(append);
2801 assert_eq!(table.state().version, 1);
2802 assert!(!temp.path().join("data").exists());
2803 Ok(())
2804 }
2805
2806 #[tokio::test(flavor = "current_thread")]
2807 async fn cancelling_append_during_batch_reads_removes_output() -> TestResult {
2808 const BATCH_COUNT: usize = 64;
2809
2810 let temp = TempDir::new()?;
2811 let mut table =
2812 TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2813 .await?;
2814 let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2815 let (reader, observations) = InstrumentedReader::new(
2816 batch.schema(),
2817 (0..BATCH_COUNT).map(|_| Ok(batch.clone())).collect(),
2818 );
2819 let mut append = Box::pin(table.append(reader));
2820
2821 tokio::select! {
2822 biased;
2823 result = &mut append => panic!("append drained the reader without yielding: {result:?}"),
2824 () = async {
2825 while observations.next_calls.get() < 2 {
2826 tokio::task::yield_now().await;
2827 }
2828 } => {}
2829 }
2830
2831 assert!(observations.next_calls.get() < BATCH_COUNT);
2832 assert_eq!(data_files(temp.path())?.len(), 1);
2833 drop(append);
2834 assert_eq!(table.state().version, 1);
2835 assert!(table.state().segments.is_empty());
2836 assert!(data_files(temp.path())?.is_empty());
2837 assert!(coverage_files(temp.path())?.is_empty());
2838 assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
2839 Ok(())
2840 }
2841
2842 #[tokio::test]
2843 async fn append_rejects_empty_sources_before_creating_data() -> TestResult {
2844 let temp = TempDir::new()?;
2845 let mut table =
2846 TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2847 .await?;
2848 let zero = time_series_batch(Vec::new(), Vec::new(), Vec::new())?;
2849 let schema = zero.schema();
2850
2851 assert!(matches!(
2852 table.append(vec![zero.clone(), zero]).await,
2853 Err(TableError::Append {
2854 source: AppendError::EmptyInput
2855 })
2856 ));
2857 let empty = RecordBatchIterator::new(Vec::<Result<RecordBatch, ArrowError>>::new(), schema);
2858 assert!(matches!(
2859 table.append(empty).await,
2860 Err(TableError::Append {
2861 source: AppendError::EmptyInput
2862 })
2863 ));
2864 assert_eq!(table.state().version, 1);
2865 assert!(table.state().segments.is_empty());
2866 assert!(!temp.path().join("data").exists());
2867
2868 let data = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2869 let leading_zero = time_series_batch(Vec::new(), Vec::new(), Vec::new())?;
2870 assert_eq!(
2871 table
2872 .append(vec![leading_zero, data])
2873 .await?
2874 .committed_version,
2875 2
2876 );
2877 assert_eq!(table.state().version, 2);
2878 assert_eq!(table.state().segments.len(), 1);
2879 Ok(())
2880 }
2881
2882 #[tokio::test]
2883 async fn append_cleans_up_after_later_schema_or_source_error() -> TestResult {
2884 for (configured, source_error) in
2885 [(false, false), (false, true), (true, false), (true, true)]
2886 {
2887 let temp = TempDir::new()?;
2888 let mut table =
2889 TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
2890 .await?;
2891 let first = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
2892 let later = if source_error {
2893 Err(ArrowError::ComputeError(
2894 "injected batch source failure".to_string(),
2895 ))
2896 } else {
2897 Ok(input_batch(vec![1])?)
2898 };
2899 let (reader, _) = InstrumentedReader::new(first.schema(), vec![Ok(first), later]);
2900
2901 let result = if configured {
2902 table
2903 .append(AppendRequest::new(reader).max_rows_per_row_group(1))
2904 .await
2905 } else {
2906 table.append(reader).await
2907 };
2908 let error = result.expect_err("later batch must fail the append");
2909 assert!(
2910 matches!(
2911 error,
2912 TableError::Append {
2913 source: AppendError::ArrowInput { .. }
2914 }
2915 ),
2916 "configured={configured}, source_error={source_error}"
2917 );
2918 assert_eq!(
2919 table.state().version,
2920 1,
2921 "configured={configured}, source_error={source_error}"
2922 );
2923 assert!(
2924 table.state().segments.is_empty(),
2925 "configured={configured}, source_error={source_error}"
2926 );
2927 assert!(
2928 data_files(temp.path())?.is_empty(),
2929 "configured={configured}, source_error={source_error}"
2930 );
2931 assert!(
2932 coverage_files(temp.path())?.is_empty(),
2933 "configured={configured}, source_error={source_error}"
2934 );
2935 assert!(
2936 !temp.path().join(layout::commit_rel_path(2)).exists(),
2937 "configured={configured}, source_error={source_error}"
2938 );
2939 }
2940 Ok(())
2941 }
2942
2943 #[tokio::test]
2944 async fn append_cleans_up_writer_lifecycle_failures() -> TestResult {
2945 #[derive(Clone, Copy, Debug)]
2946 enum Stage {
2947 Write,
2948 Close,
2949 Finish,
2950 }
2951
2952 for stage in [Stage::Write, Stage::Close, Stage::Finish] {
2953 let temp = TempDir::new()?;
2954 let mut table =
2955 TimeSeriesTable::create(TableLocation::local(temp.path()), timestamp_only_meta())
2956 .await?;
2957 let data_dir = temp.path().join("data");
2958 match stage {
2959 Stage::Write => crate::storage::inject_output_write_failure(data_dir.clone(), 2),
2960 Stage::Close => crate::storage::inject_output_write_failure(data_dir.clone(), 1),
2961 Stage::Finish => crate::storage::inject_output_finish_failure(data_dir.clone()),
2962 }
2963 let row_count = if matches!(stage, Stage::Write) {
2964 parquet::file::properties::DEFAULT_MAX_ROW_GROUP_ROW_COUNT
2965 } else {
2966 1
2967 };
2968
2969 let result = table.append(timestamp_only_batch(row_count)?).await;
2970 let Err(error) = result else {
2971 panic!("{stage:?} writer failure was not injected");
2972 };
2973 if matches!(stage, Stage::Finish) {
2974 assert!(
2975 matches!(
2976 error,
2977 TableError::Append {
2978 source: AppendError::Storage { .. }
2979 }
2980 ),
2981 "{stage:?}"
2982 );
2983 } else {
2984 assert!(
2985 matches!(
2986 error,
2987 TableError::Append {
2988 source: AppendError::ParquetWrite { .. }
2989 }
2990 ),
2991 "{stage:?}: {error}"
2992 );
2993 }
2994 assert_eq!(table.state().version, 1, "{stage:?}");
2995 assert!(table.state().segments.is_empty(), "{stage:?}");
2996 assert!(data_files(temp.path())?.is_empty(), "{stage:?}");
2997 assert!(coverage_files(temp.path())?.is_empty(), "{stage:?}");
2998 assert!(
2999 !temp.path().join(layout::commit_rel_path(2)).exists(),
3000 "{stage:?}"
3001 );
3002 }
3003 Ok(())
3004 }
3005
3006 #[tokio::test]
3007 async fn cancelling_append_before_publication_removes_owned_artifacts() -> TestResult {
3008 let temp = TempDir::new()?;
3009 let location = TableLocation::local(temp.path());
3010 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
3011 let commit_path = temp.path().join(layout::commit_rel_path(2));
3012 let current_path = temp.path().join(layout::current_rel_path());
3013 let current_temp_path = current_path.with_extension("tmp");
3014 let mut pause = crate::storage::pause_atomic_write_before_rename(current_path);
3015 let mut append = Box::pin(table.append(time_series_batch(vec![0], vec!["A"], vec![1.0])?));
3016
3017 tokio::select! {
3018 () = pause.wait_until_paused() => {}
3019 result = &mut append => panic!("append completed before cancellation: {result:?}"),
3020 }
3021
3022 assert!(commit_path.is_file());
3023 assert!(current_temp_path.is_file());
3024 assert_eq!(data_files(temp.path())?.len(), 1);
3025 assert_eq!(coverage_files(temp.path())?.len(), 2);
3026 let before_publication = TimeSeriesTable::open(location.clone()).await?;
3027 assert_eq!(before_publication.state().version, 1);
3028 assert!(before_publication.state().segments.is_empty());
3029
3030 drop(append);
3031 pause.release();
3032
3033 assert!(!commit_path.exists());
3034 assert!(!current_temp_path.exists());
3035 assert!(data_files(temp.path())?.is_empty());
3036 assert!(coverage_files(temp.path())?.is_empty());
3037 assert_eq!(table.state().version, 1);
3038 assert!(table.state().segments.is_empty());
3039 let reopened = TimeSeriesTable::open(location).await?;
3040 assert_eq!(reopened.state().version, 1);
3041 assert!(reopened.state().segments.is_empty());
3042 Ok(())
3043 }
3044
3045 #[tokio::test]
3046 async fn post_commit_observer_panic_preserves_owned_data_files() -> TestResult {
3047 let streamed_root = TempDir::new()?;
3048 let streamed_location = TableLocation::local(streamed_root.path());
3049 let mut streamed_table =
3050 TimeSeriesTable::create(streamed_location.clone(), make_basic_table_meta()).await?;
3051 let callsite_guard = panic_on_commit_close_dispatch();
3052 let outcome = std::panic::AssertUnwindSafe(
3053 streamed_table
3054 .append(time_series_batch(vec![0], vec!["A"], vec![1.0])?)
3055 .with_subscriber(panic_on_commit_close_dispatch()),
3056 )
3057 .catch_unwind()
3058 .await;
3059 drop(callsite_guard);
3060 assert!(outcome.is_err());
3061 let reopened = TimeSeriesTable::open(streamed_location).await?;
3062 assert_eq!(reopened.state().version, 2);
3063 assert_eq!(reopened.state().segments.len(), 1);
3064 assert!(
3065 reopened
3066 .state()
3067 .segments
3068 .values()
3069 .all(|segment| streamed_root.path().join(&segment.path).is_file())
3070 );
3071
3072 Ok(())
3073 }
3074
3075 #[tokio::test]
3076 async fn simultaneous_appends_publish_one_and_clean_loser() -> TestResult {
3077 let temp = TempDir::new()?;
3078 let location = TableLocation::local(temp.path());
3079 let mut winner = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
3080 let mut loser = TimeSeriesTable::open(location.clone()).await?;
3081 let loser_state_before = loser.state().clone();
3082 let commit_path = temp.path().join(layout::commit_rel_path(2));
3083 let mut pause = crate::storage::pause_atomic_write_before_rename(
3084 temp.path().join(layout::current_rel_path()),
3085 );
3086 let mut winner_append =
3087 Box::pin(winner.append(time_series_batch(vec![0], vec!["A"], vec![1.0])?));
3088
3089 tokio::select! {
3090 () = pause.wait_until_paused() => {}
3091 result = &mut winner_append => panic!("winner completed before race: {result:?}"),
3092 }
3093
3094 let data_before = data_files(temp.path())?;
3095 let coverage_before = coverage_files(temp.path())?;
3096 let commit_before = std::fs::read(&commit_path)?;
3097
3098 let loser_error = loser
3099 .append(time_series_batch(vec![120_000], vec!["A"], vec![2.0])?)
3100 .await
3101 .expect_err("second simultaneous writer must lose commit 2");
3102 assert!(matches!(
3103 loser_error,
3104 TableError::Append {
3105 source: AppendError::Commit {
3106 source: CommitError::Storage {
3107 source: StorageError::AlreadyExists { .. },
3108 },
3109 }
3110 }
3111 ));
3112 assert_eq!(loser.state(), &loser_state_before);
3113
3114 let data_after = data_files(temp.path())?;
3115 assert_eq!(data_after, data_before);
3116 assert_eq!(coverage_files(temp.path())?, coverage_before);
3117 assert_eq!(std::fs::read(&commit_path)?, commit_before);
3118
3119 pause.release();
3120 let winner_report = winner_append.await?;
3121 assert_eq!(winner_report.starting_version, 1);
3122 assert_eq!(winner_report.committed_version, 2);
3123 assert_eq!(winner.state().version, 2);
3124 assert_eq!(winner.state().segments.len(), 1);
3125 assert!(
3126 winner
3127 .state()
3128 .segments
3129 .contains_key(&winner_report.segment_path)
3130 );
3131 assert_eq!(
3132 data_files(temp.path())?,
3133 [temp.path().join(&winner_report.segment_path)]
3134 );
3135 assert!(!temp.path().join(layout::commit_rel_path(3)).exists());
3136
3137 let reopened = TimeSeriesTable::open(location).await?;
3138 assert_eq!(reopened.state().version, 2);
3139 assert_eq!(reopened.state().segments.len(), 1);
3140 assert!(
3141 reopened
3142 .state()
3143 .segments
3144 .contains_key(&winner_report.segment_path)
3145 );
3146 assert_eq!(
3147 reopened
3148 .state()
3149 .segments
3150 .values()
3151 .map(|segment| segment.row_count)
3152 .sum::<u64>(),
3153 1
3154 );
3155 Ok(())
3156 }
3157
3158 #[tokio::test]
3159 async fn guarded_output_cleans_up_parquet_writer_creation_failure() -> TestResult {
3160 let temp = TempDir::new()?;
3161 let location = StorageLocation::local(temp.path());
3162 let path = Path::new("data/create-failure.parquet");
3163 let sink = storage::open_new_output_sink(&location, path).await?;
3164 let schema = Arc::new(Schema::new(vec![Field::new(
3165 "unsupported",
3166 DataType::Struct(arrow::datatypes::Fields::empty()),
3167 true,
3168 )]));
3169
3170 assert!(ParquetArrowWriter::try_new(sink, schema, None).is_err());
3171 assert!(!temp.path().join(path).exists());
3172 Ok(())
3173 }
3174
3175 #[tokio::test]
3176 async fn append_preserves_writer_error_and_failed_cleanup_path() -> TestResult {
3177 let temp = TempDir::new()?;
3178 let mut table =
3179 TimeSeriesTable::create(TableLocation::local(temp.path()), timestamp_only_meta())
3180 .await?;
3181 let data_dir = temp.path().join("data");
3182 crate::storage::inject_output_write_failure(data_dir.clone(), 1);
3183 crate::storage::inject_cleanup_failure(data_dir.clone());
3184 let capture = TraceCapture::default();
3185
3186 let error = capture
3187 .run(table.append(timestamp_only_batch(1)?))
3188 .await
3189 .expect_err("writer and cleanup failures must fail append");
3190 let TableError::Append {
3191 source:
3192 AppendError::Rollback {
3193 source,
3194 cleanup_errors,
3195 },
3196 } = error
3197 else {
3198 panic!("unexpected error: {error}");
3199 };
3200 assert!(matches!(*source, AppendError::ParquetWrite { .. }));
3201 assert_eq!(cleanup_errors.len(), 1);
3202 let remaining = data_files(temp.path())?;
3203 assert_eq!(remaining.len(), 1);
3204 let remaining_path = &remaining[0];
3205 let remaining_name = remaining_path
3206 .file_name()
3207 .expect("remaining data filename")
3208 .to_str()
3209 .expect("UTF-8 data filename");
3210 assert!(matches!(
3211 &cleanup_errors[0],
3212 StorageError::OtherIo { path, .. } if path.ends_with(remaining_name)
3213 ));
3214 assert_eq!(table.state().version, 1);
3215 assert!(table.state().segments.is_empty());
3216 assert!(coverage_files(temp.path())?.is_empty());
3217 assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
3218
3219 let failures = capture
3220 .events()
3221 .into_iter()
3222 .filter(|event| {
3223 event.name == "table.append"
3224 && event.fields.get("outcome").map(String::as_str) == Some("failed")
3225 })
3226 .collect::<Vec<_>>();
3227 assert_eq!(failures.len(), 1);
3228 for (field, expected) in [
3229 ("phase", "writer_finalize_and_sync"),
3230 (
3231 "last_completed_phase",
3232 "source_consumption_and_parquet_write",
3233 ),
3234 ("cleanup_outcome", "failed"),
3235 ] {
3236 assert_eq!(
3237 failures[0].fields.get(field).map(String::as_str),
3238 Some(expected)
3239 );
3240 }
3241 Ok(())
3242 }
3243
3244 #[tokio::test]
3245 async fn append_retries_cleanup_after_sidecar_write_cleanup_failure() -> TestResult {
3246 for sidecar_dir in [SEGMENT_COVERAGE_DIR, TABLE_SNAPSHOT_DIR] {
3247 let temp = TempDir::new()?;
3248 let mut table =
3249 TimeSeriesTable::create(TableLocation::local(temp.path()), timestamp_only_meta())
3250 .await?;
3251 crate::storage::inject_write_new_failure(temp.path().join(sidecar_dir), true);
3252
3253 let error = table
3254 .append(timestamp_only_batch(1)?)
3255 .await
3256 .expect_err("sidecar write and its first cleanup must fail");
3257
3258 assert!(matches!(
3259 error,
3260 TableError::Append {
3261 source: AppendError::CoverageSidecar { source }
3262 } if matches!(
3263 source.as_ref(),
3264 CoverageSidecarError::Storage {
3265 source: StorageError::CleanupFailed { .. }
3266 }
3267 )
3268 ));
3269 assert_eq!(table.state().version, 1, "{sidecar_dir}");
3270 assert!(table.state().segments.is_empty(), "{sidecar_dir}");
3271 assert!(data_files(temp.path())?.is_empty(), "{sidecar_dir}");
3272 assert!(coverage_files(temp.path())?.is_empty(), "{sidecar_dir}");
3273 assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
3274 }
3275 Ok(())
3276 }
3277
3278 #[tokio::test]
3279 async fn append_conflict_cleans_attempt_and_stays_invisible() -> TestResult {
3280 let temp = TempDir::new()?;
3281 let location = TableLocation::local(temp.path());
3282 let mut winner = TimeSeriesTable::create(location.clone(), widening_table_meta()).await?;
3283 let mut loser = TimeSeriesTable::open(location.clone()).await?;
3284 let loser_state_before = loser.state().clone();
3285
3286 let winner_report = winner
3287 .append(widening_batch(
3288 vec![0],
3289 vec![1],
3290 vec![Some(1.0)],
3291 vec!["winner"],
3292 )?)
3293 .await?;
3294 let data_before = data_files(temp.path())?;
3295 let coverage_before = coverage_files(temp.path())?;
3296 assert_eq!(winner_report.starting_version, 1);
3297 assert_eq!(winner_report.committed_version, 2);
3298 assert_eq!(data_before, [temp.path().join(&winner_report.segment_path)]);
3299
3300 let error = loser
3301 .append(
3302 AppendRequest::new(widening_batch(
3303 vec![u32::MAX],
3304 vec![2],
3305 vec![Some(2.0)],
3306 vec!["loser"],
3307 )?)
3308 .max_rows_per_row_group(1),
3309 )
3310 .await
3311 .expect_err("stale streaming append must conflict");
3312
3313 assert!(matches!(
3314 &error,
3315 TableError::Append {
3316 source: AppendError::Commit {
3317 source: CommitError::Conflict {
3318 expected: 1,
3319 found: 2,
3320 ..
3321 }
3322 }
3323 }
3324 ));
3325 assert_eq!(loser.state(), &loser_state_before);
3326 assert_eq!(data_files(temp.path())?, data_before);
3327 assert_eq!(coverage_files(temp.path())?, coverage_before);
3328 assert!(!temp.path().join(layout::commit_rel_path(3)).exists());
3329
3330 let reopened = TimeSeriesTable::open(location).await?;
3331 assert_eq!(reopened.state().version, 2);
3332 assert_eq!(reopened.state().segments.len(), 1);
3333 assert_eq!(
3334 reopened
3335 .state()
3336 .segments
3337 .values()
3338 .next()
3339 .map(|segment| { temp.path().join(&segment.path) }),
3340 data_before.first().cloned()
3341 );
3342 Ok(())
3343 }
3344
3345 #[tokio::test]
3346 async fn append_ambiguous_commit_reports_and_preserves_generated_path() -> TestResult {
3347 let temp = TempDir::new()?;
3348 let location = TableLocation::local(temp.path());
3349 let mut table = TimeSeriesTable::create(location.clone(), widening_table_meta()).await?;
3350 let commit_path = temp.path().join(layout::commit_rel_path(2));
3351 crate::storage::inject_write_new_failure(commit_path.clone(), true);
3352 let capture = TraceCapture::default();
3353
3354 let error = capture
3355 .run(
3356 table.append(
3357 AppendRequest::new(widening_batch(
3358 vec![0],
3359 vec![1],
3360 vec![Some(1.0)],
3361 vec!["ambiguous"],
3362 )?)
3363 .max_rows_per_row_group(1),
3364 ),
3365 )
3366 .await
3367 .expect_err("failed commit cleanup must be ambiguous");
3368 let TableError::Append {
3369 source:
3370 AppendError::CommitAmbiguous {
3371 segment_path,
3372 source,
3373 },
3374 } = error
3375 else {
3376 panic!("unexpected error: {error}");
3377 };
3378
3379 assert!(matches!(
3380 source.as_ref(),
3381 CommitError::AmbiguousOutcome { .. }
3382 ));
3383 assert_eq!(
3384 captured_span(&capture, "table.append").target,
3385 "timeseries_table_format::table::append"
3386 );
3387 assert!(segment_path.starts_with("data/"));
3388 assert!(temp.path().join(&segment_path).is_file());
3389 assert_eq!(
3390 data_files(temp.path())?,
3391 vec![temp.path().join(&segment_path)]
3392 );
3393 assert_eq!(table.state().version, 1);
3394 assert!(table.state().segments.is_empty());
3395 assert_eq!(table.log.load_current_version().await?, 1);
3396 assert!(commit_path.exists());
3397 assert_eq!(coverage_files(temp.path())?.len(), 2);
3398 let append_span = capture
3399 .spans()
3400 .into_iter()
3401 .find(|span| span.name == "table.append")
3402 .expect("table.append span");
3403 assert_no_row_group_statistics(&append_span.fields);
3404 assert_eq!(append_span.fields.get("segment_path"), Some(&segment_path));
3405 assert_eq!(
3406 append_span.fields.get("outcome").map(String::as_str),
3407 Some("ambiguous")
3408 );
3409 assert_eq!(
3410 append_span
3411 .fields
3412 .get("cleanup_outcome")
3413 .map(String::as_str),
3414 Some("not_attempted")
3415 );
3416 let events = capture.events();
3417 for event in &events {
3418 assert_no_row_group_statistics(&event.fields);
3419 }
3420 let failures = events
3421 .into_iter()
3422 .filter(|event| {
3423 event.name == "table.append"
3424 && event.fields.get("outcome").map(String::as_str) == Some("ambiguous")
3425 })
3426 .collect::<Vec<_>>();
3427 assert_eq!(failures.len(), 1);
3428 assert_eq!(
3429 failures[0].fields.get("phase").map(String::as_str),
3430 Some("transaction_commit")
3431 );
3432 assert_eq!(
3433 failures[0]
3434 .fields
3435 .get("last_completed_phase")
3436 .map(String::as_str),
3437 Some("sidecar_write")
3438 );
3439 assert_eq!(
3440 failures[0]
3441 .fields
3442 .get("cleanup_outcome")
3443 .map(String::as_str),
3444 Some("not_attempted")
3445 );
3446 assert_eq!(failures[0].fields.get("retained_path"), Some(&segment_path));
3447 let reopened = TimeSeriesTable::open(location).await?;
3448 assert_eq!(reopened.state().version, 1);
3449 assert!(reopened.state().segments.is_empty());
3450 Ok(())
3451 }
3452
3453 #[tokio::test]
3454 async fn append_conflict_reports_every_failed_cleanup_path() -> TestResult {
3455 let temp = TempDir::new()?;
3456 let location = TableLocation::local(temp.path());
3457 let mut winner = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
3458 let mut loser = TimeSeriesTable::open(location.clone()).await?;
3459 winner
3460 .append(time_series_batch(vec![0], vec!["A"], vec![1.0])?)
3461 .await?;
3462 let data_before = data_files(temp.path())?;
3463 let coverage_before = coverage_files(temp.path())?;
3464
3465 for dir in [
3466 Path::new("data"),
3467 Path::new(SEGMENT_COVERAGE_DIR),
3468 Path::new(TABLE_SNAPSHOT_DIR),
3469 ] {
3470 crate::storage::inject_cleanup_failure(temp.path().join(dir));
3471 }
3472 let error = loser
3473 .append(time_series_batch(vec![120_000], vec!["A"], vec![2.0])?)
3474 .await
3475 .expect_err("conflict cleanup failures must be reported");
3476 let message = error.to_string();
3477
3478 let (source, cleanup_errors) = match &error {
3479 TableError::Append {
3480 source:
3481 AppendError::Rollback {
3482 source,
3483 cleanup_errors,
3484 },
3485 } => (source, cleanup_errors),
3486 other => panic!("unexpected error: {other:?}"),
3487 };
3488 assert!(matches!(
3489 source.as_ref(),
3490 AppendError::Commit {
3491 source: CommitError::Conflict { .. }
3492 }
3493 ));
3494 assert_eq!(cleanup_errors.len(), 3);
3495
3496 assert!(message.contains("Commit conflict"));
3497 let data_after = data_files(temp.path())?;
3498 assert_eq!(data_after.len(), data_before.len() + 1);
3499 let coverage_after = coverage_files(temp.path())?;
3500 assert_eq!(coverage_after.len(), coverage_before.len() + 2);
3501 let mut failed_paths = data_after
3502 .iter()
3503 .filter(|path| !data_before.contains(path))
3504 .cloned()
3505 .collect::<Vec<_>>();
3506 failed_paths.extend(
3507 coverage_after
3508 .keys()
3509 .filter(|path| !coverage_before.contains_key(*path))
3510 .map(|path| temp.path().join(path)),
3511 );
3512 for path in failed_paths {
3513 let filename = path
3514 .file_name()
3515 .expect("failed cleanup filename")
3516 .to_str()
3517 .expect("UTF-8 cleanup filename");
3518 assert!(
3519 cleanup_errors.iter().any(|error| matches!(
3520 error,
3521 StorageError::OtherIo { path, .. } if path.ends_with(filename)
3522 )),
3523 "missing typed cleanup failure for {}",
3524 path.display()
3525 );
3526 assert!(
3527 message.contains(filename),
3528 "missing cleanup diagnostic for {}",
3529 path.display()
3530 );
3531 }
3532
3533 let reopened = TimeSeriesTable::open(location).await?;
3534 assert_eq!(reopened.state().version, 2);
3535 assert_eq!(reopened.state().segments.len(), 1);
3536 Ok(())
3537 }
3538
3539 #[tokio::test]
3540 async fn append_validates_schema_before_creating_data() -> TestResult {
3541 let temp = TempDir::new()?;
3542 let mut table =
3543 TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
3544 .await?;
3545
3546 assert!(matches!(
3547 table.append(input_batch(vec![1])?).await,
3548 Err(TableError::Append {
3549 source: AppendError::SchemaValidation { .. }
3550 })
3551 ));
3552 assert_eq!(table.state().version, 1);
3553 assert!(!temp.path().join("data").exists());
3554 Ok(())
3555 }
3556
3557 #[tokio::test]
3558 async fn append_removes_owned_data_after_overlap() -> TestResult {
3559 let temp = TempDir::new()?;
3560 let mut table =
3561 TimeSeriesTable::create(TableLocation::local(temp.path()), make_basic_table_meta())
3562 .await?;
3563 let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
3564 assert_eq!(table.append(batch.clone()).await?.committed_version, 2);
3565
3566 assert!(matches!(
3567 table
3568 .append(AppendRequest::new(batch).max_rows_per_row_group(1))
3569 .await,
3570 Err(TableError::Append {
3571 source: AppendError::PersistedIndexIntervalOverlap { .. }
3572 })
3573 ));
3574 assert_eq!(table.state().version, 2);
3575 assert_eq!(data_files(temp.path())?.len(), 1);
3576 Ok(())
3577 }
3578
3579 #[tokio::test]
3580 async fn append_adopts_schema_on_first_append() -> TestResult {
3581 let temp = TempDir::new()?;
3582 let mut meta = make_basic_table_meta();
3583 meta.logical_schema = None;
3584 let mut table = TimeSeriesTable::create(TableLocation::local(temp.path()), meta).await?;
3585
3586 let batch = time_series_batch(vec![0], vec!["A"], vec![1.0])?;
3587 assert_eq!(table.append(batch).await?.committed_version, 2);
3588 assert!(table.state().table_meta.logical_schema.is_some());
3589 Ok(())
3590 }
3591
3592 #[tokio::test]
3593 async fn append_preserves_exact_schema_through_optimize() -> TestResult {
3594 let timezone = "America/Phoenix";
3595 let schema = Arc::new(Schema::new(vec![
3596 Field::new(
3597 "ts",
3598 DataType::Timestamp(ArrowTimeUnit::Millisecond, Some(timezone.into())),
3599 false,
3600 ),
3601 Field::new("symbol", DataType::Utf8, false),
3602 Field::new(
3603 "items",
3604 DataType::List(Arc::new(Field::new("item", DataType::Int64, true))),
3605 true,
3606 ),
3607 Field::new(
3608 "attrs",
3609 DataType::Map(
3610 Arc::new(Field::new(
3611 "entries",
3612 DataType::Struct(Fields::from(vec![
3613 Field::new("key", DataType::Utf8, false),
3614 Field::new("value", DataType::Binary, true),
3615 ])),
3616 false,
3617 )),
3618 true,
3619 ),
3620 true,
3621 ),
3622 ]));
3623 let expected = LogicalSchema::try_from_arrow_schema(&schema).expect("logical schema");
3624 let batch = RecordBatch::try_new(
3625 Arc::clone(&schema),
3626 vec![
3627 Arc::new(TimestampMillisecondArray::from(vec![0, 0]).with_timezone(timezone)),
3628 Arc::new(StringArray::from(vec!["A", "B"])),
3629 new_null_array(schema.field(2).data_type(), 2),
3630 new_null_array(schema.field(3).data_type(), 2),
3631 ],
3632 )?;
3633
3634 for has_canonical_schema in [true, false] {
3635 let temp = TempDir::new()?;
3636 let location = TableLocation::local(temp.path());
3637 let mut meta = TableMeta::new_time_series_with_schema(
3638 IndexSpec {
3639 column: "ts".to_string(),
3640 entity_columns: vec!["symbol".to_string()],
3641 kind: IndexKind::Timestamp {
3642 index_granularity: TimeIndexGranularity::Minutes(1),
3643 timezone: Some(timezone.to_string()),
3644 },
3645 },
3646 expected.clone(),
3647 );
3648 if !has_canonical_schema {
3649 meta.logical_schema = None;
3650 }
3651 let mut table = TimeSeriesTable::create(location.clone(), meta).await?;
3652
3653 assert_eq!(table.append(batch.clone()).await?.committed_version, 2);
3654 assert_eq!(
3655 table.state().table_meta.logical_schema.as_ref(),
3656 Some(&expected)
3657 );
3658
3659 let later_batch = RecordBatch::try_new(
3660 Arc::clone(&schema),
3661 vec![
3662 Arc::new(
3663 TimestampMillisecondArray::from(vec![120_000, 120_000])
3664 .with_timezone(timezone),
3665 ),
3666 Arc::new(StringArray::from(vec!["A", "B"])),
3667 new_null_array(schema.field(2).data_type(), 2),
3668 new_null_array(schema.field(3).data_type(), 2),
3669 ],
3670 )?;
3671 assert_eq!(table.append(later_batch).await?.committed_version, 3);
3672
3673 let report = table.optimize().await?;
3674 assert_eq!(report.committed_version, 4);
3675 assert_eq!(report.candidate_source_segments, 2);
3676 assert_eq!(report.replacement_segments_written, 4);
3677 assert_eq!(report.rows_read, 4);
3678 assert_eq!(report.rows_written, 4);
3679 assert!(
3680 table
3681 .state()
3682 .segments
3683 .values()
3684 .all(|segment| matches!(segment.entity_layout, SegmentEntityLayout::Single(_)))
3685 );
3686 let reopened = TimeSeriesTable::open(location).await?;
3687 assert_eq!(
3688 reopened.state().table_meta.logical_schema.as_ref(),
3689 Some(&expected)
3690 );
3691 }
3692 Ok(())
3693 }
3694
3695 fn registered_index(kind: IndexKind) -> IndexSpec {
3696 IndexSpec {
3697 column: "ts".to_string(),
3698 entity_columns: Vec::new(),
3699 kind,
3700 }
3701 }
3702
3703 fn write_single_index_parquet(
3704 path: &Path,
3705 data_type: DataType,
3706 values: ArrayRef,
3707 ) -> TestResult {
3708 if let Some(parent) = path.parent() {
3709 std::fs::create_dir_all(parent)?;
3710 }
3711 let schema = Arc::new(Schema::new(vec![Field::new("ts", data_type, false)]));
3712 let batch = RecordBatch::try_new(Arc::clone(&schema), vec![values])?;
3713 let mut writer = ArrowWriter::try_new(File::create(path)?, schema, None)?;
3714 writer.write(&batch)?;
3715 writer.close()?;
3716 Ok(())
3717 }
3718
3719 fn composite_entity_batch(rows: &[(i64, &str, &str, f64)]) -> Result<RecordBatch, ArrowError> {
3720 let schema = Arc::new(Schema::new(vec![
3721 Field::new(
3722 "ts",
3723 DataType::Timestamp(arrow::datatypes::TimeUnit::Millisecond, None),
3724 false,
3725 ),
3726 Field::new("symbol", DataType::Utf8, false),
3727 Field::new("venue", DataType::Utf8, false),
3728 Field::new("price", DataType::Float64, false),
3729 ]));
3730 RecordBatch::try_new(
3731 Arc::clone(&schema),
3732 vec![
3733 Arc::new(TimestampMillisecondArray::from(
3734 rows.iter().map(|row| row.0).collect::<Vec<_>>(),
3735 )),
3736 Arc::new(StringArray::from(
3737 rows.iter().map(|row| row.1).collect::<Vec<_>>(),
3738 )),
3739 Arc::new(StringArray::from(
3740 rows.iter().map(|row| row.2).collect::<Vec<_>>(),
3741 )),
3742 Arc::new(Float64Array::from(
3743 rows.iter().map(|row| row.3).collect::<Vec<_>>(),
3744 )),
3745 ],
3746 )
3747 }
3748
3749 fn write_composite_entity_parquet(path: &Path, rows: &[(i64, &str, &str, f64)]) -> TestResult {
3750 if let Some(parent) = path.parent() {
3751 std::fs::create_dir_all(parent)?;
3752 }
3753 let batch = composite_entity_batch(rows)?;
3754 let schema = batch.schema();
3755 let mut writer = ArrowWriter::try_new(File::create(path)?, schema, None)?;
3756 writer.write(&batch)?;
3757 writer.close()?;
3758 Ok(())
3759 }
3760
3761 fn coverage_files(root: &Path) -> std::io::Result<BTreeMap<PathBuf, Vec<u8>>> {
3762 let mut files = BTreeMap::new();
3763 for rel_dir in [SEGMENT_COVERAGE_DIR, TABLE_SNAPSHOT_DIR] {
3764 let dir = root.join(rel_dir);
3765 if !dir.exists() {
3766 continue;
3767 }
3768 for entry in std::fs::read_dir(dir)? {
3769 let path = entry?.path();
3770 if path.is_file() {
3771 files.insert(
3772 path.strip_prefix(root)
3773 .expect("coverage path under root")
3774 .to_owned(),
3775 std::fs::read(path)?,
3776 );
3777 }
3778 }
3779 }
3780 Ok(files)
3781 }
3782
3783 fn data_files(root: &Path) -> std::io::Result<Vec<PathBuf>> {
3784 let dir = root.join("data");
3785 if !dir.exists() {
3786 return Ok(Vec::new());
3787 }
3788 let mut pending = vec![dir];
3789 let mut files = Vec::new();
3790 while let Some(dir) = pending.pop() {
3791 for entry in std::fs::read_dir(dir)? {
3792 let entry = entry?;
3793 if entry.file_type()?.is_dir() {
3794 pending.push(entry.path());
3795 } else {
3796 files.push(entry.path());
3797 }
3798 }
3799 }
3800 files.sort();
3801 Ok(files)
3802 }
3803
3804 #[test]
3805 fn entity_layout_classification_rejects_empty_coverage() {
3806 assert!(matches!(
3807 classify_entity_layout("data/empty.parquet", &EntityCoverage::empty()),
3808 Err(AppendError::EmptySegmentEntityCoverage { segment_path })
3809 if segment_path == "data/empty.parquet"
3810 ));
3811 }
3812
3813 #[tokio::test]
3814 async fn duplicate_implicit_interval_rolls_back_first_append() -> TestResult {
3815 let temp = TempDir::new()?;
3816 let location = TableLocation::local(temp.path());
3817 let mut table = TimeSeriesTable::create(
3818 location.clone(),
3819 TableMeta::new_time_series(timestamp_only_index()),
3820 )
3821 .await?;
3822 let state_before = table.state().clone();
3823
3824 let error = table
3825 .append(timestamp_only_batch_with_millis([0, 30_000])?)
3826 .await
3827 .expect_err("duplicate implicit interval must fail");
3828
3829 let message = error.to_string();
3830 match error {
3831 TableError::Append {
3832 source: AppendError::GeneratedSegmentCoverage { source, .. },
3833 } => match *source {
3834 SegmentCoverageError::DuplicateIndexInterval {
3835 path: segment_path,
3836 example_identity: None,
3837 example_index_interval,
3838 } => {
3839 assert!(segment_path.starts_with("data/"));
3840 assert!(segment_path.ends_with(".parquet"));
3841 assert!(message.contains(&segment_path));
3842 assert!(message.contains("Duplicate ordered-index interval"));
3843 assert!(message.contains("ordered-index interval"));
3844 assert_eq!(
3845 example_index_interval.to_string(),
3846 "[1970-01-01T00:00:00Z, 1970-01-01T00:01:00Z)"
3847 );
3848 }
3849 other => panic!("expected duplicate interval source, got {other:?}"),
3850 },
3851 other => panic!("expected duplicate interval error, got {other:?}"),
3852 }
3853 assert_eq!(table.state(), &state_before);
3854 assert_eq!(table.log.load_current_version().await?, 1);
3855 assert!(data_files(temp.path())?.is_empty());
3856 assert!(coverage_files(temp.path())?.is_empty());
3857 assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
3858 assert_eq!(
3859 TimeSeriesTable::open(location).await?.state(),
3860 &state_before
3861 );
3862
3863 assert_eq!(
3864 table
3865 .append(timestamp_only_batch(1)?)
3866 .await?
3867 .committed_version,
3868 2
3869 );
3870 assert!(table.state().table_meta.logical_schema.is_some());
3871 Ok(())
3872 }
3873
3874 #[tokio::test]
3875 async fn duplicate_interval_across_input_batches_is_rejected() -> TestResult {
3876 let temp = TempDir::new()?;
3877 let location = TableLocation::local(temp.path());
3878 let mut table = TimeSeriesTable::create(location.clone(), timestamp_only_meta()).await?;
3879 let state_before = table.state().clone();
3880
3881 let error = table
3882 .append(
3883 AppendRequest::new(vec![
3884 timestamp_only_batch_with_millis([0])?,
3885 timestamp_only_batch_with_millis([30_000])?,
3886 ])
3887 .max_rows_per_row_group(1),
3888 )
3889 .await
3890 .expect_err("duplicate split across input batches and row groups must fail");
3891
3892 assert!(matches!(
3893 error,
3894 TableError::Append {
3895 source: AppendError::GeneratedSegmentCoverage { source, .. },
3896 } if matches!(source.as_ref(),
3897 SegmentCoverageError::DuplicateIndexInterval {
3898 example_identity: None,
3899 example_index_interval,
3900 ..
3901 } if example_index_interval.to_string()
3902 == "[1970-01-01T00:00:00Z, 1970-01-01T00:01:00Z)"
3903 )
3904 ));
3905 assert_eq!(table.state(), &state_before);
3906 assert!(data_files(temp.path())?.is_empty());
3907 assert!(coverage_files(temp.path())?.is_empty());
3908 assert_eq!(
3909 TimeSeriesTable::open(location).await?.state(),
3910 &state_before
3911 );
3912 Ok(())
3913 }
3914
3915 #[tokio::test]
3916 async fn composite_identity_uses_every_component_for_duplicates() -> TestResult {
3917 let temp = TempDir::new()?;
3918 let index = IndexSpec {
3919 column: "ts".to_string(),
3920 entity_columns: vec!["symbol".to_string(), "venue".to_string()],
3921 kind: IndexKind::Timestamp {
3922 index_granularity: TimeIndexGranularity::Minutes(1),
3923 timezone: None,
3924 },
3925 };
3926 let mut table = TimeSeriesTable::create(
3927 TableLocation::local(temp.path()),
3928 TableMeta::new_time_series(index),
3929 )
3930 .await?;
3931
3932 assert_eq!(
3933 table
3934 .append(composite_entity_batch(&[
3935 (0, "A", "X", 1.0),
3936 (0, "A", "Y", 2.0),
3937 ])?)
3938 .await?
3939 .committed_version,
3940 2
3941 );
3942 let state_before = table.state().clone();
3943
3944 let error = table
3945 .append(composite_entity_batch(&[
3946 (60_000, "A", "Z", 3.0),
3947 (90_000, "A", "Z", 4.0),
3948 ])?)
3949 .await
3950 .expect_err("matching composite identity must be rejected");
3951
3952 assert!(matches!(
3953 error,
3954 TableError::Append {
3955 source: AppendError::GeneratedSegmentCoverage { source, .. },
3956 } if matches!(source.as_ref(),
3957 SegmentCoverageError::DuplicateIndexInterval {
3958 example_identity: Some(example_identity),
3959 ..
3960 } if example_identity.components()
3961 == [EntityValue::from("A"), EntityValue::from("Z")]
3962 )
3963 ));
3964 assert_eq!(table.state(), &state_before);
3965 Ok(())
3966 }
3967
3968 #[tokio::test]
3969 async fn duplicate_entity_interval_leaves_nonempty_table_unchanged() -> TestResult {
3970 let temp = TempDir::new()?;
3971 let location = TableLocation::local(temp.path());
3972 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
3973 assert_eq!(
3974 table
3975 .append(time_series_batch(vec![0], vec!["A"], vec![1.0])?)
3976 .await?
3977 .committed_version,
3978 2
3979 );
3980 let state_before = table.state().clone();
3981 let data_before = data_files(temp.path())?;
3982 let coverage_before = coverage_files(temp.path())?;
3983
3984 let error = table
3985 .append(
3986 AppendRequest::new(time_series_batch(
3987 vec![60_000, 90_000],
3988 vec!["A", "A"],
3989 vec![2.0, 3.0],
3990 )?)
3991 .max_rows_per_row_group(1),
3992 )
3993 .await
3994 .expect_err("duplicate entity interval must fail");
3995 let expected_identity = EntityIdentity::try_new(vec!["A".into()])?;
3996
3997 assert!(matches!(
3998 error,
3999 TableError::Append {
4000 source: AppendError::GeneratedSegmentCoverage { source, .. },
4001 } if matches!(source.as_ref(),
4002 SegmentCoverageError::DuplicateIndexInterval {
4003 example_identity: Some(example_identity),
4004 example_index_interval,
4005 ..
4006 } if example_identity == &expected_identity
4007 && example_index_interval.to_string()
4008 == "[1970-01-01T00:01:00Z, 1970-01-01T00:02:00Z)"
4009 )
4010 ));
4011 assert_eq!(table.state(), &state_before);
4012 assert_eq!(table.log.load_current_version().await?, 2);
4013 assert_eq!(data_files(temp.path())?, data_before);
4014 assert_eq!(coverage_files(temp.path())?, coverage_before);
4015 assert!(!temp.path().join(layout::commit_rel_path(3)).exists());
4016 assert_eq!(
4017 TimeSeriesTable::open(location).await?.state(),
4018 &state_before
4019 );
4020 Ok(())
4021 }
4022
4023 #[tokio::test]
4024 async fn append_updates_state_and_log() -> TestResult {
4025 let tmp = TempDir::new()?;
4026 let location = TableLocation::local(tmp.path());
4027 let meta = make_basic_table_meta();
4028
4029 let mut table = TimeSeriesTable::create(location.clone(), meta).await?;
4030
4031 let rel_path = "data/seg1.parquet";
4032 let abs_path = tmp.path().join(rel_path);
4033 write_test_parquet(
4034 &abs_path,
4035 true,
4036 false,
4037 &[TestRow {
4038 ts_millis: 1_000,
4039 symbol: "A",
4040 price: 10.0,
4041 }],
4042 )?;
4043
4044 let new_version = append_parquet_fixture(&mut table, rel_path).await?;
4045
4046 assert_eq!(new_version, 2);
4047 assert_eq!(table.state.version, 2);
4048 let seg = table
4049 .state
4050 .segments
4051 .values()
4052 .next()
4053 .expect("segment present");
4054 assert_eq!(seg.row_count, 1);
4055 assert_eq!(
4056 seg.entity_layout,
4057 SegmentEntityLayout::Single(EntityIdentity::try_new(vec!["A".into()])?)
4058 );
4059 assert!(matches!(
4060 &seg.index_min,
4061 IndexValue::Timestamp(value) if value.timestamp_millis() == 1_000
4062 ));
4063 assert!(matches!(
4064 &seg.index_max,
4065 IndexValue::Timestamp(value) if value.timestamp_millis() == 1_000
4066 ));
4067
4068 let commit_path = tmp.path().join(layout::commit_rel_path(2));
4069 assert!(commit_path.is_file());
4070 let current =
4071 tokio::fs::read_to_string(tmp.path().join(layout::current_rel_path())).await?;
4072 assert_eq!(current.trim(), "2");
4073
4074 let reopened = TimeSeriesTable::open(location).await?;
4075 assert_eq!(reopened.state, table.state);
4076 Ok(())
4077 }
4078
4079 #[tokio::test]
4080 async fn append_without_entities_uses_global_coverage() -> TestResult {
4081 let tmp = TempDir::new()?;
4082 let location = TableLocation::local(tmp.path());
4083 let index = registered_index(IndexKind::Int64 {
4084 index_granularity: NonZeroU64::new(10).unwrap(),
4085 });
4086 let mut table =
4087 TimeSeriesTable::create(location.clone(), TableMeta::new_time_series(index.clone()))
4088 .await?;
4089 let rel_path = "data/int64.parquet";
4090 write_arrow_parquet_int_time(
4091 &tmp.path().join(rel_path),
4092 &[i64::MIN, -1, 0, i64::MAX],
4093 &["A", "A", "A", "A"],
4094 &[1.0, 2.0, 3.0, 4.0],
4095 )?;
4096
4097 assert_eq!(append_parquet_fixture(&mut table, rel_path).await?, 2);
4098
4099 let segment = table
4100 .state
4101 .segments
4102 .values()
4103 .next()
4104 .expect("segment present");
4105 assert_eq!(segment.entity_layout, SegmentEntityLayout::NotApplicable);
4106 assert_eq!(segment.index_min, IndexValue::Int64(i64::MIN));
4107 assert_eq!(segment.index_max, IndexValue::Int64(i64::MAX));
4108 let pointer = table.state.table_coverage.as_ref().expect("table coverage");
4109 assert_eq!(pointer.index_kind, index.kind);
4110 let persisted = read_coverage_sidecar(&location, Path::new(&pointer.coverage_path)).await?;
4111 let expected =
4112 compute_segment_coverage(&location, Path::new(&segment.path), table.index_spec())
4113 .await?;
4114 assert_eq!(persisted, expected);
4115 let reopened = TimeSeriesTable::open(location).await?;
4116 assert_eq!(reopened.state, table.state);
4117 Ok(())
4118 }
4119
4120 #[tokio::test]
4121 async fn int64_appends_enforce_coverage_and_exact_later_schema() -> TestResult {
4122 let tmp = TempDir::new()?;
4123 let location = TableLocation::local(tmp.path());
4124 let index = registered_index(IndexKind::Int64 {
4125 index_granularity: NonZeroU64::new(10).unwrap(),
4126 });
4127 let mut table =
4128 TimeSeriesTable::create(location, TableMeta::new_time_series(index)).await?;
4129
4130 for (path, values) in [
4131 ("data/negative.parquet", &[-25, -15][..]),
4132 ("data/positive.parquet", &[5, 15][..]),
4133 ] {
4134 write_arrow_parquet_int_time(&tmp.path().join(path), values, &["A", "A"], &[1.0, 2.0])?;
4135 append_parquet_fixture(&mut table, path).await?;
4136 }
4137 assert_eq!(table.state.version, 3);
4138
4139 let state_before = table.state.clone();
4140 let coverage_before = coverage_files(tmp.path())?;
4141 let overlap_path = "data/negative-overlap.parquet";
4142 write_arrow_parquet_int_time(&tmp.path().join(overlap_path), &[-19], &["A"], &[3.0])?;
4143 let overlap_error = append_parquet_fixture(&mut table, overlap_path)
4144 .await
4145 .expect_err("negative index interval overlap must fail");
4146 assert!(matches!(
4147 &overlap_error,
4148 TableError::Append {
4149 source: AppendError::PersistedIndexIntervalOverlap {
4150 example_identity: None,
4151 example_index_interval,
4152 ..
4153 }
4154 } if example_index_interval.to_string() == "[-20, -10)"
4155 ));
4156 assert!(
4157 overlap_error
4158 .to_string()
4159 .contains("example_index_interval=[-20, -10)")
4160 );
4161
4162 let mismatch_path = "data/schema-mismatch.parquet";
4163 write_single_index_parquet(
4164 &tmp.path().join(mismatch_path),
4165 DataType::Int64,
4166 Arc::new(Int64Array::from(vec![100])),
4167 )?;
4168 assert!(matches!(
4169 append_parquet_fixture(&mut table, mismatch_path)
4170 .await
4171 .expect_err("later schema mismatch must fail"),
4172 TableError::Append {
4173 source: AppendError::SchemaValidation { .. }
4174 }
4175 ));
4176 assert_eq!(table.state, state_before);
4177 assert_eq!(coverage_files(tmp.path())?, coverage_before);
4178 Ok(())
4179 }
4180
4181 #[tokio::test]
4182 async fn append_supports_registered_uint64_index() -> TestResult {
4183 let tmp = TempDir::new()?;
4184 let location = TableLocation::local(tmp.path());
4185 let index = registered_index(IndexKind::UInt64 {
4186 index_granularity: NonZeroU64::new(10).unwrap(),
4187 });
4188 let mut table =
4189 TimeSeriesTable::create(location.clone(), TableMeta::new_time_series(index.clone()))
4190 .await?;
4191 let rel_path = "data/uint64.parquet";
4192 write_single_index_parquet(
4193 &tmp.path().join(rel_path),
4194 DataType::UInt64,
4195 Arc::new(UInt64Array::from(vec![0, i64::MAX as u64 + 1, u64::MAX])),
4196 )?;
4197
4198 assert_eq!(append_parquet_fixture(&mut table, rel_path).await?, 2);
4199
4200 let segment = table
4201 .state
4202 .segments
4203 .values()
4204 .next()
4205 .expect("segment present");
4206 assert_eq!(segment.index_min, IndexValue::UInt64(0));
4207 assert_eq!(segment.index_max, IndexValue::UInt64(u64::MAX));
4208 assert_eq!(
4209 table
4210 .state
4211 .table_meta
4212 .logical_schema
4213 .as_ref()
4214 .expect("schema adopted")
4215 .columns()[0]
4216 .data_type,
4217 LogicalDataType::UInt64
4218 );
4219 assert_eq!(
4220 table
4221 .state
4222 .table_coverage
4223 .as_ref()
4224 .expect("table coverage")
4225 .index_kind,
4226 index.kind
4227 );
4228
4229 let non_overlap_path = "data/uint64-non-overlap.parquet";
4230 write_single_index_parquet(
4231 &tmp.path().join(non_overlap_path),
4232 DataType::UInt64,
4233 Arc::new(UInt64Array::from(vec![u64::MAX - 20])),
4234 )?;
4235 assert_eq!(
4236 append_parquet_fixture(&mut table, non_overlap_path).await?,
4237 3
4238 );
4239
4240 let state_before = table.state.clone();
4241 let coverage_before = coverage_files(tmp.path())?;
4242 let overlap_path = "data/uint64-overlap.parquet";
4243 write_single_index_parquet(
4244 &tmp.path().join(overlap_path),
4245 DataType::UInt64,
4246 Arc::new(UInt64Array::from(vec![u64::MAX - 1])),
4247 )?;
4248 let overlap_error = append_parquet_fixture(&mut table, overlap_path)
4249 .await
4250 .expect_err("large uint64 index interval overlap must fail");
4251 assert!(matches!(
4252 overlap_error,
4253 TableError::Append {
4254 source: AppendError::PersistedIndexIntervalOverlap {
4255 example_identity: None,
4256 example_index_interval,
4257 ..
4258 }
4259 } if example_index_interval.to_string()
4260 == "[18446744073709551610, 18446744073709551615]"
4261 ));
4262 assert_eq!(table.state, state_before);
4263 assert_eq!(coverage_files(tmp.path())?, coverage_before);
4264 let reopened = TimeSeriesTable::open(location).await?;
4265 assert_eq!(reopened.state, table.state);
4266 Ok(())
4267 }
4268
4269 #[tokio::test]
4270 async fn append_rejects_signed_data_for_uint64_index_without_mutation() -> TestResult {
4271 let tmp = TempDir::new()?;
4272 let location = TableLocation::local(tmp.path());
4273 let index = registered_index(IndexKind::UInt64 {
4274 index_granularity: NonZeroU64::new(1).unwrap(),
4275 });
4276 let mut table =
4277 TimeSeriesTable::create(location.clone(), TableMeta::new_time_series(index)).await?;
4278 let state_before = table.state.clone();
4279 let coverage_before = coverage_files(tmp.path())?;
4280 let rel_path = "data/signed.parquet";
4281 write_arrow_parquet_int_time(&tmp.path().join(rel_path), &[1], &["A"], &[1.0])?;
4282
4283 let error = append_parquet_fixture(&mut table, rel_path)
4284 .await
4285 .expect_err("signed data must not append to a uint64 index");
4286
4287 assert!(
4288 matches!(
4289 error,
4290 TableError::Append {
4291 source: AppendError::SchemaValidation {
4292 ref source,
4293 ..
4294 }
4295 } if matches!(
4296 source.as_ref(),
4297 crate::metadata::schema_compat::SchemaCompatibilityError::IndexKindMismatch {
4298 expected: "uint64",
4299 actual: LogicalDataType::Int64,
4300 ..
4301 }
4302 )
4303 ),
4304 "unexpected error: {error:?}"
4305 );
4306 assert_eq!(table.state, state_before);
4307 assert_eq!(table.log.load_current_version().await?, 1);
4308 assert_eq!(coverage_files(tmp.path())?, coverage_before);
4309 Ok(())
4310 }
4311
4312 #[tokio::test]
4313 async fn append_records_single_layout_for_each_entity_segment() -> TestResult {
4314 let tmp = TempDir::new()?;
4315 let location = TableLocation::local(tmp.path());
4316 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4317
4318 for (path, symbol) in [
4319 ("data/entity-a.parquet", "A"),
4320 ("data/entity-b.parquet", "B"),
4321 ] {
4322 write_test_parquet(
4323 &tmp.path().join(path),
4324 true,
4325 false,
4326 &[TestRow {
4327 ts_millis: 1_000,
4328 symbol,
4329 price: 10.0,
4330 }],
4331 )?;
4332 append_parquet_fixture(&mut table, path).await?;
4333 let expected_layout =
4334 SegmentEntityLayout::Single(EntityIdentity::try_new(vec![symbol.into()])?);
4335 assert!(
4336 table
4337 .state
4338 .segments
4339 .values()
4340 .any(|segment| segment.entity_layout == expected_layout)
4341 );
4342 }
4343
4344 assert_eq!(table.state.version, 3);
4345 let pointer = table
4346 .state
4347 .table_coverage
4348 .as_ref()
4349 .expect("table coverage pointer");
4350 let coverage =
4351 read_entity_coverage_sidecar(&location, Path::new(&pointer.coverage_path)).await?;
4352 assert_eq!(coverage.identity_count(), 2);
4353 assert_eq!(coverage.cardinality(), 2);
4354 Ok(())
4355 }
4356
4357 #[tokio::test]
4358 async fn append_records_mixed_layout_for_multiple_identities() -> TestResult {
4359 let tmp = TempDir::new()?;
4360 let location = TableLocation::local(tmp.path());
4361 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4362 let path = "data/multiple-identities.parquet";
4363 write_test_parquet(
4364 &tmp.path().join(path),
4365 true,
4366 false,
4367 &[
4368 TestRow {
4369 ts_millis: 1_000,
4370 symbol: "A",
4371 price: 10.0,
4372 },
4373 TestRow {
4374 ts_millis: 1_000,
4375 symbol: "B",
4376 price: 20.0,
4377 },
4378 ],
4379 )?;
4380
4381 append_parquet_fixture(&mut table, path).await?;
4382
4383 let segment = table
4384 .state
4385 .segments
4386 .values()
4387 .next()
4388 .expect("segment present");
4389 assert_eq!(segment.entity_layout, SegmentEntityLayout::Mixed);
4390 let coverage = read_entity_coverage_sidecar(
4391 &location,
4392 Path::new(segment.coverage_path.as_ref().expect("coverage path")),
4393 )
4394 .await?;
4395 assert_eq!(coverage.identity_count(), 2);
4396 assert_eq!(coverage.cardinality(), 2);
4397 Ok(())
4398 }
4399
4400 #[tokio::test]
4401 async fn numeric_entities_append_overlap_and_recover_with_exact_types() -> TestResult {
4402 let tmp = TempDir::new()?;
4403 let location = TableLocation::local(tmp.path());
4404 let mut table =
4405 TimeSeriesTable::create(location.clone(), make_int32_entity_table_meta()).await?;
4406
4407 let negative_path = "data/negative-device.parquet";
4408 write_int32_entity_parquet(
4409 &tmp.path().join(negative_path),
4410 &[1_000, 61_000],
4411 &[-1, -1],
4412 &[10.0, 11.0],
4413 )?;
4414 append_parquet_fixture(&mut table, negative_path).await?;
4415 let negative_identity = EntityIdentity::try_new(vec![EntityValue::Int32(-1)])?;
4416 let negative_layout = SegmentEntityLayout::Single(negative_identity.clone());
4417 assert!(
4418 table
4419 .state
4420 .segments
4421 .values()
4422 .any(|segment| segment.entity_layout == negative_layout)
4423 );
4424
4425 let maximum_path = "data/maximum-device.parquet";
4426 write_int32_entity_parquet(
4427 &tmp.path().join(maximum_path),
4428 &[1_000],
4429 &[i32::MAX],
4430 &[20.0],
4431 )?;
4432 append_parquet_fixture(&mut table, maximum_path).await?;
4433 let maximum_layout =
4434 SegmentEntityLayout::Single(EntityIdentity::try_new(vec![EntityValue::Int32(
4435 i32::MAX,
4436 )])?);
4437 assert!(
4438 table
4439 .state
4440 .segments
4441 .values()
4442 .any(|segment| segment.entity_layout == maximum_layout)
4443 );
4444
4445 let overlap_path = "data/negative-overlap.parquet";
4446 write_int32_entity_parquet(&tmp.path().join(overlap_path), &[1_500], &[-1], &[12.0])?;
4447 let error = append_parquet_fixture(&mut table, overlap_path)
4448 .await
4449 .expect_err("same typed identity and index interval must overlap");
4450 assert!(matches!(
4451 error,
4452 TableError::Append {
4453 source: AppendError::PersistedIndexIntervalOverlap {
4454 overlap_count: 1,
4455 example_identity: Some(example_identity),
4456 ..
4457 }
4458 } if example_identity == negative_identity
4459 ));
4460
4461 let snapshot = table
4462 .load_entity_coverage_with_recovery::<AppendError>()
4463 .await?;
4464 let reopened = TimeSeriesTable::open(location).await?;
4465 assert_eq!(reopened.state(), table.state());
4466 assert_eq!(
4467 reopened
4468 .recover_entity_coverage_from_segments::<AppendError>()
4469 .await?,
4470 snapshot
4471 );
4472 Ok(())
4473 }
4474
4475 #[tokio::test]
4476 async fn append_preserves_composite_identity_order_in_layout() -> TestResult {
4477 let tmp = TempDir::new()?;
4478 let location = TableLocation::local(tmp.path());
4479 let index = IndexSpec {
4480 column: "ts".to_string(),
4481 entity_columns: vec!["symbol".to_string(), "venue".to_string()],
4482 kind: IndexKind::Timestamp {
4483 index_granularity: TimeIndexGranularity::Minutes(1),
4484 timezone: None,
4485 },
4486 };
4487 let schema = LogicalSchema::new(vec![
4488 LogicalField {
4489 name: "ts".to_string(),
4490 data_type: LogicalDataType::Timestamp {
4491 unit: LogicalTimestampUnit::Millis,
4492 timezone: None,
4493 },
4494 nullable: false,
4495 },
4496 LogicalField {
4497 name: "symbol".to_string(),
4498 data_type: LogicalDataType::Utf8,
4499 nullable: false,
4500 },
4501 LogicalField {
4502 name: "venue".to_string(),
4503 data_type: LogicalDataType::Utf8,
4504 nullable: false,
4505 },
4506 LogicalField {
4507 name: "price".to_string(),
4508 data_type: LogicalDataType::Float64,
4509 nullable: false,
4510 },
4511 ])?;
4512 let mut table = TimeSeriesTable::create(
4513 location,
4514 TableMeta::new_time_series_with_schema(index, schema),
4515 )
4516 .await?;
4517
4518 for (path, venue) in [
4519 ("data/composite-x.parquet", "X"),
4520 ("data/composite-y.parquet", "Y"),
4521 ] {
4522 write_composite_entity_parquet(&tmp.path().join(path), &[(1_000, "A", venue, 10.0)])?;
4523 append_parquet_fixture(&mut table, path).await?;
4524 let expected_layout = SegmentEntityLayout::Single(EntityIdentity::try_new(vec![
4525 "A".into(),
4526 venue.into(),
4527 ])?);
4528 assert!(
4529 table
4530 .state
4531 .segments
4532 .values()
4533 .any(|segment| segment.entity_layout == expected_layout)
4534 );
4535 }
4536
4537 let overlap_path = "data/composite-x-overlap.parquet";
4538 write_composite_entity_parquet(&tmp.path().join(overlap_path), &[(1_500, "A", "X", 20.0)])?;
4539 let error = append_parquet_fixture(&mut table, overlap_path)
4540 .await
4541 .expect_err("matching composite identity and index interval must overlap");
4542 assert!(matches!(
4543 error,
4544 TableError::Append {
4545 source: AppendError::PersistedIndexIntervalOverlap {
4546 overlap_count: 1,
4547 example_identity: Some(example_identity),
4548 ..
4549 }
4550 } if example_identity.components()
4551 == [EntityValue::from("A"), EntityValue::from("X")]
4552 ));
4553 Ok(())
4554 }
4555
4556 #[tokio::test]
4557 async fn entity_with_only_null_index_values_is_rejected() -> TestResult {
4558 let tmp = TempDir::new()?;
4559 let location = TableLocation::local(tmp.path());
4560 let mut table = TimeSeriesTable::create(
4561 location,
4562 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
4563 )
4564 .await?;
4565 let state_before = table.state.clone();
4566 let path = "data/entity-without-index-coverage.parquet";
4567 write_arrow_parquet_with_unit(
4568 &tmp.path().join(path),
4569 ArrowTimeUnit::Millisecond,
4570 &[Some(1_000), None],
4571 &["A", "B"],
4572 &[10.0, 20.0],
4573 )?;
4574
4575 let error = append_parquet_fixture(&mut table, path)
4576 .await
4577 .expect_err("identity without index coverage must be rejected");
4578
4579 match error {
4580 TableError::Append {
4581 source: AppendError::EntityWithoutIndexCoverage { identity, .. },
4582 } => {
4583 assert_eq!(identity.components(), [EntityValue::from("B")]);
4584 }
4585 other => panic!("unexpected error: {other:?}"),
4586 }
4587 assert_eq!(table.state, state_before);
4588 assert!(coverage_files(tmp.path())?.is_empty());
4589 Ok(())
4590 }
4591
4592 #[tokio::test]
4593 async fn append_rejects_unsupported_entity_type_without_publication() -> TestResult {
4594 let tmp = TempDir::new()?;
4595 let location = TableLocation::local(tmp.path());
4596 let index = IndexSpec {
4597 column: "ts".to_string(),
4598 entity_columns: vec!["device_id".to_string()],
4599 kind: IndexKind::Timestamp {
4600 index_granularity: TimeIndexGranularity::Minutes(1),
4601 timezone: None,
4602 },
4603 };
4604 let mut table =
4605 TimeSeriesTable::create(location, TableMeta::new_time_series(index)).await?;
4606 let schema = Arc::new(Schema::new(vec![
4607 Field::new(
4608 "ts",
4609 DataType::Timestamp(ArrowTimeUnit::Millisecond, None),
4610 false,
4611 ),
4612 Field::new("device_id", DataType::Boolean, false),
4613 ]));
4614 let batch = RecordBatch::try_new(
4615 Arc::clone(&schema),
4616 vec![
4617 Arc::new(TimestampMillisecondArray::from(vec![1_000])),
4618 Arc::new(BooleanArray::from(vec![true])),
4619 ],
4620 )?;
4621 let state_before = table.state.clone();
4622
4623 let error = table
4624 .append(batch)
4625 .await
4626 .expect_err("Boolean entity columns must be rejected");
4627
4628 assert!(matches!(
4629 error,
4630 TableError::Append {
4631 source: AppendError::SchemaValidation {
4632 ref source,
4633 ..
4634 }
4635 } if matches!(
4636 source.as_ref(),
4637 crate::metadata::schema_compat::SchemaCompatibilityError::UnsupportedEntityColumnType {
4638 column,
4639 actual: LogicalDataType::Bool,
4640 } if column == "device_id"
4641 )
4642 ));
4643 assert_eq!(table.state, state_before);
4644 assert!(coverage_files(tmp.path())?.is_empty());
4645 assert_eq!(table.log.load_current_version().await?, 1);
4646 Ok(())
4647 }
4648
4649 #[tokio::test]
4650 async fn append_updates_snapshot() -> TestResult {
4651 let tmp = TempDir::new()?;
4652 let location = TableLocation::local(tmp.path());
4653 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4654
4655 let rel1 = "data/seg-auto-1.parquet";
4656 let rel2 = "data/seg-auto-2.parquet";
4657 let path1 = tmp.path().join(rel1);
4658 let path2 = tmp.path().join(rel2);
4659
4660 write_test_parquet(
4661 &path1,
4662 true,
4663 false,
4664 &[
4665 TestRow {
4666 ts_millis: 1_000,
4667 symbol: "A",
4668 price: 10.0,
4669 },
4670 TestRow {
4671 ts_millis: 61_000,
4672 symbol: "A",
4673 price: 20.0,
4674 },
4675 ],
4676 )?;
4677 write_test_parquet(
4678 &path2,
4679 true,
4680 false,
4681 &[
4682 TestRow {
4683 ts_millis: 120_000,
4684 symbol: "A",
4685 price: 30.0,
4686 },
4687 TestRow {
4688 ts_millis: 180_000,
4689 symbol: "A",
4690 price: 40.0,
4691 },
4692 ],
4693 )?;
4694
4695 let v2 = append_parquet_fixture(&mut table, rel1).await?;
4696 let v3 = append_parquet_fixture(&mut table, rel2).await?;
4697 assert_eq!(v2, 2);
4698 assert_eq!(v3, 3);
4699
4700 assert_eq!(table.state.segments.len(), 2);
4701 assert!(
4702 table
4703 .state
4704 .segments
4705 .values()
4706 .all(|segment| segment.coverage_path.is_some())
4707 );
4708 let expected_snapshot = table
4709 .recover_entity_coverage_from_segments::<AppendError>()
4710 .await?;
4711
4712 let ptr = table
4713 .state
4714 .table_coverage
4715 .as_ref()
4716 .expect("table snapshot pointer present after append");
4717 assert_eq!(ptr.version, v3);
4718 assert_eq!(ptr.index_kind, table.index_spec().kind);
4719
4720 let snapshot_cov =
4721 read_entity_coverage_sidecar(&location, Path::new(&ptr.coverage_path)).await?;
4722
4723 assert_eq!(snapshot_cov, expected_snapshot);
4724 Ok(())
4725 }
4726
4727 #[tokio::test]
4728 async fn append_rejects_overlap() -> TestResult {
4729 let tmp = TempDir::new()?;
4730 let location = TableLocation::local(tmp.path());
4731 let mut table = TimeSeriesTable::create(location, make_basic_table_meta()).await?;
4732
4733 let rel1 = "data/seg-overlap-a.parquet";
4734 let rel2 = "data/seg-overlap-b.parquet";
4735 let path1 = tmp.path().join(rel1);
4736 let path2 = tmp.path().join(rel2);
4737
4738 write_test_parquet(
4739 &path1,
4740 true,
4741 false,
4742 &[
4743 TestRow {
4744 ts_millis: 1_000,
4745 symbol: "B",
4746 price: 10.0,
4747 },
4748 TestRow {
4749 ts_millis: 1_000,
4750 symbol: "A",
4751 price: 20.0,
4752 },
4753 TestRow {
4754 ts_millis: 61_000,
4755 symbol: "A",
4756 price: 30.0,
4757 },
4758 ],
4759 )?;
4760 write_test_parquet(
4761 &path2,
4762 true,
4763 false,
4764 &[
4765 TestRow {
4766 ts_millis: 1_500,
4767 symbol: "B",
4768 price: 40.0,
4769 },
4770 TestRow {
4771 ts_millis: 1_500,
4772 symbol: "A",
4773 price: 50.0,
4774 },
4775 TestRow {
4776 ts_millis: 61_500,
4777 symbol: "A",
4778 price: 60.0,
4779 },
4780 ],
4781 )?;
4782
4783 append_parquet_fixture(&mut table, rel1).await?;
4784
4785 let err = append_parquet_fixture(&mut table, rel2)
4786 .await
4787 .expect_err("overlapping append should fail");
4788
4789 assert!(matches!(
4790 err,
4791 TableError::Append {
4792 source: AppendError::PersistedIndexIntervalOverlap {
4793 overlap_count: 3,
4794 example_identity: Some(example_identity),
4795 example_index_interval_id: 0x8000_0000_0000_0000,
4796 example_index_interval,
4797 ..
4798 }
4799 } if example_identity.components() == [EntityValue::from("A")]
4800 && example_index_interval.to_string()
4801 == "[1970-01-01T00:00:00Z, 1970-01-01T00:01:00Z)"
4802 ));
4803 Ok(())
4804 }
4805
4806 #[tokio::test]
4807 async fn append_snapshot_survives_reopen() -> TestResult {
4808 let tmp = TempDir::new()?;
4809 let location = TableLocation::local(tmp.path());
4810 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4811
4812 let rel1 = "data/seg-reopen-a.parquet";
4813 let rel2 = "data/seg-reopen-b.parquet";
4814 let path1 = tmp.path().join(rel1);
4815 let path2 = tmp.path().join(rel2);
4816
4817 write_test_parquet(
4818 &path1,
4819 true,
4820 false,
4821 &[TestRow {
4822 ts_millis: 1_000,
4823 symbol: "A",
4824 price: 10.0,
4825 }],
4826 )?;
4827 write_test_parquet(
4828 &path2,
4829 true,
4830 false,
4831 &[TestRow {
4832 ts_millis: 120_000,
4833 symbol: "A",
4834 price: 20.0,
4835 }],
4836 )?;
4837
4838 append_parquet_fixture(&mut table, rel1).await?;
4839 append_parquet_fixture(&mut table, rel2).await?;
4840
4841 let reopened = TimeSeriesTable::open(location.clone()).await?;
4842 let ptr = reopened
4843 .state()
4844 .table_coverage
4845 .as_ref()
4846 .expect("table snapshot pointer present after reopen");
4847
4848 assert_eq!(ptr.index_kind, reopened.index_spec().kind);
4849
4850 let expected = reopened
4851 .recover_entity_coverage_from_segments::<AppendError>()
4852 .await?;
4853
4854 let snapshot_cov =
4855 read_entity_coverage_sidecar(&location, Path::new(&ptr.coverage_path)).await?;
4856 assert_eq!(snapshot_cov, expected);
4857 Ok(())
4858 }
4859
4860 #[tokio::test]
4861 async fn load_snapshot_recovers_when_missing_file() -> TestResult {
4862 let tmp = TempDir::new()?;
4863 let location = TableLocation::local(tmp.path());
4864 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4865
4866 let rel1 = "data/seg-missing-snap-a.parquet";
4868 let rel2 = "data/seg-missing-snap-b.parquet";
4869 let path1 = tmp.path().join(rel1);
4870 let path2 = tmp.path().join(rel2);
4871 write_test_parquet(
4872 &path1,
4873 true,
4874 false,
4875 &[TestRow {
4876 ts_millis: 1_000,
4877 symbol: "A",
4878 price: 10.0,
4879 }],
4880 )?;
4881 write_test_parquet(
4882 &path2,
4883 true,
4884 false,
4885 &[TestRow {
4886 ts_millis: 120_000,
4887 symbol: "A",
4888 price: 20.0,
4889 }],
4890 )?;
4891
4892 append_parquet_fixture(&mut table, rel1).await?;
4893 append_parquet_fixture(&mut table, rel2).await?;
4894
4895 let state = table.state.clone();
4896 let ptr = state
4897 .table_coverage
4898 .as_ref()
4899 .expect("snapshot pointer present");
4900 let snapshot_abs = match &location.as_ref() {
4901 StorageLocation::Local(root) => root.join(&ptr.coverage_path),
4902 };
4903
4904 tokio::fs::remove_file(&snapshot_abs).await?;
4905
4906 let recovered = table
4907 .load_entity_coverage_with_recovery::<AppendError>()
4908 .await?;
4909
4910 let mut expected = EntityCoverage::empty();
4911 for seg in state.segments.values() {
4912 let cov_path = seg.coverage_path.as_ref().expect("coverage path");
4913 let cov = read_entity_coverage_sidecar(&location, Path::new(cov_path)).await?;
4914 expected.union_inplace(&cov);
4915 }
4916
4917 assert_eq!(recovered, expected);
4918 Ok(())
4919 }
4920
4921 #[tokio::test]
4922 async fn load_snapshot_recovers_when_corrupt_file() -> TestResult {
4923 let tmp = TempDir::new()?;
4924 let location = TableLocation::local(tmp.path());
4925 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4926
4927 let rel1 = "data/seg-corrupt-snap-a.parquet";
4928 let rel2 = "data/seg-corrupt-snap-b.parquet";
4929 let path1 = tmp.path().join(rel1);
4930 let path2 = tmp.path().join(rel2);
4931 write_test_parquet(
4932 &path1,
4933 true,
4934 false,
4935 &[TestRow {
4936 ts_millis: 1_000,
4937 symbol: "A",
4938 price: 10.0,
4939 }],
4940 )?;
4941 write_test_parquet(
4942 &path2,
4943 true,
4944 false,
4945 &[TestRow {
4946 ts_millis: 120_000,
4947 symbol: "A",
4948 price: 20.0,
4949 }],
4950 )?;
4951
4952 append_parquet_fixture(&mut table, rel1).await?;
4953 append_parquet_fixture(&mut table, rel2).await?;
4954
4955 let state = table.state.clone();
4956 let ptr = state
4957 .table_coverage
4958 .as_ref()
4959 .expect("snapshot pointer present");
4960 let snapshot_abs = match &location.as_ref() {
4961 StorageLocation::Local(root) => root.join(&ptr.coverage_path),
4962 };
4963
4964 tokio::fs::write(&snapshot_abs, b"garbage").await?;
4965
4966 let recovered = table
4967 .load_entity_coverage_with_recovery::<AppendError>()
4968 .await?;
4969
4970 let mut expected = EntityCoverage::empty();
4971 for seg in state.segments.values() {
4972 let cov_path = seg.coverage_path.as_ref().expect("coverage path");
4973 let cov = read_entity_coverage_sidecar(&location, Path::new(cov_path)).await?;
4974 expected.union_inplace(&cov);
4975 }
4976
4977 assert_eq!(recovered, expected);
4978 Ok(())
4979 }
4980
4981 #[tokio::test]
4982 async fn rejected_append_does_not_heal_corrupt_snapshot() -> TestResult {
4983 let tmp = TempDir::new()?;
4984 let location = TableLocation::local(tmp.path());
4985 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
4986
4987 let existing = "data/existing.parquet";
4988 write_test_parquet(
4989 &tmp.path().join(existing),
4990 true,
4991 false,
4992 &[TestRow {
4993 ts_millis: 1_000,
4994 symbol: "A",
4995 price: 10.0,
4996 }],
4997 )?;
4998 append_parquet_fixture(&mut table, existing).await?;
4999
5000 let snapshot_path = table
5001 .state
5002 .table_coverage
5003 .as_ref()
5004 .expect("snapshot pointer present")
5005 .coverage_path
5006 .clone();
5007 let snapshot_abs = tmp.path().join(snapshot_path);
5008 tokio::fs::write(&snapshot_abs, b"garbage").await?;
5009
5010 let overlapping = "data/overlapping.parquet";
5011 write_test_parquet(
5012 &tmp.path().join(overlapping),
5013 true,
5014 false,
5015 &[TestRow {
5016 ts_millis: 1_000,
5017 symbol: "A",
5018 price: 20.0,
5019 }],
5020 )?;
5021
5022 let err = append_parquet_fixture(&mut table, overlapping)
5023 .await
5024 .expect_err("overlap must be rejected");
5025 assert!(matches!(
5026 err,
5027 TableError::Append {
5028 source: AppendError::PersistedIndexIntervalOverlap { .. }
5029 }
5030 ));
5031 assert_eq!(tokio::fs::read(snapshot_abs).await?, b"garbage");
5032 Ok(())
5033 }
5034
5035 #[tokio::test]
5036 async fn load_snapshot_errors_when_segment_missing_coverage_path() -> TestResult {
5037 let tmp = TempDir::new()?;
5038 let location = TableLocation::local(tmp.path());
5039 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
5040
5041 let rel1 = "data/seg-missing-cov-path.parquet";
5042 let path1 = tmp.path().join(rel1);
5043 write_test_parquet(
5044 &path1,
5045 true,
5046 false,
5047 &[TestRow {
5048 ts_millis: 1_000,
5049 symbol: "A",
5050 price: 10.0,
5051 }],
5052 )?;
5053
5054 append_parquet_fixture(&mut table, rel1).await?;
5055
5056 let mut state = table.state.clone();
5057 state.table_coverage = None;
5058
5059 let segment_path = state
5060 .segments
5061 .keys()
5062 .next()
5063 .expect("segment present")
5064 .clone();
5065 state
5066 .segments
5067 .get_mut(&segment_path)
5068 .expect("segment present")
5069 .coverage_path = None;
5070
5071 table.state = state;
5073
5074 let err = table
5075 .load_entity_coverage_with_recovery::<AppendError>()
5076 .await
5077 .expect_err("missing coverage_path should error");
5078
5079 assert!(matches!(
5080 err,
5081 AppendError::ExistingSegmentMissingCoverageMetadata { .. }
5082 ));
5083 Ok(())
5084 }
5085
5086 #[tokio::test]
5087 async fn load_snapshot_errors_when_segment_sidecar_corrupt() -> TestResult {
5088 let tmp = TempDir::new()?;
5089 let location = TableLocation::local(tmp.path());
5090 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
5091
5092 let rel1 = "data/seg-corrupt-sidecar.parquet";
5093 let rel2 = "data/seg-corrupt-sidecar-ok.parquet";
5094 let path1 = tmp.path().join(rel1);
5095 let path2 = tmp.path().join(rel2);
5096 write_test_parquet(
5097 &path1,
5098 true,
5099 false,
5100 &[TestRow {
5101 ts_millis: 1_000,
5102 symbol: "A",
5103 price: 10.0,
5104 }],
5105 )?;
5106 write_test_parquet(
5107 &path2,
5108 true,
5109 false,
5110 &[TestRow {
5111 ts_millis: 120_000,
5112 symbol: "A",
5113 price: 20.0,
5114 }],
5115 )?;
5116
5117 append_parquet_fixture(&mut table, rel1).await?;
5118 append_parquet_fixture(&mut table, rel2).await?;
5119
5120 let mut state = table.state.clone();
5121 state.table_coverage = None;
5122 let (corrupt_segment_path, corrupt_cov_path) = state
5123 .segments
5124 .values()
5125 .next()
5126 .map(|meta| {
5127 (
5128 meta.path.clone(),
5129 meta.coverage_path.as_ref().expect("coverage path").clone(),
5130 )
5131 })
5132 .expect("at least one segment");
5133 table.state = state;
5134
5135 let corrupt_abs = match &location.as_ref() {
5136 StorageLocation::Local(root) => root.join(&corrupt_cov_path),
5137 };
5138 tokio::fs::write(&corrupt_abs, b"not a coverage bitmap").await?;
5139
5140 let err = table
5141 .load_entity_coverage_with_recovery::<AppendError>()
5142 .await
5143 .expect_err("corrupt sidecar should error");
5144
5145 match err {
5146 AppendError::ExistingSegmentCoverageSidecarRead {
5147 segment_path,
5148 coverage_path,
5149 ..
5150 } => {
5151 assert_eq!(segment_path, corrupt_segment_path);
5152 assert_eq!(coverage_path, corrupt_cov_path);
5153 }
5154 other => panic!("unexpected error: {other:?}"),
5155 }
5156
5157 Ok(())
5158 }
5159
5160 #[tokio::test]
5161 async fn entity_aware_stale_append_cleans_sidecars_without_state_mutation() -> TestResult {
5162 let tmp = TempDir::new()?;
5163 let location = TableLocation::local(tmp.path());
5164 let meta = make_basic_table_meta();
5165 let mut winner = TimeSeriesTable::create(location.clone(), meta).await?;
5166 let mut loser = TimeSeriesTable::open(location.clone()).await?;
5167 let loser_state_before = loser.state.clone();
5168
5169 let winner_path = "data/winner.parquet";
5170 let loser_path = "data/loser.parquet";
5171 write_test_parquet(
5172 &tmp.path().join(winner_path),
5173 true,
5174 false,
5175 &[TestRow {
5176 ts_millis: 10_000,
5177 symbol: "X",
5178 price: 100.0,
5179 }],
5180 )?;
5181 write_test_parquet(
5182 &tmp.path().join(loser_path),
5183 true,
5184 false,
5185 &[TestRow {
5186 ts_millis: 120_000,
5187 symbol: "X",
5188 price: 200.0,
5189 }],
5190 )?;
5191
5192 assert_eq!(append_parquet_fixture(&mut winner, winner_path).await?, 2);
5193 let coverage_before = coverage_files(tmp.path())?;
5194
5195 let err = append_parquet_fixture(&mut loser, loser_path)
5196 .await
5197 .expect_err("expected conflict due to stale version");
5198
5199 match err {
5200 TableError::Append {
5201 source: AppendError::Commit { source },
5202 } => {
5203 assert!(matches!(
5204 source,
5205 CommitError::Conflict {
5206 expected: 1,
5207 found: 2,
5208 ..
5209 }
5210 ));
5211 }
5212 other => panic!("unexpected error: {other:?}"),
5213 }
5214
5215 assert_eq!(loser.state, loser_state_before);
5216 assert_eq!(loser.log.load_current_version().await?, 2);
5217 let committed = loser.load_latest_state().await?;
5218 assert_eq!(committed, winner.state);
5219 assert_eq!(coverage_files(tmp.path())?, coverage_before);
5220 for bytes in coverage_before.values() {
5221 entity_coverage_from_bytes(bytes)?;
5222 }
5223 assert!(!tmp.path().join(layout::commit_rel_path(3)).exists());
5224 Ok(())
5225 }
5226
5227 #[tokio::test]
5228 async fn stale_int64_append_cleans_only_its_writer_owned_sidecars() -> TestResult {
5229 let tmp = TempDir::new()?;
5230 let location = TableLocation::local(tmp.path());
5231 let index = registered_index(IndexKind::Int64 {
5232 index_granularity: NonZeroU64::new(10).unwrap(),
5233 });
5234 let mut winner =
5235 TimeSeriesTable::create(location.clone(), TableMeta::new_time_series(index)).await?;
5236 let mut loser = TimeSeriesTable::open(location).await?;
5237 let winner_path = "data/writer-owned-winner.parquet";
5238 let loser_path = "data/writer-owned-loser.parquet";
5239
5240 write_arrow_parquet_int_time(&tmp.path().join(winner_path), &[0], &["X"], &[100.0])?;
5241 write_arrow_parquet_int_time(&tmp.path().join(loser_path), &[100], &["X"], &[200.0])?;
5242
5243 append_parquet_fixture(&mut winner, winner_path).await?;
5244 let coverage_before = coverage_files(tmp.path())?;
5245
5246 let err = append_parquet_fixture(&mut loser, loser_path)
5247 .await
5248 .expect_err("stale append should conflict");
5249
5250 assert!(matches!(
5251 err,
5252 TableError::Append {
5253 source: AppendError::Commit {
5254 source: CommitError::Conflict { .. }
5255 }
5256 }
5257 ));
5258 assert_eq!(coverage_files(tmp.path())?, coverage_before);
5259 Ok(())
5260 }
5261
5262 #[tokio::test]
5263 async fn ambiguous_int64_commit_retains_writer_owned_sidecars() -> TestResult {
5264 let tmp = TempDir::new()?;
5265 let location = TableLocation::local(tmp.path());
5266 let index = registered_index(IndexKind::Int64 {
5267 index_granularity: NonZeroU64::new(10).unwrap(),
5268 });
5269 let mut table =
5270 TimeSeriesTable::create(location, TableMeta::new_time_series(index)).await?;
5271 let state_before = table.state.clone();
5272 let coverage_before = coverage_files(tmp.path())?;
5273 let segment_path = "data/ambiguous.parquet";
5274
5275 write_arrow_parquet_int_time(&tmp.path().join(segment_path), &[10], &["X"], &[100.0])?;
5276
5277 let commit_path = tmp.path().join(layout::commit_rel_path(2));
5278 crate::storage::inject_write_new_failure(commit_path.clone(), true);
5279
5280 let err = append_parquet_fixture(&mut table, segment_path)
5281 .await
5282 .expect_err("failed commit cleanup should make the outcome ambiguous");
5283
5284 assert!(
5285 matches!(
5286 &err,
5287 TableError::Append {
5288 source: AppendError::CommitAmbiguous {
5289 source,
5290 ..
5291 }
5292 } if matches!(source.as_ref(), CommitError::AmbiguousOutcome { .. })
5293 ),
5294 "unexpected error: {err:?}"
5295 );
5296 assert_eq!(table.state, state_before);
5297 assert_eq!(table.log.load_current_version().await?, 1);
5298 assert!(commit_path.exists());
5299 assert_eq!(coverage_files(tmp.path())?.len(), coverage_before.len() + 2);
5300 Ok(())
5301 }
5302
5303 #[tokio::test]
5304 async fn entity_sidecar_cleanup_failures_preserve_error_and_reverse_order() -> TestResult {
5305 let tmp = TempDir::new()?;
5306 let location = TableLocation::local(tmp.path());
5307 let table = TimeSeriesTable::create(location, make_basic_table_meta()).await?;
5308 let sidecars = [
5309 format!("{SEGMENT_COVERAGE_DIR}/first-stuck.roar"),
5310 format!("{TABLE_SNAPSHOT_DIR}/second-stuck.roar"),
5311 ];
5312 for sidecar in &sidecars {
5313 tokio::fs::create_dir_all(tmp.path().join(sidecar)).await?;
5314 }
5315 let example_index_interval_id =
5316 crate::coverage::index_interval::index_interval_id_for_value(
5317 &table.index.kind,
5318 &IndexValue::Timestamp(utc_datetime(1970, 1, 1, 0, 0, 0)),
5319 )?;
5320 let source = AppendError::PersistedIndexIntervalOverlap {
5321 segment_path: "data/failed.parquet".to_string(),
5322 overlap_count: 1,
5323 example_identity: Some(EntityIdentity::try_new(vec!["A".into()])?),
5324 example_index_interval_id,
5325 example_index_interval: Box::new(index_interval_for_id(
5326 &table.index.kind,
5327 example_index_interval_id,
5328 )?),
5329 };
5330 let err = table.rollback_created_artifacts(&sidecars, source).await;
5331 let message = err.to_string();
5332
5333 assert!(matches!(
5334 err,
5335 AppendError::Rollback {
5336 source,
5337 cleanup_errors,
5338 } if matches!(
5339 *source,
5340 AppendError::PersistedIndexIntervalOverlap { .. }
5341 )
5342 && cleanup_errors.len() == 2
5343 && matches!(
5344 &cleanup_errors[0],
5345 StorageError::OtherIo { path, .. } if path.ends_with("second-stuck.roar")
5346 )
5347 && matches!(
5348 &cleanup_errors[1],
5349 StorageError::OtherIo { path, .. } if path.ends_with("first-stuck.roar")
5350 )
5351 ));
5352 assert!(message.contains("data/failed.parquet"));
5353 assert!(message.contains("first-stuck.roar"));
5354 assert!(message.contains("second-stuck.roar"));
5355 Ok(())
5356 }
5357
5358 #[tokio::test]
5359 async fn append_fails_when_existing_segment_missing_coverage_path() -> TestResult {
5360 let tmp = TempDir::new()?;
5361 let location = TableLocation::local(tmp.path());
5362 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
5363
5364 let rel1 = "data/seg-missing-cov.parquet";
5365 let rel2 = "data/seg-next.parquet";
5366 let path1 = tmp.path().join(rel1);
5367 let path2 = tmp.path().join(rel2);
5368
5369 write_test_parquet(
5370 &path1,
5371 true,
5372 false,
5373 &[TestRow {
5374 ts_millis: 1_000,
5375 symbol: "A",
5376 price: 10.0,
5377 }],
5378 )?;
5379 write_test_parquet(
5380 &path2,
5381 true,
5382 false,
5383 &[TestRow {
5384 ts_millis: 120_000,
5385 symbol: "A",
5386 price: 20.0,
5387 }],
5388 )?;
5389
5390 append_parquet_fixture(&mut table, rel1).await?;
5391
5392 let seg = table
5394 .state
5395 .segments
5396 .values_mut()
5397 .next()
5398 .expect("segment present");
5399 seg.coverage_path = None;
5400
5401 let err = append_parquet_fixture(&mut table, rel2)
5402 .await
5403 .expect_err("append should fail when existing segment lacks coverage");
5404
5405 assert!(matches!(
5406 err,
5407 TableError::Append {
5408 source: AppendError::ExistingSegmentMissingCoverageMetadata { .. }
5409 }
5410 ));
5411 Ok(())
5412 }
5413
5414 #[tokio::test]
5415 async fn append_recovers_when_table_snapshot_pointer_missing() -> TestResult {
5420 let tmp = TempDir::new()?;
5421 let location = TableLocation::local(tmp.path());
5422 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
5423
5424 let rel1 = "data/seg-no-pointer-a.parquet";
5425 let rel2 = "data/seg-no-pointer-b.parquet";
5426 let path1 = tmp.path().join(rel1);
5427 let path2 = tmp.path().join(rel2);
5428
5429 write_test_parquet(
5430 &path1,
5431 true,
5432 false,
5433 &[TestRow {
5434 ts_millis: 1_000,
5435 symbol: "A",
5436 price: 10.0,
5437 }],
5438 )?;
5439 write_test_parquet(
5440 &path2,
5441 true,
5442 false,
5443 &[TestRow {
5444 ts_millis: 120_000,
5445 symbol: "A",
5446 price: 20.0,
5447 }],
5448 )?;
5449
5450 append_parquet_fixture(&mut table, rel1).await?;
5451
5452 table.state.table_coverage = None;
5454
5455 append_parquet_fixture(&mut table, rel2).await?;
5456
5457 let ptr = table
5459 .state
5460 .table_coverage
5461 .as_ref()
5462 .expect("snapshot pointer restored");
5463
5464 let cov = read_entity_coverage_sidecar(&location, Path::new(&ptr.coverage_path)).await?;
5465
5466 let mut expected = EntityCoverage::empty();
5467 for seg in table.state.segments.values() {
5468 let path = seg.coverage_path.as_ref().expect("coverage path");
5469 let seg_cov = read_entity_coverage_sidecar(&location, Path::new(path)).await?;
5470 expected.union_inplace(&seg_cov);
5471 }
5472
5473 assert_eq!(cov, expected);
5474 Ok(())
5475 }
5476
5477 #[tokio::test]
5478 async fn generated_segment_metadata_failure_preserves_source_cleans_data_and_allows_retry()
5479 -> TestResult {
5480 let tmp = TempDir::new()?;
5481 let location = TableLocation::local(tmp.path());
5482 let mut table = TimeSeriesTable::create(location, timestamp_only_meta()).await?;
5483 let state_before = table.state().clone();
5484 let segment_path = "data/corrupt-generated.parquet";
5485 tokio::fs::create_dir_all(tmp.path().join("data")).await?;
5486 tokio::fs::write(tmp.path().join(segment_path), b"not parquet data").await?;
5487
5488 let mut data_guard = storage::FileCleanupGuard::new_disarmed(
5489 table.location().as_ref(),
5490 Path::new(segment_path),
5491 )?;
5492 data_guard.arm();
5493 let mut telemetry = AppendTelemetry::new(false);
5494 let source = table
5495 .publish_generated_parquet_segment(
5496 segment_path,
5497 2,
5498 &mut data_guard,
5499 &mut telemetry,
5500 AppendWriterSettings::default(),
5501 )
5502 .await
5503 .expect_err("corrupt generated Parquet must fail metadata loading");
5504 let source = table
5505 .rollback_created_artifacts(&[segment_path.to_string()], source)
5506 .await;
5507 data_guard.disarm();
5508 let error = crate::table::error::AppendSnafu.into_error(source);
5509
5510 assert!(matches!(
5511 &error,
5512 TableError::Append {
5513 source: AppendError::SegmentMetadata { source, .. }
5514 } if matches!(source.as_ref(), SegmentError::Metadata { .. })
5515 ));
5516 assert!(error.to_string().contains(segment_path));
5517 assert!(ErrorCompat::backtrace(&error).is_some());
5518 assert_eq!(table.state(), &state_before);
5519 assert_eq!(table.log.load_current_version().await?, 1);
5520 assert!(!tmp.path().join(segment_path).exists());
5521 assert!(data_files(tmp.path())?.is_empty());
5522 assert!(coverage_files(tmp.path())?.is_empty());
5523
5524 assert_eq!(
5525 table
5526 .append(timestamp_only_batch(1)?)
5527 .await?
5528 .committed_version,
5529 2
5530 );
5531 Ok(())
5532 }
5533
5534 #[tokio::test]
5535 async fn append_missing_or_corrupt_segment_sidecar_preserves_source_and_rolls_back()
5536 -> TestResult {
5537 #[derive(Debug, Clone, Copy)]
5538 enum Damage {
5539 Missing,
5540 Corrupt,
5541 }
5542
5543 for damage in [Damage::Missing, Damage::Corrupt] {
5544 let tmp = TempDir::new()?;
5545 let location = TableLocation::local(tmp.path());
5546 let mut table =
5547 TimeSeriesTable::create(location.clone(), timestamp_only_meta()).await?;
5548 assert_eq!(
5549 table
5550 .append(timestamp_only_batch(1)?)
5551 .await?
5552 .committed_version,
5553 2
5554 );
5555
5556 let (segment_path, coverage_path) = table
5557 .state()
5558 .segments
5559 .values()
5560 .next()
5561 .map(|segment| {
5562 (
5563 segment.path.clone(),
5564 segment.coverage_path.clone().expect("coverage path"),
5565 )
5566 })
5567 .expect("committed segment");
5568 let coverage_abs = tmp.path().join(&coverage_path);
5569 let original_coverage = tokio::fs::read(&coverage_abs).await?;
5570 match damage {
5571 Damage::Missing => tokio::fs::remove_file(&coverage_abs).await?,
5572 Damage::Corrupt => {
5573 tokio::fs::write(&coverage_abs, b"not a coverage sidecar").await?
5574 }
5575 }
5576 table.state.table_coverage = None;
5577
5578 let state_before = table.state().clone();
5579 let data_before = data_files(tmp.path())?;
5580 let coverage_before = coverage_files(tmp.path())?;
5581 let error = table
5582 .append(timestamp_only_batch_starting_at_minute_offset(2, 1)?)
5583 .await
5584 .expect_err("damaged segment sidecar must fail append recovery");
5585
5586 let sidecar_source = match &error {
5587 TableError::Append {
5588 source:
5589 AppendError::ExistingSegmentCoverageSidecarRead {
5590 segment_path: actual_segment_path,
5591 coverage_path: actual_coverage_path,
5592 source,
5593 },
5594 } => {
5595 assert_eq!(actual_segment_path, &segment_path);
5596 assert_eq!(actual_coverage_path, &coverage_path);
5597 source.as_ref()
5598 }
5599 other => panic!("unexpected {damage:?} error: {other:?}"),
5600 };
5601 match damage {
5602 Damage::Missing => assert!(matches!(
5603 sidecar_source,
5604 CoverageSidecarError::Storage {
5605 source: StorageError::NotFound { .. }
5606 }
5607 )),
5608 Damage::Corrupt => {
5609 assert!(matches!(sidecar_source, CoverageSidecarError::Codec { .. }))
5610 }
5611 }
5612 assert!(ErrorCompat::backtrace(&error).is_some());
5613 assert_eq!(table.state(), &state_before, "{damage:?}");
5614 assert_eq!(data_files(tmp.path())?, data_before, "{damage:?}");
5615 assert_eq!(coverage_files(tmp.path())?, coverage_before, "{damage:?}");
5616
5617 tokio::fs::write(&coverage_abs, &original_coverage).await?;
5618 assert_eq!(
5619 table
5620 .append(timestamp_only_batch_starting_at_minute_offset(2, 1)?)
5621 .await?
5622 .committed_version,
5623 3,
5624 "{damage:?}"
5625 );
5626 }
5627 Ok(())
5628 }
5629}