1use std::collections::HashMap;
18use std::sync::Arc;
19
20use api::v1::region::build_index_request;
21use common_telemetry::{debug, error, warn};
22use store_api::region_request::RegionBuildIndexRequest;
23use store_api::storage::{FileId, RegionId};
24use tokio::sync::mpsc;
25
26use crate::cache::CacheStrategy;
27use crate::error::{InvalidRequestSnafu, Result, WorkerStoppedSnafu};
28use crate::manifest::action::RegionEdit;
29use crate::metrics::INDEX_PUBLICATION_STALE_TOTAL;
30use crate::region::version::VersionRef;
31use crate::region::{IndexBuildSource, MitoRegionRef};
32use crate::request::{
33 BuildIndexRequest, IndexBuildFailed, IndexBuildFinished, IndexBuildStopped, OptionOutputTx,
34};
35use crate::sst::file::{FileHandle, FileMeta, RegionFileId, RegionIndexId};
36use crate::sst::index::{
37 IndexBuildOutcome, IndexBuildTask, IndexBuildType, IndexerBuilderImpl, ResultMpscSender,
38 cleanup_stale_index_caches,
39};
40use crate::worker::RegionWorkerLoop;
41
42impl<S> RegionWorkerLoop<S> {
43 pub(crate) fn new_index_build_task(
44 &self,
45 region: &MitoRegionRef,
46 version: &VersionRef,
47 file: FileHandle,
48 file_meta: FileMeta,
49 build_type: IndexBuildType,
50 result_sender: ResultMpscSender,
51 ) -> IndexBuildTask {
52 let access_layer = region.access_layer.clone();
53
54 let puffin_manager = if let Some(write_cache) = self.cache_manager.write_cache() {
55 write_cache.build_puffin_manager()
56 } else {
57 access_layer.build_puffin_manager()
58 };
59
60 let intermediate_manager = if let Some(write_cache) = self.cache_manager.write_cache() {
61 write_cache.intermediate_manager().clone()
62 } else {
63 access_layer.intermediate_manager().clone()
64 };
65
66 let indexer_builder_ref = Arc::new(IndexerBuilderImpl {
67 build_type: build_type.clone(),
68 metadata: version.metadata.clone(),
69 inverted_index_config: self.config.inverted_index.clone(),
70 fulltext_index_config: self.config.fulltext_index.clone(),
71 bloom_filter_index_config: self.config.bloom_filter_index.clone(),
72 index_options: version.options.index_options.clone(),
73 intermediate_manager,
74 puffin_manager,
75 write_cache_enabled: self.cache_manager.write_cache().is_some(),
76 });
77
78 IndexBuildTask {
79 region_id: region.region_id,
80 file: file.clone(),
81 target_region_metadata: version.metadata.clone(),
82 source: IndexBuildSource::new(file_meta, version.metadata.schema_version),
83 reason: build_type,
84 access_layer: access_layer.clone(),
85 listener: self.listener.clone(),
86 manifest_ctx: region.manifest_ctx.clone(),
87 write_cache: self.cache_manager.write_cache().cloned(),
88 cache_manager: Some(self.cache_manager.clone()),
89 file_purger: file.file_purger(),
90 request_sender: self.sender.clone(),
91 indexer_builder: indexer_builder_ref.clone(),
92 result_sender,
93 }
94 }
95
96 pub(crate) async fn handle_build_index_request(
98 &mut self,
99 region_id: RegionId,
100 req: RegionBuildIndexRequest,
101 sender: OptionOutputTx,
102 ) {
103 match req.options {
104 Some(build_index_request::Options::SeriesIndex(_)) => {
105 self.handle_build_series_index_request(region_id, sender);
106 }
107 None | Some(build_index_request::Options::SstIndex(_)) => {
108 self.handle_rebuild_index(
109 BuildIndexRequest {
110 region_id,
111 build_type: IndexBuildType::Manual,
112 file_metas: Vec::new(),
113 },
114 sender,
115 )
116 .await;
117 }
118 }
119 }
120
121 fn handle_build_series_index_request(&self, region_id: RegionId, mut sender: OptionOutputTx) {
123 let Some(region) = self.regions.writable_region_or(region_id, &mut sender) else {
124 return;
125 };
126 let Some(state) = &self.series_index_task_state else {
127 sender.send(
128 InvalidRequestSnafu {
129 region_id,
130 reason: "series index is disabled; enable experimental_enable_series_index",
131 }
132 .fail(),
133 );
134 return;
135 };
136 let receiver = match state.try_reconcile(region) {
137 Ok(receiver) => receiver,
138 Err(error) => {
139 sender.send(Err(error));
140 return;
141 }
142 };
143 let worker_id = self.id;
144 common_runtime::spawn_global(async move {
145 let result = receiver
146 .await
147 .unwrap_or_else(|_| WorkerStoppedSnafu { id: worker_id }.fail());
148 sender.send(result.map(|_| 0));
149 });
150 }
151
152 pub(crate) async fn handle_rebuild_index(
153 &mut self,
154 request: BuildIndexRequest,
155 mut sender: OptionOutputTx,
156 ) {
157 let region_id = request.region_id;
158 let Some(region) = self.regions.writable_region_or(region_id, &mut sender) else {
159 return;
160 };
161
162 let version_control = region.version_control.clone();
163 let version = version_control.current().version;
164 let manifest = region.manifest_ctx.manifest().await;
168 let current_file_meta = |file_id| {
169 let file_meta = manifest.files.get(&file_id);
170 if file_meta.is_none() {
171 debug!(
172 "Skipping index build because file is absent from manifest, region: {}, file_id: {}",
173 region_id, file_id
174 );
175 }
176 file_meta.cloned()
177 };
178
179 let all_files: HashMap<FileId, FileHandle> = version
180 .ssts
181 .levels()
182 .iter()
183 .flat_map(|level| level.files.iter())
184 .filter(|(_, handle)| !handle.is_deleted() && !handle.compacting())
185 .map(|(id, handle)| (*id, handle.clone()))
186 .collect();
187
188 let build_tasks = if request.file_metas.is_empty() {
189 all_files
191 .values()
192 .filter_map(|file| {
193 let file_meta = current_file_meta(file.meta_ref().file_id)?;
194 (!file_meta.is_index_consistent_with_region(&version.metadata.column_metadatas))
195 .then(|| (file.clone(), file_meta))
196 })
197 .collect::<Vec<_>>()
198 } else {
199 request
200 .file_metas
201 .iter()
202 .filter_map(|meta| {
203 let file = all_files.get(&meta.file_id)?;
204 let file_meta = current_file_meta(meta.file_id)?;
205 Some((file.clone(), file_meta))
206 })
207 .collect::<Vec<_>>()
208 };
209
210 if build_tasks.is_empty() {
211 debug!(
212 "No files need to build index for region {}, request: {:?}",
213 region_id, request
214 );
215 sender.send(Ok(0));
216 return;
217 }
218
219 let num_tasks = build_tasks.len();
220 let (tx, mut rx) = mpsc::channel::<Result<IndexBuildOutcome>>(num_tasks);
221
222 for (file_handle, file_meta) in build_tasks {
223 debug!(
224 "Scheduling index build for region {}, file_id {}",
225 region_id,
226 file_handle.meta_ref().file_id
227 );
228
229 if region.should_abort_index() {
230 warn!(
231 "Region {} is in state {:?}, abort index rebuild process for file_id {}",
232 region_id,
233 region.state(),
234 file_handle.meta_ref().file_id
235 );
236 break;
237 }
238
239 let task = self.new_index_build_task(
240 ®ion,
241 &version,
242 file_handle.clone(),
243 file_meta,
244 request.build_type.clone(),
245 tx.clone(),
246 );
247 let _ = self
248 .index_build_scheduler
249 .schedule_build(®ion.version_control, task)
250 .await;
251 }
252 common_runtime::spawn_global(async move {
254 for _ in 0..num_tasks {
255 if let Some(Err(e)) = rx.recv().await {
256 warn!(e; "Index build task failed for region: {}", region_id);
257 sender.send(Err(e));
258 return;
259 }
260 }
261 sender.send(Ok(0));
262 });
263 }
264
265 pub(crate) async fn handle_index_build_finished(
266 &mut self,
267 region_id: RegionId,
268 request: IndexBuildFinished,
269 ) {
270 let region = match self.regions.get_region(region_id) {
271 Some(region) => region,
272 None => {
273 warn!(
274 "Region not found for index build finished, region_id: {}",
275 region_id
276 );
277 return;
278 }
279 };
280
281 let file_meta = request.file_meta;
282 let region_file_id = RegionFileId::new(file_meta.region_id, file_meta.file_id);
283 let index_id = RegionIndexId::new(region_file_id, file_meta.index_version);
284
285 let cache_strategy = CacheStrategy::EnableAll(self.cache_manager.clone());
287 cache_strategy.evict_puffin_cache(index_id).await;
288
289 let manager = region.manifest_ctx.manifest_manager.read().await;
290 let manifest = manager.manifest();
291 let is_current = manifest.files.get(&file_meta.file_id) == Some(&file_meta);
292 if is_current {
293 region.version_control.apply_edit(
294 Some(RegionEdit {
295 files_to_add: vec![file_meta.clone()],
296 files_to_remove: Vec::new(),
297 timestamp_ms: None,
298 flushed_sequence: None,
299 flushed_entry_id: None,
300 committed_sequence: None,
301 compaction_time_window: None,
302 }),
303 &[],
304 region.file_purger.clone(),
305 );
306 }
307 drop(manager);
308
309 if !is_current {
310 INDEX_PUBLICATION_STALE_TOTAL
311 .with_label_values(&["worker_apply"])
312 .inc();
313 warn!(
314 "Ignores stale index build result, region: {}, file_id: {}, index_version: {}, committed manifest version: {}",
315 region_id, file_meta.file_id, file_meta.index_version, request.manifest_version
316 );
317 cleanup_stale_index_caches(
318 index_id,
319 ®ion.access_layer,
320 Some(&self.cache_manager),
321 self.cache_manager.write_cache(),
322 )
323 .await;
324 self.listener.on_index_build_abort(region_file_id).await;
325 return;
326 }
327
328 self.listener.on_index_build_finish(region_file_id).await;
329 }
330
331 pub(crate) async fn handle_index_build_failed(
332 &mut self,
333 region_id: RegionId,
334 request: IndexBuildFailed,
335 ) {
336 error!(request.err; "Index build failed for region: {}", region_id);
337 self.index_build_scheduler
338 .on_failure(region_id, request.err.clone())
339 .await;
340 }
341
342 pub(crate) async fn handle_index_build_stopped(
343 &mut self,
344 region_id: RegionId,
345 request: IndexBuildStopped,
346 ) {
347 self.index_build_scheduler
348 .on_task_stopped(region_id, request.file_id);
349 }
350}