Skip to main content

timeseries_table_format/table/
optimize.rs

1//! Entity-layout optimization for time-series tables.
2
3use 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/// Result of one entity-layout optimization operation.
15#[derive(Debug, Clone, PartialEq, Eq)]
16pub struct OptimizeReport {
17    /// Table version used to select candidates.
18    pub starting_version: u64,
19    /// Version containing the atomic replacement commit, or the starting
20    /// version for a no-op.
21    pub committed_version: u64,
22    /// Mixed live segments selected from the starting snapshot.
23    pub candidate_source_segments: u64,
24    /// Selected sources removed by the successful commit.
25    pub source_segments_replaced: u64,
26    /// Verified single-entity replacements added by the successful commit.
27    pub replacement_segments_written: u64,
28    /// Unique complete identities represented across all replacements.
29    pub distinct_identities_materialized: u64,
30    /// Logical rows in the selected source segments.
31    pub rows_read: u64,
32    /// Logical rows in the committed replacement segments.
33    pub rows_written: u64,
34    /// Whether no mixed live segments existed at the starting version.
35    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    /// Replace every live mixed-entity segment with verified single-entity
270    /// Parquet segments in one expected-version commit.
271    ///
272    /// Optimization preserves logical rows, schema, and per-entity coverage,
273    /// but may change physical row order.
274    ///
275    /// # Errors
276    ///
277    /// Returns [`TableError`] when optimization is not applicable, staging or
278    /// validation fails, the commit cannot be confirmed, or rollback fails.
279    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}