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