1use std::collections::{HashMap, VecDeque};
20use std::num::NonZeroU64;
21use std::sync::Arc;
22
23use common_telemetry::{debug, info, warn};
24use parquet::file::metadata::PageIndexPolicy;
25use snafu::ResultExt;
26use store_api::logstore::LogStore;
27use store_api::metadata::RegionMetadataRef;
28use store_api::storage::{RegionId, SequenceNumber};
29
30use crate::cache::CacheManagerRef;
31use crate::cache::file_cache::{FileType, IndexKey};
32use crate::config::IndexBuildMode;
33use crate::error::{EditRegionSnafu, RegionBusySnafu, RegionNotFoundSnafu, Result};
34use crate::manifest::action::{
35 RegionChange, RegionEdit, RegionMetaAction, RegionMetaActionList, RegionTruncate, TruncateKind,
36};
37use crate::memtable::MemtableBuilderProvider;
38use crate::metrics::WRITE_CACHE_INFLIGHT_DOWNLOAD;
39use crate::region::opener::{sanitize_region_options, version_builder_from_manifest};
40use crate::region::options::RegionOptions;
41use crate::region::version::VersionControlRef;
42use crate::region::{MitoRegionRef, RegionLeaderState, RegionRoleState};
43use crate::request::{
44 BackgroundNotify, BuildIndexRequest, DiscardUnflushedResult, OptionOutputTx,
45 RegionChangeResult, RegionEditRequest, RegionEditResult, RegionSyncRequest, TruncateResult,
46 WorkerRequest, WorkerRequestWithTime,
47};
48use crate::sst::index::IndexBuildType;
49use crate::sst::location;
50use crate::wal::EntryId;
51use crate::worker::{RegionWorkerLoop, WorkerListener};
52
53pub(crate) type RegionEditQueues = HashMap<RegionId, RegionEditQueue>;
54
55pub(crate) struct RegionEditQueue {
63 region_id: RegionId,
64 requests: VecDeque<RegionEditRequest>,
65}
66
67impl RegionEditQueue {
68 const QUEUE_MAX_LEN: usize = 128;
69
70 fn new(region_id: RegionId) -> Self {
71 Self {
72 region_id,
73 requests: VecDeque::new(),
74 }
75 }
76
77 fn enqueue(&mut self, request: RegionEditRequest) {
78 if self.requests.len() > Self::QUEUE_MAX_LEN {
79 request.waiters.reply_with(|| {
80 RegionBusySnafu {
81 region_id: self.region_id,
82 }
83 .fail()
84 });
85 return;
86 };
87 self.requests.push_back(request);
88 }
89
90 fn dequeue(&mut self) -> Option<RegionEditRequest> {
91 fn can_merge(edit: &RegionEdit) -> bool {
92 edit.files_to_add.iter().all(|f| f.sequence.is_none())
102 && edit.files_to_remove.is_empty()
103 && edit.timestamp_ms.is_none()
104 && edit.compaction_time_window.is_none()
105 && edit.flushed_entry_id.is_none()
106 && edit.flushed_sequence.is_none()
107 && edit.committed_sequence.is_none()
108 }
109
110 let mut merged = self.requests.pop_front()?;
111 if !can_merge(&merged.edit) {
112 return Some(merged);
113 }
114
115 while let Some(request) = self
116 .requests
117 .pop_front_if(|request| can_merge(&request.edit))
118 {
119 merged.edit.files_to_add.extend(request.edit.files_to_add);
120 merged.waiters.merge(request.waiters);
121 }
122 debug!(
123 "the files to add: [{}] are merged in one edit",
124 merged
125 .edit
126 .files_to_add
127 .iter()
128 .map(|x| x.file_id.to_string())
129 .collect::<Vec<_>>()
130 .join(", ")
131 );
132 Some(merged)
133 }
134
135 fn is_empty(&self) -> bool {
136 self.requests.is_empty()
137 }
138
139 fn reject_all_as_not_found(mut self) {
140 while let Some(request) = self.requests.pop_front() {
141 request.waiters.reply_with(|| {
142 RegionNotFoundSnafu {
143 region_id: self.region_id,
144 }
145 .fail()
146 });
147 }
148 }
149}
150
151impl<S: LogStore> RegionWorkerLoop<S> {
152 pub(crate) fn reject_region_edit_queue_as_not_found(&mut self, region_id: RegionId) {
154 if let Some(edit_queue) = self.region_edit_queues.remove(®ion_id) {
155 edit_queue.reject_all_as_not_found();
156 }
157 }
158
159 pub(crate) async fn handle_manifest_region_change_result(
161 &mut self,
162 change_result: RegionChangeResult,
163 ) {
164 let region = match self.regions.get_region(change_result.region_id) {
165 Some(region) => region,
166 None => {
167 self.reject_region_stalled_requests(&change_result.region_id);
168 change_result.sender.send(
169 RegionNotFoundSnafu {
170 region_id: change_result.region_id,
171 }
172 .fail(),
173 );
174 return;
175 }
176 };
177
178 if change_result.result.is_ok() {
179 Self::update_region_version(
181 ®ion.version_control,
182 change_result.new_meta,
183 change_result.new_options,
184 &self.memtable_builder_provider,
185 );
186 }
187
188 region.switch_state_to_writable(RegionLeaderState::Altering);
190 change_result.sender.send(change_result.result.map(|_| 0));
192
193 if self.config.index.build_mode == IndexBuildMode::Async && change_result.need_index {
195 self.handle_rebuild_index(
196 BuildIndexRequest {
197 region_id: region.region_id,
198 build_type: IndexBuildType::SchemaChange,
199 file_metas: Vec::new(),
200 },
201 OptionOutputTx::new(None),
202 )
203 .await;
204 }
205 self.handle_region_stalled_requests(&change_result.region_id, true)
207 .await;
208 }
209
210 pub(crate) async fn handle_region_sync(&mut self, request: RegionSyncRequest) {
215 let region_id = request.region_id;
216 let sender = request.sender;
217 let region = match self.regions.follower_region(region_id) {
218 Ok(region) => region,
219 Err(e) => {
220 let _ = sender.send(Err(e));
221 return;
222 }
223 };
224
225 let original_manifest_version = region.manifest_ctx.manifest_version().await;
226 let manifest = match region
227 .manifest_ctx
228 .install_manifest_to(request.manifest_version)
229 .await
230 {
231 Ok(manifest) => manifest,
232 Err(e) => {
233 let _ = sender.send(Err(e));
234 return;
235 }
236 };
237 let version = region.version();
238 let mut region_options = version.options.clone();
239 let old_format = region_options.sst_format.unwrap_or_default();
240 sanitize_region_options(&manifest, &mut region_options);
242 if !version.memtables.is_empty() {
243 let current = region.version_control.current();
244 warn!(
245 "Region {} memtables is not empty, which should not happen, manifest version: {}, last entry id: {}",
246 region.region_id, manifest.manifest_version, current.last_entry_id
247 );
248 }
249
250 let memtable_builder = if old_format != region_options.sst_format.unwrap_or_default() {
252 Some(
254 self.memtable_builder_provider
255 .builder_for_options(®ion_options),
256 )
257 } else {
258 None
259 };
260 let new_mutable = version
261 .memtables
262 .mutable
263 .new_with_part_duration(version.compaction_time_window, memtable_builder);
264 let metadata = manifest.metadata.clone();
266
267 let version_builder = version_builder_from_manifest(
268 &manifest,
269 metadata,
270 region.file_purger.clone(),
271 new_mutable,
272 region_options,
273 );
274 let version = version_builder.build();
275 region.version_control.overwrite_current(Arc::new(version));
276
277 let updated = manifest.manifest_version > original_manifest_version;
278 let _ = sender.send(Ok((manifest.manifest_version, updated)));
279 }
280}
281
282impl<S: LogStore> RegionWorkerLoop<S> {
283 pub(crate) fn handle_region_edit(&mut self, request: RegionEditRequest) {
285 let region_id = request.region_id;
286 let Some(region) = self.regions.get_region(region_id) else {
287 request
288 .waiters
289 .reply_with(|| RegionNotFoundSnafu { region_id }.fail());
290 return;
291 };
292
293 if !region.is_writable() {
294 if region.state() == RegionRoleState::Leader(RegionLeaderState::Editing) {
295 self.region_edit_queues
296 .entry(region_id)
297 .or_insert_with(|| RegionEditQueue::new(region_id))
298 .enqueue(request);
299 } else {
300 request
301 .waiters
302 .reply_with(|| RegionBusySnafu { region_id }.fail());
303 }
304 return;
305 }
306
307 let RegionEditRequest {
308 region_id: _,
309 mut edit,
310 waiters,
311 preload_sst_cache,
312 } = request;
313 let file_sequence = region.version_control.committed_sequence() + 1;
314 edit.committed_sequence = Some(file_sequence);
315
316 for file in &mut edit.files_to_add {
322 file.sequence = NonZeroU64::new(file_sequence);
323 file.preserve_row_sequence = false;
324 }
325
326 let is_staging = region.is_staging();
328 let expect_state = if is_staging {
329 RegionLeaderState::Staging
330 } else {
331 RegionLeaderState::Writable
332 };
333 if let Err(e) = region.set_editing(expect_state) {
335 let e = Arc::new(e);
336 waiters.reply_with(|| Err(e.clone()).context(EditRegionSnafu { region_id }));
337 return;
338 }
339
340 let request_sender = self.sender.clone();
341 let cache_manager = self.cache_manager.clone();
342 let listener = self.listener.clone();
343 common_runtime::spawn_global(async move {
346 let result = edit_region(
347 ®ion,
348 edit.clone(),
349 cache_manager,
350 listener,
351 is_staging,
352 preload_sst_cache,
353 )
354 .await
355 .map_err(Arc::new);
356 let notify = WorkerRequest::Background {
357 region_id,
358 notify: BackgroundNotify::RegionEdit(RegionEditResult {
359 region_id,
360 waiters,
361 edit,
362 result,
363 update_region_state: true,
365 is_staging,
366 }),
367 };
368
369 if let Err(res) = request_sender
371 .send(WorkerRequestWithTime::new(notify))
372 .await
373 {
374 warn!(
375 "Failed to send region edit result back to the worker, region_id: {}, res: {:?}",
376 region_id, res
377 );
378 }
379 });
380 }
381
382 pub(crate) async fn handle_region_edit_result(&mut self, edit_result: RegionEditResult) {
384 let region = match self.regions.get_region(edit_result.region_id) {
385 Some(region) => region,
386 None => {
387 self.fail_region_stalled_requests_as_not_found(&edit_result.region_id);
390 self.reject_region_edit_queue_as_not_found(edit_result.region_id);
391
392 edit_result.waiters.reply_with(|| {
393 RegionNotFoundSnafu {
394 region_id: edit_result.region_id,
395 }
396 .fail()
397 });
398 return;
399 }
400 };
401
402 let need_compaction = if edit_result.is_staging {
403 if edit_result.update_region_state {
404 region.switch_state_to_staging(RegionLeaderState::Editing);
407 }
408
409 false
410 } else {
411 let need_compaction = self.config.schedule_compaction_after_edit
412 && edit_result.result.is_ok()
413 && !edit_result.edit.files_to_add.is_empty();
414
415 if edit_result.result.is_ok() {
417 region.version_control.apply_edit(
419 Some(edit_result.edit),
420 &[],
421 region.file_purger.clone(),
422 );
423 }
424 if edit_result.update_region_state {
425 region.switch_state_to_writable(RegionLeaderState::Editing);
426 }
427
428 need_compaction
429 };
430
431 edit_result
432 .waiters
433 .reply_with(|| match &edit_result.result {
434 Ok(()) => Ok(()),
435 Err(e) => Err(e.clone()).context(EditRegionSnafu {
436 region_id: edit_result.region_id,
437 }),
438 });
439
440 if edit_result.update_region_state {
441 self.handle_region_stalled_requests(&edit_result.region_id, false)
444 .await;
445 }
446
447 let next_request =
448 if let Some(edit_queue) = self.region_edit_queues.get_mut(&edit_result.region_id) {
449 let request = edit_queue.dequeue();
450 if edit_queue.is_empty() {
451 self.region_edit_queues.remove(&edit_result.region_id);
452 }
453 request
454 } else {
455 None
456 };
457 if let Some(request) = next_request {
458 self.handle_region_edit(request);
459 }
460
461 if need_compaction {
462 self.schedule_compaction(®ion).await;
463 }
464 }
465
466 pub(crate) fn handle_manifest_truncate_action(
468 &self,
469 region: MitoRegionRef,
470 truncate: RegionTruncate,
471 sender: OptionOutputTx,
472 ) {
473 if let Err(e) = region.set_truncating() {
476 sender.send(Err(e));
477 return;
478 }
479 let request_sender = self.sender.clone();
482 let manifest_ctx = region.manifest_ctx.clone();
483 let is_staging = region.is_staging();
484 let durability_barrier = self
485 .wal
486 .durability_barrier(region.region_id, ®ion.provider);
487
488 common_runtime::spawn_global(async move {
490 let durable = match &truncate.kind {
493 TruncateKind::All {
494 truncated_entry_id, ..
495 } => durability_barrier.wait(*truncated_entry_id).await,
496 _ => Ok(()),
497 };
498 let action_list =
500 RegionMetaActionList::with_action(RegionMetaAction::Truncate(truncate.clone()));
501
502 let result = match durable {
503 Ok(()) => manifest_ctx
504 .update_manifest(RegionLeaderState::Truncating, action_list, is_staging)
505 .await
506 .map(|_| ()),
507 Err(e) => Err(e),
508 };
509
510 let truncate_result = TruncateResult {
512 region_id: truncate.region_id,
513 sender,
514 result,
515 kind: truncate.kind,
516 };
517 let _ = request_sender
518 .send(WorkerRequestWithTime::new(WorkerRequest::Background {
519 region_id: truncate.region_id,
520 notify: BackgroundNotify::Truncate(truncate_result),
521 }))
522 .await
523 .inspect_err(|_| warn!("failed to send truncate result"));
524 });
525 }
526
527 pub(crate) fn handle_manifest_discard_unflushed_action(
529 &self,
530 region: MitoRegionRef,
531 discarded_entry_id: EntryId,
532 discarded_sequence: SequenceNumber,
533 discarded_rows: u64,
534 discarded_bytes: u64,
535 sender: OptionOutputTx,
536 ) {
537 if let Err(e) = region.set_truncating() {
538 sender.send(Err(e));
539 return;
540 }
541
542 let region_id = region.region_id;
543 let request_sender = self.sender.clone();
544 let manifest_ctx = region.manifest_ctx.clone();
545 let durability_barrier = self.wal.durability_barrier(region_id, ®ion.provider);
546
547 common_runtime::spawn_global(async move {
548 let edit = RegionEdit {
552 files_to_add: Vec::new(),
553 files_to_remove: Vec::new(),
554 timestamp_ms: None,
555 compaction_time_window: None,
556 flushed_entry_id: Some(discarded_entry_id),
557 flushed_sequence: Some(discarded_sequence),
558 committed_sequence: None,
559 };
560 let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit));
561 let result = match durability_barrier.wait(discarded_entry_id).await {
562 Ok(()) => manifest_ctx
563 .update_manifest(RegionLeaderState::Truncating, action_list, false)
564 .await
565 .map(|_| ()),
566 Err(e) => Err(e),
567 };
568
569 let result = DiscardUnflushedResult {
570 region_id,
571 sender,
572 result,
573 discarded_entry_id,
574 discarded_sequence,
575 discarded_rows,
576 discarded_bytes,
577 };
578 let _ = request_sender
579 .send(WorkerRequestWithTime::new(WorkerRequest::Background {
580 region_id,
581 notify: BackgroundNotify::DiscardUnflushed(result),
582 }))
583 .await
584 .inspect_err(|_| warn!("failed to send discard unflushed result"));
585 });
586 }
587
588 pub(crate) fn handle_manifest_region_change(
590 &self,
591 region: MitoRegionRef,
592 change: RegionChange,
593 need_index: bool,
594 new_options: Option<RegionOptions>,
595 sender: OptionOutputTx,
596 ) {
597 if let Err(e) = region.set_altering() {
599 sender.send(Err(e));
600 return;
601 }
602 let listener = self.listener.clone();
603 let request_sender = self.sender.clone();
604 let is_staging = region.is_staging();
605 common_runtime::spawn_global(async move {
607 let new_meta = change.metadata.clone();
608 let action_list = RegionMetaActionList::with_action(RegionMetaAction::Change(change));
609
610 let result = region
611 .manifest_ctx
612 .update_manifest(RegionLeaderState::Altering, action_list, is_staging)
613 .await
614 .map(|_| ());
615 let notify = WorkerRequest::Background {
616 region_id: region.region_id,
617 notify: BackgroundNotify::RegionChange(RegionChangeResult {
618 region_id: region.region_id,
619 sender,
620 result,
621 new_meta,
622 need_index,
623 new_options,
624 }),
625 };
626 listener
627 .on_notify_region_change_result_begin(region.region_id)
628 .await;
629
630 if let Err(res) = request_sender
631 .send(WorkerRequestWithTime::new(notify))
632 .await
633 {
634 warn!(
635 "Failed to send region change result back to the worker, region_id: {}, res: {:?}",
636 region.region_id, res
637 );
638 }
639 });
640 }
641
642 fn update_region_version(
643 version_control: &VersionControlRef,
644 new_meta: RegionMetadataRef,
645 new_options: Option<RegionOptions>,
646 memtable_builder_provider: &MemtableBuilderProvider,
647 ) {
648 let options_changed = new_options.is_some();
649 let region_id = new_meta.region_id;
650 if let Some(new_options) = new_options {
651 let new_memtable_builder = memtable_builder_provider.builder_for_options(&new_options);
654 version_control.alter_schema_and_format(new_meta, new_options, new_memtable_builder);
655 } else {
656 version_control.alter_schema(new_meta);
658 }
659
660 let version_data = version_control.current();
661 let version = version_data.version;
662 info!(
663 "Region {} is altered, metadata is {:?}, options: {:?}, options_changed: {}",
664 region_id, version.metadata, version.options, options_changed,
665 );
666 }
667}
668
669async fn edit_region(
671 region: &MitoRegionRef,
672 edit: RegionEdit,
673 cache_manager: CacheManagerRef,
674 listener: WorkerListener,
675 is_staging: bool,
676 preload_sst_cache: bool,
677) -> Result<()> {
678 let region_id = region.region_id;
679 if let Some(write_cache) = cache_manager.write_cache()
680 && preload_sst_cache
681 {
682 for file_meta in &edit.files_to_add {
683 let write_cache = write_cache.clone();
684 let layer = region.access_layer.clone();
685 let listener = listener.clone();
686
687 let index_key = IndexKey::new(region_id, file_meta.file_id, FileType::Parquet);
688 let remote_path =
689 location::sst_file_path(layer.table_dir(), file_meta.file_id(), layer.path_type());
690
691 let is_index_exist = file_meta.exists_index();
692 let index_file_size = file_meta.index_file_size();
693
694 let index_file_index_key = IndexKey::new(
695 region_id,
696 file_meta.index_id().file_id.file_id(),
697 FileType::Puffin(file_meta.index_version),
698 );
699 let index_remote_path = location::index_file_path(
700 layer.table_dir(),
701 file_meta.index_id(),
702 layer.path_type(),
703 );
704
705 let file_size = file_meta.file_size;
706 common_runtime::spawn_global(async move {
707 WRITE_CACHE_INFLIGHT_DOWNLOAD.add(1);
708
709 let parquet_cached = write_cache
710 .download_if_absent(index_key, &remote_path, layer.object_store(), file_size)
711 .await;
712
713 if parquet_cached.is_ok() {
714 let mut cache_metrics = Default::default();
717 let _ = write_cache
718 .file_cache()
719 .get_parquet_meta_data(
720 index_key,
721 &mut cache_metrics,
722 PageIndexPolicy::Optional,
723 )
724 .await;
725
726 if matches!(parquet_cached, Ok(true)) {
727 listener.on_file_cache_filled(index_key.file_id);
728 }
729 }
730 if is_index_exist {
731 if let Err(err) = write_cache
733 .download(
734 index_file_index_key,
735 &index_remote_path,
736 layer.object_store(),
737 index_file_size,
738 )
739 .await
740 {
741 common_telemetry::error!(
742 err; "Failed to download puffin file, region_id: {}, index_file_index_key: {:?}, index_remote_path: {}", region_id, index_file_index_key, index_remote_path
743 );
744 }
745 }
746
747 WRITE_CACHE_INFLIGHT_DOWNLOAD.sub(1);
748 });
749 }
750 }
751
752 info!(
753 "Applying {edit:?} to region {}, is_staging: {}",
754 region_id, is_staging
755 );
756
757 let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit));
758 region
759 .manifest_ctx
760 .update_manifest(RegionLeaderState::Editing, action_list, is_staging)
761 .await
762 .map(|_| ())
763}