1use common_telemetry::{debug, info, warn};
18use store_api::logstore::LogStore;
19use store_api::region_request::RegionTruncateRequest;
20use store_api::storage::RegionId;
21
22use crate::error::RegionNotFoundSnafu;
23use crate::manifest::action::{RegionTruncate, TruncateKind};
24use crate::region::RegionLeaderState;
25use crate::request::{DdlRequest, DiscardUnflushedResult, OptionOutputTx, TruncateResult};
26use crate::worker::RegionWorkerLoop;
27
28impl<S: LogStore> RegionWorkerLoop<S> {
29 pub(crate) async fn handle_truncate_request(
30 &mut self,
31 region_id: RegionId,
32 req: RegionTruncateRequest,
33 sender: OptionOutputTx,
34 ) {
35 let region = match self.regions.writable_non_staging_region(region_id) {
36 Ok(region) => region,
37 Err(e) => {
38 sender.send(Err(e));
39 return;
40 }
41 };
42
43 let (sender, req) = match self.flush_scheduler.try_cancel_and_add_ddl(
44 region_id,
45 sender,
46 req,
47 DdlRequest::Truncate,
48 ) {
49 Ok(()) => {
50 self.listener.on_flush_cancel_requested(region_id);
51 return;
52 }
53 Err(request) => request,
54 };
55
56 let (sender, req) = match self.compaction_scheduler.try_cancel_and_add_ddl(
57 region_id,
58 sender,
59 req,
60 DdlRequest::Truncate,
61 ) {
62 Ok(()) => {
63 self.listener.on_compaction_cancel_requested(region_id);
64 return;
65 }
66 Err(request) => request,
67 };
68
69 let version_data = region.version_control.current();
70
71 match req {
72 RegionTruncateRequest::All => {
73 info!("Try to fully truncate region {}", region_id);
74
75 let truncated_entry_id = version_data.last_entry_id;
76 let truncated_sequence = version_data.committed_sequence;
77
78 let truncate = RegionTruncate {
80 region_id,
81 kind: TruncateKind::All {
82 truncated_entry_id,
83 truncated_sequence,
84 },
85 timestamp_ms: None,
86 };
87
88 self.handle_manifest_truncate_action(region, truncate, sender);
89 }
90 RegionTruncateRequest::Unflushed => {
91 let memtables = &version_data.version.memtables;
92 if memtables.is_empty()
93 && version_data.version.flushed_entry_id == version_data.last_entry_id
94 && version_data.version.flushed_sequence == version_data.committed_sequence
95 {
96 self.update_topic_latest_entry_id(®ion);
97 let result = self
98 .wal
99 .obsolete(
100 region_id,
101 version_data.version.flushed_entry_id,
102 ®ion.provider,
103 )
104 .await
105 .map(|_| 0);
106 sender.send(result);
107 return;
108 }
109
110 let discarded_rows = memtables.num_rows();
111 let discarded_bytes =
112 (memtables.mutable_usage() + memtables.immutables_usage()) as u64;
113 warn!(
114 "Discarding unflushed data from region {}, entry_id: {}, sequence: {}, estimated rows: {}, estimated bytes: {}",
115 region_id,
116 version_data.last_entry_id,
117 version_data.committed_sequence,
118 discarded_rows,
119 discarded_bytes,
120 );
121
122 self.handle_manifest_discard_unflushed_action(
123 region,
124 version_data.last_entry_id,
125 version_data.committed_sequence,
126 discarded_rows,
127 discarded_bytes,
128 sender,
129 );
130 }
131 RegionTruncateRequest::ByTimeRanges { time_ranges } => {
132 info!(
133 "Try to partially truncate region {} by time ranges: {:?}",
134 region_id, time_ranges
135 );
136 let mut files_to_truncate = Vec::new();
138 for level in version_data.version.ssts.levels() {
139 for file in level.files() {
140 let file_time_range = file.time_range();
141 let is_subset = time_ranges.iter().any(|(start, end)| {
144 file_time_range.0 >= *start && file_time_range.1 <= *end
145 });
146
147 if is_subset {
148 files_to_truncate.push(file.meta_ref().clone());
149 }
150 }
151 }
152 debug!(
153 "Found {} files to partially truncate in region {}",
154 files_to_truncate.len(),
155 region_id
156 );
157
158 if files_to_truncate.is_empty() {
160 info!("No files to truncate in region {}", region_id);
161 sender.send(Ok(0));
164 return;
165 }
166 let truncate = RegionTruncate {
167 region_id,
168 kind: TruncateKind::Partial {
169 files_to_remove: files_to_truncate,
170 },
171 timestamp_ms: None,
172 };
173 self.handle_manifest_truncate_action(region, truncate, sender);
174 }
175 };
176 }
177
178 pub(crate) async fn handle_truncate_result(&mut self, truncate_result: TruncateResult) {
180 let region_id = truncate_result.region_id;
181 let Some(region) = self.regions.get_region(region_id) else {
182 truncate_result.sender.send(
183 RegionNotFoundSnafu {
184 region_id: truncate_result.region_id,
185 }
186 .fail(),
187 );
188 return;
189 };
190
191 region.switch_state_to_writable(RegionLeaderState::Truncating);
193
194 match truncate_result.result {
195 Ok(()) => {
196 region
198 .version_control
199 .truncate(truncate_result.kind.clone());
200 }
201 Err(e) => {
202 truncate_result.sender.send(Err(e));
204 return;
205 }
206 }
207
208 self.flush_scheduler.on_region_truncated(region_id);
210 self.compaction_scheduler.on_region_truncated(region_id);
212 self.index_build_scheduler
214 .on_region_truncated(region_id)
215 .await;
216
217 if let TruncateKind::All {
218 truncated_entry_id,
219 truncated_sequence: _,
220 } = &truncate_result.kind
221 {
222 if let Err(e) = self
224 .wal
225 .obsolete(region_id, *truncated_entry_id, ®ion.provider)
226 .await
227 {
228 truncate_result.sender.send(Err(e));
229 return;
230 }
231 }
232
233 info!(
234 "Complete truncating region: {}, kind: {:?}.",
235 region_id, truncate_result.kind
236 );
237
238 truncate_result.sender.send(Ok(0));
239 }
240
241 pub(crate) async fn handle_discard_unflushed_result(&mut self, result: DiscardUnflushedResult) {
243 let region_id = result.region_id;
244 let Some(region) = self.regions.get_region(region_id) else {
245 result.sender.send(RegionNotFoundSnafu { region_id }.fail());
246 return;
247 };
248
249 region.switch_state_to_writable(RegionLeaderState::Truncating);
250
251 if let Err(e) = result.result {
252 result.sender.send(Err(e));
253 return;
254 }
255
256 region
257 .version_control
258 .discard_unflushed(result.discarded_entry_id, result.discarded_sequence);
259
260 self.flush_scheduler.on_region_truncated(region_id);
266 self.compaction_scheduler.on_region_truncated(region_id);
267 self.index_build_scheduler
268 .on_region_truncated(region_id)
269 .await;
270 self.update_topic_latest_entry_id(®ion);
271
272 self.notify_group();
276 self.handle_stalled_requests().await;
277
278 if let Err(e) = self
279 .wal
280 .obsolete(region_id, result.discarded_entry_id, ®ion.provider)
281 .await
282 {
283 result.sender.send(Err(e));
284 return;
285 }
286
287 warn!(
288 "Discarded unflushed data from region {}, entry_id: {}, sequence: {}, estimated rows: {}, estimated bytes: {}",
289 region_id,
290 result.discarded_entry_id,
291 result.discarded_sequence,
292 result.discarded_rows,
293 result.discarded_bytes,
294 );
295 result.sender.send(Ok(0));
296 }
297}