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