1mod buckets;
16pub mod compactor;
17mod json2;
18mod last_non_null;
19pub mod memory_manager;
20mod overlap;
21pub mod picker;
22mod reader;
23pub mod run;
24mod scheduler;
25mod task;
26#[cfg(test)]
27mod test_util;
28mod twcs;
29mod window;
30
31use std::collections::HashMap;
32
33use common_meta::key::SchemaMetadataManagerRef;
34use common_telemetry::{debug, error};
35use common_time::TimeToLive;
36use common_time::range::TimestampRange;
37pub(crate) use json2::{
38 Json2RewritePlans, collect_json2_rewrite_plans, rewrite_json2_batch, rewrite_json2_schema,
39};
40pub use scheduler::CompactionRequest;
41pub(crate) use scheduler::{
42 CompactionExecution, CompactionPickFinished, CompactionScheduler, CompactionTransition,
43};
44use serde::{Deserialize, Serialize};
45use snafu::ResultExt;
46use store_api::mito_engine_options::{TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM, TWCS_TRIGGER_FILE_NUM};
47use store_api::storage::RegionId;
48
49use crate::error::{GetSchemaMetadataSnafu, Result, TimeoutSnafu};
50use crate::sst::file::{FileHandle, FileMeta, Level};
51
52async fn find_dynamic_options(
54 region_id: RegionId,
55 region_options: &crate::region::options::RegionOptions,
56 schema_metadata_manager: &SchemaMetadataManagerRef,
57) -> Result<(crate::region::options::CompactionOptions, TimeToLive)> {
58 let table_id = region_id.table_id();
59 if let (true, Some(ttl)) = (region_options.compaction_override, region_options.ttl) {
60 debug!(
61 "Use region options directly for table {}: compaction={:?}, ttl={:?}",
62 table_id, region_options.compaction, region_options.ttl
63 );
64 return Ok((region_options.compaction.clone(), ttl));
65 }
66
67 let db_options = tokio::time::timeout(
68 crate::config::FETCH_OPTION_TIMEOUT,
69 schema_metadata_manager.get_schema_options_by_table_id(table_id),
70 )
71 .await
72 .context(TimeoutSnafu)?
73 .context(GetSchemaMetadataSnafu)?;
74
75 let ttl = if let Some(ttl) = region_options.ttl {
76 debug!(
77 "Use region TTL directly for table {}: ttl={:?}",
78 table_id, region_options.ttl
79 );
80 ttl
81 } else {
82 db_options
83 .as_ref()
84 .and_then(|options| options.ttl)
85 .unwrap_or_default()
86 .into()
87 };
88
89 let compaction = if !region_options.compaction_override {
90 if let Some(schema_opts) = db_options {
91 let mut map: HashMap<String, String> = schema_opts
92 .extra_options
93 .iter()
94 .filter_map(|(k, v)| {
95 if k.starts_with("compaction.") {
96 Some((k.clone(), v.clone()))
97 } else {
98 None
99 }
100 })
101 .collect();
102 if map.contains_key(TWCS_ACTIVE_WINDOW_TRIGGER_FILE_NUM) {
104 map.remove(TWCS_TRIGGER_FILE_NUM);
105 }
106 if map.is_empty() {
107 region_options.compaction.clone()
108 } else {
109 crate::region::options::RegionOptions::try_from_options(region_id, &map)
110 .map(|o| o.compaction)
111 .unwrap_or_else(|e| {
112 error!(e; "Failed to create RegionOptions from map");
113 region_options.compaction.clone()
114 })
115 }
116 } else {
117 debug!(
118 "DB options is None for table {}, use region compaction: compaction={:?}",
119 table_id, region_options.compaction
120 );
121 region_options.compaction.clone()
122 }
123 } else {
124 debug!(
125 "No schema options for table {}, use region compaction: compaction={:?}",
126 table_id, region_options.compaction
127 );
128 region_options.compaction.clone()
129 };
130
131 debug!(
132 "Resolved dynamic options for table {}: compaction={:?}, ttl={:?}",
133 table_id, compaction, ttl
134 );
135 Ok((compaction, ttl))
136}
137
138#[derive(Debug, Clone)]
139pub struct CompactionOutput {
140 pub output_level: Level,
142 pub inputs: Vec<FileHandle>,
144 pub filter_deleted: bool,
146 pub output_time_range: Option<TimestampRange>,
148}
149
150#[derive(Debug, Clone, Serialize, Deserialize)]
152pub struct SerializedCompactionOutput {
153 output_level: Level,
154 inputs: Vec<FileMeta>,
155 filter_deleted: bool,
156 output_time_range: Option<TimestampRange>,
157}