1use std::{collections::HashSet, path::Path};
4
5use crate::{
6 coverage::{EntityCoverage, EntityIdentity, io::read_entity_coverage_sidecar},
7 formats::parquet::{StagedEntityRewrite, rewrite_mixed_parquet_segment},
8 metadata::{segments::SegmentEntityLayout, table_metadata::IndexValueError},
9 storage::{StorageLocation, normalize_relative_storage_path, remove_file_if_exists},
10 table::{TableError, TimeSeriesTable},
11 transaction_log::{CommitError, LogAction, SegmentMeta, TableState},
12};
13
14#[derive(Debug, Clone, PartialEq, Eq)]
16pub struct OptimizeReport {
17 pub starting_version: u64,
19 pub committed_version: u64,
22 pub candidate_source_segments: u64,
24 pub source_segments_replaced: u64,
26 pub replacement_segments_written: u64,
28 pub distinct_identities_materialized: u64,
30 pub rows_read: u64,
32 pub rows_written: u64,
34 pub no_op: bool,
36}
37
38impl OptimizeReport {
39 fn no_op(version: u64) -> Self {
40 Self {
41 starting_version: version,
42 committed_version: version,
43 candidate_source_segments: 0,
44 source_segments_replaced: 0,
45 replacement_segments_written: 0,
46 distinct_identities_materialized: 0,
47 rows_read: 0,
48 rows_written: 0,
49 no_op: true,
50 }
51 }
52}
53
54struct PlanCounts {
55 candidates: u64,
56 replacements: u64,
57 identities: u64,
58 rows: u64,
59}
60
61fn mixed_segment_candidates(state: &TableState) -> Result<Vec<SegmentMeta>, IndexValueError> {
62 Ok(state
63 .segments_sorted_by_index()?
64 .into_iter()
65 .filter(|segment| segment.entity_layout == SegmentEntityLayout::Mixed)
66 .cloned()
67 .collect())
68}
69
70fn count(field: &'static str, value: usize) -> Result<u64, TableError> {
71 value
72 .try_into()
73 .map_err(|_| TableError::OptimizeCountOverflow { field })
74}
75
76fn add(field: &'static str, total: &mut u64, value: u64) -> Result<(), TableError> {
77 *total = total
78 .checked_add(value)
79 .ok_or(TableError::OptimizeCountOverflow { field })?;
80 Ok(())
81}
82
83fn invalid_plan(reason: impl Into<String>) -> TableError {
84 TableError::OptimizeInvariant {
85 reason: reason.into(),
86 }
87}
88
89fn ensure_canonical(path: &str, description: &str) -> Result<(), TableError> {
90 let (canonical, _) = normalize_relative_storage_path(Path::new(path))
91 .map_err(|error| invalid_plan(format!("invalid {description} path {path:?}: {error}")))?;
92 if canonical != path {
93 return Err(invalid_plan(format!(
94 "{description} path {path:?} is not canonical; expected {canonical:?}"
95 )));
96 }
97 Ok(())
98}
99
100async fn validate_staged_plan(
101 table: &TimeSeriesTable,
102 candidates: &[SegmentMeta],
103 staged: &[StagedEntityRewrite],
104) -> Result<PlanCounts, TableError> {
105 if candidates.len() != staged.len() {
106 return Err(invalid_plan(format!(
107 "staged {} rewrites for {} candidates",
108 staged.len(),
109 candidates.len()
110 )));
111 }
112
113 let mut candidate_paths = HashSet::new();
114 let mut object_paths = HashSet::new();
115 let mut live_object_paths = HashSet::new();
116 for segment in table.state.segments.values() {
117 live_object_paths.insert(segment.path.as_str());
118 if let Some(path) = segment.coverage_path.as_deref() {
119 live_object_paths.insert(path);
120 }
121 }
122 if let Some(coverage) = &table.state.table_coverage {
123 live_object_paths.insert(coverage.coverage_path.as_str());
124 }
125 let mut source_coverage = EntityCoverage::empty();
126 let mut replacement_coverage = EntityCoverage::empty();
127 let mut identities = HashSet::<EntityIdentity>::new();
128 let mut source_rows = 0u64;
129 let mut replacement_rows = 0u64;
130 let mut replacement_count = 0u64;
131
132 for (candidate, rewrite) in candidates.iter().zip(staged) {
133 if rewrite.source_path != candidate.path {
134 return Err(invalid_plan(format!(
135 "candidate {} is not represented exactly once by its staged rewrite",
136 candidate.path
137 )));
138 }
139 if !candidate_paths.insert(&candidate.path) {
140 return Err(invalid_plan(format!(
141 "candidate path {} appears more than once",
142 candidate.path
143 )));
144 }
145 if table.state.segments.get(&candidate.path) != Some(candidate) {
146 return Err(invalid_plan(format!(
147 "candidate {} is not live in the starting snapshot",
148 candidate.path
149 )));
150 }
151 add("rows_read", &mut source_rows, candidate.row_count)?;
152
153 let source_coverage_path = candidate.coverage_path.as_deref().ok_or_else(|| {
154 invalid_plan(format!(
155 "candidate {} has no coverage sidecar",
156 candidate.path
157 ))
158 })?;
159 let candidate_coverage =
160 read_entity_coverage_sidecar(table.location(), Path::new(source_coverage_path))
161 .await
162 .map_err(|source| TableError::CoverageSidecar { source })?;
163 let mut candidate_replacement_coverage = EntityCoverage::empty();
164 let mut candidate_rows = 0u64;
165
166 for replacement in &rewrite.replacements {
167 let coverage_path = replacement.meta.coverage_path.as_deref().ok_or_else(|| {
168 invalid_plan(format!(
169 "replacement {} has no coverage sidecar",
170 replacement.meta.path
171 ))
172 })?;
173 for (path, description) in [
174 (replacement.meta.path.as_str(), "replacement data"),
175 (coverage_path, "replacement coverage"),
176 ] {
177 ensure_canonical(path, description)?;
178 if !object_paths.insert(path) {
179 return Err(invalid_plan(format!(
180 "staged object path {path} appears more than once"
181 )));
182 }
183 if live_object_paths.contains(path) {
184 return Err(invalid_plan(format!(
185 "staged object path {path} conflicts with a live table object"
186 )));
187 }
188 }
189 if replacement.meta.entity_layout
190 != SegmentEntityLayout::Single(replacement.identity.clone())
191 || replacement.coverage.identity_count() != 1
192 || replacement.coverage.get(&replacement.identity).is_none()
193 {
194 return Err(invalid_plan(format!(
195 "replacement {} is not truthful Single metadata",
196 replacement.meta.path
197 )));
198 }
199 add("replacement_segments_written", &mut replacement_count, 1)?;
200 add(
201 "rows_written",
202 &mut candidate_rows,
203 replacement.meta.row_count,
204 )?;
205 identities.insert(replacement.identity.clone());
206 candidate_replacement_coverage.union_inplace(&replacement.coverage);
207 }
208
209 if candidate_rows != candidate.row_count {
210 return Err(invalid_plan(format!(
211 "replacement rows {candidate_rows} do not equal source rows {} for {}",
212 candidate.row_count, candidate.path
213 )));
214 }
215 if candidate_replacement_coverage != candidate_coverage {
216 return Err(invalid_plan(format!(
217 "replacement coverage does not equal source coverage for {}",
218 candidate.path
219 )));
220 }
221 add("rows_written", &mut replacement_rows, candidate_rows)?;
222 source_coverage.union_inplace(&candidate_coverage);
223 replacement_coverage.union_inplace(&candidate_replacement_coverage);
224 }
225
226 if source_rows != replacement_rows {
227 return Err(invalid_plan(format!(
228 "replacement rows {replacement_rows} do not equal source rows {source_rows}"
229 )));
230 }
231 if source_coverage != replacement_coverage {
232 return Err(invalid_plan(
233 "replacement coverage does not reconstruct selected source coverage",
234 ));
235 }
236
237 Ok(PlanCounts {
238 candidates: count("candidate_source_segments", candidates.len())?,
239 replacements: replacement_count,
240 identities: count("distinct_identities_materialized", identities.len())?,
241 rows: source_rows,
242 })
243}
244
245impl TimeSeriesTable {
246 async fn rollback_optimization(
247 &self,
248 staged_paths: &[String],
249 source: TableError,
250 ) -> TableError {
251 let mut cleanup_errors = Vec::new();
252 for path in staged_paths.iter().rev() {
253 if let Err(error) =
254 remove_file_if_exists(self.location().as_ref(), Path::new(path)).await
255 {
256 cleanup_errors.push(format!("{path}: {error}"));
257 }
258 }
259 if cleanup_errors.is_empty() {
260 source
261 } else {
262 TableError::OptimizeRollback {
263 source: Box::new(source),
264 cleanup_errors,
265 }
266 }
267 }
268
269 pub async fn optimize(&mut self) -> Result<OptimizeReport, TableError> {
280 if self.index.entity_columns.is_empty() {
281 let table_root = match self.location().as_ref() {
282 StorageLocation::Local(root) => root.display().to_string(),
283 };
284 return Err(TableError::OptimizeNotApplicable { table_root });
285 }
286
287 let starting_version = self.state.version;
288 let candidates = mixed_segment_candidates(&self.state)
289 .map_err(|source| TableError::InvalidSegmentBounds { source })?;
290 if candidates.is_empty() {
291 return Ok(OptimizeReport::no_op(starting_version));
292 }
293 let committed_version =
294 starting_version
295 .checked_add(1)
296 .ok_or(TableError::OptimizeCountOverflow {
297 field: "committed_version",
298 })?;
299 let table_schema = self.state.table_meta.logical_schema.clone().ok_or(
300 TableError::MissingCanonicalSchema {
301 version: starting_version,
302 },
303 )?;
304
305 let mut staged = Vec::with_capacity(candidates.len());
306 let mut staged_paths = Vec::new();
307 for source in &candidates {
308 match rewrite_mixed_parquet_segment(self.location(), &table_schema, &self.index, source)
309 .await
310 {
311 Ok(rewrite) => {
312 staged_paths.extend(rewrite.staged_object_paths.iter().cloned());
313 staged.push(rewrite);
314 }
315 Err(source) => {
316 let error = TableError::OptimizeRewrite { source };
317 return Err(self.rollback_optimization(&staged_paths, error).await);
318 }
319 }
320 }
321
322 let counts = match validate_staged_plan(self, &candidates, &staged).await {
323 Ok(counts) => counts,
324 Err(source) => {
325 return Err(self.rollback_optimization(&staged_paths, source).await);
326 }
327 };
328
329 let mut actions = Vec::new();
330 for source in &candidates {
331 actions.push(LogAction::RemoveSegment {
332 path: source.path.clone(),
333 });
334 }
335 for rewrite in &staged {
336 actions.extend(
337 rewrite
338 .replacements
339 .iter()
340 .map(|replacement| LogAction::AddSegment(replacement.meta.clone())),
341 );
342 }
343
344 let new_version = match self
345 .log
346 .commit_with_expected_version(starting_version, actions)
347 .await
348 {
349 Ok(version) => version,
350 Err(source @ CommitError::AmbiguousOutcome { .. }) => {
351 return Err(TableError::TransactionLog { source });
352 }
353 Err(source) => {
354 let error = TableError::TransactionLog { source };
355 return Err(self.rollback_optimization(&staged_paths, error).await);
356 }
357 };
358 assert_eq!(
359 new_version, committed_version,
360 "transaction log returned unexpected optimize version"
361 );
362
363 for source in &candidates {
364 self.state.segments.remove(&source.path);
365 }
366 for rewrite in staged {
367 for replacement in rewrite.replacements {
368 self.state
369 .segments
370 .insert(replacement.meta.path.clone(), replacement.meta);
371 }
372 }
373 self.state.version = new_version;
374
375 Ok(OptimizeReport {
376 starting_version,
377 committed_version: new_version,
378 candidate_source_segments: counts.candidates,
379 source_segments_replaced: counts.candidates,
380 replacement_segments_written: counts.replacements,
381 distinct_identities_materialized: counts.identities,
382 rows_read: counts.rows,
383 rows_written: counts.rows,
384 no_op: false,
385 })
386 }
387}
388
389#[cfg(test)]
390mod tests {
391 use super::*;
392 use std::{
393 collections::{BTreeSet, HashMap},
394 path::{Path, PathBuf},
395 };
396
397 use arrow::datatypes::TimeUnit;
398 use futures::StreamExt;
399 use tempfile::TempDir;
400
401 use crate::{
402 coverage::{EntityIdentity, EntityValue},
403 metadata::{
404 logical_schema::LogicalTimestampUnit,
405 segments::{FileFormat, SegmentEntityLayout},
406 table_metadata::IndexValue,
407 },
408 storage::{TableLocation, layout},
409 table::test_util::{
410 make_int32_entity_table_meta, make_table_meta_with_unit, utc_datetime,
411 write_arrow_parquet_with_unit, write_int32_entity_parquet,
412 },
413 transaction_log::TableKind,
414 };
415
416 fn segment(path: &str, layout: SegmentEntityLayout, minute: u32) -> SegmentMeta {
417 SegmentMeta {
418 path: path.to_string(),
419 format: FileFormat::Parquet,
420 entity_layout: layout,
421 index_min: IndexValue::Timestamp(utc_datetime(2025, 1, 1, 0, minute, 0)),
422 index_max: IndexValue::Timestamp(utc_datetime(2025, 1, 1, 0, minute + 1, 0)),
423 row_count: 1,
424 file_size: Some(1),
425 coverage_path: Some(format!("_coverage/segments/{minute}.roar")),
426 }
427 }
428
429 fn state(segments: impl IntoIterator<Item = SegmentMeta>) -> TableState {
430 TableState {
431 version: 7,
432 table_meta: make_table_meta_with_unit(LogicalTimestampUnit::Millis),
433 segments: segments
434 .into_iter()
435 .map(|segment| (segment.path.clone(), segment))
436 .collect::<HashMap<_, _>>(),
437 table_coverage: None,
438 }
439 }
440
441 fn single_identity() -> SegmentEntityLayout {
442 SegmentEntityLayout::Single(
443 EntityIdentity::try_new(vec!["A".into()]).expect("valid identity"),
444 )
445 }
446
447 fn files_below(root: &Path) -> std::io::Result<BTreeSet<PathBuf>> {
448 if !root.exists() {
449 return Ok(BTreeSet::new());
450 }
451 let mut files = BTreeSet::new();
452 let mut directories = vec![root.to_owned()];
453 while let Some(directory) = directories.pop() {
454 for entry in std::fs::read_dir(directory)? {
455 let entry = entry?;
456 if entry.file_type()?.is_dir() {
457 directories.push(entry.path());
458 } else {
459 files.insert(entry.path());
460 }
461 }
462 }
463 Ok(files)
464 }
465
466 fn optimization_objects(root: &Path) -> std::io::Result<BTreeSet<PathBuf>> {
467 let mut files = BTreeSet::new();
468 for relative in ["data/_staged", layout::SEGMENT_COVERAGE_DIR] {
469 files.extend(files_below(&root.join(relative))?);
470 }
471 Ok(files)
472 }
473
474 async fn append_mixed_source(
475 table: &mut TimeSeriesTable,
476 root: &Path,
477 path: &str,
478 start_millis: i64,
479 ) -> Result<(), TableError> {
480 write_arrow_parquet_with_unit(
481 &root.join(path),
482 TimeUnit::Millisecond,
483 &[
484 Some(start_millis + 1_000),
485 Some(start_millis + 2_000),
486 Some(start_millis + 3_000),
487 Some(start_millis + 4_000),
488 ],
489 &["A", "B", "A", "B"],
490 &[10.0, 20.0, 11.0, 21.0],
491 )
492 .expect("write mixed source");
493 table.append_parquet_segment(path).await?;
494 assert_eq!(
495 table.state().segments[path].entity_layout,
496 SegmentEntityLayout::Mixed
497 );
498 Ok(())
499 }
500
501 #[test]
502 fn discovery_selects_all_and_only_mixed_segments() {
503 let state = state([
504 segment("data/mixed-late.parquet", SegmentEntityLayout::Mixed, 2),
505 segment("data/single.parquet", single_identity(), 0),
506 segment(
507 "data/not-applicable.parquet",
508 SegmentEntityLayout::NotApplicable,
509 0,
510 ),
511 segment("data/mixed.parquet", SegmentEntityLayout::Mixed, 1),
512 ]);
513
514 let paths = mixed_segment_candidates(&state)
515 .expect("valid segment bounds")
516 .into_iter()
517 .map(|segment| segment.path)
518 .collect::<Vec<_>>();
519
520 assert_eq!(paths, ["data/mixed.parquet", "data/mixed-late.parquet"]);
521 }
522
523 #[test]
524 fn discovery_order_is_independent_of_hash_map_insertion_order() {
525 let segments = [
526 segment("data/b.parquet", SegmentEntityLayout::Mixed, 1),
527 segment("data/a.parquet", SegmentEntityLayout::Mixed, 1),
528 segment("data/later.parquet", SegmentEntityLayout::Mixed, 2),
529 ];
530 let forward = state(segments.clone());
531 let reverse = state(segments.into_iter().rev());
532
533 let paths = |state: &TableState| {
534 mixed_segment_candidates(state)
535 .expect("valid segment bounds")
536 .into_iter()
537 .map(|segment| segment.path)
538 .collect::<Vec<_>>()
539 };
540
541 assert_eq!(paths(&forward), paths(&reverse));
542 assert_eq!(
543 paths(&forward),
544 ["data/a.parquet", "data/b.parquet", "data/later.parquet"]
545 );
546 }
547
548 #[tokio::test]
549 async fn optimize_without_mixed_segments_is_a_zero_write_no_op() -> Result<(), TableError> {
550 let temp = TempDir::new().expect("temp directory");
551 let mut table = TimeSeriesTable::create(
552 TableLocation::local(temp.path()),
553 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
554 )
555 .await?;
556 let starting_version = table.state().version;
557
558 let report = table.optimize().await?;
559
560 assert_eq!(report, OptimizeReport::no_op(starting_version));
561 assert!(!temp.path().join("data/_staged").exists());
562 let reopened = TimeSeriesTable::open(TableLocation::local(temp.path())).await?;
563 assert_eq!(reopened.state().version, starting_version);
564 Ok(())
565 }
566
567 #[tokio::test]
568 async fn optimize_rejects_a_table_without_entity_columns() -> Result<(), TableError> {
569 let temp = TempDir::new().expect("temp directory");
570 let mut table_meta = make_table_meta_with_unit(LogicalTimestampUnit::Millis);
571 let TableKind::TimeSeries(index) = &mut table_meta.kind else {
572 unreachable!("test table is time-series");
573 };
574 index.entity_columns.clear();
575 let mut table =
576 TimeSeriesTable::create(TableLocation::local(temp.path()), table_meta).await?;
577
578 let error = table
579 .optimize()
580 .await
581 .expect_err("entity-free table must be rejected");
582
583 assert!(matches!(
584 error,
585 TableError::OptimizeNotApplicable { table_root }
586 if table_root == temp.path().display().to_string()
587 ));
588 assert!(!temp.path().join("data/_staged").exists());
589 Ok(())
590 }
591
592 #[tokio::test]
593 async fn optimize_atomically_replaces_one_mixed_source() -> Result<(), TableError> {
594 let temp = TempDir::new().expect("temp directory");
595 let location = TableLocation::local(temp.path());
596 let mut table = TimeSeriesTable::create(
597 location.clone(),
598 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
599 )
600 .await?;
601 let source_path = "data/mixed.parquet";
602 write_arrow_parquet_with_unit(
603 &temp.path().join(source_path),
604 TimeUnit::Millisecond,
605 &[Some(1_000), Some(2_000), Some(3_000), Some(4_000)],
606 &["A", "B", "A", "B"],
607 &[10.0, 20.0, 11.0, 21.0],
608 )
609 .expect("write mixed source");
610 table.append_parquet_segment(source_path).await?;
611 let source = table
612 .state()
613 .segments
614 .get(source_path)
615 .expect("committed source")
616 .clone();
617 assert_eq!(source.entity_layout, SegmentEntityLayout::Mixed);
618 let source_bytes = std::fs::read(temp.path().join(source_path)).expect("source bytes");
619 let source_coverage_path = source.coverage_path.as_deref().expect("source coverage");
620 let source_coverage_bytes =
621 std::fs::read(temp.path().join(source_coverage_path)).expect("source coverage bytes");
622 let table_coverage = table.state().table_coverage.clone();
623 let starting_version = table.state().version;
624
625 let report = table.optimize().await?;
626
627 assert_eq!(report.starting_version, starting_version);
628 assert_eq!(report.committed_version, starting_version + 1);
629 assert_eq!(report.candidate_source_segments, 1);
630 assert_eq!(report.source_segments_replaced, 1);
631 assert_eq!(report.replacement_segments_written, 2);
632 assert_eq!(report.distinct_identities_materialized, 2);
633 assert_eq!(report.rows_read, 4);
634 assert_eq!(report.rows_written, 4);
635 assert!(!report.no_op);
636 assert!(!table.state().segments.contains_key(source_path));
637 assert_eq!(table.state().segments.len(), 2);
638 assert!(
639 table
640 .state()
641 .segments
642 .values()
643 .all(|segment| matches!(segment.entity_layout, SegmentEntityLayout::Single(_)))
644 );
645 assert_eq!(table.state().table_coverage, table_coverage);
646 assert_eq!(
647 std::fs::read(temp.path().join(source_path)).expect("source remains"),
648 source_bytes
649 );
650 assert_eq!(
651 std::fs::read(temp.path().join(source_coverage_path)).expect("source coverage remains"),
652 source_coverage_bytes
653 );
654 Ok(())
655 }
656
657 #[tokio::test]
658 async fn optimize_rewrites_numeric_entities_with_typed_single_layouts() -> Result<(), TableError>
659 {
660 let temp = TempDir::new().expect("temp directory");
661 let location = TableLocation::local(temp.path());
662 let mut table =
663 TimeSeriesTable::create(location.clone(), make_int32_entity_table_meta()).await?;
664 let source_path = "data/numeric-mixed.parquet";
665 write_int32_entity_parquet(
666 &temp.path().join(source_path),
667 &[1_000, 2_000, 3_000, 4_000],
668 &[-1, i32::MAX, -1, i32::MAX],
669 &[10.0, 20.0, 11.0, 21.0],
670 )
671 .expect("write numeric mixed source");
672 table.append_parquet_segment(source_path).await?;
673 assert_eq!(
674 table.state().segments[source_path].entity_layout,
675 SegmentEntityLayout::Mixed
676 );
677 let table_meta = table.state().table_meta.clone();
678
679 let report = table.optimize().await?;
680
681 assert_eq!(report.replacement_segments_written, 2);
682 assert_eq!(report.rows_written, 4);
683 assert_eq!(table.state().table_meta, table_meta);
684 let identities = table
685 .state()
686 .segments
687 .values()
688 .map(|segment| match &segment.entity_layout {
689 SegmentEntityLayout::Single(identity) => identity.clone(),
690 layout => panic!("expected typed single-entity layout, found {layout:?}"),
691 })
692 .collect::<BTreeSet<_>>();
693 assert_eq!(
694 identities,
695 BTreeSet::from([
696 EntityIdentity::try_new(vec![EntityValue::Int32(-1)]).expect("negative identity"),
697 EntityIdentity::try_new(vec![EntityValue::Int32(i32::MAX)])
698 .expect("maximum identity"),
699 ])
700 );
701
702 let reopened = TimeSeriesTable::open(location).await?;
703 assert_eq!(reopened.state(), table.state());
704 Ok(())
705 }
706
707 #[tokio::test]
708 async fn later_staging_failure_cleans_every_earlier_rewrite() -> Result<(), TableError> {
709 let temp = TempDir::new().expect("temp directory");
710 let mut table = TimeSeriesTable::create(
711 TableLocation::local(temp.path()),
712 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
713 )
714 .await?;
715 append_mixed_source(&mut table, temp.path(), "data/first.parquet", 0).await?;
716 append_mixed_source(&mut table, temp.path(), "data/broken.parquet", 60_000).await?;
717 let state_before = table.state().clone();
718 let objects_before = optimization_objects(temp.path()).expect("optimization objects");
719 std::fs::remove_file(temp.path().join("data/broken.parquet")).expect("remove later source");
720
721 let error = table.optimize().await.expect_err("later staging must fail");
722
723 assert!(matches!(error, TableError::OptimizeRewrite { .. }));
724 assert_eq!(table.state(), &state_before);
725 assert_eq!(
726 optimization_objects(temp.path()).expect("optimization objects"),
727 objects_before
728 );
729 Ok(())
730 }
731
732 #[tokio::test]
733 async fn occ_conflict_cleans_staged_objects_and_preserves_state() -> Result<(), TableError> {
734 let temp = TempDir::new().expect("temp directory");
735 let location = TableLocation::local(temp.path());
736 let mut table = TimeSeriesTable::create(
737 location.clone(),
738 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
739 )
740 .await?;
741 append_mixed_source(&mut table, temp.path(), "data/candidate.parquet", 0).await?;
742 let state_before = table.state().clone();
743 let mut concurrent = TimeSeriesTable::open(location).await?;
744 append_mixed_source(
745 &mut concurrent,
746 temp.path(),
747 "data/concurrent.parquet",
748 60_000,
749 )
750 .await?;
751 let objects_before = optimization_objects(temp.path()).expect("optimization objects");
752
753 let error = table.optimize().await.expect_err("stale commit must fail");
754
755 assert!(matches!(
756 error,
757 TableError::TransactionLog {
758 source: CommitError::Conflict { .. }
759 }
760 ));
761 assert_eq!(table.state(), &state_before);
762 assert_eq!(
763 table
764 .log
765 .load_current_version()
766 .await
767 .expect("current version"),
768 state_before.version + 1
769 );
770 assert_eq!(
771 optimization_objects(temp.path()).expect("optimization objects"),
772 objects_before
773 );
774 Ok(())
775 }
776
777 #[tokio::test]
778 async fn ambiguous_commit_retains_staged_objects_and_preserves_state() -> Result<(), TableError>
779 {
780 let temp = TempDir::new().expect("temp directory");
781 let mut table = TimeSeriesTable::create(
782 TableLocation::local(temp.path()),
783 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
784 )
785 .await?;
786 append_mixed_source(&mut table, temp.path(), "data/candidate.parquet", 0).await?;
787 let state_before = table.state().clone();
788 let objects_before = optimization_objects(temp.path()).expect("optimization objects");
789 let commit_path = temp
790 .path()
791 .join(layout::commit_rel_path(state_before.version + 1));
792 crate::storage::inject_write_new_failure(commit_path.clone(), true);
793
794 let error = table
795 .optimize()
796 .await
797 .expect_err("commit outcome must be ambiguous");
798
799 assert!(matches!(
800 error,
801 TableError::TransactionLog {
802 source: CommitError::AmbiguousOutcome { .. }
803 }
804 ));
805 assert_eq!(table.state(), &state_before);
806 assert_eq!(
807 table
808 .log
809 .load_current_version()
810 .await
811 .expect("current version"),
812 state_before.version
813 );
814 assert!(commit_path.exists());
815 let objects_after = optimization_objects(temp.path()).expect("optimization objects");
816 assert_eq!(objects_after.difference(&objects_before).count(), 4);
817 assert_eq!(
818 files_below(&temp.path().join("data/_staged"))
819 .expect("staged data")
820 .len(),
821 2
822 );
823 Ok(())
824 }
825
826 #[tokio::test]
827 async fn rollback_reports_every_cleanup_failure_in_reverse_order() -> Result<(), TableError> {
828 let temp = TempDir::new().expect("temp directory");
829 let table = TimeSeriesTable::create(
830 TableLocation::local(temp.path()),
831 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
832 )
833 .await?;
834 let paths = [
835 "data/_staged/entity-rewrite/first.parquet".to_string(),
836 format!("{}/second.roar", layout::SEGMENT_COVERAGE_DIR),
837 ];
838 for path in &paths {
839 let absolute = temp.path().join(path);
840 std::fs::create_dir_all(absolute.parent().expect("object parent"))
841 .expect("create object parent");
842 std::fs::write(&absolute, b"staged").expect("write staged object");
843 crate::storage::inject_cleanup_failure(absolute);
844 }
845
846 let error = table
847 .rollback_optimization(&paths, invalid_plan("primary failure"))
848 .await;
849 let message = error.to_string();
850
851 assert!(matches!(
852 error,
853 TableError::OptimizeRollback {
854 source,
855 cleanup_errors,
856 } if matches!(*source, TableError::OptimizeInvariant { .. })
857 && cleanup_errors.len() == 2
858 && cleanup_errors[0].contains("second.roar")
859 && cleanup_errors[1].contains("first.parquet")
860 ));
861 assert!(message.contains("primary failure"));
862 assert!(paths.iter().all(|path| temp.path().join(path).exists()));
863 Ok(())
864 }
865
866 #[tokio::test]
867 async fn multiple_sources_reopen_recover_and_repeat_as_a_no_op() -> Result<(), TableError> {
868 let temp = TempDir::new().expect("temp directory");
869 let location = TableLocation::local(temp.path());
870 let mut table = TimeSeriesTable::create(
871 location.clone(),
872 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
873 )
874 .await?;
875 let source_paths = ["data/first.parquet", "data/second.parquet"];
876 append_mixed_source(&mut table, temp.path(), source_paths[0], 0).await?;
877 append_mixed_source(&mut table, temp.path(), source_paths[1], 60_000).await?;
878 let sources = source_paths.map(|path| table.state().segments[path].clone());
879 let expected_coverage = table.load_table_entity_snapshot_coverage_readonly().await?;
880 let coverage_pointer = table
881 .state()
882 .table_coverage
883 .clone()
884 .expect("table coverage pointer");
885 let coverage_bytes = std::fs::read(temp.path().join(&coverage_pointer.coverage_path))
886 .expect("table coverage bytes");
887 let snapshot_files =
888 files_below(&temp.path().join(layout::TABLE_SNAPSHOT_DIR)).expect("snapshot files");
889 let table_meta = table.state().table_meta.clone();
890 let starting_version = table.state().version;
891
892 let report = table.optimize().await?;
893
894 assert_eq!(
895 report,
896 OptimizeReport {
897 starting_version,
898 committed_version: starting_version + 1,
899 candidate_source_segments: 2,
900 source_segments_replaced: 2,
901 replacement_segments_written: 4,
902 distinct_identities_materialized: 2,
903 rows_read: 8,
904 rows_written: 8,
905 no_op: false,
906 }
907 );
908 assert_eq!(table.state().segments.len(), 4);
909 assert_eq!(table.state().table_meta, table_meta);
910 assert!(
911 table
912 .state()
913 .segments
914 .values()
915 .all(|segment| matches!(segment.entity_layout, SegmentEntityLayout::Single(_)))
916 );
917 assert_eq!(table.state().table_coverage, Some(coverage_pointer.clone()));
918 assert_eq!(
919 std::fs::read(temp.path().join(&coverage_pointer.coverage_path))
920 .expect("table coverage bytes"),
921 coverage_bytes
922 );
923 assert_eq!(
924 files_below(&temp.path().join(layout::TABLE_SNAPSHOT_DIR)).expect("snapshot files"),
925 snapshot_files
926 );
927 for source in &sources {
928 assert!(temp.path().join(&source.path).exists());
929 assert!(
930 temp.path()
931 .join(source.coverage_path.as_deref().expect("source coverage"))
932 .exists()
933 );
934 }
935
936 let commit = table
937 .log
938 .load_commit(report.committed_version)
939 .await
940 .expect("optimization commit");
941 assert_eq!(commit.base_version, starting_version);
942 assert_eq!(commit.actions.len(), 6);
943 assert!(
944 commit.actions[..2]
945 .iter()
946 .zip(source_paths)
947 .all(|(action, expected)| matches!(
948 action,
949 LogAction::RemoveSegment { path } if path == expected
950 ))
951 );
952 assert!(
953 commit.actions[2..]
954 .iter()
955 .all(|action| matches!(action, LogAction::AddSegment(_)))
956 );
957
958 let state_after_first = table.state().clone();
959 let objects_after_first = optimization_objects(temp.path()).expect("optimization objects");
960 let second_report = table.optimize().await?;
961 assert_eq!(
962 second_report,
963 OptimizeReport::no_op(report.committed_version)
964 );
965 assert_eq!(table.state(), &state_after_first);
966 assert_eq!(
967 optimization_objects(temp.path()).expect("optimization objects"),
968 objects_after_first
969 );
970 assert_eq!(
971 table
972 .log
973 .load_current_version()
974 .await
975 .expect("current version"),
976 report.committed_version
977 );
978
979 let reopened = TimeSeriesTable::open(location).await?;
980 assert_eq!(reopened.state(), table.state());
981 assert_eq!(
982 reopened
983 .recover_table_entity_coverage_from_segments()
984 .await?,
985 expected_coverage
986 );
987 let mut scan = reopened
988 .scan_range(
989 chrono::DateTime::from_timestamp_millis(0).expect("range start"),
990 chrono::DateTime::from_timestamp_millis(120_000).expect("range end"),
991 )
992 .await?;
993 let mut rows = 0;
994 while let Some(batch) = scan.next().await {
995 rows += batch?.num_rows();
996 }
997 assert_eq!(rows, 8);
998 Ok(())
999 }
1000
1001 #[test]
1002 fn accumulated_report_counts_do_not_wrap() {
1003 let mut total = u64::MAX;
1004
1005 let error = add("rows_written", &mut total, 1).expect_err("count overflow must fail");
1006
1007 assert!(matches!(
1008 error,
1009 TableError::OptimizeCountOverflow {
1010 field: "rows_written"
1011 }
1012 ));
1013 assert_eq!(total, u64::MAX);
1014 }
1015
1016 #[tokio::test]
1017 async fn version_overflow_fails_before_staging() -> Result<(), TableError> {
1018 let temp = TempDir::new().expect("temp directory");
1019 let mut table = TimeSeriesTable::create(
1020 TableLocation::local(temp.path()),
1021 make_table_meta_with_unit(LogicalTimestampUnit::Millis),
1022 )
1023 .await?;
1024 append_mixed_source(&mut table, temp.path(), "data/candidate.parquet", 0).await?;
1025 table.state.version = u64::MAX;
1026 let objects_before = optimization_objects(temp.path()).expect("optimization objects");
1027
1028 let error = table
1029 .optimize()
1030 .await
1031 .expect_err("version overflow must fail");
1032
1033 assert!(matches!(
1034 error,
1035 TableError::OptimizeCountOverflow {
1036 field: "committed_version"
1037 }
1038 ));
1039 assert_eq!(table.state().version, u64::MAX);
1040 assert_eq!(
1041 optimization_objects(temp.path()).expect("optimization objects"),
1042 objects_before
1043 );
1044 Ok(())
1045 }
1046}