Skip to main content

mito2/
flush.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Flush related utilities and structs.
16
17use std::collections::HashMap;
18use std::num::NonZeroU64;
19use std::sync::atomic::{AtomicUsize, Ordering};
20use std::sync::{Arc, Mutex};
21use std::time::Instant;
22
23use bytes::Bytes;
24use common_base::cancellation::CancellableFuture;
25use common_telemetry::{debug, error, info};
26use datatypes::arrow::datatypes::SchemaRef;
27use datatypes::extension::json::is_json2_extension_type;
28use partition::expr::PartitionExpr;
29use smallvec::{SmallVec, smallvec};
30use snafu::ResultExt;
31use store_api::region_request::RegionFlushReason;
32use store_api::storage::{RegionId, SequenceNumber};
33use strum::IntoStaticStr;
34use tokio::sync::{Semaphore, mpsc, watch};
35
36use crate::access_layer::{
37    AccessLayerRef, Metrics, OperationType, SstInfoArray, SstWriteRequest, WriteType,
38};
39use crate::cache::CacheManagerRef;
40use crate::config::MitoConfig;
41use crate::engine::region_hook::SstFileInfo;
42use crate::error::{
43    Error, FlushCancelledSnafu, FlushRegionSnafu, JoinSnafu, RegionBusySnafu, RegionClosedSnafu,
44    RegionDroppedSnafu, RegionTruncatedSnafu, Result,
45};
46use crate::manifest::action::{RegionEdit, RegionMetaAction, RegionMetaActionList};
47use crate::memtable::bulk::ENCODE_ROW_THRESHOLD;
48use crate::memtable::bulk::json_align::Json2Aligner;
49use crate::memtable::{BoxedRecordBatchIterator, EncodedRange, MemtableRanges, RangesOptions};
50use crate::metrics::{
51    FLUSH_BYTES_TOTAL, FLUSH_ELAPSED, FLUSH_FAILURE_TOTAL, FLUSH_FILE_TOTAL, FLUSH_REQUESTS_TOTAL,
52    INFLIGHT_FLUSH_COUNT,
53};
54use crate::read::FlatSource;
55use crate::read::flat_dedup::{FlatDedupIterator, FlatLastNonNull, FlatLastRow};
56use crate::read::flat_merge::FlatMergeIterator;
57use crate::region::options::{IndexOptions, MergeMode, RegionOptions};
58use crate::region::version::{VersionControlData, VersionControlRef, VersionRef};
59use crate::region::{ManifestContextRef, RegionLeaderState, RegionRoleState, parse_partition_expr};
60use crate::request::{
61    BackgroundNotify, DdlRequest, FlushFailed, FlushFinished, OnFailure, OptionOutputTx, OutputTx,
62    SenderBulkRequest, SenderDdlRequest, SenderWriteRequest, WorkerRequest, WorkerRequestWithTime,
63};
64use crate::schedule::CancellableTaskState;
65use crate::schedule::scheduler::{Job, SchedulerRef};
66use crate::sst::file::{FileMeta, UncommittedSsts};
67use crate::sst::parquet::metadata::extract_primary_key_range;
68use crate::sst::parquet::{
69    DEFAULT_READ_BATCH_SIZE, DEFAULT_ROW_GROUP_SIZE, SstInfo, WriteOptions, flat_format,
70};
71use crate::sst::{FlatSchemaOptions, FormatType, to_flat_sst_arrow_schema};
72use crate::worker::WorkerListener;
73
74/// Global write buffer (memtable) manager.
75///
76/// Tracks write buffer (memtable) usages and decide whether the engine needs to flush.
77pub trait WriteBufferManager: Send + Sync + std::fmt::Debug {
78    /// Returns whether to trigger the engine.
79    fn should_flush_engine(&self) -> bool;
80
81    /// Returns whether to stall write requests.
82    fn should_stall(&self) -> bool;
83
84    /// Reserves `mem` bytes.
85    fn reserve_mem(&self, mem: usize);
86
87    /// Tells the manager we are freeing `mem` bytes.
88    ///
89    /// We are in the process of freeing `mem` bytes, so it is not considered
90    /// when checking the soft limit.
91    fn schedule_free_mem(&self, mem: usize);
92
93    /// We have freed `mem` bytes.
94    fn free_mem(&self, mem: usize);
95
96    /// Returns the total memory used by memtables.
97    fn memory_usage(&self) -> usize;
98
99    /// Returns the mutable memtable memory limit.
100    ///
101    /// The write buffer manager should flush memtables when the mutable memory usage
102    /// exceeds this limit.
103    fn flush_limit(&self) -> usize;
104}
105
106pub type WriteBufferManagerRef = Arc<dyn WriteBufferManager>;
107
108/// Default [WriteBufferManager] implementation.
109///
110/// Inspired by RocksDB's WriteBufferManager.
111/// <https://github.com/facebook/rocksdb/blob/main/include/rocksdb/write_buffer_manager.h>
112#[derive(Debug)]
113pub struct WriteBufferManagerImpl {
114    /// Write buffer size for the engine.
115    global_write_buffer_size: usize,
116    /// Mutable memtable memory size limit.
117    mutable_limit: usize,
118    /// Memory in used (e.g. used by mutable and immutable memtables).
119    memory_used: AtomicUsize,
120    /// Memory that hasn't been scheduled to free (e.g. used by mutable memtables).
121    memory_active: AtomicUsize,
122    /// Optional notifier.
123    /// The manager can wake up the worker once we free the write buffer.
124    notifier: Option<watch::Sender<()>>,
125}
126
127impl WriteBufferManagerImpl {
128    /// Returns a new manager with specific `global_write_buffer_size`.
129    pub fn new(global_write_buffer_size: usize) -> Self {
130        Self {
131            global_write_buffer_size,
132            mutable_limit: Self::get_mutable_limit(global_write_buffer_size),
133            memory_used: AtomicUsize::new(0),
134            memory_active: AtomicUsize::new(0),
135            notifier: None,
136        }
137    }
138
139    /// Attaches a notifier to the manager.
140    pub fn with_notifier(mut self, notifier: watch::Sender<()>) -> Self {
141        self.notifier = Some(notifier);
142        self
143    }
144
145    /// Returns memory usage of mutable memtables.
146    pub fn mutable_usage(&self) -> usize {
147        self.memory_active.load(Ordering::Relaxed)
148    }
149
150    /// Returns the size limit for mutable memtables.
151    fn get_mutable_limit(global_write_buffer_size: usize) -> usize {
152        // Reserves half of the write buffer for mutable memtable.
153        global_write_buffer_size / 2
154    }
155}
156
157impl WriteBufferManager for WriteBufferManagerImpl {
158    fn should_flush_engine(&self) -> bool {
159        let mutable_memtable_memory_usage = self.memory_active.load(Ordering::Relaxed);
160        if mutable_memtable_memory_usage >= self.mutable_limit {
161            debug!(
162                "Engine should flush (over mutable limit), mutable_usage: {}, memory_usage: {}, mutable_limit: {}, global_limit: {}",
163                mutable_memtable_memory_usage,
164                self.memory_usage(),
165                self.mutable_limit,
166                self.global_write_buffer_size,
167            );
168            return true;
169        }
170
171        let memory_usage = self.memory_used.load(Ordering::Relaxed);
172        if memory_usage >= self.global_write_buffer_size {
173            return true;
174        }
175
176        false
177    }
178
179    fn should_stall(&self) -> bool {
180        self.memory_usage() >= self.global_write_buffer_size
181    }
182
183    fn reserve_mem(&self, mem: usize) {
184        self.memory_used.fetch_add(mem, Ordering::Relaxed);
185        self.memory_active.fetch_add(mem, Ordering::Relaxed);
186    }
187
188    fn schedule_free_mem(&self, mem: usize) {
189        self.memory_active.fetch_sub(mem, Ordering::Relaxed);
190    }
191
192    fn free_mem(&self, mem: usize) {
193        self.memory_used.fetch_sub(mem, Ordering::Relaxed);
194        if let Some(notifier) = &self.notifier {
195            // Notifies the worker after the memory usage is decreased. When we drop the memtable
196            // outside of the worker, the worker may still stall requests because the memory usage
197            // is not updated. So we need to notify the worker to handle stalled requests again.
198            let _ = notifier.send(());
199        }
200    }
201
202    fn memory_usage(&self) -> usize {
203        self.memory_used.load(Ordering::Relaxed)
204    }
205
206    fn flush_limit(&self) -> usize {
207        self.mutable_limit
208    }
209}
210
211/// Reason of a flush task.
212#[derive(Debug, IntoStaticStr, Clone, Copy, PartialEq, Eq)]
213pub enum FlushReason {
214    /// Engine reaches flush threshold.
215    EngineFull,
216    /// Region reaches its write buffer threshold.
217    RegionFull,
218    /// Manual flush.
219    Manual,
220    /// Flush to alter table.
221    Alter,
222    /// Flush periodically.
223    Periodically,
224    /// Flush memtable during downgrading state.
225    Downgrading,
226    /// Enter staging mode.
227    EnterStaging,
228    /// Flush triggered before region migration.
229    RegionMigration,
230    /// Flush triggered by repartition procedure.
231    Repartition,
232    /// Flush triggered by remote WAL pruning.
233    RemoteWalPrune,
234    /// Flush before closing a Noop WAL region.
235    Closing,
236}
237
238impl FlushReason {
239    /// Get flush reason as static str.
240    fn as_str(&self) -> &'static str {
241        self.into()
242    }
243}
244
245impl From<RegionFlushReason> for FlushReason {
246    fn from(reason: RegionFlushReason) -> Self {
247        match reason {
248            RegionFlushReason::RegionMigration => FlushReason::RegionMigration,
249            RegionFlushReason::Repartition => FlushReason::Repartition,
250            RegionFlushReason::RemoteWalPrune => FlushReason::RemoteWalPrune,
251            RegionFlushReason::Closing => FlushReason::Closing,
252            RegionFlushReason::Downgrading => FlushReason::Downgrading,
253        }
254    }
255}
256
257/// Task to flush a region.
258pub(crate) struct RegionFlushTask {
259    /// Region to flush.
260    pub(crate) region_id: RegionId,
261    /// Reason to flush.
262    pub(crate) reason: FlushReason,
263    /// Flush result senders.
264    pub(crate) senders: Vec<OutputTx>,
265    /// Request sender to notify the worker.
266    pub(crate) request_sender: mpsc::Sender<WorkerRequestWithTime>,
267
268    pub(crate) access_layer: AccessLayerRef,
269    pub(crate) listener: WorkerListener,
270    pub(crate) engine_config: Arc<MitoConfig>,
271    pub(crate) row_group_size: Option<usize>,
272    pub(crate) cache_manager: CacheManagerRef,
273    pub(crate) manifest_ctx: ManifestContextRef,
274
275    /// Index options for the region.
276    pub(crate) index_options: IndexOptions,
277    /// Semaphore to control flush concurrency.
278    pub(crate) flush_semaphore: Arc<Semaphore>,
279    /// Whether the region is in staging mode.
280    pub(crate) is_staging: bool,
281    /// Partition expression of the region.
282    ///
283    /// This is used to generate the file meta.
284    pub(crate) partition_expr: Option<String>,
285}
286
287struct FlushTaskWaiters {
288    region_id: RegionId,
289    senders: Mutex<Vec<OutputTx>>,
290}
291
292impl FlushTaskWaiters {
293    fn new(region_id: RegionId, senders: Vec<OutputTx>) -> Self {
294        Self {
295            region_id,
296            senders: Mutex::new(senders),
297        }
298    }
299
300    fn take(&self) -> Vec<OutputTx> {
301        std::mem::take(&mut *self.senders.lock().unwrap())
302    }
303
304    fn on_failure(&self, err: Arc<Error>) {
305        for sender in self.take() {
306            sender.send(Err(err.clone()).context(FlushRegionSnafu {
307                region_id: self.region_id,
308            }));
309        }
310    }
311}
312
313impl Drop for FlushTaskWaiters {
314    fn drop(&mut self) {
315        self.on_failure(Arc::new(
316            RegionBusySnafu {
317                region_id: self.region_id,
318            }
319            .build(),
320        ));
321    }
322}
323
324impl RegionFlushTask {
325    /// Push the sender if it is not none.
326    pub(crate) fn push_sender(&mut self, mut sender: OptionOutputTx) {
327        if let Some(sender) = sender.take_inner() {
328            self.senders.push(sender);
329        }
330    }
331
332    /// Consumes the task and notify the sender the job is success.
333    fn on_success(self) {
334        for sender in self.senders {
335            sender.send(Ok(0));
336        }
337    }
338
339    /// Send flush error to waiter.
340    fn on_failure(&mut self, err: Arc<Error>) {
341        for sender in self.senders.drain(..) {
342            sender.send(Err(err.clone()).context(FlushRegionSnafu {
343                region_id: self.region_id,
344            }));
345        }
346    }
347
348    /// Converts the flush task into a background job.
349    ///
350    /// We must call this in the region worker.
351    fn into_flush_job(
352        mut self,
353        version_control: &VersionControlRef,
354        state: CancellableTaskState,
355    ) -> (Job, Arc<FlushTaskWaiters>) {
356        // Get a version of this region before creating a job to get current
357        // wal entry id, sequence and immutable memtables.
358        let version_data = version_control.current();
359        let waiters = Arc::new(FlushTaskWaiters::new(
360            self.region_id,
361            std::mem::take(&mut self.senders),
362        ));
363        let job_waiters = waiters.clone();
364
365        let job = Box::pin(async move {
366            self.senders = job_waiters.take();
367            INFLIGHT_FLUSH_COUNT.inc();
368            self.do_flush(version_data, state).await;
369            INFLIGHT_FLUSH_COUNT.dec();
370        });
371        (job, waiters)
372    }
373
374    /// Runs the flush task.
375    async fn do_flush(&mut self, version_data: VersionControlData, state: CancellableTaskState) {
376        let timer = FLUSH_ELAPSED.with_label_values(&["total"]).start_timer();
377        let uncommitted = UncommittedSsts::new(
378            self.region_id,
379            self.access_layer.clone(),
380            Some(self.cache_manager.clone()),
381        );
382        self.listener.on_flush_begin(self.region_id).await;
383        let flush_result = if state.is_cancelled() {
384            FlushCancelledSnafu.fail()
385        } else {
386            self.flush_memtables(&version_data, &state, &uncommitted)
387                .await
388        };
389
390        let worker_request = match flush_result {
391            Ok(edit) => {
392                let memtables_to_remove = version_data
393                    .version
394                    .memtables
395                    .immutables()
396                    .iter()
397                    .map(|m| m.id())
398                    .collect();
399                let flush_finished = FlushFinished {
400                    region_id: self.region_id,
401                    flush_reason: self.reason,
402                    // The last entry has been flushed.
403                    flushed_entry_id: version_data.last_entry_id,
404                    senders: std::mem::take(&mut self.senders),
405                    _timer: timer,
406                    edit,
407                    memtables_to_remove,
408                    is_staging: self.is_staging,
409                };
410                WorkerRequest::Background {
411                    region_id: self.region_id,
412                    notify: BackgroundNotify::FlushFinished(flush_finished),
413                }
414            }
415            Err(e) => {
416                let err = Arc::new(e);
417                let failed = FlushFailed { err: err.clone() };
418                if failed.is_cancelled() {
419                    info!("Flush cancelled for region {}", self.region_id);
420                } else {
421                    error!(err; "Failed to flush region {}", self.region_id);
422                }
423                // Discard the timer.
424                timer.stop_and_discard();
425                uncommitted.cleanup().await;
426
427                self.on_failure(err.clone());
428                WorkerRequest::Background {
429                    region_id: self.region_id,
430                    notify: BackgroundNotify::FlushFailed(failed),
431                }
432            }
433        };
434        self.send_worker_request(worker_request).await;
435    }
436
437    /// Flushes memtables to level 0 SSTs and updates the manifest.
438    /// Returns the [RegionEdit] to apply.
439    async fn flush_memtables(
440        &self,
441        version_data: &VersionControlData,
442        state: &CancellableTaskState,
443        uncommitted: &UncommittedSsts,
444    ) -> Result<RegionEdit> {
445        // We must use the immutable memtables list and entry ids from the `version_data`
446        // for consistency as others might already modify the version in the `version_control`.
447        let version = &version_data.version;
448        let timer = FLUSH_ELAPSED
449            .with_label_values(&["flush_memtables"])
450            .start_timer();
451
452        let mut write_opts = WriteOptions {
453            write_buffer_size: self.engine_config.sst_write_buffer_size,
454            ..Default::default()
455        };
456        if let Some(row_group_size) = self.row_group_size {
457            write_opts.row_group_size = row_group_size;
458        }
459
460        let DoFlushMemtablesResult {
461            file_metas,
462            flushed_bytes,
463            series_count,
464            encoded_part_count,
465            flush_metrics,
466            sst_infos,
467        } = self
468            .do_flush_memtables(version, write_opts, state, uncommitted)
469            .await?;
470
471        if !file_metas.is_empty() {
472            FLUSH_BYTES_TOTAL.inc_by(flushed_bytes);
473        }
474
475        let mut file_ids = Vec::with_capacity(file_metas.len());
476        let mut total_rows = 0;
477        let mut total_bytes = 0;
478        for meta in &file_metas {
479            file_ids.push(meta.file_id);
480            total_rows += meta.num_rows;
481            total_bytes += meta.file_size;
482        }
483        info!(
484            "Successfully flush memtables, region: {}, reason: {}, files: {:?}, series count: {}, total_rows: {}, total_bytes: {}, cost: {:?}, encoded_part_count: {}, metrics: {:?}",
485            self.region_id,
486            self.reason.as_str(),
487            file_ids,
488            series_count,
489            total_rows,
490            total_bytes,
491            timer.stop_and_record(),
492            encoded_part_count,
493            flush_metrics,
494        );
495        flush_metrics.observe();
496
497        let hook = self.manifest_ctx.hook();
498        if let Some(hook) = &hook {
499            let files: Vec<SstFileInfo<'_>> = sst_infos
500                .iter()
501                .zip(file_metas.iter())
502                .map(|(sst_info, file_meta)| SstFileInfo {
503                    sst_info_ref: sst_info,
504                    file_meta,
505                })
506                .collect();
507            hook.on_sst_files_written(self.region_id, &version.metadata, &files)
508                .await;
509        }
510
511        let edit = RegionEdit {
512            files_to_add: file_metas,
513            files_to_remove: Vec::new(),
514            timestamp_ms: Some(chrono::Utc::now().timestamp_millis()),
515            compaction_time_window: None,
516            // The last entry has been flushed.
517            flushed_entry_id: Some(version_data.last_entry_id),
518            flushed_sequence: Some(version_data.committed_sequence),
519            committed_sequence: None,
520        };
521        info!(
522            "Applying {edit:?} to region {}, is_staging: {}",
523            self.region_id, self.is_staging
524        );
525
526        let action_list = RegionMetaActionList::with_action(RegionMetaAction::Edit(edit.clone()));
527
528        // Stop accepting cancellation once the flush is about to publish its manifest edit.
529        if !state.mark_commit_started() {
530            return FlushCancelledSnafu.fail();
531        }
532        self.listener.on_flush_commit_begin(self.region_id).await;
533
534        let expected_state = if matches!(self.reason, FlushReason::Downgrading) {
535            RegionLeaderState::Downgrading
536        } else {
537            // Check if region is in staging mode
538            let current_state = self.manifest_ctx.current_state();
539            if current_state == RegionRoleState::Leader(RegionLeaderState::Staging) {
540                RegionLeaderState::Staging
541            } else {
542                RegionLeaderState::Writable
543            }
544        };
545        let manifest_version = match self
546            .manifest_ctx
547            .update_manifest(expected_state, action_list, self.is_staging)
548            .await
549        {
550            Ok(manifest_version) => {
551                uncommitted.disarm_cleanup();
552                manifest_version
553            }
554            Err(e) => {
555                if e.may_have_persisted_manifest_update() {
556                    uncommitted.disarm_cleanup();
557                } else {
558                    info!(
559                        "Cleaning uncommitted SSTs because the manifest update was not persisted, region: {}, job: flush, error: {:?}",
560                        self.region_id, e
561                    );
562                    uncommitted.cleanup().await;
563                }
564                return Err(e);
565            }
566        };
567        info!(
568            "Successfully update manifest version to {manifest_version}, region: {}, is_staging: {}, reason: {}",
569            self.region_id,
570            self.is_staging,
571            self.reason.as_str()
572        );
573
574        Ok(edit)
575    }
576
577    async fn do_flush_memtables(
578        &self,
579        version: &VersionRef,
580        write_opts: WriteOptions,
581        state: &CancellableTaskState,
582        uncommitted: &UncommittedSsts,
583    ) -> Result<DoFlushMemtablesResult> {
584        let memtables = version.memtables.immutables();
585        let mut file_metas = Vec::with_capacity(memtables.len());
586        let mut flushed_bytes = 0;
587        let mut series_count = 0;
588        let mut encoded_part_count = 0;
589        let mut flush_metrics = Metrics::new(WriteType::Flush);
590        let partition_expr = parse_partition_expr(self.partition_expr.as_deref())?;
591        let hook = self.manifest_ctx.hook();
592        let mut all_sst_infos = Vec::new();
593        for mem in memtables {
594            if mem.is_empty() {
595                // Skip empty memtables.
596                continue;
597            }
598
599            // Compact the memtable first, this waits the background compaction to finish.
600            let compact_start = std::time::Instant::now();
601            if let Err(e) = mem.compact(true) {
602                common_telemetry::error!(e; "Failed to compact memtable before flush");
603            }
604            let compact_cost = compact_start.elapsed();
605            flush_metrics.compact_memtable += compact_cost;
606
607            let mem_stats = mem.stats();
608            let batch_size = crate::batch_size::estimate_batch_size([(
609                mem_stats.num_rows() as u64,
610                mem_stats.bytes_allocated() as u64,
611            )]);
612            // Sets `for_flush` flag and propagates the reader batch size.
613            let mem_ranges =
614                mem.ranges(None, RangesOptions::for_flush().with_batch_size(batch_size))?;
615            let num_mem_ranges = mem_ranges.ranges.len();
616
617            // Aggregate stats from all ranges
618            let num_mem_rows = mem_ranges.num_rows();
619            let memtable_series_count = mem_ranges.series_count();
620            let memtable_id = mem.id();
621            // Increases series count for each mem range. We consider each mem range has different series so
622            // the counter may have more series than the actual series count.
623            series_count += memtable_series_count;
624
625            let flush_start = Instant::now();
626            let FlushFlatMemResult {
627                num_encoded,
628                num_sources,
629                results,
630            } = self
631                .flush_flat_mem_ranges(version, &write_opts, mem_ranges, state, uncommitted)
632                .await?;
633            encoded_part_count += num_encoded;
634            for (source_idx, result) in results.into_iter().enumerate() {
635                let (max_sequence, ssts_written, metrics) = result?;
636                if ssts_written.is_empty() {
637                    // No data written.
638                    continue;
639                }
640
641                common_telemetry::debug!(
642                    "Region {} flush one memtable {} {}/{}, metrics: {:?}",
643                    self.region_id,
644                    memtable_id,
645                    source_idx,
646                    num_sources,
647                    metrics
648                );
649
650                flush_metrics = flush_metrics.merge(metrics);
651
652                for sst_info in &ssts_written {
653                    flushed_bytes += sst_info.file_size;
654                    let pk_range = sst_info
655                        .file_metadata
656                        .as_ref()
657                        .and_then(|meta| extract_primary_key_range(meta, &version.metadata));
658                    file_metas.push(Self::new_file_meta(
659                        self.region_id,
660                        max_sequence,
661                        sst_info,
662                        partition_expr.clone(),
663                        pk_range,
664                    ));
665                }
666                if hook.is_some() {
667                    all_sst_infos.extend(ssts_written);
668                }
669            }
670
671            common_telemetry::debug!(
672                "Region {} flush {} memtables for {}, num_mem_ranges: {}, num_encoded: {}, num_rows: {}, flush_cost: {:?}, compact_cost: {:?}",
673                self.region_id,
674                num_sources,
675                memtable_id,
676                num_mem_ranges,
677                num_encoded,
678                num_mem_rows,
679                flush_start.elapsed(),
680                compact_cost,
681            );
682        }
683
684        Ok(DoFlushMemtablesResult {
685            file_metas,
686            flushed_bytes,
687            series_count,
688            encoded_part_count,
689            flush_metrics,
690            sst_infos: all_sst_infos,
691        })
692    }
693
694    async fn flush_flat_mem_ranges(
695        &self,
696        version: &VersionRef,
697        write_opts: &WriteOptions,
698        mem_ranges: MemtableRanges,
699        state: &CancellableTaskState,
700        uncommitted: &UncommittedSsts,
701    ) -> Result<FlushFlatMemResult> {
702        let batch_schema = to_flat_sst_arrow_schema(
703            &version.metadata,
704            &FlatSchemaOptions::from_encoding(version.metadata.primary_key_encoding),
705        );
706        let field_column_start =
707            flat_format::field_column_start(&version.metadata, batch_schema.fields().len());
708        let flat_sources = memtable_flat_sources(
709            batch_schema,
710            mem_ranges,
711            &version.options,
712            field_column_start,
713        )?;
714        let mut tasks = Vec::with_capacity(flat_sources.encoded.len() + flat_sources.sources.len());
715        let num_encoded = flat_sources.encoded.len();
716        for (source, max_sequence) in flat_sources.sources {
717            let write_request = self.new_write_request(version, max_sequence, source);
718            let access_layer = self.access_layer.clone();
719            let write_opts = write_opts.clone();
720            let semaphore = self.flush_semaphore.clone();
721            let uncommitted = uncommitted.clone();
722            let task = common_runtime::spawn_global(async move {
723                let _permit = semaphore.acquire().await.unwrap();
724                let mut metrics = Metrics::new(WriteType::Flush);
725                let ssts = access_layer
726                    .write_sst(write_request, &write_opts, &mut metrics)
727                    .await?;
728                uncommitted.track(&ssts);
729                FLUSH_FILE_TOTAL.inc_by(ssts.len() as u64);
730                Ok((max_sequence, ssts, metrics))
731            });
732            tasks.push(task);
733        }
734        for (encoded, max_sequence) in flat_sources.encoded {
735            let access_layer = self.access_layer.clone();
736            let cache_manager = self.cache_manager.clone();
737            let region_id = version.metadata.region_id;
738            let semaphore = self.flush_semaphore.clone();
739            let uncommitted = uncommitted.clone();
740            let task = common_runtime::spawn_global(async move {
741                let _permit = semaphore.acquire().await.unwrap();
742                let metrics = access_layer
743                    .put_sst(&encoded.data, region_id, &encoded.sst_info, &cache_manager)
744                    .await?;
745                uncommitted.track(std::slice::from_ref(&encoded.sst_info));
746                FLUSH_FILE_TOTAL.inc();
747                Ok((max_sequence, smallvec![encoded.sst_info], metrics))
748            });
749            tasks.push(task);
750        }
751        let num_sources = tasks.len();
752        let abort_handles = tasks
753            .iter()
754            .map(|task| task.abort_handle())
755            .collect::<Vec<_>>();
756        let join_all = futures::future::join_all(tasks);
757        tokio::pin!(join_all);
758        let results = match CancellableFuture::new(join_all.as_mut(), state.cancel_handle()).await {
759            Ok(results) => results
760                .into_iter()
761                .map(|result| result.context(JoinSnafu))
762                .collect::<Result<Vec<_>>>()?,
763            Err(_) => {
764                for handle in abort_handles {
765                    handle.abort();
766                }
767                // Wait until every writer observes the abort so cleanup cannot race with a late
768                // finalized output.
769                let _ = join_all.await;
770                return FlushCancelledSnafu.fail();
771            }
772        };
773        Ok(FlushFlatMemResult {
774            num_encoded,
775            num_sources,
776            results,
777        })
778    }
779
780    fn new_file_meta(
781        region_id: RegionId,
782        max_sequence: u64,
783        sst_info: &SstInfo,
784        partition_expr: Option<PartitionExpr>,
785        primary_key_range: Option<(Bytes, Bytes)>,
786    ) -> FileMeta {
787        let (primary_key_min, primary_key_max) = match primary_key_range {
788            Some((min, max)) => (Some(min), Some(max)),
789            None => (None, None),
790        };
791        FileMeta {
792            region_id,
793            file_id: sst_info.file_id,
794            time_range: sst_info.time_range,
795            level: 0,
796            file_size: sst_info.file_size,
797            max_row_group_uncompressed_size: sst_info.max_row_group_uncompressed_size,
798            available_indexes: sst_info.index_metadata.build_available_indexes(),
799            indexes: sst_info.index_metadata.build_indexes(),
800            index_file_size: sst_info.index_metadata.file_size,
801            index_version: 0,
802            num_rows: sst_info.num_rows as u64,
803            num_row_groups: sst_info.num_row_groups,
804            sequence: NonZeroU64::new(max_sequence),
805            partition_expr,
806            num_series: sst_info.num_series,
807            primary_key_min,
808            primary_key_max,
809        }
810    }
811
812    fn new_write_request(
813        &self,
814        version: &VersionRef,
815        max_sequence: u64,
816        source: FlatSource,
817    ) -> SstWriteRequest {
818        let flat_format = version
819            .options
820            .sst_format
821            .map(|f| f == FormatType::Flat)
822            .unwrap_or(self.engine_config.default_flat_format);
823        SstWriteRequest {
824            op_type: OperationType::Flush,
825            metadata: version.metadata.clone(),
826            source,
827            cache_manager: self.cache_manager.clone(),
828            storage: version.options.storage.clone(),
829            max_sequence: Some(max_sequence),
830            sst_write_format: if flat_format {
831                FormatType::Flat
832            } else {
833                FormatType::PrimaryKey
834            },
835            index_options: self.index_options.clone(),
836            index_config: self.engine_config.index.clone(),
837            inverted_index_config: self.engine_config.inverted_index.clone(),
838            fulltext_index_config: self.engine_config.fulltext_index.clone(),
839            bloom_filter_index_config: self.engine_config.bloom_filter_index.clone(),
840            #[cfg(feature = "vector_index")]
841            vector_index_config: self.engine_config.vector_index.clone(),
842        }
843    }
844
845    /// Notify flush job status.
846    pub(crate) async fn send_worker_request(&self, request: WorkerRequest) {
847        if let Err(e) = self
848            .request_sender
849            .send(WorkerRequestWithTime::new(request))
850            .await
851        {
852            let request = e.0.request;
853            error!(
854                "Failed to notify flush job status for region {}, request: {:?}",
855                self.region_id, request
856            );
857            if let WorkerRequest::Background {
858                notify: BackgroundNotify::FlushFinished(mut finished),
859                ..
860            } = request
861            {
862                finished.on_failure(
863                    RegionClosedSnafu {
864                        region_id: self.region_id,
865                    }
866                    .build(),
867                );
868            }
869        }
870    }
871
872    /// Merge two flush tasks.
873    fn merge(&mut self, mut other: RegionFlushTask) {
874        assert_eq!(self.region_id, other.region_id);
875        // Now we only merge senders. They share the same flush reason.
876        self.senders.append(&mut other.senders);
877    }
878}
879
880struct FlushFlatMemResult {
881    num_encoded: usize,
882    num_sources: usize,
883    results: Vec<Result<(SequenceNumber, SstInfoArray, Metrics)>>,
884}
885
886struct DoFlushMemtablesResult {
887    file_metas: Vec<FileMeta>,
888    flushed_bytes: u64,
889    series_count: usize,
890    encoded_part_count: usize,
891    flush_metrics: Metrics,
892    sst_infos: Vec<SstInfo>,
893}
894
895struct FlatSources {
896    sources: SmallVec<[(FlatSource, SequenceNumber); 4]>,
897    encoded: SmallVec<[(EncodedRange, SequenceNumber); 4]>,
898}
899
900/// Returns the max sequence and [FlatSource] for the given memtable.
901fn memtable_flat_sources(
902    schema: SchemaRef,
903    mem_ranges: MemtableRanges,
904    options: &RegionOptions,
905    field_column_start: usize,
906) -> Result<FlatSources> {
907    let MemtableRanges { ranges } = mem_ranges;
908    let mut flat_sources = FlatSources {
909        sources: SmallVec::new(),
910        encoded: SmallVec::new(),
911    };
912
913    if ranges.len() == 1 {
914        debug!("Flushing single flat range");
915
916        let only_range = ranges.into_values().next().unwrap();
917        let max_sequence = only_range.stats().max_sequence();
918        if let Some(encoded) = only_range.encoded() {
919            flat_sources.encoded.push((encoded, max_sequence));
920        } else {
921            let iter = only_range.build_record_batch_iter(None, None)?;
922            // Dedup according to append mode and merge mode.
923            // Even single range may have duplicate rows.
924            let iter = maybe_dedup_one(
925                options.append_mode,
926                options.merge_mode(),
927                field_column_start,
928                iter,
929            );
930            flat_sources
931                .sources
932                .push((FlatSource::new_iter(schema, iter), max_sequence));
933        };
934    } else {
935        let min_flush_rows = *ENCODE_ROW_THRESHOLD;
936        // Calculate total rows from non-encoded ranges.
937        let total_rows: usize = ranges
938            .values()
939            .filter(|r| r.encoded().is_none())
940            .map(|r| r.num_rows())
941            .sum();
942        debug!(
943            "Flushing multiple flat ranges, total_rows: {}, min_flush_rows: {}, num_ranges: {}",
944            total_rows,
945            min_flush_rows,
946            ranges.len()
947        );
948        let mut rows_remaining = total_rows;
949        let mut last_iter_rows = 0;
950        let num_ranges = ranges.len();
951        let mut input_iters = Vec::with_capacity(num_ranges);
952        let mut current_ranges = Vec::new();
953
954        let has_json2 = schema.fields().iter().any(is_json2_extension_type);
955        let mut json_align_schemas = if has_json2 {
956            Some(Vec::with_capacity(num_ranges))
957        } else {
958            None
959        };
960
961        for (_range_id, range) in ranges {
962            if let Some(encoded) = range.encoded() {
963                let max_sequence = range.stats().max_sequence();
964                flat_sources.encoded.push((encoded, max_sequence));
965                continue;
966            }
967
968            // Collect schemas if has json2 field.
969            if let Some(schemas) = json_align_schemas.as_mut() {
970                let schema = range
971                    .record_batch_schema_hint()
972                    .unwrap_or_else(|| schema.clone());
973                schemas.push(schema);
974            }
975
976            let iter = range.build_record_batch_iter(None, None)?;
977            input_iters.push(iter);
978            let range_rows = range.num_rows();
979            last_iter_rows += range_rows;
980            rows_remaining -= range_rows;
981            current_ranges.push(range);
982
983            // Flush if we have enough rows, but don't flush if the remaining rows
984            // would be less than DEFAULT_ROW_GROUP_SIZE (to avoid small last files).
985            if last_iter_rows >= min_flush_rows
986                && (rows_remaining == 0 || rows_remaining >= DEFAULT_ROW_GROUP_SIZE)
987            {
988                debug!(
989                    "Flush batch ready, rows: {}, min_rows: {}, num_iters: {}, remaining: {}",
990                    last_iter_rows,
991                    min_flush_rows,
992                    input_iters.len(),
993                    rows_remaining
994                );
995
996                // Calculate max_sequence from all merged ranges
997                let max_sequence = current_ranges
998                    .iter()
999                    .map(|r| r.stats().max_sequence())
1000                    .max()
1001                    .unwrap_or(0);
1002                let batch_size =
1003                    crate::batch_size::estimate_batch_size(current_ranges.iter().map(|range| {
1004                        let stats = range.stats();
1005                        (stats.num_rows() as u64, stats.bytes_allocated() as u64)
1006                    }));
1007
1008                let input_iters =
1009                    std::mem::replace(&mut input_iters, Vec::with_capacity(num_ranges));
1010                let (schema, input_iters) = maybe_align_json2_iters(
1011                    schema.clone(),
1012                    json_align_schemas.take(),
1013                    input_iters,
1014                )?;
1015
1016                let maybe_dedup = merge_and_dedup_with_batch_size(
1017                    &schema,
1018                    options.append_mode,
1019                    options.merge_mode(),
1020                    field_column_start,
1021                    input_iters,
1022                    batch_size,
1023                )?;
1024
1025                flat_sources
1026                    .sources
1027                    .push((FlatSource::new_iter(schema, maybe_dedup), max_sequence));
1028                last_iter_rows = 0;
1029                current_ranges.clear();
1030
1031                json_align_schemas = if has_json2 {
1032                    Some(Vec::with_capacity(num_ranges))
1033                } else {
1034                    None
1035                };
1036            }
1037        }
1038
1039        // Handle remaining iters.
1040        if !input_iters.is_empty() {
1041            debug!(
1042                "Flush remaining batch, rows: {}, min_rows: {}, num_iters: {}, remaining: {}",
1043                last_iter_rows,
1044                min_flush_rows,
1045                input_iters.len(),
1046                rows_remaining
1047            );
1048
1049            let (schema, input_iters) =
1050                maybe_align_json2_iters(schema, json_align_schemas, input_iters)?;
1051
1052            let max_sequence = current_ranges
1053                .iter()
1054                .map(|r| r.stats().max_sequence())
1055                .max()
1056                .unwrap_or(0);
1057            let batch_size =
1058                crate::batch_size::estimate_batch_size(current_ranges.iter().map(|range| {
1059                    let stats = range.stats();
1060                    (stats.num_rows() as u64, stats.bytes_allocated() as u64)
1061                }));
1062
1063            let maybe_dedup = merge_and_dedup_with_batch_size(
1064                &schema,
1065                options.append_mode,
1066                options.merge_mode(),
1067                field_column_start,
1068                input_iters,
1069                batch_size,
1070            )?;
1071
1072            flat_sources
1073                .sources
1074                .push((FlatSource::new_iter(schema, maybe_dedup), max_sequence));
1075        }
1076    }
1077
1078    Ok(flat_sources)
1079}
1080
1081fn maybe_align_json2_iters(
1082    schema: SchemaRef,
1083    schemas: Option<Vec<SchemaRef>>,
1084    input_iters: Vec<BoxedRecordBatchIterator>,
1085) -> Result<(SchemaRef, Vec<BoxedRecordBatchIterator>)> {
1086    let Some(schemas) = schemas else {
1087        return Ok((schema, input_iters));
1088    };
1089
1090    let aligner = Json2Aligner::try_new(schemas)?;
1091    let input_iters = input_iters
1092        .into_iter()
1093        .map(|input_iter| aligner.wrap_iter(input_iter))
1094        .collect();
1095
1096    Ok((aligner.schema().clone(), input_iters))
1097}
1098
1099/// Merges multiple record batch iterators and applies deduplication based on the specified mode.
1100///
1101/// This function is used during the flush process to combine data from multiple memtable ranges
1102/// into a single stream while handling duplicate records according to the configured merge strategy.
1103///
1104/// # Arguments
1105///
1106/// * `schema` - The Arrow schema reference that defines the structure of the record batches
1107/// * `append_mode` - When true, no deduplication is performed and all records are preserved.
1108///                  This is used for append-only workloads where duplicate handling is not required.
1109/// * `merge_mode` - The strategy used for deduplication when not in append mode:
1110///   - `MergeMode::LastRow`: Keeps the last record for each primary key
1111///   - `MergeMode::LastNonNull`: Keeps the last non-null values for each field
1112/// * `field_column_start` - The starting column index for fields in the record batch.
1113///                          Used when `MergeMode::LastNonNull` to identify which columns
1114///                          contain field values versus primary key columns.
1115/// * `input_iters` - A vector of record batch iterators to be merged and deduplicated
1116///
1117/// # Returns
1118///
1119/// Returns a boxed record batch iterator that yields the merged and potentially deduplicated
1120/// record batches.
1121///
1122/// # Behavior
1123///
1124/// 1. Creates a `FlatMergeIterator` to merge all input iterators in sorted order based on
1125///    primary key and timestamp
1126/// 2. If `append_mode` is true, returns the merge iterator directly without deduplication
1127/// 3. If `append_mode` is false, wraps the merge iterator with a `FlatDedupIterator` that
1128///    applies the specified merge mode:
1129///    - `LastRow`: Removes duplicate rows, keeping only the last one
1130///    - `LastNonNull`: Removes duplicates but preserves the last non-null value for each field
1131///
1132/// # Examples
1133///
1134/// ```ignore
1135/// let merged_iter = merge_and_dedup(
1136///     &schema,
1137///     false,  // not append mode, apply dedup
1138///     MergeMode::LastRow,
1139///     2,  // fields start at column 2 after primary key columns
1140///     vec![iter1, iter2, iter3],
1141/// )?;
1142/// ```
1143pub fn merge_and_dedup(
1144    schema: &SchemaRef,
1145    append_mode: bool,
1146    merge_mode: MergeMode,
1147    field_column_start: usize,
1148    input_iters: Vec<BoxedRecordBatchIterator>,
1149) -> Result<BoxedRecordBatchIterator> {
1150    merge_and_dedup_with_batch_size(
1151        schema,
1152        append_mode,
1153        merge_mode,
1154        field_column_start,
1155        input_iters,
1156        DEFAULT_READ_BATCH_SIZE,
1157    )
1158}
1159
1160/// Merges and optionally deduplicates record batch iterators with an explicit output batch size.
1161///
1162/// `batch_size` controls the target number of rows in batches assembled by the merge iterator and
1163/// is clamped to at least one. The other arguments have the same meaning as in
1164/// [`merge_and_dedup`].
1165pub fn merge_and_dedup_with_batch_size(
1166    schema: &SchemaRef,
1167    append_mode: bool,
1168    merge_mode: MergeMode,
1169    field_column_start: usize,
1170    input_iters: Vec<BoxedRecordBatchIterator>,
1171    batch_size: usize,
1172) -> Result<BoxedRecordBatchIterator> {
1173    let batch_size = batch_size.max(1);
1174    let merge_iter = FlatMergeIterator::new(schema.clone(), input_iters, batch_size)?;
1175    let maybe_dedup = if append_mode {
1176        // No dedup in append mode
1177        Box::new(merge_iter) as _
1178    } else {
1179        // Dedup according to merge mode.
1180        match merge_mode {
1181            MergeMode::LastRow => {
1182                Box::new(FlatDedupIterator::new(merge_iter, FlatLastRow::new(false))) as _
1183            }
1184            MergeMode::LastNonNull => Box::new(FlatDedupIterator::new(
1185                merge_iter,
1186                FlatLastNonNull::new(field_column_start, false),
1187            )) as _,
1188        }
1189    };
1190    Ok(maybe_dedup)
1191}
1192
1193pub fn maybe_dedup_one(
1194    append_mode: bool,
1195    merge_mode: MergeMode,
1196    field_column_start: usize,
1197    input_iter: BoxedRecordBatchIterator,
1198) -> BoxedRecordBatchIterator {
1199    if append_mode {
1200        // No dedup in append mode
1201        input_iter
1202    } else {
1203        // Dedup according to merge mode.
1204        match merge_mode {
1205            MergeMode::LastRow => {
1206                Box::new(FlatDedupIterator::new(input_iter, FlatLastRow::new(false)))
1207            }
1208            MergeMode::LastNonNull => Box::new(FlatDedupIterator::new(
1209                input_iter,
1210                FlatLastNonNull::new(field_column_start, false),
1211            )),
1212        }
1213    }
1214}
1215
1216/// Manages background flushes of a worker.
1217pub(crate) struct FlushScheduler {
1218    /// Tracks regions need to flush.
1219    region_status: HashMap<RegionId, FlushStatus>,
1220    /// Background job scheduler.
1221    scheduler: SchedulerRef,
1222}
1223
1224impl FlushScheduler {
1225    /// Creates a new flush scheduler.
1226    pub(crate) fn new(scheduler: SchedulerRef) -> FlushScheduler {
1227        FlushScheduler {
1228            region_status: HashMap::new(),
1229            scheduler,
1230        }
1231    }
1232
1233    /// Returns true if the region already requested flush.
1234    pub(crate) fn is_flush_requested(&self, region_id: RegionId) -> bool {
1235        self.region_status.contains_key(&region_id)
1236    }
1237
1238    fn schedule_flush_task(
1239        &mut self,
1240        version_control: &VersionControlRef,
1241        mut task: RegionFlushTask,
1242    ) -> Result<CancellableTaskState> {
1243        let region_id = task.region_id;
1244
1245        // If current region doesn't have flush status, we can flush the region directly.
1246        if let Err(e) = version_control.freeze_mutable() {
1247            error!(e; "Failed to freeze the mutable memtable for region {}", region_id);
1248            task.on_failure(Arc::new(RegionBusySnafu { region_id }.build()));
1249
1250            return Err(e);
1251        }
1252        // Submit a flush job.
1253        let state = CancellableTaskState::new();
1254        let (job, waiters) = task.into_flush_job(version_control, state.clone());
1255        if let Err(e) = self.scheduler.schedule(job) {
1256            error!(e; "Failed to schedule flush job for region {}", region_id);
1257            waiters.on_failure(Arc::new(RegionBusySnafu { region_id }.build()));
1258
1259            return Err(e);
1260        }
1261        Ok(state)
1262    }
1263
1264    /// Schedules a flush `task` for specific `region`.
1265    pub(crate) fn schedule_flush(
1266        &mut self,
1267        region_id: RegionId,
1268        version_control: &VersionControlRef,
1269        mut task: RegionFlushTask,
1270    ) -> Result<()> {
1271        debug_assert_eq!(region_id, task.region_id);
1272
1273        let version = version_control.current().version;
1274        if version.memtables.is_empty() {
1275            debug_assert!(!self.region_status.contains_key(&region_id));
1276            // The region has nothing to flush.
1277            task.on_success();
1278            return Ok(());
1279        }
1280
1281        // Don't increase the counter if a region has nothing to flush.
1282        FLUSH_REQUESTS_TOTAL
1283            .with_label_values(&[task.reason.as_str()])
1284            .inc();
1285
1286        // If current region has flush status, merge the task.
1287        if let Some(flush_status) = self.region_status.get_mut(&region_id) {
1288            if flush_status.has_pending_lifecycle_ddl() {
1289                task.on_failure(Arc::new(FlushCancelledSnafu.build()));
1290                return Ok(());
1291            }
1292            // Checks whether we can flush the region now.
1293            debug!("Merging flush task for region {}", region_id);
1294            flush_status.merge_task(task);
1295            return Ok(());
1296        }
1297
1298        let closing = task.reason == FlushReason::Closing;
1299        let state = self.schedule_flush_task(version_control, task)?;
1300
1301        // Add this region to status map.
1302        let _ = self.region_status.insert(
1303            region_id,
1304            FlushStatus::new(region_id, version_control.clone(), state, closing),
1305        );
1306
1307        Ok(())
1308    }
1309
1310    /// Notifies the scheduler that the flush job is finished.
1311    ///
1312    /// Returns all pending requests if the region doesn't need to flush again.
1313    pub(crate) fn on_flush_success(
1314        &mut self,
1315        region_id: RegionId,
1316    ) -> Option<(
1317        Vec<SenderDdlRequest>,
1318        Vec<SenderWriteRequest>,
1319        Vec<SenderBulkRequest>,
1320    )> {
1321        let flush_status = self.region_status.get_mut(&region_id)?;
1322        // If region doesn't have any pending flush task, we need to remove it from the status.
1323        if flush_status.pending_task.is_none() {
1324            // The region doesn't have any pending flush task.
1325            // Safety: The flush status must exist.
1326            debug!(
1327                "Region {} doesn't have any pending flush task, removing it from the status",
1328                region_id
1329            );
1330            let flush_status = self.region_status.remove(&region_id).unwrap();
1331            return Some((
1332                flush_status.pending_ddls,
1333                flush_status.pending_writes,
1334                flush_status.pending_bulk_writes,
1335            ));
1336        }
1337
1338        // If region has pending task, but has nothing to flush, we need to remove it from the status.
1339        let version_data = flush_status.version_control.current();
1340        if version_data.version.memtables.is_empty() {
1341            // The region has nothing to flush, we also need to remove it from the status.
1342            // Safety: The pending task is not None.
1343            let task = flush_status.pending_task.take().unwrap();
1344            // The region has nothing to flush. We can notify pending task.
1345            task.on_success();
1346            debug!(
1347                "Region {} has nothing to flush, removing it from the status",
1348                region_id
1349            );
1350            // Safety: The flush status must exist.
1351            let flush_status = self.region_status.remove(&region_id).unwrap();
1352            return Some((
1353                flush_status.pending_ddls,
1354                flush_status.pending_writes,
1355                flush_status.pending_bulk_writes,
1356            ));
1357        }
1358
1359        // If region has pending task and has something to flush, we need to schedule it.
1360        debug!("Scheduling pending flush task for region {}", region_id);
1361        // Safety: The flush status must exist.
1362        let task = flush_status.pending_task.take().unwrap();
1363        let version_control = flush_status.version_control.clone();
1364        match self.schedule_flush_task(&version_control, task) {
1365            Ok(state) => {
1366                self.region_status.get_mut(&region_id).unwrap().state = state;
1367            }
1368            Err(err) => {
1369                error!(
1370                    err;
1371                    "Flush succeeded for region {region_id}, but failed to schedule next flush for it."
1372                );
1373                let flush_status = self.region_status.remove(&region_id).unwrap();
1374                flush_status.fail_all(Arc::new(RegionBusySnafu { region_id }.build()));
1375                return None;
1376            }
1377        }
1378        // We can flush the region again, keep it in the region status.
1379        None
1380    }
1381
1382    /// Notifies the scheduler that the flush job failed.
1383    ///
1384    /// Returns pending drop and truncate requests in their original order. All other pending
1385    /// requests are failed with `err`.
1386    pub(crate) fn on_flush_failed(
1387        &mut self,
1388        region_id: RegionId,
1389        err: Arc<Error>,
1390    ) -> Vec<SenderDdlRequest> {
1391        if matches!(err.as_ref(), Error::FlushCancelled { .. }) {
1392            info!("Region {} flush was cancelled", region_id);
1393        } else {
1394            error!(err; "Region {} failed to flush, cancel all pending tasks", region_id);
1395            FLUSH_FAILURE_TOTAL.inc();
1396        }
1397
1398        // Remove this region.
1399        let Some(flush_status) = self.region_status.remove(&region_id) else {
1400            return Vec::new();
1401        };
1402
1403        flush_status.on_failure(err)
1404    }
1405
1406    /// Cancels the running flush and queues its dependent lifecycle DDL atomically.
1407    pub(crate) fn try_cancel_and_add_ddl<T>(
1408        &mut self,
1409        region_id: RegionId,
1410        sender: OptionOutputTx,
1411        request: T,
1412        into_ddl_request: impl FnOnce(T) -> DdlRequest,
1413    ) -> std::result::Result<(), (OptionOutputTx, T)> {
1414        let Some(status) = self.region_status.get_mut(&region_id) else {
1415            return Err((sender, request));
1416        };
1417
1418        let cancel_result = status.state.request_cancel();
1419        debug!(
1420            "Requested flush cancellation for region {}, result: {:?}",
1421            region_id, cancel_result
1422        );
1423        status.pending_ddls.push(SenderDdlRequest {
1424            region_id,
1425            sender,
1426            request: into_ddl_request(request),
1427        });
1428        Ok(())
1429    }
1430
1431    /// Notifies the scheduler that the region is dropped.
1432    pub(crate) fn on_region_dropped(&mut self, region_id: RegionId) {
1433        self.remove_region_on_failure(
1434            region_id,
1435            Arc::new(RegionDroppedSnafu { region_id }.build()),
1436        );
1437    }
1438
1439    /// Notifies the scheduler that the region is closed.
1440    pub(crate) fn on_region_closed(&mut self, region_id: RegionId) {
1441        let Some(flush_status) = self.region_status.remove(&region_id) else {
1442            return;
1443        };
1444
1445        flush_status.on_region_closed(Arc::new(RegionClosedSnafu { region_id }.build()));
1446    }
1447
1448    /// Notifies the scheduler that the region is truncated.
1449    pub(crate) fn on_region_truncated(&mut self, region_id: RegionId) {
1450        self.remove_region_on_failure(
1451            region_id,
1452            Arc::new(RegionTruncatedSnafu { region_id }.build()),
1453        );
1454    }
1455
1456    fn remove_region_on_failure(&mut self, region_id: RegionId, err: Arc<Error>) {
1457        // Remove this region.
1458        let Some(flush_status) = self.region_status.remove(&region_id) else {
1459            return;
1460        };
1461
1462        // Notifies all pending tasks.
1463        flush_status.fail_all(err);
1464    }
1465
1466    /// Add ddl request to pending queue.
1467    ///
1468    /// # Panics
1469    /// Panics if region didn't request flush.
1470    pub(crate) fn add_ddl_request_to_pending(&mut self, request: SenderDdlRequest) {
1471        let status = self.region_status.get_mut(&request.region_id).unwrap();
1472        status.pending_ddls.push(request);
1473    }
1474
1475    /// Add write request to pending queue.
1476    ///
1477    /// # Panics
1478    /// Panics if region didn't request flush.
1479    pub(crate) fn add_write_request_to_pending(&mut self, request: SenderWriteRequest) {
1480        let status = self
1481            .region_status
1482            .get_mut(&request.request.region_id)
1483            .unwrap();
1484        status.pending_writes.push(request);
1485    }
1486
1487    /// Add bulk write request to pending queue.
1488    ///
1489    /// # Panics
1490    /// Panics if region didn't request flush.
1491    pub(crate) fn add_bulk_request_to_pending(&mut self, request: SenderBulkRequest) {
1492        let status = self.region_status.get_mut(&request.region_id).unwrap();
1493        status.pending_bulk_writes.push(request);
1494    }
1495
1496    /// Returns true if the region has pending DDLs or a close-time flush.
1497    pub(crate) fn has_pending_ddls(&self, region_id: RegionId) -> bool {
1498        self.region_status
1499            .get(&region_id)
1500            .map(|status| !status.pending_ddls.is_empty() || status.closing)
1501            .unwrap_or(false)
1502    }
1503}
1504
1505impl Drop for FlushScheduler {
1506    fn drop(&mut self) {
1507        for (region_id, flush_status) in self.region_status.drain() {
1508            // We are shutting down so notify all pending tasks.
1509            flush_status.fail_all(Arc::new(RegionClosedSnafu { region_id }.build()));
1510        }
1511    }
1512}
1513
1514/// Flush status of a region scheduled by the [FlushScheduler].
1515///
1516/// Tracks running and pending flush tasks and all pending requests of a region.
1517struct FlushStatus {
1518    /// Current region.
1519    region_id: RegionId,
1520    /// Version control of the region.
1521    version_control: VersionControlRef,
1522    /// Cancellation state of the running flush.
1523    state: CancellableTaskState,
1524    /// Task waiting for next flush.
1525    pending_task: Option<RegionFlushTask>,
1526    /// Whether a close-time flush is in progress or pending.
1527    closing: bool,
1528    /// Pending ddl requests.
1529    pending_ddls: Vec<SenderDdlRequest>,
1530    /// Requests waiting to write after altering the region.
1531    pending_writes: Vec<SenderWriteRequest>,
1532    /// Bulk requests waiting to write after altering the region.
1533    pending_bulk_writes: Vec<SenderBulkRequest>,
1534}
1535
1536impl FlushStatus {
1537    fn new(
1538        region_id: RegionId,
1539        version_control: VersionControlRef,
1540        state: CancellableTaskState,
1541        closing: bool,
1542    ) -> FlushStatus {
1543        FlushStatus {
1544            region_id,
1545            version_control,
1546            state,
1547            pending_task: None,
1548            closing,
1549            pending_ddls: Vec::new(),
1550            pending_writes: Vec::new(),
1551            pending_bulk_writes: Vec::new(),
1552        }
1553    }
1554
1555    /// Merges the task to pending task.
1556    fn merge_task(&mut self, task: RegionFlushTask) {
1557        self.closing |= task.reason == FlushReason::Closing;
1558        if let Some(pending) = &mut self.pending_task {
1559            pending.merge(task);
1560        } else {
1561            self.pending_task = Some(task);
1562        }
1563    }
1564
1565    fn has_pending_lifecycle_ddl(&self) -> bool {
1566        self.pending_ddls
1567            .iter()
1568            .any(|ddl| matches!(ddl.request, DdlRequest::Drop(_) | DdlRequest::Truncate(_)))
1569    }
1570
1571    /// Fails pending requests except drop and truncate, which the worker must revalidate.
1572    fn on_failure(self, err: Arc<Error>) -> Vec<SenderDdlRequest> {
1573        if let Some(mut task) = self.pending_task {
1574            task.on_failure(err.clone());
1575        }
1576        let mut lifecycle_ddls = Vec::new();
1577        for ddl in self.pending_ddls {
1578            if matches!(ddl.request, DdlRequest::Drop(_) | DdlRequest::Truncate(_)) {
1579                lifecycle_ddls.push(ddl);
1580            } else {
1581                ddl.sender.send(Err(err.clone()).context(FlushRegionSnafu {
1582                    region_id: self.region_id,
1583                }));
1584            }
1585        }
1586        for write_req in self.pending_writes {
1587            write_req
1588                .sender
1589                .send(Err(err.clone()).context(FlushRegionSnafu {
1590                    region_id: self.region_id,
1591                }));
1592        }
1593        for bulk_req in self.pending_bulk_writes {
1594            bulk_req
1595                .sender
1596                .send(Err(err.clone()).context(FlushRegionSnafu {
1597                    region_id: self.region_id,
1598                }));
1599        }
1600        lifecycle_ddls
1601    }
1602
1603    fn fail_all(self, err: Arc<Error>) {
1604        let region_id = self.region_id;
1605        let ddls = self.on_failure(err.clone());
1606        for ddl in ddls {
1607            ddl.sender
1608                .send(Err(err.clone()).context(FlushRegionSnafu { region_id }));
1609        }
1610    }
1611
1612    fn on_region_closed(self, err: Arc<Error>) {
1613        if let Some(mut task) = self.pending_task {
1614            task.on_failure(err.clone());
1615        }
1616
1617        for ddl in self.pending_ddls {
1618            if matches!(ddl.request, DdlRequest::Close(_)) {
1619                ddl.sender.send(Ok(0));
1620            } else {
1621                ddl.sender.send(Err(err.clone()).context(FlushRegionSnafu {
1622                    region_id: self.region_id,
1623                }));
1624            }
1625        }
1626
1627        for write_req in self.pending_writes {
1628            write_req
1629                .sender
1630                .send(Err(err.clone()).context(FlushRegionSnafu {
1631                    region_id: self.region_id,
1632                }));
1633        }
1634        for bulk_req in self.pending_bulk_writes {
1635            bulk_req
1636                .sender
1637                .send(Err(err.clone()).context(FlushRegionSnafu {
1638                    region_id: self.region_id,
1639                }));
1640        }
1641    }
1642}
1643
1644#[cfg(test)]
1645mod tests {
1646    use api::v1::{OpType, Rows};
1647    use common_error::ext::ErrorExt;
1648    use common_error::status_code::StatusCode;
1649    use mito_codec::row_converter::build_primary_key_codec;
1650    use tokio::sync::oneshot;
1651
1652    use super::*;
1653    use crate::cache::CacheManager;
1654    use crate::error::InvalidSchedulerStateSnafu;
1655    use crate::memtable::bulk::part::BulkPartConverter;
1656    use crate::memtable::time_series::TimeSeriesMemtableBuilder;
1657    use crate::memtable::{Memtable, RangesOptions};
1658    use crate::request::WriteRequest;
1659    use crate::schedule::scheduler::Scheduler;
1660    use crate::sst::{FlatSchemaOptions, to_flat_sst_arrow_schema};
1661    use crate::test_util::memtable_util::{build_key_values_with_ts_seq_values, metadata_for_test};
1662    use crate::test_util::scheduler_util::{SchedulerEnv, VecScheduler};
1663    use crate::test_util::version_util::{VersionControlBuilder, write_rows_to_version};
1664
1665    struct FailingScheduler;
1666
1667    #[async_trait::async_trait]
1668    impl Scheduler for FailingScheduler {
1669        fn schedule(&self, _job: Job) -> Result<()> {
1670            InvalidSchedulerStateSnafu.fail()
1671        }
1672
1673        async fn stop(&self, _await_termination: bool) -> Result<()> {
1674            Ok(())
1675        }
1676    }
1677
1678    fn new_test_flush_task(
1679        env: &SchedulerEnv,
1680        region_id: RegionId,
1681        reason: FlushReason,
1682        request_sender: mpsc::Sender<WorkerRequestWithTime>,
1683        manifest_ctx: ManifestContextRef,
1684    ) -> RegionFlushTask {
1685        RegionFlushTask {
1686            region_id,
1687            reason,
1688            senders: Vec::new(),
1689            request_sender,
1690            access_layer: env.access_layer.clone(),
1691            listener: WorkerListener::default(),
1692            engine_config: Arc::new(MitoConfig::default()),
1693            row_group_size: None,
1694            cache_manager: Arc::new(CacheManager::default()),
1695            manifest_ctx,
1696            index_options: IndexOptions::default(),
1697            flush_semaphore: Arc::new(Semaphore::new(2)),
1698            is_staging: false,
1699            partition_expr: None,
1700        }
1701    }
1702
1703    fn new_test_bulk_request(
1704        region_id: RegionId,
1705    ) -> (
1706        SenderBulkRequest,
1707        oneshot::Receiver<Result<store_api::region_request::AffectedRows>>,
1708    ) {
1709        let metadata = metadata_for_test();
1710        let schema = to_flat_sst_arrow_schema(
1711            &metadata,
1712            &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
1713        );
1714        let pk_codec = build_primary_key_codec(&metadata);
1715        let mut converter = BulkPartConverter::new(&metadata, schema, 16, pk_codec, true);
1716        let kvs = build_key_values_with_ts_seq_values(
1717            &metadata,
1718            "bulk_key".to_string(),
1719            1,
1720            std::iter::once(1000i64),
1721            std::iter::once(Some(1.0f64)),
1722            1,
1723        );
1724        converter.append_key_values(&kvs).unwrap();
1725        let (sender, receiver) = oneshot::channel();
1726
1727        (
1728            SenderBulkRequest {
1729                sender: OptionOutputTx::from(sender),
1730                region_id,
1731                request: converter.convert().unwrap(),
1732                region_metadata: Some(metadata),
1733                partition_expr_version: None,
1734            },
1735            receiver,
1736        )
1737    }
1738
1739    fn new_test_write_request(
1740        region_id: RegionId,
1741    ) -> (
1742        SenderWriteRequest,
1743        oneshot::Receiver<Result<store_api::region_request::AffectedRows>>,
1744    ) {
1745        let (sender, receiver) = oneshot::channel();
1746        let request = WriteRequest::new(region_id, OpType::Put, Rows::default(), None).unwrap();
1747        (
1748            SenderWriteRequest {
1749                sender: OptionOutputTx::from(sender),
1750                request,
1751            },
1752            receiver,
1753        )
1754    }
1755
1756    #[test]
1757    fn test_get_mutable_limit() {
1758        assert_eq!(4, WriteBufferManagerImpl::get_mutable_limit(8));
1759        assert_eq!(5, WriteBufferManagerImpl::get_mutable_limit(10));
1760        assert_eq!(32, WriteBufferManagerImpl::get_mutable_limit(64));
1761        assert_eq!(0, WriteBufferManagerImpl::get_mutable_limit(0));
1762    }
1763
1764    #[test]
1765    fn test_over_mutable_limit() {
1766        // Mutable limit is 500.
1767        let manager = WriteBufferManagerImpl::new(1000);
1768        manager.reserve_mem(400);
1769        assert!(!manager.should_flush_engine());
1770        assert!(!manager.should_stall());
1771
1772        // More than mutable limit.
1773        manager.reserve_mem(400);
1774        assert!(manager.should_flush_engine());
1775
1776        // Freezes mutable.
1777        manager.schedule_free_mem(400);
1778        assert!(!manager.should_flush_engine());
1779        assert_eq!(800, manager.memory_used.load(Ordering::Relaxed));
1780        assert_eq!(400, manager.memory_active.load(Ordering::Relaxed));
1781
1782        // Releases immutable.
1783        manager.free_mem(400);
1784        assert_eq!(400, manager.memory_used.load(Ordering::Relaxed));
1785        assert_eq!(400, manager.memory_active.load(Ordering::Relaxed));
1786    }
1787
1788    #[test]
1789    fn test_over_global() {
1790        // Mutable limit is 500.
1791        let manager = WriteBufferManagerImpl::new(1000);
1792        manager.reserve_mem(1100);
1793        assert!(manager.should_stall());
1794        // Global usage is still 1100.
1795        manager.schedule_free_mem(200);
1796        assert!(manager.should_flush_engine());
1797        assert!(manager.should_stall());
1798
1799        // More than global limit, mutable (1100-200-450=450) is less than mutable limit (< 500).
1800        manager.schedule_free_mem(450);
1801        assert!(manager.should_flush_engine());
1802        assert!(manager.should_stall());
1803
1804        // Now mutable is enough.
1805        manager.reserve_mem(50);
1806        assert!(manager.should_flush_engine());
1807        manager.reserve_mem(100);
1808        assert!(manager.should_flush_engine());
1809    }
1810
1811    #[test]
1812    fn test_manager_notify() {
1813        let (sender, receiver) = watch::channel(());
1814        let manager = WriteBufferManagerImpl::new(1000).with_notifier(sender);
1815        manager.reserve_mem(500);
1816        assert!(!receiver.has_changed().unwrap());
1817        manager.schedule_free_mem(500);
1818        assert!(!receiver.has_changed().unwrap());
1819        manager.free_mem(500);
1820        assert!(receiver.has_changed().unwrap());
1821    }
1822
1823    #[tokio::test]
1824    async fn test_schedule_empty() {
1825        let env = SchedulerEnv::new().await;
1826        let (tx, _rx) = mpsc::channel(4);
1827        let mut scheduler = env.mock_flush_scheduler();
1828        let builder = VersionControlBuilder::new();
1829
1830        let version_control = Arc::new(builder.build());
1831        let (output_tx, output_rx) = oneshot::channel();
1832        let mut task = RegionFlushTask {
1833            region_id: builder.region_id(),
1834            reason: FlushReason::Manual,
1835            senders: Vec::new(),
1836            request_sender: tx,
1837            access_layer: env.access_layer.clone(),
1838            listener: WorkerListener::default(),
1839            engine_config: Arc::new(MitoConfig::default()),
1840            row_group_size: None,
1841            cache_manager: Arc::new(CacheManager::default()),
1842            manifest_ctx: env
1843                .mock_manifest_context(version_control.current().version.metadata.clone())
1844                .await,
1845            index_options: IndexOptions::default(),
1846            flush_semaphore: Arc::new(Semaphore::new(2)),
1847            is_staging: false,
1848            partition_expr: None,
1849        };
1850        task.push_sender(OptionOutputTx::from(output_tx));
1851        scheduler
1852            .schedule_flush(builder.region_id(), &version_control, task)
1853            .unwrap();
1854        assert!(scheduler.region_status.is_empty());
1855        let output = output_rx.await.unwrap().unwrap();
1856        assert_eq!(output, 0);
1857        assert!(scheduler.region_status.is_empty());
1858    }
1859
1860    #[tokio::test]
1861    async fn test_schedule_flush_failure_notifies_waiter() {
1862        let env = SchedulerEnv::new()
1863            .await
1864            .scheduler(Arc::new(FailingScheduler));
1865        let (tx, _rx) = mpsc::channel(4);
1866        let mut scheduler = env.mock_flush_scheduler();
1867        let mut builder = VersionControlBuilder::new();
1868        builder.set_memtable_builder(Arc::new(TimeSeriesMemtableBuilder::default()));
1869        let version_control = Arc::new(builder.build());
1870        let version_data = version_control.current();
1871        write_rows_to_version(&version_data.version, "host0", 0, 10);
1872        let manifest_ctx = env
1873            .mock_manifest_context(version_data.version.metadata.clone())
1874            .await;
1875        let (output_tx, output_rx) = oneshot::channel();
1876        let mut task = new_test_flush_task(
1877            &env,
1878            builder.region_id(),
1879            FlushReason::Manual,
1880            tx,
1881            manifest_ctx,
1882        );
1883        task.push_sender(OptionOutputTx::from(output_tx));
1884
1885        scheduler
1886            .schedule_flush(builder.region_id(), &version_control, task)
1887            .unwrap_err();
1888
1889        let err = output_rx
1890            .await
1891            .expect("waiter must receive explicit error")
1892            .unwrap_err();
1893        assert_eq!(err.status_code(), StatusCode::RegionBusy);
1894    }
1895
1896    #[tokio::test]
1897    async fn test_send_worker_request_failure_notifies_flush_finished_waiter() {
1898        let env = SchedulerEnv::new().await;
1899        let (tx, rx) = mpsc::channel(1);
1900        drop(rx);
1901        let builder = VersionControlBuilder::new();
1902        let version_control = Arc::new(builder.build());
1903        let manifest_ctx = env
1904            .mock_manifest_context(version_control.current().version.metadata.clone())
1905            .await;
1906        let task = new_test_flush_task(
1907            &env,
1908            builder.region_id(),
1909            FlushReason::Manual,
1910            tx,
1911            manifest_ctx,
1912        );
1913        let (output_tx, output_rx) = oneshot::channel();
1914        let request = WorkerRequest::Background {
1915            region_id: builder.region_id(),
1916            notify: BackgroundNotify::FlushFinished(FlushFinished {
1917                region_id: builder.region_id(),
1918                flush_reason: FlushReason::Manual,
1919                flushed_entry_id: 0,
1920                senders: vec![OutputTx::new(output_tx)],
1921                _timer: FLUSH_ELAPSED.with_label_values(&["total"]).start_timer(),
1922                edit: RegionEdit {
1923                    files_to_add: Vec::new(),
1924                    files_to_remove: Vec::new(),
1925                    timestamp_ms: None,
1926                    compaction_time_window: None,
1927                    flushed_entry_id: None,
1928                    flushed_sequence: None,
1929                    committed_sequence: None,
1930                },
1931                memtables_to_remove: smallvec![],
1932                is_staging: false,
1933            }),
1934        };
1935
1936        task.send_worker_request(request).await;
1937
1938        let output = output_rx.await.expect("waiter must receive explicit error");
1939        assert!(output.is_err());
1940    }
1941
1942    #[tokio::test]
1943    async fn test_flush_waiters_drop_notifies_waiter() {
1944        let region_id = RegionId::new(1, 1);
1945        let (output_tx, output_rx) = oneshot::channel();
1946        let waiters = FlushTaskWaiters::new(region_id, vec![OutputTx::new(output_tx)]);
1947
1948        drop(waiters);
1949
1950        let err = output_rx
1951            .await
1952            .expect("waiter must receive explicit error")
1953            .unwrap_err();
1954        assert_eq!(err.status_code(), StatusCode::RegionBusy);
1955    }
1956
1957    #[tokio::test]
1958    async fn test_flush_failure_notifies_pending_bulk_writes() {
1959        let region_id = RegionId::new(1, 1);
1960        let version_control = Arc::new(VersionControlBuilder::new().build());
1961        let (bulk_req, output_rx) = new_test_bulk_request(region_id);
1962        let status = FlushStatus {
1963            region_id,
1964            version_control,
1965            state: CancellableTaskState::new(),
1966            pending_task: None,
1967            closing: false,
1968            pending_ddls: Vec::new(),
1969            pending_writes: Vec::new(),
1970            pending_bulk_writes: vec![bulk_req],
1971        };
1972
1973        let pending_ddls = status.on_failure(Arc::new(RegionClosedSnafu { region_id }.build()));
1974        assert!(pending_ddls.is_empty());
1975
1976        let err = output_rx
1977            .await
1978            .expect("pending bulk write must receive explicit error")
1979            .unwrap_err();
1980        assert_eq!(err.status_code(), StatusCode::Cancelled);
1981    }
1982
1983    #[tokio::test]
1984    async fn test_flush_failure_retains_lifecycle_ddls_and_fails_pending_writes() {
1985        let region_id = RegionId::new(1, 1);
1986        let version_control = Arc::new(VersionControlBuilder::new().build());
1987        let (write_req, write_rx) = new_test_write_request(region_id);
1988        let (bulk_req, bulk_rx) = new_test_bulk_request(region_id);
1989        let (truncate_tx, mut truncate_rx) = oneshot::channel();
1990        let (drop_tx, mut drop_rx) = oneshot::channel();
1991        let status = FlushStatus {
1992            region_id,
1993            version_control,
1994            state: CancellableTaskState::new(),
1995            pending_task: None,
1996            closing: false,
1997            pending_ddls: vec![
1998                SenderDdlRequest {
1999                    region_id,
2000                    sender: OptionOutputTx::from(truncate_tx),
2001                    request: DdlRequest::Truncate(
2002                        store_api::region_request::RegionTruncateRequest::All,
2003                    ),
2004                },
2005                SenderDdlRequest {
2006                    region_id,
2007                    sender: OptionOutputTx::from(drop_tx),
2008                    request: DdlRequest::Drop(store_api::region_request::RegionDropRequest {
2009                        fast_path: false,
2010                        force: false,
2011                        partial_drop: false,
2012                    }),
2013                },
2014            ],
2015            pending_writes: vec![write_req],
2016            pending_bulk_writes: vec![bulk_req],
2017        };
2018
2019        let ddls = status.on_failure(Arc::new(RegionBusySnafu { region_id }.build()));
2020        assert_eq!(2, ddls.len());
2021        assert!(matches!(ddls[0].request, DdlRequest::Truncate(_)));
2022        assert!(matches!(ddls[1].request, DdlRequest::Drop(_)));
2023
2024        let write_err = write_rx
2025            .await
2026            .expect("pending write must receive explicit error")
2027            .unwrap_err();
2028        assert_eq!(write_err.status_code(), StatusCode::RegionBusy);
2029        let bulk_err = bulk_rx
2030            .await
2031            .expect("pending bulk write must receive explicit error")
2032            .unwrap_err();
2033        assert_eq!(bulk_err.status_code(), StatusCode::RegionBusy);
2034
2035        assert!(truncate_rx.try_recv().is_err());
2036        assert!(drop_rx.try_recv().is_err());
2037        for ddl in ddls {
2038            ddl.sender.send(Ok(0));
2039        }
2040        assert_eq!(0, truncate_rx.await.unwrap().unwrap());
2041        assert_eq!(0, drop_rx.await.unwrap().unwrap());
2042    }
2043
2044    #[tokio::test]
2045    async fn test_uncommitted_ssts_cleanup_finalized_file() {
2046        let env = SchedulerEnv::new().await;
2047        let region_id = RegionId::new(1, 1);
2048        let file_id = store_api::storage::FileId::random();
2049        let path = crate::sst::location::sst_file_path(
2050            env.access_layer.table_dir(),
2051            crate::sst::file::RegionFileId::new(region_id, file_id),
2052            env.access_layer.path_type(),
2053        );
2054        env.access_layer
2055            .object_store()
2056            .write(&path, Bytes::from_static(b"sst"))
2057            .await
2058            .unwrap();
2059        assert!(env.access_layer.object_store().exists(&path).await.unwrap());
2060
2061        let uncommitted = UncommittedSsts::new(region_id, env.access_layer.clone(), None);
2062        uncommitted.track(&[SstInfo {
2063            file_id,
2064            ..Default::default()
2065        }]);
2066        assert_eq!(1, uncommitted.num_tracked_files());
2067        uncommitted.cleanup_for_test().await.unwrap();
2068        assert_eq!(0, uncommitted.num_tracked_files());
2069
2070        // The object-store wrapper caches successful stat results, so list the directory instead
2071        // of using `exists()` again to verify the deletion.
2072        let entries = env
2073            .access_layer
2074            .object_store()
2075            .list(&env.access_layer.build_region_dir(region_id))
2076            .await
2077            .unwrap();
2078        assert!(entries.iter().all(|entry| entry.path() != path));
2079    }
2080
2081    #[tokio::test]
2082    async fn test_region_closed_notifies_pending_bulk_writes() {
2083        let region_id = RegionId::new(1, 1);
2084        let version_control = Arc::new(VersionControlBuilder::new().build());
2085        let (bulk_req, output_rx) = new_test_bulk_request(region_id);
2086        let status = FlushStatus {
2087            region_id,
2088            version_control,
2089            state: CancellableTaskState::new(),
2090            pending_task: None,
2091            closing: false,
2092            pending_ddls: Vec::new(),
2093            pending_writes: Vec::new(),
2094            pending_bulk_writes: vec![bulk_req],
2095        };
2096
2097        status.on_region_closed(Arc::new(RegionClosedSnafu { region_id }.build()));
2098
2099        let err = output_rx
2100            .await
2101            .expect("pending bulk write must receive explicit error")
2102            .unwrap_err();
2103        assert_eq!(err.status_code(), StatusCode::Cancelled);
2104    }
2105
2106    #[tokio::test]
2107    async fn test_schedule_pending_request() {
2108        let job_scheduler = Arc::new(VecScheduler::default());
2109        let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone());
2110        let (tx, _rx) = mpsc::channel(4);
2111        let mut scheduler = env.mock_flush_scheduler();
2112        let mut builder = VersionControlBuilder::new();
2113        // Overwrites the empty memtable builder.
2114        builder.set_memtable_builder(Arc::new(TimeSeriesMemtableBuilder::default()));
2115        let version_control = Arc::new(builder.build());
2116        // Writes data to the memtable so it is not empty.
2117        let version_data = version_control.current();
2118        write_rows_to_version(&version_data.version, "host0", 0, 10);
2119        let manifest_ctx = env
2120            .mock_manifest_context(version_data.version.metadata.clone())
2121            .await;
2122        // Creates 3 tasks.
2123        let mut tasks: Vec<_> = (0..3)
2124            .map(|_| RegionFlushTask {
2125                region_id: builder.region_id(),
2126                reason: FlushReason::Manual,
2127                senders: Vec::new(),
2128                request_sender: tx.clone(),
2129                access_layer: env.access_layer.clone(),
2130                listener: WorkerListener::default(),
2131                engine_config: Arc::new(MitoConfig::default()),
2132                row_group_size: None,
2133                cache_manager: Arc::new(CacheManager::default()),
2134                manifest_ctx: manifest_ctx.clone(),
2135                index_options: IndexOptions::default(),
2136                flush_semaphore: Arc::new(Semaphore::new(2)),
2137                is_staging: false,
2138                partition_expr: None,
2139            })
2140            .collect();
2141        // Schedule first task.
2142        let task = tasks.pop().unwrap();
2143        scheduler
2144            .schedule_flush(builder.region_id(), &version_control, task)
2145            .unwrap();
2146        // Should schedule 1 flush.
2147        assert_eq!(1, scheduler.region_status.len());
2148        assert_eq!(1, job_scheduler.num_jobs());
2149        // Check the new version.
2150        let version_data = version_control.current();
2151        assert_eq!(0, version_data.version.memtables.immutables()[0].id());
2152        // Schedule remaining tasks.
2153        let output_rxs: Vec<_> = tasks
2154            .into_iter()
2155            .map(|mut task| {
2156                let (output_tx, output_rx) = oneshot::channel();
2157                task.push_sender(OptionOutputTx::from(output_tx));
2158                scheduler
2159                    .schedule_flush(builder.region_id(), &version_control, task)
2160                    .unwrap();
2161                output_rx
2162            })
2163            .collect();
2164        // Assumes the flush job is finished.
2165        version_control.apply_edit(
2166            Some(RegionEdit {
2167                files_to_add: Vec::new(),
2168                files_to_remove: Vec::new(),
2169                timestamp_ms: None,
2170                compaction_time_window: None,
2171                flushed_entry_id: None,
2172                flushed_sequence: None,
2173                committed_sequence: None,
2174            }),
2175            &[0],
2176            builder.file_purger(),
2177        );
2178        scheduler.on_flush_success(builder.region_id());
2179        // No new flush task.
2180        assert_eq!(1, job_scheduler.num_jobs());
2181        // The flush status is cleared.
2182        assert!(scheduler.region_status.is_empty());
2183        for output_rx in output_rxs {
2184            let output = output_rx.await.unwrap().unwrap();
2185            assert_eq!(output, 0);
2186        }
2187    }
2188
2189    // Verifies single-range flat flush path respects append_mode (no dedup) vs dedup when disabled.
2190    #[test]
2191    fn test_memtable_flat_sources_single_range_append_mode_behavior() {
2192        // Build test metadata and flat schema
2193        let metadata = metadata_for_test();
2194        let schema = to_flat_sst_arrow_schema(
2195            &metadata,
2196            &FlatSchemaOptions::from_encoding(metadata.primary_key_encoding),
2197        );
2198
2199        // Prepare a bulk part containing duplicate rows for the same PK and timestamp
2200        // Two rows with identical keys and timestamps (ts = 1000), different field values
2201        let capacity = 16;
2202        let pk_codec = build_primary_key_codec(&metadata);
2203        let mut converter =
2204            BulkPartConverter::new(&metadata, schema.clone(), capacity, pk_codec, true);
2205        let kvs = build_key_values_with_ts_seq_values(
2206            &metadata,
2207            "dup_key".to_string(),
2208            1,
2209            vec![1000i64, 1000i64].into_iter(),
2210            vec![Some(1.0f64), Some(2.0f64)].into_iter(),
2211            1,
2212        );
2213        converter.append_key_values(&kvs).unwrap();
2214        let part = converter.convert().unwrap();
2215
2216        // Helper to build MemtableRanges with a single range from one bulk part.
2217        // We use BulkMemtable directly because it produces record batch iterators.
2218        let build_ranges = |append_mode: bool| -> MemtableRanges {
2219            let memtable = crate::memtable::bulk::BulkMemtable::new(
2220                1,
2221                crate::memtable::bulk::BulkMemtableConfig::default(),
2222                metadata.clone(),
2223                None,
2224                None,
2225                append_mode,
2226                MergeMode::LastRow,
2227            );
2228            memtable.write_bulk(part.clone()).unwrap();
2229            memtable.ranges(None, RangesOptions::for_flush()).unwrap()
2230        };
2231
2232        // Case 1: append_mode = false => dedup happens, total rows should be 1
2233        {
2234            let mem_ranges = build_ranges(false);
2235            assert_eq!(1, mem_ranges.ranges.len());
2236
2237            let options = RegionOptions {
2238                append_mode: false,
2239                merge_mode: Some(MergeMode::LastRow),
2240                ..Default::default()
2241            };
2242
2243            let flat_sources = memtable_flat_sources(
2244                schema.clone(),
2245                mem_ranges,
2246                &options,
2247                metadata.primary_key.len(),
2248            )
2249            .unwrap();
2250            assert!(flat_sources.encoded.is_empty());
2251            assert_eq!(1, flat_sources.sources.len());
2252
2253            // Consume the iterator and count rows
2254            let mut total_rows = 0usize;
2255            for (source, _sequence) in flat_sources.sources {
2256                total_rows += source
2257                    .take_iter()
2258                    .map(|x| x.unwrap().num_rows())
2259                    .sum::<usize>();
2260            }
2261            assert_eq!(1, total_rows, "dedup should keep a single row");
2262        }
2263
2264        // Case 2: append_mode = true => no dedup, total rows should be 2
2265        {
2266            let mem_ranges = build_ranges(true);
2267            assert_eq!(1, mem_ranges.ranges.len());
2268
2269            let options = RegionOptions {
2270                append_mode: true,
2271                ..Default::default()
2272            };
2273
2274            let flat_sources =
2275                memtable_flat_sources(schema, mem_ranges, &options, metadata.primary_key.len())
2276                    .unwrap();
2277            assert!(flat_sources.encoded.is_empty());
2278            assert_eq!(1, flat_sources.sources.len());
2279
2280            let mut total_rows = 0usize;
2281            for (source, _sequence) in flat_sources.sources {
2282                total_rows += source
2283                    .take_iter()
2284                    .map(|x| x.unwrap().num_rows())
2285                    .sum::<usize>();
2286            }
2287            assert_eq!(2, total_rows, "append_mode should preserve duplicates");
2288        }
2289    }
2290
2291    #[tokio::test]
2292    async fn test_schedule_pending_request_on_flush_success() {
2293        common_telemetry::init_default_ut_logging();
2294        let job_scheduler = Arc::new(VecScheduler::default());
2295        let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone());
2296        let (tx, _rx) = mpsc::channel(4);
2297        let mut scheduler = env.mock_flush_scheduler();
2298        let mut builder = VersionControlBuilder::new();
2299        // Overwrites the empty memtable builder.
2300        builder.set_memtable_builder(Arc::new(TimeSeriesMemtableBuilder::default()));
2301        let version_control = Arc::new(builder.build());
2302        // Writes data to the memtable so it is not empty.
2303        let version_data = version_control.current();
2304        write_rows_to_version(&version_data.version, "host0", 0, 10);
2305        let manifest_ctx = env
2306            .mock_manifest_context(version_data.version.metadata.clone())
2307            .await;
2308        // Creates 2 tasks.
2309        let mut tasks: Vec<_> = (0..2)
2310            .map(|_| RegionFlushTask {
2311                region_id: builder.region_id(),
2312                reason: FlushReason::Manual,
2313                senders: Vec::new(),
2314                request_sender: tx.clone(),
2315                access_layer: env.access_layer.clone(),
2316                listener: WorkerListener::default(),
2317                engine_config: Arc::new(MitoConfig::default()),
2318                row_group_size: None,
2319                cache_manager: Arc::new(CacheManager::default()),
2320                manifest_ctx: manifest_ctx.clone(),
2321                index_options: IndexOptions::default(),
2322                flush_semaphore: Arc::new(Semaphore::new(2)),
2323                is_staging: false,
2324                partition_expr: None,
2325            })
2326            .collect();
2327        // Schedule first task.
2328        let task = tasks.pop().unwrap();
2329        scheduler
2330            .schedule_flush(builder.region_id(), &version_control, task)
2331            .unwrap();
2332        // Should schedule 1 flush.
2333        assert_eq!(1, scheduler.region_status.len());
2334        assert_eq!(1, job_scheduler.num_jobs());
2335        // Schedule second task.
2336        let task = tasks.pop().unwrap();
2337        scheduler
2338            .schedule_flush(builder.region_id(), &version_control, task)
2339            .unwrap();
2340        assert!(
2341            scheduler
2342                .region_status
2343                .get(&builder.region_id())
2344                .unwrap()
2345                .pending_task
2346                .is_some()
2347        );
2348
2349        // Check the new version.
2350        let version_data = version_control.current();
2351        assert_eq!(0, version_data.version.memtables.immutables()[0].id());
2352        // Assumes the flush job is finished.
2353        version_control.apply_edit(
2354            Some(RegionEdit {
2355                files_to_add: Vec::new(),
2356                files_to_remove: Vec::new(),
2357                timestamp_ms: None,
2358                compaction_time_window: None,
2359                flushed_entry_id: None,
2360                flushed_sequence: None,
2361                committed_sequence: None,
2362            }),
2363            &[0],
2364            builder.file_purger(),
2365        );
2366        write_rows_to_version(&version_data.version, "host1", 0, 10);
2367        scheduler.on_flush_success(builder.region_id());
2368        assert_eq!(2, job_scheduler.num_jobs());
2369        // The pending task is cleared.
2370        assert!(
2371            scheduler
2372                .region_status
2373                .get(&builder.region_id())
2374                .unwrap()
2375                .pending_task
2376                .is_none()
2377        );
2378    }
2379
2380    #[tokio::test]
2381    async fn test_schedule_pending_request_failure_drains_pending_ddls() {
2382        common_telemetry::init_default_ut_logging();
2383        let job_scheduler = Arc::new(VecScheduler::default());
2384        let env = SchedulerEnv::new().await.scheduler(job_scheduler.clone());
2385        let (tx, _rx) = mpsc::channel(4);
2386        let mut scheduler = env.mock_flush_scheduler();
2387        let mut builder = VersionControlBuilder::new();
2388        builder.set_memtable_builder(Arc::new(TimeSeriesMemtableBuilder::default()));
2389        let version_control = Arc::new(builder.build());
2390
2391        let version_data = version_control.current();
2392        write_rows_to_version(&version_data.version, "host0", 0, 10);
2393        let manifest_ctx = env
2394            .mock_manifest_context(version_data.version.metadata.clone())
2395            .await;
2396
2397        let task = new_test_flush_task(
2398            &env,
2399            builder.region_id(),
2400            FlushReason::Manual,
2401            tx.clone(),
2402            manifest_ctx.clone(),
2403        );
2404        scheduler
2405            .schedule_flush(builder.region_id(), &version_control, task)
2406            .unwrap();
2407
2408        let task = new_test_flush_task(
2409            &env,
2410            builder.region_id(),
2411            FlushReason::Closing,
2412            tx,
2413            manifest_ctx,
2414        );
2415        scheduler
2416            .schedule_flush(builder.region_id(), &version_control, task)
2417            .unwrap();
2418
2419        let (sender, receiver) = oneshot::channel();
2420        scheduler.add_ddl_request_to_pending(SenderDdlRequest {
2421            sender: OptionOutputTx::from(sender),
2422            region_id: builder.region_id(),
2423            request: DdlRequest::Close(store_api::region_request::RegionCloseRequest::default()),
2424        });
2425
2426        let version_data = version_control.current();
2427        version_control.apply_edit(
2428            Some(RegionEdit {
2429                files_to_add: Vec::new(),
2430                files_to_remove: Vec::new(),
2431                timestamp_ms: None,
2432                compaction_time_window: None,
2433                flushed_entry_id: None,
2434                flushed_sequence: None,
2435                committed_sequence: None,
2436            }),
2437            &[0],
2438            builder.file_purger(),
2439        );
2440        write_rows_to_version(&version_data.version, "host1", 0, 10);
2441
2442        scheduler.scheduler = Arc::new(FailingScheduler);
2443        scheduler.on_flush_success(builder.region_id());
2444
2445        assert!(scheduler.region_status.is_empty());
2446        let err = receiver
2447            .await
2448            .expect("pending DDL must be notified")
2449            .unwrap_err();
2450        assert_eq!(err.status_code(), StatusCode::RegionBusy);
2451    }
2452}