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