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