Skip to main content

timeseries_table_format/table/operations/
optimize.rs

1//! Entity-layout optimization for time-series tables.
2
3use 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/// Errors owned by an entity-layout optimization operation.
26#[derive(Debug, Snafu)]
27#[snafu(module, visibility(pub(crate)))]
28#[non_exhaustive]
29pub enum OptimizeError {
30    /// The table protocol does not permit this client to optimize.
31    #[snafu(context(false), display("Table protocol error: {source}"))]
32    Protocol {
33        /// Complete table protocol failure.
34        #[snafu(source)]
35        source: TableProtocolError,
36        /// Backtrace captured at the optimization boundary.
37        backtrace: Backtrace,
38    },
39
40    /// Entity-layout optimization requires at least one entity column.
41    #[snafu(display(
42        "Entity-layout optimization is not applicable to table {table_root}: no entity columns are configured"
43    ))]
44    NotApplicable {
45        /// User-facing table root.
46        table_root: String,
47    },
48
49    /// Staging replacements for one mixed segment failed.
50    #[snafu(context(false), display("Mixed-segment rewrite failed: {source}"))]
51    MixedSegmentRewrite {
52        /// Complete mixed-segment rewrite failure.
53        #[snafu(source(from(EntityRewriteError, Box::new)), backtrace)]
54        source: Box<EntityRewriteError>,
55    },
56
57    /// Live segment bounds cannot be ordered in one native index domain.
58    #[snafu(
59        context(false),
60        display("Invalid segment ordered-index bounds: {source}")
61    )]
62    InvalidSegmentBounds {
63        /// Complete bounds validation failure.
64        #[snafu(source)]
65        source: IndexValueError,
66        /// Backtrace captured because index value validation does not own one.
67        backtrace: Backtrace,
68    },
69
70    /// Optimization cannot use the table's canonical schema.
71    #[snafu(
72        context(false),
73        display("Optimization schema validation failed: {source}")
74    )]
75    SchemaValidation {
76        /// Complete schema compatibility failure.
77        #[snafu(source(from(SchemaCompatibilityError, Box::new)), backtrace)]
78        source: Box<SchemaCompatibilityError>,
79    },
80
81    /// A coverage sidecar required to validate a staged plan could not be read.
82    #[snafu(
83        context(false),
84        display("Optimization coverage sidecar error: {source}")
85    )]
86    CoverageSidecar {
87        /// Complete coverage sidecar failure.
88        #[snafu(source(from(CoverageSidecarError, Box::new)), backtrace)]
89        source: Box<CoverageSidecarError>,
90    },
91
92    /// A staged replacement path failed table-relative storage validation.
93    #[snafu(display("Invalid {description} path {path:?}: {source}"))]
94    InvalidStagedPath {
95        /// Role of the rejected path in the staged replacement.
96        description: &'static str,
97        /// Rejected table-relative path.
98        path: String,
99        /// Complete storage path validation failure.
100        #[snafu(source(from(StorageError, Box::new)), backtrace)]
101        source: Box<StorageError>,
102    },
103
104    /// A staged optimization plan violated an atomic publication invariant.
105    #[snafu(display("Invalid staged entity-layout optimization plan: {reason}"))]
106    InvalidStagedPlan {
107        /// Failed plan invariant.
108        reason: String,
109        /// Backtrace captured at the failed internal invariant.
110        backtrace: Backtrace,
111    },
112
113    /// An optimization count could not be represented without wrapping.
114    #[snafu(display("Entity-layout optimization count overflow: {field}"))]
115    CountOverflow {
116        /// Report or version field that overflowed.
117        field: &'static str,
118        /// Backtrace captured at the failed internal arithmetic boundary.
119        backtrace: Backtrace,
120    },
121
122    /// Publishing the optimization transaction failed.
123    #[snafu(context(false), display("Optimization commit failed: {source}"))]
124    Commit {
125        /// Complete transaction-log failure.
126        #[snafu(source, backtrace)]
127        source: CommitError,
128    },
129
130    /// Optimization failed and one or more owned staged objects could not be removed.
131    #[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        /// Original optimization failure that triggered rollback.
141        #[snafu(source, backtrace)]
142        source: Box<OptimizeError>,
143        /// Typed cleanup failure for every staged object that could not be removed.
144        cleanup_errors: Vec<StorageError>,
145    },
146}
147
148/// Result of one entity-layout optimization operation.
149#[derive(Debug, Clone, PartialEq, Eq)]
150pub struct OptimizeReport {
151    /// Table version used to select candidates.
152    pub starting_version: u64,
153    /// Version containing the atomic replacement commit, or the starting
154    /// version for a no-op.
155    pub committed_version: u64,
156    /// Mixed live segments selected from the starting snapshot.
157    pub candidate_source_segments: u64,
158    /// Selected sources removed by the successful commit.
159    pub source_segments_replaced: u64,
160    /// Verified single-entity replacements added by the successful commit.
161    pub replacement_segments_written: u64,
162    /// Unique complete identities represented across all replacements.
163    pub distinct_identities_materialized: u64,
164    /// Logical rows in the selected source segments.
165    pub rows_read: u64,
166    /// Logical rows in the committed replacement segments.
167    pub rows_written: u64,
168    /// Whether no mixed live segments existed at the starting version.
169    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    /// Replace every live mixed-entity segment with verified single-entity
404    /// Parquet segments in one expected-version commit.
405    ///
406    /// Optimization preserves logical rows, schema, and per-entity coverage,
407    /// but may change physical row order.
408    ///
409    /// # Errors
410    ///
411    /// Returns [`TableError`] when optimization is not applicable, staging or
412    /// validation fails, the commit cannot be confirmed, or rollback fails.
413    #[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}