1use std::fmt::Debug;
16use std::sync::Arc;
17
18use api::v1::region::compact_request;
19use common_time::range::TimestampRange;
20use common_time::{TimeToLive, Timestamp};
21use serde::{Deserialize, Serialize};
22
23use crate::compaction::compactor::CompactionRegion;
24use crate::compaction::last_non_null::LastNonNullPicker;
25use crate::compaction::twcs::TwcsPicker;
26use crate::compaction::window::WindowedCompactionPicker;
27use crate::compaction::{CompactionOutput, SerializedCompactionOutput};
28use crate::error::Result;
29use crate::region::options::{CompactionOptions, MergeMode, RegionOptions};
30use crate::sst::file::{FileHandle, FileMeta};
31use crate::sst::file_purger::FilePurger;
32use crate::sst::version::LevelMeta;
33
34#[async_trait::async_trait]
35pub(crate) trait CompactionTask: Debug + Send + Sync + 'static {
36 async fn run(&mut self);
37}
38
39#[async_trait::async_trait]
42pub trait Picker: Debug + Send + Sync + 'static {
43 async fn pick(&self, compaction_region: &CompactionRegion) -> Result<Option<PickerOutput>>;
45}
46
47#[derive(Default, Clone, Debug)]
50pub struct PickerOutput {
51 pub outputs: Vec<CompactionOutput>,
52 pub expired_ssts: Vec<FileHandle>,
53 pub time_window_size: i64,
54 pub max_file_size: Option<usize>,
56}
57
58#[derive(Default, Clone, Debug, Serialize, Deserialize)]
60pub struct SerializedPickerOutput {
61 pub outputs: Vec<SerializedCompactionOutput>,
62 pub expired_ssts: Vec<FileMeta>,
63 pub time_window_size: i64,
64 pub max_file_size: Option<usize>,
65}
66
67impl From<&PickerOutput> for SerializedPickerOutput {
68 fn from(input: &PickerOutput) -> Self {
69 let outputs = input
70 .outputs
71 .iter()
72 .map(|output| SerializedCompactionOutput {
73 output_level: output.output_level,
74 inputs: output.inputs.iter().map(|s| s.meta_ref().clone()).collect(),
75 filter_deleted: output.filter_deleted,
76 output_time_range: output.output_time_range,
77 })
78 .collect();
79 let expired_ssts = input
80 .expired_ssts
81 .iter()
82 .map(|s| s.meta_ref().clone())
83 .collect();
84 Self {
85 outputs,
86 expired_ssts,
87 time_window_size: input.time_window_size,
88 max_file_size: input.max_file_size,
89 }
90 }
91}
92
93impl PickerOutput {
94 pub fn from_serialized(
97 input: SerializedPickerOutput,
98 file_purger: Arc<dyn FilePurger>,
99 ) -> Self {
100 let outputs = input
101 .outputs
102 .into_iter()
103 .map(|output| CompactionOutput {
104 output_level: output.output_level,
105 inputs: output
106 .inputs
107 .into_iter()
108 .map(|file_meta| FileHandle::new(file_meta, file_purger.clone()))
109 .collect(),
110 filter_deleted: output.filter_deleted,
111 output_time_range: output.output_time_range,
112 })
113 .collect();
114
115 let expired_ssts = input
116 .expired_ssts
117 .into_iter()
118 .map(|file_meta| FileHandle::new(file_meta, file_purger.clone()))
119 .collect();
120
121 Self {
122 outputs,
123 expired_ssts,
124 time_window_size: input.time_window_size,
125 max_file_size: input.max_file_size,
126 }
127 }
128}
129
130pub fn new_picker(
132 compact_request_options: &compact_request::Options,
133 region_options: &RegionOptions,
134 max_background_tasks: Option<usize>,
135 time_range: Option<TimestampRange>,
136) -> Arc<dyn Picker> {
137 if let compact_request::Options::StrictWindow(window) = compact_request_options {
138 let window = if window.window_seconds == 0 {
139 None
140 } else {
141 Some(window.window_seconds)
142 };
143 Arc::new(WindowedCompactionPicker::new(window).with_time_range(time_range)) as Arc<_>
146 } else {
147 let picker = match ®ion_options.compaction {
148 CompactionOptions::Twcs(twcs_opts) => TwcsPicker {
149 trigger_file_num: twcs_opts.active_window_trigger_file_num,
150 active_window_l1_merge_trigger: twcs_opts.active_window_l1_merge_trigger,
151 inactive_window_trigger_file_num: twcs_opts.inactive_window_trigger_file_num,
152 inactive_window_l1_merge_trigger: twcs_opts.inactive_window_l1_merge_trigger,
153 time_window_seconds: twcs_opts.time_window_seconds(),
154 max_output_file_size: twcs_opts.max_output_file_size.map(|r| r.as_bytes()),
155 append_mode: region_options.append_mode,
156 max_background_tasks,
157 time_range,
158 },
159 };
160 if region_options.merge_mode() == MergeMode::LastNonNull {
161 Arc::new(LastNonNullPicker::new(picker))
164 } else {
165 Arc::new(picker)
166 }
167 }
168}
169
170pub(super) fn get_expired_ssts(
172 levels: &[LevelMeta],
173 ttl: Option<TimeToLive>,
174 now: Timestamp,
175) -> Vec<FileHandle> {
176 let Some(ttl) = ttl else {
177 return vec![];
178 };
179
180 levels
181 .iter()
182 .flat_map(|l| l.get_expired_files(&now, &ttl).into_iter())
183 .collect()
184}
185
186#[cfg(test)]
187mod tests {
188 use store_api::storage::FileId;
189
190 use super::*;
191 use crate::compaction::test_util::new_file_handle;
192 use crate::test_util::new_noop_file_purger;
193
194 #[test]
195 fn test_picker_output_serialization() {
196 let inputs_file_handle = vec![
197 new_file_handle(FileId::random(), 0, 999, 0),
198 new_file_handle(FileId::random(), 0, 999, 0),
199 new_file_handle(FileId::random(), 0, 999, 0),
200 ];
201 let expired_ssts_file_handle = vec![
202 new_file_handle(FileId::random(), 0, 999, 0),
203 new_file_handle(FileId::random(), 0, 999, 0),
204 ];
205
206 let picker_output = PickerOutput {
207 outputs: vec![
208 CompactionOutput {
209 output_level: 0,
210 inputs: inputs_file_handle.clone(),
211 filter_deleted: false,
212 output_time_range: None,
213 },
214 CompactionOutput {
215 output_level: 0,
216 inputs: inputs_file_handle.clone(),
217 filter_deleted: false,
218 output_time_range: None,
219 },
220 ],
221 expired_ssts: expired_ssts_file_handle.clone(),
222 time_window_size: 1000,
223 max_file_size: None,
224 };
225
226 let picker_output_str =
227 serde_json::to_string(&SerializedPickerOutput::from(&picker_output)).unwrap();
228 let serialized_picker_output: SerializedPickerOutput =
229 serde_json::from_str(&picker_output_str).unwrap();
230 let picker_output_from_serialized =
231 PickerOutput::from_serialized(serialized_picker_output, new_noop_file_purger());
232
233 picker_output
234 .expired_ssts
235 .iter()
236 .zip(picker_output_from_serialized.expired_ssts.iter())
237 .for_each(|(expected, actual)| {
238 assert_eq!(expected.meta_ref(), actual.meta_ref());
239 });
240
241 picker_output
242 .outputs
243 .iter()
244 .zip(picker_output_from_serialized.outputs.iter())
245 .for_each(|(expected, actual)| {
246 assert_eq!(expected.output_level, actual.output_level);
247 expected
248 .inputs
249 .iter()
250 .zip(actual.inputs.iter())
251 .for_each(|(expected, actual)| {
252 assert_eq!(expected.meta_ref(), actual.meta_ref());
253 });
254 assert_eq!(expected.filter_deleted, actual.filter_deleted);
255 assert_eq!(expected.output_time_range, actual.output_time_range);
256 });
257 }
258}