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