Skip to main content

mito2/sst/
file.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Structures to describe metadata of files.
16
17use std::collections::HashMap;
18use std::fmt;
19use std::fmt::{Debug, Formatter};
20use std::num::NonZeroU64;
21use std::sync::atomic::{AtomicBool, Ordering};
22use std::sync::{Arc, Mutex, RwLock};
23
24use base64::prelude::{BASE64_STANDARD, Engine};
25use bytes::Bytes;
26use common_base::readable_size::ReadableSize;
27use common_telemetry::{debug, error};
28use common_time::Timestamp;
29use partition::expr::PartitionExpr;
30use serde::{Deserialize, Serialize};
31use smallvec::SmallVec;
32use store_api::metadata::ColumnMetadata;
33use store_api::region_request::PathType;
34use store_api::storage::{ColumnId, FileId, IndexVersion, RegionId};
35
36use crate::access_layer::AccessLayerRef;
37use crate::cache::CacheManagerRef;
38use crate::cache::file_cache::{FileType, IndexKey};
39use crate::sst::file_purger::FilePurgerRef;
40use crate::sst::location;
41use crate::sst::parquet::SstInfo;
42
43/// Custom serde functions for Bytes fields serialized as base64 strings.
44fn serialize_bytes_option<S>(bytes: &Option<Bytes>, serializer: S) -> Result<S::Ok, S::Error>
45where
46    S: serde::Serializer,
47{
48    match bytes {
49        None => serializer.serialize_none(),
50        Some(b) => serializer.serialize_some(&BASE64_STANDARD.encode(b)),
51    }
52}
53
54fn deserialize_bytes_option<'de, D>(deserializer: D) -> Result<Option<Bytes>, D::Error>
55where
56    D: serde::Deserializer<'de>,
57{
58    let opt: Option<String> = Option::deserialize(deserializer)?;
59    match opt {
60        None => Ok(None),
61        Some(s) => {
62            let decoded = BASE64_STANDARD
63                .decode(&s)
64                .map_err(serde::de::Error::custom)?;
65            Ok(Some(Bytes::from(decoded)))
66        }
67    }
68}
69
70/// Custom serde functions for partition_expr field in FileMeta
71fn serialize_partition_expr<S>(
72    partition_expr: &Option<PartitionExpr>,
73    serializer: S,
74) -> Result<S::Ok, S::Error>
75where
76    S: serde::Serializer,
77{
78    use serde::ser::Error;
79
80    match partition_expr {
81        None => serializer.serialize_none(),
82        Some(expr) => {
83            let json_str = expr.as_json_str().map_err(S::Error::custom)?;
84            serializer.serialize_some(&json_str)
85        }
86    }
87}
88
89fn deserialize_partition_expr<'de, D>(deserializer: D) -> Result<Option<PartitionExpr>, D::Error>
90where
91    D: serde::Deserializer<'de>,
92{
93    use serde::de::Error;
94
95    let opt_json_str: Option<String> = Option::deserialize(deserializer)?;
96    match opt_json_str {
97        None => Ok(None),
98        Some(json_str) => {
99            if json_str.is_empty() {
100                // Empty string represents explicit "single-region/no-partition" designation
101                Ok(None)
102            } else {
103                // Parse the JSON string to PartitionExpr
104                PartitionExpr::from_json_str(&json_str).map_err(D::Error::custom)
105            }
106        }
107    }
108}
109
110/// Type to store SST level.
111pub type Level = u8;
112/// Maximum level of SSTs.
113pub const MAX_LEVEL: Level = 2;
114/// Type to store index types for a column.
115pub type IndexTypes = SmallVec<[IndexType; 4]>;
116
117/// Cross-region file id.
118///
119/// It contains a region id and a file id. The string representation is `{region_id}/{file_id}`.
120#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
121pub struct RegionFileId {
122    /// The region that creates the file.
123    region_id: RegionId,
124    /// The id of the file.
125    file_id: FileId,
126}
127
128impl RegionFileId {
129    /// Creates a new [RegionFileId] from `region_id` and `file_id`.
130    pub fn new(region_id: RegionId, file_id: FileId) -> Self {
131        Self { region_id, file_id }
132    }
133
134    /// Gets the region id.
135    pub fn region_id(&self) -> RegionId {
136        self.region_id
137    }
138
139    /// Gets the file id.
140    pub fn file_id(&self) -> FileId {
141        self.file_id
142    }
143}
144
145impl fmt::Display for RegionFileId {
146    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
147        write!(f, "{}/{}", self.region_id, self.file_id)
148    }
149}
150
151/// Unique identifier for an index file, combining the SST file ID and the index version.
152#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
153pub struct RegionIndexId {
154    pub file_id: RegionFileId,
155    pub version: IndexVersion,
156}
157
158impl RegionIndexId {
159    pub fn new(file_id: RegionFileId, version: IndexVersion) -> Self {
160        Self { file_id, version }
161    }
162
163    pub fn region_id(&self) -> RegionId {
164        self.file_id.region_id
165    }
166
167    pub fn file_id(&self) -> FileId {
168        self.file_id.file_id
169    }
170}
171
172impl fmt::Display for RegionIndexId {
173    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
174        if self.version == 0 {
175            write!(f, "{}/{}", self.file_id.region_id, self.file_id.file_id)
176        } else {
177            write!(
178                f,
179                "{}/{}.{}",
180                self.file_id.region_id, self.file_id.file_id, self.version
181            )
182        }
183    }
184}
185
186/// Time range (min and max timestamps) of a SST file.
187/// Both min and max are inclusive.
188pub type FileTimeRange = (Timestamp, Timestamp);
189
190/// Checks if two inclusive timestamp ranges overlap with each other.
191pub(crate) fn overlaps(l: &FileTimeRange, r: &FileTimeRange) -> bool {
192    let (l, r) = if l.0 <= r.0 { (l, r) } else { (r, l) };
193    let (_, l_end) = l;
194    let (r_start, _) = r;
195
196    r_start <= l_end
197}
198
199/// Metadata of a SST file.
200#[derive(Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
201#[serde(default)]
202pub struct FileMeta {
203    /// Region that created the file. The region id may not be the id of the current region.
204    pub region_id: RegionId,
205    /// Compared to normal file names, FileId ignore the extension
206    pub file_id: FileId,
207    /// Timestamp range of file. The timestamps have the same time unit as the
208    /// data in the SST.
209    pub time_range: FileTimeRange,
210    /// SST level of the file.
211    pub level: Level,
212    /// Size of the file.
213    pub file_size: u64,
214    /// Maximum uncompressed row group size of the file. 0 means unknown.
215    pub max_row_group_uncompressed_size: u64,
216    /// Available indexes of the file.
217    pub available_indexes: IndexTypes,
218    /// Created indexes of the file for each column.
219    ///
220    /// This is essentially a more granular, column-level version of `available_indexes`,
221    /// primarily used for manual index building in the asynchronous index construction mode.
222    ///
223    /// For backward compatibility, older `FileMeta` versions might only contain `available_indexes`.
224    /// In such cases, we cannot deduce specific column index information from `available_indexes` alone.
225    /// Therefore, defaulting this `indexes` field to an empty list during deserialization is a
226    /// reasonable and necessary step to ensure column information consistency.
227    pub indexes: Vec<ColumnIndexMetadata>,
228    /// Size of the index file.
229    pub index_file_size: u64,
230    /// Version of the index file.
231    /// Used to generate the index file name: "{file_id}.{index_version}.puffin".
232    /// Default is 0 (which maps to "{file_id}.puffin" for compatibility).
233    pub index_version: u64,
234    /// Number of rows in the file.
235    ///
236    /// For historical reasons, this field might be missing in old files. Thus
237    /// the default value `0` doesn't means the file doesn't contains any rows,
238    /// but instead means the number of rows is unknown.
239    pub num_rows: u64,
240    /// Number of row groups in the file.
241    ///
242    /// For historical reasons, this field might be missing in old files. Thus
243    /// the default value `0` doesn't means the file doesn't contains any rows,
244    /// but instead means the number of rows is unknown.
245    pub num_row_groups: u64,
246    /// Sequence in this file.
247    ///
248    /// This sequence is the only sequence in this file. And it's retrieved from the max
249    /// sequence of the rows on generating this file.
250    pub sequence: Option<NonZeroU64>,
251    /// Partition expression from the region metadata when the file is created.
252    ///
253    /// This is stored as a PartitionExpr object in memory for convenience,
254    /// but serialized as JSON string for manifest compatibility.
255    /// Compatibility behavior:
256    /// - None: no partition expr was set when the file was created (legacy files).
257    /// - Some(expr): partition expression from region metadata.
258    #[serde(
259        serialize_with = "serialize_partition_expr",
260        deserialize_with = "deserialize_partition_expr"
261    )]
262    pub partition_expr: Option<PartitionExpr>,
263    /// Number of series in the file.
264    ///
265    /// The number is 0 if the series number is not available.
266    pub num_series: u64,
267    /// Minimum primary key value in the file, encoded as bytes.
268    /// `None` if the primary key range is not available (e.g., legacy files).
269    #[serde(
270        default,
271        skip_serializing_if = "Option::is_none",
272        serialize_with = "serialize_bytes_option",
273        deserialize_with = "deserialize_bytes_option"
274    )]
275    pub primary_key_min: Option<Bytes>,
276    /// Maximum primary key value in the file, encoded as bytes.
277    /// `None` if the primary key range is not available (e.g., legacy files).
278    #[serde(
279        default,
280        skip_serializing_if = "Option::is_none",
281        serialize_with = "serialize_bytes_option",
282        deserialize_with = "deserialize_bytes_option"
283    )]
284    pub primary_key_max: Option<Bytes>,
285    /// Whether the file preserves per-row sequence numbers usable for exact
286    /// row-level sequence filtering.
287    #[serde(default, skip_serializing_if = "is_false")]
288    pub preserve_row_sequence: bool,
289}
290
291fn is_false(value: &bool) -> bool {
292    !*value
293}
294
295impl Debug for FileMeta {
296    fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
297        let mut debug_struct = f.debug_struct("FileMeta");
298        debug_struct
299            .field("region_id", &self.region_id)
300            .field_with("file_id", |f| write!(f, "{} ", self.file_id))
301            .field_with("time_range", |f| {
302                write!(
303                    f,
304                    "({}, {}) ",
305                    self.time_range.0.to_iso8601_string(),
306                    self.time_range.1.to_iso8601_string()
307                )
308            })
309            .field("level", &self.level)
310            .field("file_size", &ReadableSize(self.file_size))
311            .field(
312                "max_row_group_uncompressed_size",
313                &ReadableSize(self.max_row_group_uncompressed_size),
314            );
315        if !self.available_indexes.is_empty() {
316            debug_struct
317                .field("available_indexes", &self.available_indexes)
318                .field("indexes", &self.indexes)
319                .field("index_file_size", &ReadableSize(self.index_file_size));
320        }
321        debug_struct
322            .field("num_rows", &self.num_rows)
323            .field("num_row_groups", &self.num_row_groups)
324            .field_with("sequence", |f| match self.sequence {
325                None => {
326                    write!(f, "None")
327                }
328                Some(seq) => {
329                    write!(f, "{}", seq)
330                }
331            })
332            .field("partition_expr", &self.partition_expr)
333            .field("num_series", &self.num_series);
334        if self.primary_key_min.is_some() || self.primary_key_max.is_some() {
335            debug_struct
336                .field(
337                    "primary_key_min",
338                    &self.primary_key_min.as_ref().map(|b| b.len()),
339                )
340                .field(
341                    "primary_key_max",
342                    &self.primary_key_max.as_ref().map(|b| b.len()),
343                );
344        }
345        debug_struct.finish()
346    }
347}
348
349/// Type of index.
350#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
351pub enum IndexType {
352    /// Inverted index.
353    InvertedIndex,
354    /// Full-text index.
355    FulltextIndex,
356    /// Bloom Filter index
357    BloomFilterIndex,
358    /// Vector index (HNSW).
359    #[cfg(feature = "vector_index")]
360    VectorIndex,
361}
362
363/// Metadata of indexes created for a specific column in an SST file.
364///
365/// This structure tracks which index types have been successfully created for a column.
366/// It provides more granular, column-level index information compared to the file-level
367/// `available_indexes` field in [`FileMeta`].
368///
369/// This is primarily used for:
370/// - Manual index building in asynchronous index construction mode
371/// - Verifying index consistency between files and region metadata
372#[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize, Default)]
373#[serde(default)]
374pub struct ColumnIndexMetadata {
375    /// The column ID this index metadata applies to.
376    pub column_id: ColumnId,
377    /// List of index types that have been successfully created for this column.
378    pub created_indexes: IndexTypes,
379}
380
381impl FileMeta {
382    /// Returns the primary key range if both min and max are present.
383    pub fn primary_key_range(&self) -> Option<(Bytes, Bytes)> {
384        match (&self.primary_key_min, &self.primary_key_max) {
385            (Some(min), Some(max)) => Some((min.clone(), max.clone())),
386            _ => None,
387        }
388    }
389
390    pub fn exists_index(&self) -> bool {
391        !self.available_indexes.is_empty()
392    }
393
394    pub fn index_version(&self) -> Option<IndexVersion> {
395        if self.exists_index() {
396            Some(self.index_version)
397        } else {
398            None
399        }
400    }
401
402    /// Whether the index file is up-to-date comparing to another file meta.
403    pub fn is_index_up_to_date(&self, other: &FileMeta) -> bool {
404        self.exists_index() && other.exists_index() && self.index_version >= other.index_version
405    }
406
407    /// Returns true if the file has an inverted index
408    pub fn inverted_index_available(&self) -> bool {
409        self.available_indexes.contains(&IndexType::InvertedIndex)
410    }
411
412    /// Returns true if the file has a fulltext index
413    pub fn fulltext_index_available(&self) -> bool {
414        self.available_indexes.contains(&IndexType::FulltextIndex)
415    }
416
417    /// Returns true if the file has a bloom filter index.
418    pub fn bloom_filter_index_available(&self) -> bool {
419        self.available_indexes
420            .contains(&IndexType::BloomFilterIndex)
421    }
422
423    /// Returns true if the file has a vector index.
424    #[cfg(feature = "vector_index")]
425    pub fn vector_index_available(&self) -> bool {
426        self.available_indexes.contains(&IndexType::VectorIndex)
427    }
428
429    pub fn index_file_size(&self) -> u64 {
430        self.index_file_size
431    }
432
433    /// Check whether the file index is consistent with the given region metadata.
434    pub fn is_index_consistent_with_region(&self, metadata: &[ColumnMetadata]) -> bool {
435        let id_to_indexes = self
436            .indexes
437            .iter()
438            .map(|index| (index.column_id, index.created_indexes.clone()))
439            .collect::<std::collections::HashMap<_, _>>();
440        for column in metadata {
441            if !column.column_schema.is_indexed() {
442                continue;
443            }
444            if let Some(indexes) = id_to_indexes.get(&column.column_id) {
445                if column.column_schema.is_inverted_indexed()
446                    && !indexes.contains(&IndexType::InvertedIndex)
447                {
448                    return false;
449                }
450                if column.column_schema.is_fulltext_indexed()
451                    && !indexes.contains(&IndexType::FulltextIndex)
452                {
453                    return false;
454                }
455                if column.column_schema.is_skipping_indexed()
456                    && !indexes.contains(&IndexType::BloomFilterIndex)
457                {
458                    return false;
459                }
460            } else {
461                return false;
462            }
463        }
464        true
465    }
466
467    /// Returns the cross-region file id.
468    pub fn file_id(&self) -> RegionFileId {
469        RegionFileId::new(self.region_id, self.file_id)
470    }
471
472    /// Returns the RegionIndexId for this file.
473    pub fn index_id(&self) -> RegionIndexId {
474        RegionIndexId::new(self.file_id(), self.index_version)
475    }
476}
477
478/// Handle to a SST file.
479#[derive(Clone)]
480pub struct FileHandle {
481    inner: Arc<FileHandleInner>,
482}
483
484impl fmt::Debug for FileHandle {
485    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
486        f.debug_struct("FileHandle")
487            .field("meta", self.meta_ref())
488            .field("compacting", &self.compacting())
489            .field("deleted", &self.inner.deleted.load(Ordering::Relaxed))
490            .finish()
491    }
492}
493
494impl FileHandle {
495    pub fn new(meta: FileMeta, file_purger: FilePurgerRef) -> FileHandle {
496        let pk_range = meta.primary_key_range();
497        FileHandle {
498            inner: Arc::new(FileHandleInner::new(meta, file_purger, pk_range)),
499        }
500    }
501
502    #[cfg(test)]
503    pub fn new_with_primary_key_range(
504        meta: FileMeta,
505        file_purger: FilePurgerRef,
506        primary_key_range: Option<(Bytes, Bytes)>,
507    ) -> FileHandle {
508        FileHandle {
509            inner: Arc::new(FileHandleInner::new(meta, file_purger, primary_key_range)),
510        }
511    }
512
513    /// Returns the region id of the file.
514    pub fn region_id(&self) -> RegionId {
515        self.inner.meta.region_id
516    }
517
518    /// Returns whether this file's row sequences are trusted in the target region.
519    ///
520    /// Foreign files use their target-local sequence barrier; local files require
521    /// the preserve marker because their physical sequences belong to this region.
522    pub(crate) fn is_effective_target_sequence_trusted(&self, target_region_id: RegionId) -> bool {
523        if self.region_id() != target_region_id {
524            self.meta_ref().sequence.is_some()
525        } else {
526            self.meta_ref().preserve_row_sequence
527        }
528    }
529
530    /// Returns the cross-region file id.
531    pub fn file_id(&self) -> RegionFileId {
532        RegionFileId::new(self.inner.meta.region_id, self.inner.meta.file_id)
533    }
534
535    /// Returns the RegionIndexId for this file.
536    pub fn index_id(&self) -> RegionIndexId {
537        RegionIndexId::new(self.file_id(), self.inner.meta.index_version)
538    }
539
540    /// Returns the complete file path of the file.
541    pub fn file_path(&self, table_dir: &str, path_type: PathType) -> String {
542        location::sst_file_path(table_dir, self.file_id(), path_type)
543    }
544
545    /// Returns the time range of the file.
546    pub fn time_range(&self) -> FileTimeRange {
547        self.inner.meta.time_range
548    }
549
550    /// Mark the file as deleted and will delete it on drop asynchronously
551    pub fn mark_deleted(&self) {
552        self.inner.deleted.store(true, Ordering::Relaxed);
553    }
554
555    pub fn compacting(&self) -> bool {
556        self.inner.compacting.load(Ordering::Relaxed)
557    }
558
559    pub fn set_compacting(&self, compacting: bool) {
560        self.inner.compacting.store(compacting, Ordering::Relaxed);
561    }
562
563    /// Atomically marks this file as compacting if it is currently available.
564    pub fn try_set_compacting(&self) -> bool {
565        self.inner
566            .compacting
567            .compare_exchange(false, true, Ordering::Relaxed, Ordering::Relaxed)
568            .is_ok()
569    }
570
571    pub fn index_outdated(&self) -> bool {
572        self.inner.index_outdated.load(Ordering::Relaxed)
573    }
574
575    pub fn set_index_outdated(&self, index_outdated: bool) {
576        self.inner
577            .index_outdated
578            .store(index_outdated, Ordering::Relaxed);
579    }
580
581    /// Returns a reference to the [FileMeta].
582    pub fn meta_ref(&self) -> &FileMeta {
583        &self.inner.meta
584    }
585
586    pub fn file_purger(&self) -> FilePurgerRef {
587        self.inner.file_purger.clone()
588    }
589
590    pub fn size(&self) -> u64 {
591        self.inner.meta.file_size
592    }
593
594    pub fn index_size(&self) -> u64 {
595        self.inner.meta.index_file_size
596    }
597
598    pub fn num_rows(&self) -> usize {
599        self.inner.meta.num_rows as usize
600    }
601
602    pub fn level(&self) -> Level {
603        self.inner.meta.level
604    }
605
606    pub fn is_deleted(&self) -> bool {
607        self.inner.deleted.load(Ordering::Relaxed)
608    }
609
610    pub fn primary_key_range(&self) -> Option<(Bytes, Bytes)> {
611        self.inner.primary_key_range.read().unwrap().clone()
612    }
613
614    pub(crate) fn set_primary_key_range(&self, primary_key_range: (Bytes, Bytes)) {
615        *self.inner.primary_key_range.write().unwrap() = Some(primary_key_range);
616    }
617}
618
619/// Inner data of [FileHandle].
620///
621/// Contains meta of the file, and other mutable info like whether the file is compacting.
622struct FileHandleInner {
623    meta: FileMeta,
624    compacting: AtomicBool,
625    deleted: AtomicBool,
626    index_outdated: AtomicBool,
627    primary_key_range: RwLock<Option<(Bytes, Bytes)>>,
628    file_purger: FilePurgerRef,
629}
630
631impl Drop for FileHandleInner {
632    fn drop(&mut self) {
633        self.file_purger.remove_file(
634            self.meta.clone(),
635            self.deleted.load(Ordering::Acquire),
636            self.index_outdated.load(Ordering::Acquire),
637        );
638    }
639}
640
641impl FileHandleInner {
642    /// There should only be one `FileHandleInner` for each file on a datanode
643    fn new(
644        meta: FileMeta,
645        file_purger: FilePurgerRef,
646        primary_key_range: Option<(Bytes, Bytes)>,
647    ) -> FileHandleInner {
648        file_purger.new_file(&meta);
649        FileHandleInner {
650            meta,
651            compacting: AtomicBool::new(false),
652            deleted: AtomicBool::new(false),
653            index_outdated: AtomicBool::new(false),
654            primary_key_range: RwLock::new(primary_key_range),
655            file_purger,
656        }
657    }
658}
659
660/// Delete files for a region.
661/// - `region_id`: Region id.
662/// - `file_ids`: List of (file id, index version) tuples to delete.
663/// - `delete_index`: Whether to delete the index file from the cache.
664/// - `access_layer`: Access layer to delete files.
665/// - `cache_manager`: Cache manager to remove files from cache.
666pub async fn delete_files(
667    region_id: RegionId,
668    file_ids: &[(FileId, u64)],
669    delete_index: bool,
670    access_layer: &AccessLayerRef,
671    cache_manager: &Option<CacheManagerRef>,
672) -> crate::error::Result<()> {
673    // Remove meta of the file from cache.
674    if let Some(cache) = &cache_manager {
675        for (file_id, _) in file_ids {
676            cache.remove_parquet_meta_data(RegionFileId::new(region_id, *file_id));
677        }
678    }
679    let mut attempted_files = Vec::with_capacity(file_ids.len());
680    let mut index_ids = Vec::new();
681
682    for (file_id, index_version) in file_ids {
683        let region_file_id = RegionFileId::new(region_id, *file_id);
684        attempted_files.push(*file_id);
685        index_ids.extend(
686            (0..=*index_version).map(|version| RegionIndexId::new(region_file_id, version)),
687        );
688    }
689
690    access_layer
691        .delete_ssts(region_id, &attempted_files)
692        .await?;
693    access_layer.delete_indexes(&index_ids).await?;
694
695    debug!(
696        "Attempted to delete {} files for region {}: {:?}",
697        attempted_files.len(),
698        region_id,
699        attempted_files
700    );
701
702    for (file_id, index_version) in file_ids {
703        purge_index_cache_stager(
704            region_id,
705            delete_index,
706            access_layer,
707            cache_manager,
708            *file_id,
709            *index_version,
710        )
711        .await;
712    }
713    Ok(())
714}
715
716/// Tracks finalized SSTs until their manifest edit is committed.
717///
718/// Flush and local compaction jobs use this to remove files that were written successfully but
719/// abandoned by cancellation or a later failure.
720#[derive(Clone)]
721pub(crate) struct UncommittedSsts {
722    region_id: RegionId,
723    files: Arc<Mutex<HashMap<FileId, (u64, bool)>>>,
724    access_layer: AccessLayerRef,
725    cache_manager: Option<CacheManagerRef>,
726}
727
728impl UncommittedSsts {
729    pub(crate) fn new(
730        region_id: RegionId,
731        access_layer: AccessLayerRef,
732        cache_manager: Option<CacheManagerRef>,
733    ) -> Self {
734        Self {
735            region_id,
736            files: Arc::new(Mutex::new(HashMap::new())),
737            access_layer,
738            cache_manager,
739        }
740    }
741
742    /// Tracks newly finalized SSTs before they become visible in the manifest.
743    pub(crate) fn track(&self, ssts: &[SstInfo]) {
744        let mut files = self.files.lock().unwrap();
745        for sst in ssts {
746            files.insert(
747                sst.file_id,
748                (sst.index_metadata.version, sst.index_metadata.file_size > 0),
749            );
750        }
751    }
752
753    /// Disarms cleanup after the manifest edit has committed or may have committed.
754    ///
755    /// A manifest update error does not guarantee that the edit was not persisted. Once an edit
756    /// may become visible, its SSTs must be retained to avoid deleting referenced files.
757    pub(crate) fn disarm_cleanup(&self) {
758        self.files.lock().unwrap().clear();
759    }
760
761    #[cfg(test)]
762    pub(crate) fn num_tracked_files(&self) -> usize {
763        self.files.lock().unwrap().len()
764    }
765
766    /// Removes all finalized SSTs still owned by this job.
767    pub(crate) async fn cleanup(&self) {
768        if let Err(err) = self.try_cleanup().await {
769            error!(err; "Failed to clean uncommitted SSTs for region {}", self.region_id);
770        }
771    }
772
773    async fn try_cleanup(&self) -> crate::error::Result<()> {
774        let files = std::mem::take(&mut *self.files.lock().unwrap());
775        if files.is_empty() {
776            return Ok(());
777        }
778
779        let delete_index = files.values().any(|(_, exists_index)| *exists_index);
780        let file_ids = files
781            .into_iter()
782            .map(|(file_id, (index_version, _))| (file_id, index_version))
783            .collect::<Vec<_>>();
784        delete_files(
785            self.region_id,
786            &file_ids,
787            delete_index,
788            &self.access_layer,
789            &self.cache_manager,
790        )
791        .await
792    }
793
794    #[cfg(test)]
795    pub(crate) async fn cleanup_for_test(&self) -> crate::error::Result<()> {
796        self.try_cleanup().await
797    }
798}
799
800pub async fn delete_index(
801    region_index_id: RegionIndexId,
802    access_layer: &AccessLayerRef,
803    cache_manager: &Option<CacheManagerRef>,
804) -> crate::error::Result<()> {
805    delete_index_and_purge(region_index_id, access_layer, cache_manager).await?;
806
807    Ok(())
808}
809
810pub async fn delete_indexes(
811    index_ids: &[RegionIndexId],
812    access_layer: &AccessLayerRef,
813    cache_manager: &Option<CacheManagerRef>,
814) -> crate::error::Result<()> {
815    if index_ids.is_empty() {
816        return Ok(());
817    }
818
819    if let Err(e) = access_layer.delete_indexes(index_ids).await {
820        error!(e; "Failed to batch delete index files");
821
822        for index_id in index_ids {
823            delete_index_and_purge(*index_id, access_layer, cache_manager).await?;
824        }
825
826        return Ok(());
827    }
828
829    purge_indexes(index_ids, access_layer, cache_manager).await;
830
831    Ok(())
832}
833
834async fn delete_index_and_purge(
835    index_id: RegionIndexId,
836    access_layer: &AccessLayerRef,
837    cache_manager: &Option<CacheManagerRef>,
838) -> crate::error::Result<()> {
839    access_layer.delete_index(index_id).await?;
840    purge_index_cache_stager(
841        index_id.region_id(),
842        true,
843        access_layer,
844        cache_manager,
845        index_id.file_id(),
846        index_id.version,
847    )
848    .await;
849    Ok(())
850}
851
852async fn purge_indexes(
853    index_ids: &[RegionIndexId],
854    access_layer: &AccessLayerRef,
855    cache_manager: &Option<CacheManagerRef>,
856) {
857    for index_id in index_ids {
858        purge_index_cache_stager(
859            index_id.region_id(),
860            true,
861            access_layer,
862            cache_manager,
863            index_id.file_id(),
864            index_id.version,
865        )
866        .await;
867    }
868}
869
870async fn purge_index_cache_stager(
871    region_id: RegionId,
872    delete_index: bool,
873    access_layer: &AccessLayerRef,
874    cache_manager: &Option<CacheManagerRef>,
875    file_id: FileId,
876    index_version: u64,
877) {
878    if let Some(write_cache) = cache_manager.as_ref().and_then(|cache| cache.write_cache()) {
879        // Removes index file from the cache.
880        if delete_index {
881            write_cache
882                .remove(IndexKey::new(
883                    region_id,
884                    file_id,
885                    FileType::Puffin(index_version),
886                ))
887                .await;
888        }
889
890        // Remove the SST file from the cache.
891        write_cache
892            .remove(IndexKey::new(region_id, file_id, FileType::Parquet))
893            .await;
894    }
895
896    // Purges index content in the stager.
897    if let Err(e) = access_layer
898        .puffin_manager_factory()
899        .purge_stager(RegionIndexId::new(
900            RegionFileId::new(region_id, file_id),
901            index_version,
902        ))
903        .await
904    {
905        error!(e; "Failed to purge stager with index file, file_id: {}, index_version: {}, region: {}",
906                file_id, index_version, region_id);
907    }
908}
909
910#[cfg(test)]
911mod tests {
912    use std::str::FromStr;
913
914    use datatypes::prelude::ConcreteDataType;
915    use datatypes::schema::{
916        ColumnSchema, FulltextAnalyzer, FulltextBackend, FulltextOptions, SkippingIndexOptions,
917    };
918    use datatypes::value::Value;
919    use partition::expr::{PartitionExpr, col};
920
921    use super::*;
922
923    fn create_file_meta(file_id: FileId, level: Level) -> FileMeta {
924        FileMeta {
925            region_id: 0.into(),
926            file_id,
927            time_range: FileTimeRange::default(),
928            level,
929            file_size: 0,
930            max_row_group_uncompressed_size: 0,
931            available_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
932            indexes: vec![ColumnIndexMetadata {
933                column_id: 0,
934                created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
935            }],
936            index_file_size: 0,
937            index_version: 0,
938            num_rows: 0,
939            num_row_groups: 0,
940            sequence: None,
941            partition_expr: None,
942            num_series: 0,
943            ..Default::default()
944        }
945    }
946
947    #[test]
948    fn test_deserialize_file_meta() {
949        let file_meta = create_file_meta(FileId::random(), 0);
950        let serialized_file_meta = serde_json::to_string(&file_meta).unwrap();
951        let deserialized_file_meta = serde_json::from_str(&serialized_file_meta);
952        assert_eq!(file_meta, deserialized_file_meta.unwrap());
953    }
954
955    #[test]
956    fn test_deserialize_from_string() {
957        let json_file_meta = "{\"region_id\":0,\"file_id\":\"bc5896ec-e4d8-4017-a80d-f2de73188d55\",\
958        \"time_range\":[{\"value\":0,\"unit\":\"Millisecond\"},{\"value\":0,\"unit\":\"Millisecond\"}],\
959        \"available_indexes\":[\"InvertedIndex\"],\"indexes\":[{\"column_id\": 0, \"created_indexes\": [\"InvertedIndex\"]}],\"level\":0}";
960        let file_meta = create_file_meta(
961            FileId::from_str("bc5896ec-e4d8-4017-a80d-f2de73188d55").unwrap(),
962            0,
963        );
964        let deserialized_file_meta: FileMeta = serde_json::from_str(json_file_meta).unwrap();
965        assert_eq!(file_meta, deserialized_file_meta);
966    }
967
968    #[test]
969    fn test_file_meta_with_partition_expr() {
970        let file_id = FileId::random();
971        let partition_expr = PartitionExpr::new(
972            col("a"),
973            partition::expr::RestrictedOp::GtEq,
974            Value::UInt32(10).into(),
975        );
976
977        let file_meta_with_partition = FileMeta {
978            region_id: 0.into(),
979            file_id,
980            time_range: FileTimeRange::default(),
981            level: 0,
982            file_size: 0,
983            max_row_group_uncompressed_size: 0,
984            available_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
985            indexes: vec![ColumnIndexMetadata {
986                column_id: 0,
987                created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
988            }],
989            index_file_size: 0,
990            index_version: 0,
991            num_rows: 0,
992            num_row_groups: 0,
993            sequence: None,
994            partition_expr: Some(partition_expr.clone()),
995            num_series: 0,
996            ..Default::default()
997        };
998
999        // Test serialization/deserialization
1000        let serialized = serde_json::to_string(&file_meta_with_partition).unwrap();
1001        let deserialized: FileMeta = serde_json::from_str(&serialized).unwrap();
1002        assert_eq!(file_meta_with_partition, deserialized);
1003
1004        // Verify the serialized JSON contains the expected partition expression string
1005        let serialized_value: serde_json::Value = serde_json::from_str(&serialized).unwrap();
1006        assert!(serialized_value["partition_expr"].as_str().is_some());
1007        let partition_expr_json = serialized_value["partition_expr"].as_str().unwrap();
1008        assert!(partition_expr_json.contains("\"Column\":\"a\""));
1009        assert!(partition_expr_json.contains("\"op\":\"GtEq\""));
1010
1011        // Test with None (legacy files)
1012        let file_meta_none = FileMeta {
1013            partition_expr: None,
1014            ..file_meta_with_partition.clone()
1015        };
1016        let serialized_none = serde_json::to_string(&file_meta_none).unwrap();
1017        let deserialized_none: FileMeta = serde_json::from_str(&serialized_none).unwrap();
1018        assert_eq!(file_meta_none, deserialized_none);
1019    }
1020
1021    #[test]
1022    fn test_file_meta_partition_expr_backward_compatibility() {
1023        // Test that we can deserialize old JSON format with partition_expr as string
1024        let json_with_partition_expr = r#"{
1025            "region_id": 0,
1026            "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1027            "time_range": [
1028                {"value": 0, "unit": "Millisecond"},
1029                {"value": 0, "unit": "Millisecond"}
1030            ],
1031            "level": 0,
1032            "file_size": 0,
1033            "available_indexes": ["InvertedIndex"],
1034            "index_file_size": 0,
1035            "num_rows": 0,
1036            "num_row_groups": 0,
1037            "sequence": null,
1038            "partition_expr": "{\"Expr\":{\"lhs\":{\"Column\":\"a\"},\"op\":\"GtEq\",\"rhs\":{\"Value\":{\"UInt32\":10}}}}"
1039        }"#;
1040
1041        let file_meta: FileMeta = serde_json::from_str(json_with_partition_expr).unwrap();
1042        assert!(file_meta.partition_expr.is_some());
1043        let expr = file_meta.partition_expr.unwrap();
1044        assert_eq!(format!("{}", expr), "a >= 10");
1045
1046        // Test empty partition expression string
1047        let json_with_empty_expr = r#"{
1048            "region_id": 0,
1049            "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1050            "time_range": [
1051                {"value": 0, "unit": "Millisecond"},
1052                {"value": 0, "unit": "Millisecond"}
1053            ],
1054            "level": 0,
1055            "file_size": 0,
1056            "available_indexes": [],
1057            "index_file_size": 0,
1058            "num_rows": 0,
1059            "num_row_groups": 0,
1060            "sequence": null,
1061            "partition_expr": ""
1062        }"#;
1063
1064        let file_meta_empty: FileMeta = serde_json::from_str(json_with_empty_expr).unwrap();
1065        assert!(file_meta_empty.partition_expr.is_none());
1066
1067        // Test null partition expression
1068        let json_with_null_expr = r#"{
1069            "region_id": 0,
1070            "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1071            "time_range": [
1072                {"value": 0, "unit": "Millisecond"},
1073                {"value": 0, "unit": "Millisecond"}
1074            ],
1075            "level": 0,
1076            "file_size": 0,
1077            "available_indexes": [],
1078            "index_file_size": 0,
1079            "num_rows": 0,
1080            "num_row_groups": 0,
1081            "sequence": null,
1082            "partition_expr": null
1083        }"#;
1084
1085        let file_meta_null: FileMeta = serde_json::from_str(json_with_null_expr).unwrap();
1086        assert!(file_meta_null.partition_expr.is_none());
1087
1088        // Test partition expression doesn't exist
1089        let json_with_empty_expr = r#"{
1090            "region_id": 0,
1091            "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1092            "time_range": [
1093                {"value": 0, "unit": "Millisecond"},
1094                {"value": 0, "unit": "Millisecond"}
1095            ],
1096            "level": 0,
1097            "file_size": 0,
1098            "available_indexes": [],
1099            "index_file_size": 0,
1100            "num_rows": 0,
1101            "num_row_groups": 0,
1102            "sequence": null
1103        }"#;
1104
1105        let file_meta_empty: FileMeta = serde_json::from_str(json_with_empty_expr).unwrap();
1106        assert!(file_meta_empty.partition_expr.is_none());
1107    }
1108
1109    #[test]
1110    fn test_file_meta_indexes_backward_compatibility() {
1111        // Old FileMeta format without the 'indexes' field
1112        let json_old_file_meta = r#"{
1113            "region_id": 0,
1114            "file_id": "bc5896ec-e4d8-4017-a80d-f2de73188d55",
1115            "time_range": [
1116                {"value": 0, "unit": "Millisecond"},
1117                {"value": 0, "unit": "Millisecond"}
1118            ],
1119            "available_indexes": ["InvertedIndex"],
1120            "level": 0,
1121            "file_size": 0,
1122            "index_file_size": 0,
1123            "num_rows": 0,
1124            "num_row_groups": 0
1125        }"#;
1126
1127        let deserialized_file_meta: FileMeta = serde_json::from_str(json_old_file_meta).unwrap();
1128
1129        // Verify backward compatibility: indexes field should default to empty vec
1130        assert_eq!(deserialized_file_meta.indexes, vec![]);
1131
1132        let expected_indexes: IndexTypes = SmallVec::from_iter([IndexType::InvertedIndex]);
1133        assert_eq!(deserialized_file_meta.available_indexes, expected_indexes);
1134
1135        assert_eq!(
1136            deserialized_file_meta.file_id,
1137            FileId::from_str("bc5896ec-e4d8-4017-a80d-f2de73188d55").unwrap()
1138        );
1139        assert!(!deserialized_file_meta.preserve_row_sequence);
1140    }
1141
1142    #[test]
1143    fn test_file_meta_preserve_row_sequence_serde() {
1144        let file_meta = FileMeta {
1145            preserve_row_sequence: true,
1146            ..Default::default()
1147        };
1148
1149        let serialized = serde_json::to_string(&file_meta).unwrap();
1150        let value: serde_json::Value = serde_json::from_str(&serialized).unwrap();
1151        assert_eq!(value["preserve_row_sequence"], true);
1152
1153        let deserialized: FileMeta = serde_json::from_str(&serialized).unwrap();
1154        assert_eq!(file_meta, deserialized);
1155
1156        let file_meta_false = FileMeta {
1157            preserve_row_sequence: false,
1158            ..file_meta.clone()
1159        };
1160        let serialized_false = serde_json::to_string(&file_meta_false).unwrap();
1161        let value_false: serde_json::Value = serde_json::from_str(&serialized_false).unwrap();
1162        assert!(value_false.get("preserve_row_sequence").is_none());
1163    }
1164    #[test]
1165    fn test_is_index_consistent_with_region() {
1166        fn new_column_meta(
1167            id: ColumnId,
1168            name: &str,
1169            inverted: bool,
1170            fulltext: bool,
1171            skipping: bool,
1172        ) -> ColumnMetadata {
1173            let mut column_schema =
1174                ColumnSchema::new(name, ConcreteDataType::string_datatype(), true);
1175            if inverted {
1176                column_schema = column_schema.with_inverted_index(true);
1177            }
1178            if fulltext {
1179                column_schema = column_schema
1180                    .with_fulltext_options(FulltextOptions::new_unchecked(
1181                        true,
1182                        FulltextAnalyzer::English,
1183                        false,
1184                        FulltextBackend::Bloom,
1185                        1000,
1186                        0.01,
1187                    ))
1188                    .unwrap();
1189            }
1190            if skipping {
1191                column_schema = column_schema
1192                    .with_skipping_options(SkippingIndexOptions::new_unchecked(
1193                        1024,
1194                        0.01,
1195                        datatypes::schema::SkippingIndexType::BloomFilter,
1196                    ))
1197                    .unwrap();
1198            }
1199
1200            ColumnMetadata {
1201                column_schema,
1202                semantic_type: api::v1::SemanticType::Tag,
1203                column_id: id,
1204            }
1205        }
1206
1207        // Case 1: Perfect match. File has exactly the required indexes.
1208        let mut file_meta = FileMeta {
1209            indexes: vec![ColumnIndexMetadata {
1210                column_id: 1,
1211                created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1212            }],
1213            ..Default::default()
1214        };
1215        let region_meta = vec![new_column_meta(1, "tag1", true, false, false)];
1216        assert!(file_meta.is_index_consistent_with_region(&region_meta));
1217
1218        // Case 2: Superset match. File has more indexes than required.
1219        file_meta.indexes = vec![ColumnIndexMetadata {
1220            column_id: 1,
1221            created_indexes: SmallVec::from_iter([
1222                IndexType::InvertedIndex,
1223                IndexType::BloomFilterIndex,
1224            ]),
1225        }];
1226        let region_meta = vec![new_column_meta(1, "tag1", true, false, false)];
1227        assert!(file_meta.is_index_consistent_with_region(&region_meta));
1228
1229        // Case 3: Missing index type. File has the column but lacks the required index type.
1230        file_meta.indexes = vec![ColumnIndexMetadata {
1231            column_id: 1,
1232            created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1233        }];
1234        let region_meta = vec![new_column_meta(1, "tag1", true, true, false)]; // Requires fulltext too
1235        assert!(!file_meta.is_index_consistent_with_region(&region_meta));
1236
1237        // Case 4: Missing column. Region requires an index on a column not in the file's index list.
1238        file_meta.indexes = vec![ColumnIndexMetadata {
1239            column_id: 2, // File only has index for column 2
1240            created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1241        }];
1242        let region_meta = vec![new_column_meta(1, "tag1", true, false, false)]; // Requires index on column 1
1243        assert!(!file_meta.is_index_consistent_with_region(&region_meta));
1244
1245        // Case 5: No indexes required by region. Should always be consistent.
1246        file_meta.indexes = vec![ColumnIndexMetadata {
1247            column_id: 1,
1248            created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1249        }];
1250        let region_meta = vec![new_column_meta(1, "tag1", false, false, false)]; // No index required
1251        assert!(file_meta.is_index_consistent_with_region(&region_meta));
1252
1253        // Case 6: Empty file indexes. Region requires an index.
1254        file_meta.indexes = vec![];
1255        let region_meta = vec![new_column_meta(1, "tag1", true, false, false)];
1256        assert!(!file_meta.is_index_consistent_with_region(&region_meta));
1257
1258        // Case 7: Multiple columns, one is inconsistent.
1259        file_meta.indexes = vec![
1260            ColumnIndexMetadata {
1261                column_id: 1,
1262                created_indexes: SmallVec::from_iter([IndexType::InvertedIndex]),
1263            },
1264            ColumnIndexMetadata {
1265                column_id: 2, // Column 2 is missing the required BloomFilterIndex
1266                created_indexes: SmallVec::from_iter([IndexType::FulltextIndex]),
1267            },
1268        ];
1269        let region_meta = vec![
1270            new_column_meta(1, "tag1", true, false, false),
1271            new_column_meta(2, "tag2", false, true, true), // Requires Fulltext and BloomFilter
1272        ];
1273        assert!(!file_meta.is_index_consistent_with_region(&region_meta));
1274    }
1275}