1use std::{collections::HashSet, path::Path};
4
5use snafu::{Backtrace, ResultExt, Snafu};
6
7use crate::{
8 coverage::{
9 EntityCoverage, EntityIdentity,
10 io::{CoverageSidecarError, read_entity_coverage_sidecar},
11 },
12 formats::parquet::{EntityRewriteError, StagedEntityRewrite, rewrite_mixed_parquet_segment},
13 metadata::{
14 index::IndexValueError, protocol::TableProtocolError,
15 schema_compat::SchemaCompatibilityError, segments::SegmentEntityLayout,
16 },
17 storage::{
18 StorageError, StorageLocation, ensure_canonical_relative_storage_path,
19 remove_file_if_exists,
20 },
21 table::{TableError, TimeSeriesTable},
22 transaction_log::{CommitError, LogAction, SegmentMeta, TableState},
23};
24
25#[derive(Debug, Snafu)]
27#[snafu(module, visibility(pub(crate)))]
28#[non_exhaustive]
29pub enum OptimizeError {
30 #[snafu(context(false), display("Table protocol error: {source}"))]
32 Protocol {
33 #[snafu(source)]
35 source: TableProtocolError,
36 backtrace: Backtrace,
38 },
39
40 #[snafu(display(
42 "Entity-layout optimization is not applicable to table {table_root}: no entity columns are configured"
43 ))]
44 NotApplicable {
45 table_root: String,
47 },
48
49 #[snafu(context(false), display("Mixed-segment rewrite failed: {source}"))]
51 MixedSegmentRewrite {
52 #[snafu(source(from(EntityRewriteError, Box::new)), backtrace)]
54 source: Box<EntityRewriteError>,
55 },
56
57 #[snafu(
59 context(false),
60 display("Invalid segment ordered-index bounds: {source}")
61 )]
62 InvalidSegmentBounds {
63 #[snafu(source)]
65 source: IndexValueError,
66 backtrace: Backtrace,
68 },
69
70 #[snafu(
72 context(false),
73 display("Optimization schema validation failed: {source}")
74 )]
75 SchemaValidation {
76 #[snafu(source(from(SchemaCompatibilityError, Box::new)), backtrace)]
78 source: Box<SchemaCompatibilityError>,
79 },
80
81 #[snafu(
83 context(false),
84 display("Optimization coverage sidecar error: {source}")
85 )]
86 CoverageSidecar {
87 #[snafu(source(from(CoverageSidecarError, Box::new)), backtrace)]
89 source: Box<CoverageSidecarError>,
90 },
91
92 #[snafu(display("Invalid {description} path {path:?}: {source}"))]
94 InvalidStagedPath {
95 description: &'static str,
97 path: String,
99 #[snafu(source(from(StorageError, Box::new)), backtrace)]
101 source: Box<StorageError>,
102 },
103
104 #[snafu(display("Invalid staged entity-layout optimization plan: {reason}"))]
106 InvalidStagedPlan {
107 reason: String,
109 backtrace: Backtrace,
111 },
112
113 #[snafu(display("Entity-layout optimization count overflow: {field}"))]
115 CountOverflow {
116 field: &'static str,
118 backtrace: Backtrace,
120 },
121
122 #[snafu(context(false), display("Optimization commit failed: {source}"))]
124 Commit {
125 #[snafu(source, backtrace)]
127 source: CommitError,
128 },
129
130 #[snafu(display(
132 "{source}; staged-object rollback also failed: [{}]",
133 cleanup_errors
134 .iter()
135 .map(ToString::to_string)
136 .collect::<Vec<_>>()
137 .join("; ")
138 ))]
139 Rollback {
140 #[snafu(source, backtrace)]
142 source: Box<OptimizeError>,
143 cleanup_errors: Vec<StorageError>,
145 },
146}
147
148#[derive(Debug, Clone, PartialEq, Eq)]
150pub struct OptimizeReport {
151 pub starting_version: u64,
153 pub committed_version: u64,
156 pub candidate_source_segments: u64,
158 pub source_segments_replaced: u64,
160 pub replacement_segments_written: u64,
162 pub distinct_identities_materialized: u64,
164 pub rows_read: u64,
166 pub rows_written: u64,
168 pub no_op: bool,
170}
171
172impl OptimizeReport {
173 fn no_op(version: u64) -> Self {
174 Self {
175 starting_version: version,
176 committed_version: version,
177 candidate_source_segments: 0,
178 source_segments_replaced: 0,
179 replacement_segments_written: 0,
180 distinct_identities_materialized: 0,
181 rows_read: 0,
182 rows_written: 0,
183 no_op: true,
184 }
185 }
186}
187
188struct PlanCounts {
189 candidates: u64,
190 replacements: u64,
191 identities: u64,
192 rows: u64,
193}
194
195fn mixed_segment_candidates(state: &TableState) -> Result<Vec<SegmentMeta>, IndexValueError> {
196 Ok(state
197 .segments_sorted_by_index()?
198 .into_iter()
199 .filter(|segment| segment.entity_layout == SegmentEntityLayout::Mixed)
200 .cloned()
201 .collect())
202}
203
204fn count(field: &'static str, value: usize) -> Result<u64, OptimizeError> {
205 value.try_into().map_err(|_| OptimizeError::CountOverflow {
206 field,
207 backtrace: Backtrace::capture(),
208 })
209}
210
211fn add(field: &'static str, total: &mut u64, value: u64) -> Result<(), OptimizeError> {
212 *total = total
213 .checked_add(value)
214 .ok_or_else(|| OptimizeError::CountOverflow {
215 field,
216 backtrace: Backtrace::capture(),
217 })?;
218 Ok(())
219}
220
221fn invalid_plan(reason: impl Into<String>) -> OptimizeError {
222 OptimizeError::InvalidStagedPlan {
223 reason: reason.into(),
224 backtrace: Backtrace::capture(),
225 }
226}
227
228async fn validate_staged_plan(
229 table: &TimeSeriesTable,
230 candidates: &[SegmentMeta],
231 staged: &[StagedEntityRewrite],
232) -> Result<PlanCounts, OptimizeError> {
233 if candidates.len() != staged.len() {
234 return Err(invalid_plan(format!(
235 "staged {} rewrites for {} candidates",
236 staged.len(),
237 candidates.len()
238 )));
239 }
240
241 let mut candidate_paths = HashSet::new();
242 let mut object_paths = HashSet::new();
243 let mut live_object_paths = HashSet::new();
244 for segment in table.state.segments.values() {
245 live_object_paths.insert(segment.path.as_str());
246 if let Some(path) = segment.coverage_path.as_deref() {
247 live_object_paths.insert(path);
248 }
249 }
250 if let Some(coverage) = &table.state.table_coverage {
251 live_object_paths.insert(coverage.coverage_path.as_str());
252 }
253 let mut source_coverage = EntityCoverage::empty();
254 let mut replacement_coverage = EntityCoverage::empty();
255 let mut identities = HashSet::<EntityIdentity>::new();
256 let mut source_rows = 0u64;
257 let mut replacement_rows = 0u64;
258 let mut replacement_count = 0u64;
259
260 for (candidate, rewrite) in candidates.iter().zip(staged) {
261 if rewrite.source_path != candidate.path {
262 return Err(invalid_plan(format!(
263 "candidate {} is not represented exactly once by its staged rewrite",
264 candidate.path
265 )));
266 }
267 if !candidate_paths.insert(&candidate.path) {
268 return Err(invalid_plan(format!(
269 "candidate path {} appears more than once",
270 candidate.path
271 )));
272 }
273 if table.state.segments.get(&candidate.path) != Some(candidate) {
274 return Err(invalid_plan(format!(
275 "candidate {} is not live in the starting snapshot",
276 candidate.path
277 )));
278 }
279 add("rows_read", &mut source_rows, candidate.row_count)?;
280
281 let source_coverage_path = candidate.coverage_path.as_deref().ok_or_else(|| {
282 invalid_plan(format!(
283 "candidate {} has no coverage sidecar",
284 candidate.path
285 ))
286 })?;
287 let candidate_coverage =
288 read_entity_coverage_sidecar(table.location(), Path::new(source_coverage_path))
289 .await
290 .map_err(OptimizeError::from)?;
291 let mut candidate_replacement_coverage = EntityCoverage::empty();
292 let mut candidate_rows = 0u64;
293
294 for replacement in &rewrite.replacements {
295 let coverage_path = replacement.meta.coverage_path.as_deref().ok_or_else(|| {
296 invalid_plan(format!(
297 "replacement {} has no coverage sidecar",
298 replacement.meta.path
299 ))
300 })?;
301 for (path, description) in [
302 (replacement.meta.path.as_str(), "replacement data"),
303 (coverage_path, "replacement coverage"),
304 ] {
305 ensure_canonical_relative_storage_path(path).map_err(|source| {
306 OptimizeError::InvalidStagedPath {
307 description,
308 path: path.to_string(),
309 source: Box::new(source),
310 }
311 })?;
312 if !object_paths.insert(path) {
313 return Err(invalid_plan(format!(
314 "staged object path {path} appears more than once"
315 )));
316 }
317 if live_object_paths.contains(path) {
318 return Err(invalid_plan(format!(
319 "staged object path {path} conflicts with a live table object"
320 )));
321 }
322 }
323 if replacement.meta.entity_layout
324 != SegmentEntityLayout::Single(replacement.identity.clone())
325 || replacement.coverage.identity_count() != 1
326 || replacement.coverage.get(&replacement.identity).is_none()
327 {
328 return Err(invalid_plan(format!(
329 "replacement {} is not truthful Single metadata",
330 replacement.meta.path
331 )));
332 }
333 add("replacement_segments_written", &mut replacement_count, 1)?;
334 add(
335 "rows_written",
336 &mut candidate_rows,
337 replacement.meta.row_count,
338 )?;
339 identities.insert(replacement.identity.clone());
340 candidate_replacement_coverage.union_inplace(&replacement.coverage);
341 }
342
343 if candidate_rows != candidate.row_count {
344 return Err(invalid_plan(format!(
345 "replacement rows {candidate_rows} do not equal source rows {} for {}",
346 candidate.row_count, candidate.path
347 )));
348 }
349 if candidate_replacement_coverage != candidate_coverage {
350 return Err(invalid_plan(format!(
351 "replacement coverage does not equal source coverage for {}",
352 candidate.path
353 )));
354 }
355 add("rows_written", &mut replacement_rows, candidate_rows)?;
356 source_coverage.union_inplace(&candidate_coverage);
357 replacement_coverage.union_inplace(&candidate_replacement_coverage);
358 }
359
360 if source_rows != replacement_rows {
361 return Err(invalid_plan(format!(
362 "replacement rows {replacement_rows} do not equal source rows {source_rows}"
363 )));
364 }
365 if source_coverage != replacement_coverage {
366 return Err(invalid_plan(
367 "replacement coverage does not reconstruct selected source coverage",
368 ));
369 }
370
371 Ok(PlanCounts {
372 candidates: count("candidate_source_segments", candidates.len())?,
373 replacements: replacement_count,
374 identities: count("distinct_identities_materialized", identities.len())?,
375 rows: source_rows,
376 })
377}
378
379impl TimeSeriesTable {
380 async fn rollback_optimization(
381 &self,
382 staged_paths: &[String],
383 source: OptimizeError,
384 ) -> OptimizeError {
385 let mut cleanup_errors = Vec::new();
386 for path in staged_paths.iter().rev() {
387 if let Err(error) =
388 remove_file_if_exists(self.location().as_ref(), Path::new(path)).await
389 {
390 cleanup_errors.push(error);
391 }
392 }
393 if cleanup_errors.is_empty() {
394 source
395 } else {
396 OptimizeError::Rollback {
397 source: Box::new(source),
398 cleanup_errors,
399 }
400 }
401 }
402
403 #[tracing::instrument(
414 name = "table.optimize",
415 target = "timeseries_table_format::table::optimize",
416 level = "debug",
417 skip_all,
418 fields(
419 starting_version = self.state.version,
420 candidate_source_segments = tracing::field::Empty,
421 replacement_segments_written = tracing::field::Empty,
422 distinct_identities_materialized = tracing::field::Empty,
423 rows_read = tracing::field::Empty,
424 rows_written = tracing::field::Empty,
425 committed_version = tracing::field::Empty,
426 no_op = tracing::field::Empty,
427 outcome = tracing::field::Empty
428 )
429 )]
430 pub async fn optimize(&mut self) -> Result<OptimizeReport, TableError> {
431 let result: Result<OptimizeReport, OptimizeError> = async {
432 self.ensure_write_compatible()
433 .map_err(OptimizeError::from)?;
434
435 if self.index.entity_columns.is_empty() {
436 let table_root = match self.location().as_ref() {
437 StorageLocation::Local(root) => root.display().to_string(),
438 };
439 return Err(OptimizeError::NotApplicable { table_root });
440 }
441
442 let starting_version = self.state.version;
443 let candidates = mixed_segment_candidates(&self.state).map_err(|source| {
444 OptimizeError::InvalidSegmentBounds {
445 source,
446 backtrace: Backtrace::capture(),
447 }
448 })?;
449 tracing::Span::current().record("candidate_source_segments", candidates.len());
450 if candidates.is_empty() {
451 return Ok(OptimizeReport::no_op(starting_version));
452 }
453 let committed_version =
454 starting_version
455 .checked_add(1)
456 .ok_or_else(|| OptimizeError::CountOverflow {
457 field: "committed_version",
458 backtrace: Backtrace::capture(),
459 })?;
460 let table_schema = self
461 .state
462 .table_meta
463 .logical_schema
464 .clone()
465 .ok_or_else(|| OptimizeError::from(SchemaCompatibilityError::MissingTableSchema))?;
466
467 let mut staged = Vec::with_capacity(candidates.len());
468 let mut staged_paths = Vec::new();
469 for source in &candidates {
470 match rewrite_mixed_parquet_segment(
471 self.location(),
472 &table_schema,
473 &self.index,
474 source,
475 )
476 .await
477 {
478 Ok(rewrite) => {
479 staged_paths.extend(rewrite.staged_object_paths.iter().cloned());
480 staged.push(rewrite);
481 }
482 Err(source) => {
483 let error = OptimizeError::from(source);
484 return Err(self.rollback_optimization(&staged_paths, error).await);
485 }
486 }
487 }
488
489 let counts = match validate_staged_plan(self, &candidates, &staged).await {
490 Ok(counts) => counts,
491 Err(source) => {
492 return Err(self.rollback_optimization(&staged_paths, source).await);
493 }
494 };
495 let span = tracing::Span::current();
496 span.record("candidate_source_segments", counts.candidates);
497 span.record("replacement_segments_written", counts.replacements);
498 span.record("distinct_identities_materialized", counts.identities);
499 span.record("rows_read", counts.rows);
500 span.record("rows_written", counts.rows);
501
502 let mut actions = Vec::new();
503 for source in &candidates {
504 actions.push(LogAction::RemoveSegment {
505 path: source.path.clone(),
506 });
507 }
508 for rewrite in &staged {
509 actions.extend(
510 rewrite
511 .replacements
512 .iter()
513 .map(|replacement| LogAction::AddSegment(replacement.meta.clone())),
514 );
515 }
516
517 let new_version = match self
518 .log
519 .commit_with_expected_version(starting_version, actions)
520 .await
521 {
522 Ok(version) => version,
523 Err(source @ CommitError::AmbiguousOutcome { .. }) => {
524 return Err(OptimizeError::from(source));
525 }
526 Err(source) => {
527 let error = OptimizeError::from(source);
528 return Err(self.rollback_optimization(&staged_paths, error).await);
529 }
530 };
531 assert_eq!(
532 new_version, committed_version,
533 "transaction log returned unexpected optimize version"
534 );
535
536 for source in &candidates {
537 self.state.segments.remove(&source.path);
538 }
539 for rewrite in staged {
540 for replacement in rewrite.replacements {
541 self.state
542 .segments
543 .insert(replacement.meta.path.clone(), replacement.meta);
544 }
545 }
546 self.state.version = new_version;
547
548 Ok(OptimizeReport {
549 starting_version,
550 committed_version: new_version,
551 candidate_source_segments: counts.candidates,
552 source_segments_replaced: counts.candidates,
553 replacement_segments_written: counts.replacements,
554 distinct_identities_materialized: counts.identities,
555 rows_read: counts.rows,
556 rows_written: counts.rows,
557 no_op: false,
558 })
559 }
560 .await;
561
562 let span = tracing::Span::current();
563 match &result {
564 Ok(report) => {
565 span.record(
566 "candidate_source_segments",
567 report.candidate_source_segments,
568 );
569 span.record(
570 "replacement_segments_written",
571 report.replacement_segments_written,
572 );
573 span.record(
574 "distinct_identities_materialized",
575 report.distinct_identities_materialized,
576 );
577 span.record("rows_read", report.rows_read);
578 span.record("rows_written", report.rows_written);
579 span.record("committed_version", report.committed_version);
580 span.record("no_op", report.no_op);
581 if report.no_op {
582 span.record("outcome", "no_op");
583 } else {
584 span.record("outcome", "succeeded");
585 tracing::info!(
586 name: "table.optimize",
587 target: "timeseries_table_format::table::optimize",
588 starting_version = report.starting_version,
589 committed_version = report.committed_version,
590 candidate_source_segments = report.candidate_source_segments,
591 replacement_segments_written = report.replacement_segments_written,
592 distinct_identities_materialized = report.distinct_identities_materialized,
593 rows_read = report.rows_read,
594 rows_written = report.rows_written,
595 outcome = "succeeded",
596 "Optimized time-series table"
597 );
598 }
599 }
600 Err(OptimizeError::Commit {
601 source: CommitError::AmbiguousOutcome { .. },
602 }) => {
603 span.record("outcome", "ambiguous");
604 }
605 Err(OptimizeError::Rollback { .. }) => {
606 span.record("outcome", "cleanup_failed");
607 }
608 Err(_) => {
609 span.record("outcome", "failed");
610 }
611 }
612 result.context(crate::table::error::OptimizeSnafu)
613 }
614}
615
616#[cfg(test)]
617mod tests {
618 use std::error::Error as _;
619
620 use super::*;
621 use crate::table::AppendError;
622 use std::{
623 collections::{BTreeSet, HashMap},
624 path::{Path, PathBuf},
625 };
626
627 use arrow::datatypes::TimeUnit;
628 use futures::StreamExt;
629 use snafu::ErrorCompat;
630 use tempfile::TempDir;
631
632 use crate::{
633 coverage::{
634 EntityIdentity, EntityValue,
635 io::CoverageSidecarError,
636 layout::{SEGMENT_COVERAGE_DIR, TABLE_SNAPSHOT_DIR},
637 },
638 formats::parquet::EntityRewriteError,
639 metadata::{
640 index::IndexValue,
641 logical_schema::LogicalTimestampUnit,
642 protocol::TableProtocolError,
643 segments::{FileFormat, SegmentEntityLayout},
644 },
645 storage::{StorageError, TableLocation, layout},
646 table::test_util::{
647 CapturedSpan, TraceCapture, append_parquet_fixture, make_int32_entity_table_meta,
648 make_table_meta_with_unit, utc_datetime, write_arrow_parquet_with_unit,
649 write_int32_entity_parquet,
650 },
651 transaction_log::TableKind,
652 };
653
654 fn segment(path: &str, layout: SegmentEntityLayout, minute: u32) -> SegmentMeta {
655 SegmentMeta {
656 path: path.to_string(),
657 format: FileFormat::Parquet,
658 entity_layout: layout,
659 index_min: IndexValue::Timestamp(utc_datetime(2025, 1, 1, 0, minute, 0)),
660 index_max: IndexValue::Timestamp(utc_datetime(2025, 1, 1, 0, minute + 1, 0)),
661 row_count: 1,
662 file_size: Some(1),
663 coverage_path: Some(format!("_coverage/segments/{minute}.roar")),
664 }
665 }
666
667 fn state(segments: impl IntoIterator<Item = SegmentMeta>) -> TableState {
668 TableState {
669 version: 7,
670 table_meta: make_table_meta_with_unit(LogicalTimestampUnit::Millis),
671 segments: segments
672 .into_iter()
673 .map(|segment| (segment.path.clone(), segment))
674 .collect::<HashMap<_, _>>(),
675 table_coverage: None,
676 }
677 }
678
679 fn single_identity() -> SegmentEntityLayout {
680 SegmentEntityLayout::Single(
681 EntityIdentity::try_new(vec!["A".into()]).expect("valid identity"),
682 )
683 }
684
685 fn files_below(root: &Path) -> std::io::Result<BTreeSet<PathBuf>> {
686 if !root.exists() {
687 return Ok(BTreeSet::new());
688 }
689 let mut files = BTreeSet::new();
690 let mut directories = vec![root.to_owned()];
691 while let Some(directory) = directories.pop() {
692 for entry in std::fs::read_dir(directory)? {
693 let entry = entry?;
694 if entry.file_type()?.is_dir() {
695 directories.push(entry.path());
696 } else {
697 files.insert(entry.path());
698 }
699 }
700 }
701 Ok(files)
702 }
703
704 fn optimization_objects(root: &Path) -> std::io::Result<BTreeSet<PathBuf>> {
705 let mut files = BTreeSet::new();
706 for relative in ["data/_staged", SEGMENT_COVERAGE_DIR] {
707 files.extend(files_below(&root.join(relative))?);
708 }
709 Ok(files)
710 }
711
712 fn captured_optimize_span(capture: &TraceCapture) -> CapturedSpan {
713 let mut spans: Vec<_> = capture
714 .spans()
715 .into_iter()
716 .filter(|span| span.name == "table.optimize")
717 .collect();
718 assert_eq!(spans.len(), 1, "expected one table.optimize span");
719 spans.pop().expect("captured optimize span")
720 }
721
722 fn assert_no_optimize_event(capture: &TraceCapture) {
723 assert!(
724 !capture
725 .events()
726 .iter()
727 .any(|event| event.name == "table.optimize")
728 );
729 }
730
731 async fn append_mixed_source(
732 table: &mut TimeSeriesTable,
733 root: &Path,
734 path: &str,
735 start_millis: i64,
736 ) -> Result<String, TableError> {
737 write_arrow_parquet_with_unit(
738 &root.join(path),
739 TimeUnit::Millisecond,
740 &[
741 Some(start_millis + 1_000),
742 Some(start_millis + 2_000),
743 Some(start_millis + 61_000),
744 Some(start_millis + 62_000),
745 ],
746 &["A", "B", "A", "B"],
747 &[10.0, 20.0, 11.0, 21.0],
748 )
749 .expect("write mixed source");
750 let existing_paths = table
751 .state()
752 .segments
753 .keys()
754 .cloned()
755 .collect::<BTreeSet<_>>();
756 append_parquet_fixture(table, path).await?;
757 let committed_path = table
758 .state()
759 .segments
760 .keys()
761 .find(|path| !existing_paths.contains(*path))
762 .expect("append added a segment")
763 .clone();
764 assert_eq!(
765 table.state().segments[&committed_path].entity_layout,
766 SegmentEntityLayout::Mixed
767 );
768 Ok(committed_path)
769 }
770
771 #[test]
772 fn discovery_selects_all_and_only_mixed_segments() {
773 let state = state([
774 segment("data/mixed-late.parquet", SegmentEntityLayout::Mixed, 2),
775 segment("data/single.parquet", single_identity(), 0),
776 segment(
777 "data/not-applicable.parquet",
778 SegmentEntityLayout::NotApplicable,
779 0,
780 ),
781 segment("data/mixed.parquet", SegmentEntityLayout::Mixed, 1),
782 ]);
783
784 let paths = mixed_segment_candidates(&state)
785 .expect("valid segment bounds")
786 .into_iter()
787 .map(|segment| segment.path)
788 .collect::<Vec<_>>();
789
790 assert_eq!(paths, ["data/mixed.parquet", "data/mixed-late.parquet"]);
791 }
792
793 #[test]
794 fn discovery_order_is_independent_of_hash_map_insertion_order() {
795 let segments = [
796 segment("data/b.parquet", SegmentEntityLayout::Mixed, 1),
797 segment("data/a.parquet", SegmentEntityLayout::Mixed, 1),
798 segment("data/later.parquet", SegmentEntityLayout::Mixed, 2),
799 ];
800 let forward = state(segments.clone());
801 let reverse = state(segments.into_iter().rev());
802
803 let paths = |state: &TableState| {
804 mixed_segment_candidates(state)
805 .expect("valid segment bounds")
806 .into_iter()
807 .map(|segment| segment.path)
808 .collect::<Vec<_>>()
809 };
810
811 assert_eq!(paths(&forward), paths(&reverse));
812 assert_eq!(
813 paths(&forward),
814 ["data/a.parquet", "data/b.parquet", "data/later.parquet"]
815 );
816 }
817
818 #[test]
819 fn invalid_staged_path_preserves_storage_source_and_backtrace() {
820 let path = "../outside.parquet";
821 let source =
822 ensure_canonical_relative_storage_path(path).expect_err("parent traversal must fail");
823 let error = OptimizeError::InvalidStagedPath {
824 description: "replacement data",
825 path: path.to_string(),
826 source: Box::new(source),
827 };
828 let storage = error
829 .source()
830 .and_then(|source| source.downcast_ref::<Box<StorageError>>())
831 .map(Box::as_ref)
832 .expect("storage source");
833
834 assert!(matches!(error, OptimizeError::InvalidStagedPath { .. }));
835 assert!(std::ptr::eq(
836 ErrorCompat::backtrace(&error).expect("optimization backtrace"),
837 ErrorCompat::backtrace(storage).expect("storage backtrace")
838 ));
839 }
840
841 #[tokio::test]
842 async fn optimize_without_mixed_segments_is_a_zero_write_no_op() -> Result<(), TableError> {
843 let temp = TempDir::new().expect("temp directory");
844 let mut table = TimeSeriesTable::create(
845 TableLocation::local(temp.path()),
846 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
847 )
848 .await?;
849 let starting_version = table.state().version;
850
851 let capture = TraceCapture::default();
852 let report = capture.run(table.optimize()).await?;
853
854 assert_eq!(report, OptimizeReport::no_op(starting_version));
855 let span = captured_optimize_span(&capture);
856 assert_eq!(span.level, tracing::Level::DEBUG);
857 for (field, expected) in [
858 ("starting_version", starting_version.to_string()),
859 ("candidate_source_segments", "0".to_string()),
860 ("replacement_segments_written", "0".to_string()),
861 ("distinct_identities_materialized", "0".to_string()),
862 ("rows_read", "0".to_string()),
863 ("rows_written", "0".to_string()),
864 ("committed_version", starting_version.to_string()),
865 ("no_op", "true".to_string()),
866 ("outcome", "no_op".to_string()),
867 ] {
868 assert_eq!(span.fields.get(field), Some(&expected));
869 }
870 assert_no_optimize_event(&capture);
871 assert!(!temp.path().join("data/_staged").exists());
872 let reopened = TimeSeriesTable::open(TableLocation::local(temp.path())).await?;
873 assert_eq!(reopened.state().version, starting_version);
874 Ok(())
875 }
876
877 #[tokio::test]
878 async fn optimize_rejects_unsupported_writer_features_before_staging() -> Result<(), TableError>
879 {
880 let temp = TempDir::new().expect("temp directory");
881 let mut table = TimeSeriesTable::create(
882 TableLocation::local(temp.path()),
883 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
884 )
885 .await?;
886 table
887 .state
888 .table_meta
889 .required_writer_features
890 .insert("future_writer".to_string());
891 let state_before = table.state().clone();
892 let objects_before = optimization_objects(temp.path()).expect("optimization objects");
893
894 let error = table
895 .optimize()
896 .await
897 .expect_err("unsupported writer feature must reject optimize");
898
899 assert!(matches!(
900 error,
901 TableError::Optimize {
902 source: OptimizeError::Protocol {
903 source: TableProtocolError::UnsupportedWriterFeatures { features },
904 ..
905 }
906 } if features == ["future_writer"]
907 ));
908 assert_eq!(table.state(), &state_before);
909 assert_eq!(table.current_version().await?, 1);
910 assert_eq!(
911 optimization_objects(temp.path()).expect("optimization objects"),
912 objects_before
913 );
914 assert!(!temp.path().join(layout::commit_rel_path(2)).exists());
915 Ok(())
916 }
917
918 #[tokio::test]
919 async fn optimize_rejects_a_table_without_entity_columns() -> Result<(), TableError> {
920 let temp = TempDir::new().expect("temp directory");
921 let mut table_meta = make_table_meta_with_unit(LogicalTimestampUnit::Millis);
922 let TableKind::TimeSeries(index) = &mut table_meta.kind else {
923 unreachable!("test table is time-series");
924 };
925 index.entity_columns.clear();
926 let mut table =
927 TimeSeriesTable::create(TableLocation::local(temp.path()), table_meta).await?;
928
929 let capture = TraceCapture::default();
930 let error = capture
931 .run(table.optimize())
932 .await
933 .expect_err("entity-free table must be rejected");
934
935 assert!(matches!(
936 error,
937 TableError::Optimize {
938 source: OptimizeError::NotApplicable { table_root }
939 }
940 if table_root == temp.path().display().to_string()
941 ));
942 let span = captured_optimize_span(&capture);
943 assert_eq!(
944 span.fields.get("outcome").map(String::as_str),
945 Some("failed")
946 );
947 assert!(
948 span.fields
949 .values()
950 .all(|value| !value.contains(&temp.path().display().to_string()))
951 );
952 assert_no_optimize_event(&capture);
953 assert!(!temp.path().join("data/_staged").exists());
954 Ok(())
955 }
956
957 #[tokio::test]
958 async fn optimize_rejects_missing_canonical_schema_before_staging() -> Result<(), TableError> {
959 let temp = TempDir::new().expect("temp directory");
960 let mut table = TimeSeriesTable::create(
961 TableLocation::local(temp.path()),
962 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
963 )
964 .await?;
965 append_mixed_source(&mut table, temp.path(), "data/mixed.parquet", 0).await?;
966 table.state.table_meta.logical_schema = None;
967 let state_before = table.state().clone();
968 let objects_before = optimization_objects(temp.path()).expect("optimization objects");
969
970 let error = table
971 .optimize()
972 .await
973 .expect_err("missing canonical schema must fail");
974
975 assert!(matches!(
976 error,
977 TableError::Optimize {
978 source: OptimizeError::SchemaValidation { source, .. }
979 } if matches!(*source, SchemaCompatibilityError::MissingTableSchema)
980 ));
981 assert_eq!(table.state(), &state_before);
982 assert_eq!(
983 optimization_objects(temp.path()).expect("optimization objects"),
984 objects_before
985 );
986 Ok(())
987 }
988
989 #[tokio::test]
990 async fn optimize_rejects_invalid_segment_bounds_before_staging() -> Result<(), TableError> {
991 let temp = TempDir::new().expect("temp directory");
992 let mut table = TimeSeriesTable::create(
993 TableLocation::local(temp.path()),
994 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
995 )
996 .await?;
997 let source_path =
998 append_mixed_source(&mut table, temp.path(), "data/mixed.parquet", 0).await?;
999 table
1000 .state
1001 .segments
1002 .get_mut(&source_path)
1003 .expect("mixed source")
1004 .index_max = IndexValue::Int64(1);
1005 let state_before = table.state().clone();
1006 let objects_before = optimization_objects(temp.path()).expect("optimization objects");
1007
1008 let error = table
1009 .optimize()
1010 .await
1011 .expect_err("mixed ordered-index domains must fail");
1012
1013 assert!(matches!(
1014 error,
1015 TableError::Optimize {
1016 source: OptimizeError::InvalidSegmentBounds {
1017 source: IndexValueError::DomainMismatch { .. },
1018 ..
1019 }
1020 }
1021 ));
1022 assert_eq!(table.state(), &state_before);
1023 assert_eq!(
1024 optimization_objects(temp.path()).expect("optimization objects"),
1025 objects_before
1026 );
1027 Ok(())
1028 }
1029
1030 #[tokio::test]
1031 async fn optimize_preserves_a_missing_source_coverage_error() -> Result<(), TableError> {
1032 let temp = TempDir::new().expect("temp directory");
1033 let mut table = TimeSeriesTable::create(
1034 TableLocation::local(temp.path()),
1035 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
1036 )
1037 .await?;
1038 let source_path =
1039 append_mixed_source(&mut table, temp.path(), "data/mixed.parquet", 0).await?;
1040 let coverage_path = table.state().segments[&source_path]
1041 .coverage_path
1042 .as_deref()
1043 .expect("source coverage")
1044 .to_string();
1045 std::fs::remove_file(temp.path().join(&coverage_path)).expect("remove source coverage");
1046 let state_before = table.state().clone();
1047 let objects_before = optimization_objects(temp.path()).expect("optimization objects");
1048
1049 let error = table
1050 .optimize()
1051 .await
1052 .expect_err("missing source coverage must fail");
1053
1054 assert!(matches!(
1055 error,
1056 TableError::Optimize {
1057 source: OptimizeError::MixedSegmentRewrite { source }
1058 } if matches!(
1059 *source,
1060 EntityRewriteError::CoverageSidecar {
1061 source: CoverageSidecarError::Storage {
1062 source: StorageError::NotFound { .. }
1063 }
1064 }
1065 )
1066 ));
1067 assert_eq!(table.state(), &state_before);
1068 assert_eq!(
1069 optimization_objects(temp.path()).expect("optimization objects"),
1070 objects_before
1071 );
1072 Ok(())
1073 }
1074
1075 #[tokio::test]
1076 async fn optimize_atomically_replaces_one_mixed_source() -> Result<(), TableError> {
1077 let temp = TempDir::new().expect("temp directory");
1078 let location = TableLocation::local(temp.path());
1079 let mut table = TimeSeriesTable::create(
1080 location.clone(),
1081 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
1082 )
1083 .await?;
1084 let fixture_path = "data/mixed.parquet";
1085 write_arrow_parquet_with_unit(
1086 &temp.path().join(fixture_path),
1087 TimeUnit::Millisecond,
1088 &[Some(1_000), Some(2_000), Some(61_000), Some(62_000)],
1089 &["A", "B", "A", "B"],
1090 &[10.0, 20.0, 11.0, 21.0],
1091 )
1092 .expect("write mixed source");
1093 append_parquet_fixture(&mut table, fixture_path).await?;
1094 let source_path = table
1095 .state()
1096 .segments
1097 .keys()
1098 .next()
1099 .expect("committed source")
1100 .clone();
1101 let source = table
1102 .state()
1103 .segments
1104 .get(&source_path)
1105 .expect("committed source")
1106 .clone();
1107 assert_eq!(source.entity_layout, SegmentEntityLayout::Mixed);
1108 let source_bytes = std::fs::read(temp.path().join(&source_path)).expect("source bytes");
1109 let source_coverage_path = source.coverage_path.as_deref().expect("source coverage");
1110 let source_coverage_bytes =
1111 std::fs::read(temp.path().join(source_coverage_path)).expect("source coverage bytes");
1112 let table_coverage = table.state().table_coverage.clone();
1113 let starting_version = table.state().version;
1114
1115 let capture = TraceCapture::default();
1116 let report = capture.run(table.optimize()).await?;
1117
1118 assert_eq!(report.starting_version, starting_version);
1119 assert_eq!(report.committed_version, starting_version + 1);
1120 assert_eq!(report.candidate_source_segments, 1);
1121 assert_eq!(report.source_segments_replaced, 1);
1122 assert_eq!(report.replacement_segments_written, 2);
1123 assert_eq!(report.distinct_identities_materialized, 2);
1124 assert_eq!(report.rows_read, 4);
1125 assert_eq!(report.rows_written, 4);
1126 assert!(!report.no_op);
1127 let span = captured_optimize_span(&capture);
1128 assert_eq!(span.target, "timeseries_table_format::table::optimize");
1129 assert_eq!(span.level, tracing::Level::DEBUG);
1130 for (field, expected) in [
1131 ("starting_version", report.starting_version.to_string()),
1132 (
1133 "candidate_source_segments",
1134 report.candidate_source_segments.to_string(),
1135 ),
1136 (
1137 "replacement_segments_written",
1138 report.replacement_segments_written.to_string(),
1139 ),
1140 (
1141 "distinct_identities_materialized",
1142 report.distinct_identities_materialized.to_string(),
1143 ),
1144 ("rows_read", report.rows_read.to_string()),
1145 ("rows_written", report.rows_written.to_string()),
1146 ("committed_version", report.committed_version.to_string()),
1147 ("no_op", "false".to_string()),
1148 ("outcome", "succeeded".to_string()),
1149 ] {
1150 assert_eq!(span.fields.get(field), Some(&expected));
1151 }
1152 let events: Vec<_> = capture
1153 .events()
1154 .into_iter()
1155 .filter(|event| event.name == "table.optimize")
1156 .collect();
1157 assert_eq!(events.len(), 1, "expected one table.optimize event");
1158 assert_eq!(events[0].target, "timeseries_table_format::table::optimize");
1159 assert_eq!(events[0].level, tracing::Level::INFO);
1160 for (field, expected) in [
1161 ("starting_version", report.starting_version.to_string()),
1162 ("committed_version", report.committed_version.to_string()),
1163 (
1164 "candidate_source_segments",
1165 report.candidate_source_segments.to_string(),
1166 ),
1167 (
1168 "replacement_segments_written",
1169 report.replacement_segments_written.to_string(),
1170 ),
1171 (
1172 "distinct_identities_materialized",
1173 report.distinct_identities_materialized.to_string(),
1174 ),
1175 ("rows_read", report.rows_read.to_string()),
1176 ("rows_written", report.rows_written.to_string()),
1177 ("outcome", "succeeded".to_string()),
1178 ] {
1179 assert_eq!(events[0].fields.get(field), Some(&expected));
1180 }
1181 assert!(
1182 events[0]
1183 .fields
1184 .get("message")
1185 .is_some_and(|message| message.contains("Optimized time-series table"))
1186 );
1187 for value in span.fields.into_values().chain(
1188 events
1189 .into_iter()
1190 .flat_map(|event| event.fields.into_values()),
1191 ) {
1192 assert!(!value.contains(&temp.path().display().to_string()));
1193 assert!(!value.contains("EntityIdentity"));
1194 assert!(!value.contains("LogicalSchema"));
1195 assert_ne!(value, "A");
1196 assert_ne!(value, "B");
1197 }
1198 assert!(!table.state().segments.contains_key(&source_path));
1199 assert_eq!(table.state().segments.len(), 2);
1200 assert!(
1201 table
1202 .state()
1203 .segments
1204 .values()
1205 .all(|segment| matches!(segment.entity_layout, SegmentEntityLayout::Single(_)))
1206 );
1207 assert_eq!(table.state().table_coverage, table_coverage);
1208 assert_eq!(
1209 std::fs::read(temp.path().join(&source_path)).expect("source remains"),
1210 source_bytes
1211 );
1212 assert_eq!(
1213 std::fs::read(temp.path().join(source_coverage_path)).expect("source coverage remains"),
1214 source_coverage_bytes
1215 );
1216 Ok(())
1217 }
1218
1219 #[tokio::test]
1220 async fn optimize_rewrites_numeric_entities_with_typed_single_layouts() -> Result<(), TableError>
1221 {
1222 let temp = TempDir::new().expect("temp directory");
1223 let location = TableLocation::local(temp.path());
1224 let mut table =
1225 TimeSeriesTable::create(location.clone(), make_int32_entity_table_meta()).await?;
1226 let fixture_path = "data/numeric-mixed.parquet";
1227 write_int32_entity_parquet(
1228 &temp.path().join(fixture_path),
1229 &[1_000, 2_000, 61_000, 62_000],
1230 &[-1, i32::MAX, -1, i32::MAX],
1231 &[10.0, 20.0, 11.0, 21.0],
1232 )
1233 .expect("write numeric mixed source");
1234 append_parquet_fixture(&mut table, fixture_path).await?;
1235 let source_path = table
1236 .state()
1237 .segments
1238 .keys()
1239 .next()
1240 .expect("committed source");
1241 assert_eq!(
1242 table.state().segments[source_path].entity_layout,
1243 SegmentEntityLayout::Mixed
1244 );
1245 let table_meta = table.state().table_meta.clone();
1246
1247 let report = table.optimize().await?;
1248
1249 assert_eq!(report.replacement_segments_written, 2);
1250 assert_eq!(report.rows_written, 4);
1251 assert_eq!(table.state().table_meta, table_meta);
1252 let identities = table
1253 .state()
1254 .segments
1255 .values()
1256 .map(|segment| match &segment.entity_layout {
1257 SegmentEntityLayout::Single(identity) => identity.clone(),
1258 layout => panic!("expected typed single-entity layout, found {layout:?}"),
1259 })
1260 .collect::<BTreeSet<_>>();
1261 assert_eq!(
1262 identities,
1263 BTreeSet::from([
1264 EntityIdentity::try_new(vec![EntityValue::Int32(-1)]).expect("negative identity"),
1265 EntityIdentity::try_new(vec![EntityValue::Int32(i32::MAX)])
1266 .expect("maximum identity"),
1267 ])
1268 );
1269
1270 let reopened = TimeSeriesTable::open(location).await?;
1271 assert_eq!(reopened.state(), table.state());
1272 Ok(())
1273 }
1274
1275 #[tokio::test]
1276 async fn later_staging_failure_cleans_every_earlier_rewrite() -> Result<(), TableError> {
1277 let temp = TempDir::new().expect("temp directory");
1278 let mut table = TimeSeriesTable::create(
1279 TableLocation::local(temp.path()),
1280 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
1281 )
1282 .await?;
1283 append_mixed_source(&mut table, temp.path(), "data/first.parquet", 0).await?;
1284 let broken_path =
1285 append_mixed_source(&mut table, temp.path(), "data/broken.parquet", 120_000).await?;
1286 let state_before = table.state().clone();
1287 let objects_before = optimization_objects(temp.path()).expect("optimization objects");
1288 std::fs::remove_file(temp.path().join(broken_path)).expect("remove later source");
1289
1290 let error = table.optimize().await.expect_err("later staging must fail");
1291
1292 assert!(matches!(
1293 error,
1294 TableError::Optimize {
1295 source: OptimizeError::MixedSegmentRewrite { .. }
1296 }
1297 ));
1298 assert_eq!(table.state(), &state_before);
1299 assert_eq!(
1300 optimization_objects(temp.path()).expect("optimization objects"),
1301 objects_before
1302 );
1303 Ok(())
1304 }
1305
1306 #[tokio::test]
1307 async fn occ_conflict_cleans_staged_objects_and_preserves_state() -> Result<(), TableError> {
1308 let temp = TempDir::new().expect("temp directory");
1309 let location = TableLocation::local(temp.path());
1310 let mut table = TimeSeriesTable::create(
1311 location.clone(),
1312 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
1313 )
1314 .await?;
1315 append_mixed_source(&mut table, temp.path(), "data/candidate.parquet", 0).await?;
1316 let state_before = table.state().clone();
1317 let mut concurrent = TimeSeriesTable::open(location).await?;
1318 append_mixed_source(
1319 &mut concurrent,
1320 temp.path(),
1321 "data/concurrent.parquet",
1322 120_000,
1323 )
1324 .await?;
1325 let objects_before = optimization_objects(temp.path()).expect("optimization objects");
1326
1327 let error = table.optimize().await.expect_err("stale commit must fail");
1328
1329 assert!(matches!(
1330 error,
1331 TableError::Optimize {
1332 source: OptimizeError::Commit {
1333 source: CommitError::Conflict { .. }
1334 }
1335 }
1336 ));
1337 assert_eq!(table.state(), &state_before);
1338 assert_eq!(
1339 table
1340 .log
1341 .load_current_version()
1342 .await
1343 .expect("current version"),
1344 state_before.version + 1
1345 );
1346 assert_eq!(
1347 optimization_objects(temp.path()).expect("optimization objects"),
1348 objects_before
1349 );
1350 Ok(())
1351 }
1352
1353 #[tokio::test]
1354 async fn ambiguous_commit_retains_staged_objects_and_preserves_state() -> Result<(), TableError>
1355 {
1356 let temp = TempDir::new().expect("temp directory");
1357 let mut table = TimeSeriesTable::create(
1358 TableLocation::local(temp.path()),
1359 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
1360 )
1361 .await?;
1362 append_mixed_source(&mut table, temp.path(), "data/candidate.parquet", 0).await?;
1363 let state_before = table.state().clone();
1364 let objects_before = optimization_objects(temp.path()).expect("optimization objects");
1365 let commit_path = temp
1366 .path()
1367 .join(layout::commit_rel_path(state_before.version + 1));
1368 crate::storage::inject_write_new_failure(commit_path.clone(), true);
1369
1370 let capture = TraceCapture::default();
1371 let error = capture
1372 .run(table.optimize())
1373 .await
1374 .expect_err("commit outcome must be ambiguous");
1375
1376 assert!(matches!(
1377 error,
1378 TableError::Optimize {
1379 source: OptimizeError::Commit {
1380 source: CommitError::AmbiguousOutcome { .. }
1381 }
1382 }
1383 ));
1384 let span = captured_optimize_span(&capture);
1385 assert_eq!(
1386 span.fields.get("outcome").map(String::as_str),
1387 Some("ambiguous")
1388 );
1389 assert_no_optimize_event(&capture);
1390 assert_eq!(table.state(), &state_before);
1391 assert_eq!(
1392 table
1393 .log
1394 .load_current_version()
1395 .await
1396 .expect("current version"),
1397 state_before.version
1398 );
1399 assert!(commit_path.exists());
1400 let objects_after = optimization_objects(temp.path()).expect("optimization objects");
1401 assert_eq!(objects_after.difference(&objects_before).count(), 4);
1402 assert_eq!(
1403 files_below(&temp.path().join("data/_staged"))
1404 .expect("staged data")
1405 .len(),
1406 2
1407 );
1408 Ok(())
1409 }
1410
1411 #[tokio::test]
1412 async fn rollback_reports_every_cleanup_failure_in_reverse_order() -> Result<(), TableError> {
1413 let temp = TempDir::new().expect("temp directory");
1414 let table = TimeSeriesTable::create(
1415 TableLocation::local(temp.path()),
1416 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
1417 )
1418 .await?;
1419 let paths = [
1420 "data/_staged/entity-rewrite/first.parquet".to_string(),
1421 format!("{SEGMENT_COVERAGE_DIR}/second.roar"),
1422 ];
1423 for path in &paths {
1424 let absolute = temp.path().join(path);
1425 std::fs::create_dir_all(absolute.parent().expect("object parent"))
1426 .expect("create object parent");
1427 std::fs::write(&absolute, b"staged").expect("write staged object");
1428 crate::storage::inject_cleanup_failure(absolute);
1429 }
1430
1431 let error = table
1432 .rollback_optimization(&paths, invalid_plan("primary failure"))
1433 .await;
1434 let message = error.to_string();
1435
1436 assert!(matches!(
1437 error,
1438 OptimizeError::Rollback {
1439 source,
1440 cleanup_errors,
1441 } if matches!(*source, OptimizeError::InvalidStagedPlan { .. })
1442 && cleanup_errors.len() == 2
1443 && cleanup_errors[0].to_string().contains("second.roar")
1444 && cleanup_errors[1].to_string().contains("first.parquet")
1445 ));
1446 assert!(message.contains("primary failure"));
1447 assert!(paths.iter().all(|path| temp.path().join(path).exists()));
1448 Ok(())
1449 }
1450
1451 #[tokio::test]
1452 async fn multiple_sources_reopen_recover_and_repeat_as_a_no_op() -> Result<(), TableError> {
1453 let temp = TempDir::new().expect("temp directory");
1454 let location = TableLocation::local(temp.path());
1455 let mut table = TimeSeriesTable::create(
1456 location.clone(),
1457 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
1458 )
1459 .await?;
1460 let source_paths = [
1461 append_mixed_source(&mut table, temp.path(), "data/first.parquet", 0).await?,
1462 append_mixed_source(&mut table, temp.path(), "data/second.parquet", 120_000).await?,
1463 ];
1464 let sources = source_paths
1465 .iter()
1466 .map(|path| table.state().segments[path].clone())
1467 .collect::<Vec<_>>();
1468 let expected_coverage = table
1469 .load_entity_coverage_with_recovery::<AppendError>()
1470 .await
1471 .context(crate::table::error::AppendSnafu)?;
1472 let coverage_pointer = table
1473 .state()
1474 .table_coverage
1475 .clone()
1476 .expect("table coverage pointer");
1477 let coverage_bytes = std::fs::read(temp.path().join(&coverage_pointer.coverage_path))
1478 .expect("table coverage bytes");
1479 let snapshot_files =
1480 files_below(&temp.path().join(TABLE_SNAPSHOT_DIR)).expect("snapshot files");
1481 let table_meta = table.state().table_meta.clone();
1482 let starting_version = table.state().version;
1483
1484 let report = table.optimize().await?;
1485
1486 assert_eq!(
1487 report,
1488 OptimizeReport {
1489 starting_version,
1490 committed_version: starting_version + 1,
1491 candidate_source_segments: 2,
1492 source_segments_replaced: 2,
1493 replacement_segments_written: 4,
1494 distinct_identities_materialized: 2,
1495 rows_read: 8,
1496 rows_written: 8,
1497 no_op: false,
1498 }
1499 );
1500 assert_eq!(table.state().segments.len(), 4);
1501 assert_eq!(table.state().table_meta, table_meta);
1502 assert!(
1503 table
1504 .state()
1505 .segments
1506 .values()
1507 .all(|segment| matches!(segment.entity_layout, SegmentEntityLayout::Single(_)))
1508 );
1509 assert_eq!(table.state().table_coverage, Some(coverage_pointer.clone()));
1510 assert_eq!(
1511 std::fs::read(temp.path().join(&coverage_pointer.coverage_path))
1512 .expect("table coverage bytes"),
1513 coverage_bytes
1514 );
1515 assert_eq!(
1516 files_below(&temp.path().join(TABLE_SNAPSHOT_DIR)).expect("snapshot files"),
1517 snapshot_files
1518 );
1519 for source in &sources {
1520 assert!(temp.path().join(&source.path).exists());
1521 assert!(
1522 temp.path()
1523 .join(source.coverage_path.as_deref().expect("source coverage"))
1524 .exists()
1525 );
1526 }
1527
1528 let commit = table
1529 .log
1530 .load_commit(report.committed_version)
1531 .await
1532 .expect("optimization commit");
1533 assert_eq!(commit.base_version, starting_version);
1534 assert_eq!(commit.actions.len(), 6);
1535 assert!(
1536 commit.actions[..2]
1537 .iter()
1538 .zip(&source_paths)
1539 .all(|(action, expected)| matches!(
1540 action,
1541 LogAction::RemoveSegment { path } if path == expected
1542 ))
1543 );
1544 assert!(
1545 commit.actions[2..]
1546 .iter()
1547 .all(|action| matches!(action, LogAction::AddSegment(_)))
1548 );
1549
1550 let state_after_first = table.state().clone();
1551 let objects_after_first = optimization_objects(temp.path()).expect("optimization objects");
1552 let second_report = table.optimize().await?;
1553 assert_eq!(
1554 second_report,
1555 OptimizeReport::no_op(report.committed_version)
1556 );
1557 assert_eq!(table.state(), &state_after_first);
1558 assert_eq!(
1559 optimization_objects(temp.path()).expect("optimization objects"),
1560 objects_after_first
1561 );
1562 assert_eq!(
1563 table
1564 .log
1565 .load_current_version()
1566 .await
1567 .expect("current version"),
1568 report.committed_version
1569 );
1570
1571 let reopened = TimeSeriesTable::open(location).await?;
1572 assert_eq!(reopened.state(), table.state());
1573 assert_eq!(
1574 reopened
1575 .recover_entity_coverage_from_segments::<AppendError>()
1576 .await
1577 .context(crate::table::error::AppendSnafu)?,
1578 expected_coverage
1579 );
1580 let mut scan = reopened
1581 .scan_range(
1582 chrono::DateTime::from_timestamp_millis(0).expect("range start"),
1583 chrono::DateTime::from_timestamp_millis(240_000).expect("range end"),
1584 )
1585 .await?;
1586 let mut rows = 0;
1587 while let Some(batch) = scan.next().await {
1588 rows += batch?.num_rows();
1589 }
1590 assert_eq!(rows, 8);
1591 Ok(())
1592 }
1593
1594 #[test]
1595 fn accumulated_report_counts_do_not_wrap() {
1596 let mut total = u64::MAX;
1597
1598 let error = add("rows_written", &mut total, 1).expect_err("count overflow must fail");
1599
1600 assert!(matches!(
1601 error,
1602 OptimizeError::CountOverflow {
1603 field: "rows_written",
1604 ..
1605 }
1606 ));
1607 assert_eq!(total, u64::MAX);
1608 }
1609
1610 #[tokio::test]
1611 async fn version_overflow_fails_before_staging() -> Result<(), TableError> {
1612 let temp = TempDir::new().expect("temp directory");
1613 let mut table = TimeSeriesTable::create(
1614 TableLocation::local(temp.path()),
1615 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
1616 )
1617 .await?;
1618 append_mixed_source(&mut table, temp.path(), "data/candidate.parquet", 0).await?;
1619 table.state.version = u64::MAX;
1620 let objects_before = optimization_objects(temp.path()).expect("optimization objects");
1621
1622 let error = table
1623 .optimize()
1624 .await
1625 .expect_err("version overflow must fail");
1626
1627 assert!(matches!(
1628 error,
1629 TableError::Optimize {
1630 source: OptimizeError::CountOverflow {
1631 field: "committed_version",
1632 ..
1633 }
1634 }
1635 ));
1636 assert_eq!(table.state().version, u64::MAX);
1637 assert_eq!(
1638 optimization_objects(temp.path()).expect("optimization objects"),
1639 objects_before
1640 );
1641 Ok(())
1642 }
1643}