Skip to main content

risingwave_storage/hummock/event_handler/
hummock_event_handler.rs

1// Copyright 2022 RisingWave Labs
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
15use std::cell::RefCell;
16use std::collections::{HashMap, HashSet, VecDeque};
17use std::pin::pin;
18use std::sync::atomic::AtomicUsize;
19use std::sync::atomic::Ordering::Relaxed;
20use std::sync::{Arc, LazyLock};
21use std::time::Duration;
22
23use arc_swap::ArcSwap;
24use await_tree::{InstrumentAwait, SpanExt};
25use futures::FutureExt;
26use itertools::Itertools;
27use parking_lot::RwLock;
28use prometheus::{Histogram, IntGauge, IntGaugeVec};
29use risingwave_common::bitmap::Bitmap;
30use risingwave_common::catalog::TableId;
31use risingwave_common::config::Role;
32use risingwave_common::config::streaming::CacheRefillPolicy;
33use risingwave_common::metrics::UintGauge;
34use risingwave_hummock_sdk::compaction_group::hummock_version_ext::SstDeltaInfo;
35use risingwave_hummock_sdk::{HummockEpoch, SyncResult};
36use risingwave_pb::meta::table_cache_refill_policies::table_cache_refill_policy::PbCacheRefillPolicy;
37use risingwave_pb::meta::{PbServingTableVnodeMappings, PbTableRefillRuntimeConfig};
38use tokio::spawn;
39use tokio::sync::mpsc::error::SendError;
40use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender, unbounded_channel};
41use tokio::sync::oneshot;
42use tracing::{debug, error, info, trace, warn};
43
44use super::refiller::{CacheRefillConfig, CacheRefiller};
45use super::{LocalInstanceGuard, LocalInstanceId, ReadVersionMappingType};
46use crate::compaction_catalog_manager::CompactionCatalogManagerRef;
47use crate::hummock::compactor::{CompactorContext, await_tree_key, compact};
48use crate::hummock::event_handler::refiller::{CacheRefillerEvent, SpawnRefillTask};
49use crate::hummock::event_handler::uploader::{
50    HummockUploader, SpawnUploadTask, SyncedData, UploadTaskOutput,
51};
52use crate::hummock::event_handler::{
53    HummockEvent, HummockObserverEvent, HummockReadVersionRef, HummockVersionUpdate,
54    ReadOnlyReadVersionMapping, ReadOnlyRwLockRef,
55};
56use crate::hummock::local_version::pinned_version::PinnedVersion;
57use crate::hummock::local_version::recent_versions::RecentVersions;
58use crate::hummock::store::version::{HummockReadVersion, StagingSstableInfo, VersionUpdate};
59use crate::hummock::{HummockResult, MemoryLimiter, ObjectIdManager, SstableStoreRef};
60use crate::mem_table::ImmutableMemtable;
61use crate::monitor::HummockStateStoreMetrics;
62use crate::opts::StorageOpts;
63
64#[derive(Clone)]
65pub(crate) struct BufferTracker {
66    flush_threshold: usize,
67    min_batch_flush_size: usize,
68    global_uploading_memory_limiter: Arc<MemoryLimiter>,
69    uploader_imm_size: UintGauge,
70    uploader_uploading_task_size: UintGauge,
71}
72
73impl BufferTracker {
74    pub fn from_storage_opts(config: &StorageOpts, metrics: &HummockStateStoreMetrics) -> Self {
75        let capacity = config.shared_buffer_capacity_mb * (1 << 20);
76        let flush_threshold = (capacity as f32 * config.shared_buffer_flush_ratio) as usize;
77        let min_batch_flush_size = config.shared_buffer_min_batch_flush_size_mb * (1 << 20);
78        assert!(
79            flush_threshold < capacity,
80            "flush_threshold {} should be less or equal to capacity {}",
81            flush_threshold,
82            capacity
83        );
84        Self {
85            flush_threshold,
86            min_batch_flush_size,
87            global_uploading_memory_limiter: Arc::new(MemoryLimiter::new(capacity as u64)),
88            uploader_imm_size: metrics.uploader_imm_size.clone(),
89            uploader_uploading_task_size: metrics.uploader_uploading_task_size.clone(),
90        }
91    }
92
93    #[cfg(test)]
94    fn for_test_with_config(flush_threshold: usize, min_batch_flush_size: usize) -> Self {
95        Self {
96            flush_threshold,
97            min_batch_flush_size,
98            ..Self::for_test()
99        }
100    }
101
102    #[cfg(test)]
103    pub fn for_test() -> Self {
104        Self::from_storage_opts(&StorageOpts::default(), &HummockStateStoreMetrics::unused())
105    }
106
107    pub fn get_memory_limiter(&self) -> &Arc<MemoryLimiter> {
108        &self.global_uploading_memory_limiter
109    }
110
111    pub fn global_upload_task_size(&self) -> &UintGauge {
112        &self.uploader_uploading_task_size
113    }
114
115    /// Return true when the buffer size minus current upload task size is still greater than the
116    /// flush threshold.
117    pub fn need_flush(&self) -> bool {
118        self.uploader_imm_size.get()
119            > self.flush_threshold as u64 + self.uploader_uploading_task_size.get()
120    }
121
122    pub fn need_more_flush(&self, curr_batch_flush_size: usize) -> bool {
123        curr_batch_flush_size < self.min_batch_flush_size || self.need_flush()
124    }
125
126    #[cfg(test)]
127    pub(crate) fn flush_threshold(&self) -> usize {
128        self.flush_threshold
129    }
130}
131
132#[derive(Clone)]
133pub struct HummockEventSender {
134    inner: UnboundedSender<HummockEvent>,
135    event_count: IntGaugeVec,
136}
137
138pub fn event_channel(event_count: IntGaugeVec) -> (HummockEventSender, HummockEventReceiver) {
139    let (tx, rx) = unbounded_channel();
140    (
141        HummockEventSender {
142            inner: tx,
143            event_count: event_count.clone(),
144        },
145        HummockEventReceiver {
146            inner: rx,
147            event_count,
148        },
149    )
150}
151
152impl HummockEventSender {
153    pub fn send(&self, event: HummockEvent) -> Result<(), SendError<HummockEvent>> {
154        let event_type = event.event_name();
155        self.inner.send(event)?;
156        get_event_pending_gauge(&self.event_count, event_type).inc();
157        Ok(())
158    }
159}
160
161pub struct HummockEventReceiver {
162    inner: UnboundedReceiver<HummockEvent>,
163    event_count: IntGaugeVec,
164}
165
166impl HummockEventReceiver {
167    async fn recv(&mut self) -> Option<HummockEvent> {
168        let event = self.inner.recv().await?;
169        let event_type = event.event_name();
170        get_event_pending_gauge(&self.event_count, event_type).dec();
171        Some(event)
172    }
173}
174
175thread_local! {
176    static EVENT_PENDING_GAUGE_CACHE: RefCell<HashMap<&'static str, IntGauge>> = RefCell::new(HashMap::new());
177}
178
179fn get_event_pending_gauge(event_count: &IntGaugeVec, event_type: &'static str) -> IntGauge {
180    EVENT_PENDING_GAUGE_CACHE.with(|cache| {
181        cache
182            .borrow_mut()
183            .entry(event_type)
184            .or_insert_with(|| event_count.with_label_values(&[event_type]))
185            .clone()
186    })
187}
188
189struct HummockEventHandlerMetrics {
190    event_handler_on_upload_finish_latency: Histogram,
191    event_handler_on_apply_version_update: Histogram,
192    event_handler_on_recv_version_update: Histogram,
193    event_handler_on_spiller: Histogram,
194    event_handler_on_apply_table_refill_runtime_config: Histogram,
195}
196
197fn table_cache_refill_policy_from_protobuf(
198    pb_policy: risingwave_pb::meta::table_cache_refill_policies::PbTableCacheRefillPolicy,
199) -> Option<(TableId, CacheRefillPolicy)> {
200    let table_id = TableId::from(pb_policy.table_id);
201    let policy = match PbCacheRefillPolicy::try_from(pb_policy.policy) {
202        Ok(policy) => CacheRefillPolicy::from_protobuf(policy)?,
203        Err(_) => {
204            tracing::warn!(
205                table_id = pb_policy.table_id,
206                policy = pb_policy.policy,
207                "ignore unknown table cache refill policy"
208            );
209            return None;
210        }
211    };
212    Some((table_id, policy))
213}
214
215fn serving_table_vnode_mappings_to_map(
216    mappings: PbServingTableVnodeMappings,
217) -> HashMap<TableId, Bitmap> {
218    mappings
219        .mappings
220        .into_iter()
221        .filter_map(|mapping| {
222            let Some(bitmap) = mapping.bitmap else {
223                tracing::warn!(
224                    table_id = mapping.table_id,
225                    "ignore serving table vnode mapping without bitmap"
226                );
227                return None;
228            };
229            Some((TableId::from(mapping.table_id), Bitmap::from(bitmap)))
230        })
231        .collect()
232}
233
234#[derive(Default)]
235struct LocalReadVersions {
236    by_instance: HashMap<LocalInstanceId, (TableId, HummockReadVersionRef)>,
237    by_table: HashMap<TableId, HashMap<LocalInstanceId, HummockReadVersionRef>>,
238}
239
240impl LocalReadVersions {
241    fn insert(
242        &mut self,
243        instance_id: LocalInstanceId,
244        table_id: TableId,
245        read_version: HummockReadVersionRef,
246    ) {
247        assert!(
248            !self.by_instance.contains_key(&instance_id),
249            "duplicated local read version instance_id {}",
250            instance_id
251        );
252        assert!(
253            !self
254                .by_table
255                .get(&table_id)
256                .is_some_and(|read_versions| read_versions.contains_key(&instance_id)),
257            "duplicated local read version table_id {} instance_id {}",
258            table_id,
259            instance_id
260        );
261        self.by_instance
262            .insert(instance_id, (table_id, read_version.clone()));
263        self.by_table
264            .entry(table_id)
265            .or_default()
266            .insert(instance_id, read_version);
267    }
268
269    fn remove(&mut self, instance_id: LocalInstanceId) -> (TableId, HummockReadVersionRef) {
270        let (table_id, read_version) = self.by_instance.remove(&instance_id).unwrap_or_else(|| {
271            panic!(
272                "DestroyHummockInstance nonexistent instance instance_id {}",
273                instance_id
274            )
275        });
276        let table_read_versions = self.by_table.get_mut(&table_id).unwrap_or_else(|| {
277            panic!(
278                "DestroyHummockInstance table_id {} instance_id {} fail in local mapping",
279                table_id, instance_id
280            )
281        });
282        let removed = table_read_versions.remove(&instance_id).unwrap_or_else(|| {
283            panic!(
284                "DestroyHummockInstance nonexistent instance in local mapping table_id {} instance_id {}",
285                table_id, instance_id
286            )
287        });
288        debug_assert!(Arc::ptr_eq(&read_version, &removed));
289        if table_read_versions.is_empty() {
290            self.by_table.remove(&table_id);
291        }
292        (table_id, read_version)
293    }
294
295    fn get(&self, instance_id: &LocalInstanceId) -> Option<&(TableId, HummockReadVersionRef)> {
296        self.by_instance.get(instance_id)
297    }
298
299    fn contains_instance(&self, instance_id: &LocalInstanceId) -> bool {
300        self.by_instance.contains_key(instance_id)
301    }
302
303    fn instance_ids(&self) -> impl Iterator<Item = LocalInstanceId> + '_ {
304        self.by_instance.keys().copied()
305    }
306
307    fn is_empty(&self) -> bool {
308        self.by_instance.is_empty() && self.by_table.is_empty()
309    }
310
311    fn table_ids(&self) -> impl Iterator<Item = TableId> + '_ {
312        self.by_table.keys().copied()
313    }
314
315    fn table_vnodes(&self, table_id: TableId) -> Option<Bitmap> {
316        let read_versions = self.by_table.get(&table_id)?;
317        let mut vnodes = None;
318        for read_version in read_versions.values() {
319            let read_version_vnodes = read_version.read().vnodes();
320            let current_vnodes =
321                vnodes.get_or_insert_with(|| Bitmap::zeros(read_version_vnodes.as_ref().len()));
322            *current_vnodes |= read_version_vnodes.as_ref();
323        }
324        vnodes
325    }
326}
327
328pub(crate) struct HummockEventHandler {
329    hummock_event_tx: HummockEventSender,
330    hummock_event_rx: HummockEventReceiver,
331    observer_event_rx: UnboundedReceiver<HummockObserverEvent>,
332    read_version_mapping: Arc<RwLock<ReadVersionMappingType>>,
333    /// A copy of `read_version_mapping` but owned by event handler
334    local_read_versions: LocalReadVersions,
335
336    version_update_notifier_tx: Arc<tokio::sync::watch::Sender<PinnedVersion>>,
337    recent_versions: Arc<ArcSwap<RecentVersions>>,
338
339    uploader: HummockUploader,
340    refiller: CacheRefiller,
341
342    last_instance_id: LocalInstanceId,
343
344    metrics: HummockEventHandlerMetrics,
345}
346
347async fn flush_imms(
348    payload: Vec<ImmutableMemtable>,
349    compactor_context: CompactorContext,
350    compaction_catalog_manager_ref: CompactionCatalogManagerRef,
351    object_id_manager: Arc<ObjectIdManager>,
352) -> HummockResult<UploadTaskOutput> {
353    compact(
354        compactor_context,
355        object_id_manager,
356        payload,
357        compaction_catalog_manager_ref,
358    )
359    .instrument_await("shared_buffer_compact".verbose())
360    .await
361}
362
363impl HummockEventHandler {
364    // Table refill runtime config applies by field presence for both snapshot
365    // and update notifications.
366    fn apply_table_refill_runtime_config(
367        refiller: &mut CacheRefiller,
368        config: PbTableRefillRuntimeConfig,
369    ) {
370        if let Some(policies) = config.table_cache_refill_policies {
371            let policies = policies
372                .table_policies
373                .into_iter()
374                .chain(policies.internal_table_policies)
375                .filter_map(table_cache_refill_policy_from_protobuf)
376                .collect();
377            refiller.replace_table_cache_refill_policies(policies);
378        }
379
380        if let Some(mappings) = config.serving_table_vnode_mappings {
381            refiller
382                .replace_serving_table_vnode_mapping(serving_table_vnode_mappings_to_map(mappings));
383        }
384    }
385
386    pub fn new(
387        role: Role,
388        observer_event_rx: UnboundedReceiver<HummockObserverEvent>,
389        pinned_version: PinnedVersion,
390        compactor_context: CompactorContext,
391        compaction_catalog_manager_ref: CompactionCatalogManagerRef,
392        object_id_manager: Arc<ObjectIdManager>,
393        state_store_metrics: Arc<HummockStateStoreMetrics>,
394    ) -> Self {
395        let upload_compactor_context = compactor_context.clone();
396        let upload_task_latency = state_store_metrics.uploader_upload_task_latency.clone();
397        let wait_poll_latency = state_store_metrics.uploader_wait_poll_latency.clone();
398        let recent_versions = RecentVersions::new(
399            pinned_version,
400            compactor_context
401                .storage_opts
402                .max_cached_recent_versions_number,
403            state_store_metrics.clone(),
404        );
405        let buffer_tracker =
406            BufferTracker::from_storage_opts(&compactor_context.storage_opts, &state_store_metrics);
407        Self::new_inner(
408            role,
409            observer_event_rx,
410            compactor_context.sstable_store.clone(),
411            state_store_metrics,
412            CacheRefillConfig::from_storage_opts(&compactor_context.storage_opts),
413            recent_versions,
414            buffer_tracker,
415            Arc::new(move |payload, task_info| {
416                static NEXT_UPLOAD_TASK_ID: LazyLock<AtomicUsize> =
417                    LazyLock::new(|| AtomicUsize::new(0));
418                let tree_root = upload_compactor_context.await_tree_reg.as_ref().map(|reg| {
419                    let upload_task_id = NEXT_UPLOAD_TASK_ID.fetch_add(1, Relaxed);
420                    reg.register(
421                        await_tree_key::SpawnUploadTask { id: upload_task_id },
422                        format!("Spawn Upload Task: {}", task_info),
423                    )
424                });
425                let upload_task_latency = upload_task_latency.clone();
426                let wait_poll_latency = wait_poll_latency.clone();
427                let upload_compactor_context = upload_compactor_context.clone();
428                let compaction_catalog_manager_ref = compaction_catalog_manager_ref.clone();
429                let object_id_manager = object_id_manager.clone();
430                spawn({
431                    let future = async move {
432                        let _timer = upload_task_latency.start_timer();
433                        let mut output = flush_imms(
434                            payload
435                                .into_values()
436                                .flat_map(|imms| imms.into_iter())
437                                .collect(),
438                            upload_compactor_context.clone(),
439                            compaction_catalog_manager_ref.clone(),
440                            object_id_manager.clone(),
441                        )
442                        .await?;
443                        assert!(
444                            output
445                                .wait_poll_timer
446                                .replace(wait_poll_latency.start_timer())
447                                .is_none(),
448                            "should not set timer before"
449                        );
450                        Ok(output)
451                    };
452                    if let Some(tree_root) = tree_root {
453                        tree_root.instrument(future).left_future()
454                    } else {
455                        future.right_future()
456                    }
457                })
458            }),
459            CacheRefiller::default_spawn_refill_task(),
460        )
461    }
462
463    fn new_inner(
464        role: Role,
465        observer_event_rx: UnboundedReceiver<HummockObserverEvent>,
466        sstable_store: SstableStoreRef,
467        state_store_metrics: Arc<HummockStateStoreMetrics>,
468        refill_config: CacheRefillConfig,
469        recent_versions: RecentVersions,
470        buffer_tracker: BufferTracker,
471        spawn_upload_task: SpawnUploadTask,
472        spawn_refill_task: SpawnRefillTask,
473    ) -> Self {
474        let (hummock_event_tx, hummock_event_rx) =
475            event_channel(state_store_metrics.event_handler_pending_event.clone());
476        let (version_update_notifier_tx, _) =
477            tokio::sync::watch::channel(recent_versions.latest_version().clone());
478        let version_update_notifier_tx = Arc::new(version_update_notifier_tx);
479        let read_version_mapping = Arc::new(RwLock::new(HashMap::default()));
480
481        let metrics = HummockEventHandlerMetrics {
482            event_handler_on_upload_finish_latency: state_store_metrics
483                .event_handler_latency
484                .with_label_values(&["on_upload_finish"]),
485            event_handler_on_apply_version_update: state_store_metrics
486                .event_handler_latency
487                .with_label_values(&["apply_version"]),
488            event_handler_on_recv_version_update: state_store_metrics
489                .event_handler_latency
490                .with_label_values(&["recv_version_update"]),
491            event_handler_on_spiller: state_store_metrics
492                .event_handler_latency
493                .with_label_values(&["spiller"]),
494            event_handler_on_apply_table_refill_runtime_config: state_store_metrics
495                .event_handler_latency
496                .with_label_values(&["apply_table_refill_runtime_config"]),
497        };
498
499        let uploader = HummockUploader::new(
500            state_store_metrics,
501            recent_versions.latest_version().clone(),
502            spawn_upload_task,
503            buffer_tracker,
504        );
505        let refiller = CacheRefiller::new(role, refill_config, sstable_store, spawn_refill_task);
506
507        Self {
508            hummock_event_tx,
509            hummock_event_rx,
510            observer_event_rx,
511            version_update_notifier_tx,
512            recent_versions: Arc::new(ArcSwap::from_pointee(recent_versions)),
513            read_version_mapping,
514            local_read_versions: Default::default(),
515            uploader,
516            refiller,
517            last_instance_id: 0,
518            metrics,
519        }
520    }
521
522    pub fn version_update_notifier_tx(&self) -> Arc<tokio::sync::watch::Sender<PinnedVersion>> {
523        self.version_update_notifier_tx.clone()
524    }
525
526    pub fn recent_versions(&self) -> Arc<ArcSwap<RecentVersions>> {
527        self.recent_versions.clone()
528    }
529
530    pub fn read_version_mapping(&self) -> ReadOnlyReadVersionMapping {
531        ReadOnlyRwLockRef::new(self.read_version_mapping.clone())
532    }
533
534    pub fn event_sender(&self) -> HummockEventSender {
535        self.hummock_event_tx.clone()
536    }
537
538    pub fn buffer_tracker(&self) -> &BufferTracker {
539        self.uploader.buffer_tracker()
540    }
541}
542
543// Handler for different events
544impl HummockEventHandler {
545    /// Applies the function to local read versions selected by instance ids.
546    /// Each read version is protected by its own write lock.
547    fn for_each_read_version(
548        &self,
549        instances: impl IntoIterator<Item = LocalInstanceId>,
550        mut f: impl FnMut(LocalInstanceId, &mut HummockReadVersion),
551    ) {
552        let instances = {
553            #[cfg(debug_assertions)]
554            {
555                // check duplication on debug_mode
556                let mut id_set = std::collections::HashSet::new();
557                for instance in instances {
558                    assert!(id_set.insert(instance));
559                }
560                id_set
561            }
562            #[cfg(not(debug_assertions))]
563            {
564                instances
565            }
566        };
567        let mut pending = VecDeque::new();
568        let mut total_count = 0;
569        for instance_id in instances {
570            let Some((_, read_version)) = self.local_read_versions.get(&instance_id) else {
571                continue;
572            };
573            total_count += 1;
574            match read_version.try_write() {
575                Some(mut write_guard) => {
576                    f(instance_id, &mut write_guard);
577                }
578                _ => {
579                    pending.push_back(instance_id);
580                }
581            }
582        }
583        if !pending.is_empty() {
584            if pending.len() * 10 > total_count {
585                // Only print warn log when failed to acquire more than 10%
586                warn!(
587                    pending_count = pending.len(),
588                    total_count, "cannot acquire lock for all read version"
589                );
590            } else {
591                debug!(
592                    pending_count = pending.len(),
593                    total_count, "cannot acquire lock for all read version"
594                );
595            }
596        }
597
598        const TRY_LOCK_TIMEOUT: Duration = Duration::from_millis(1);
599
600        while let Some(instance_id) = pending.pop_front() {
601            let (_, read_version) = self
602                .local_read_versions
603                .get(&instance_id)
604                .expect("have checked exist before");
605            match read_version.try_write_for(TRY_LOCK_TIMEOUT) {
606                Some(mut write_guard) => {
607                    f(instance_id, &mut write_guard);
608                }
609                _ => {
610                    warn!(instance_id, "failed to get lock again for instance");
611                    pending.push_back(instance_id);
612                }
613            }
614        }
615    }
616
617    fn handle_uploaded_ssts_inner(&mut self, ssts: Vec<Arc<StagingSstableInfo>>) {
618        match ssts.as_slice() {
619            [] => {
620                if cfg!(debug_assertions) {
621                    panic!("empty ssts")
622                }
623            }
624            [staging_sstable_info] => {
625                trace!("data_flushed. SST size {}", staging_sstable_info.imm_size());
626                self.for_each_read_version(
627                    staging_sstable_info.imm_ids().keys().cloned(),
628                    |_, read_version| {
629                        read_version.update(VersionUpdate::Sst(staging_sstable_info.clone()))
630                    },
631                )
632            }
633            ssts => {
634                warn!(
635                    batch_size = ssts.len(),
636                    "handle multiple uploaded ssts in batch"
637                );
638                let affected_instances: HashSet<_> = ssts
639                    .iter()
640                    .flat_map(|sst| {
641                        trace!("data_flushed. SST size {}", sst.imm_size());
642                        sst.imm_ids().keys()
643                    })
644                    .copied()
645                    .collect();
646                self.for_each_read_version(affected_instances, |instance_id, read_version| {
647                    for sst in ssts {
648                        if sst.imm_ids().contains_key(&instance_id) {
649                            read_version.update(VersionUpdate::Sst(sst.clone()));
650                        }
651                    }
652                })
653            }
654        }
655    }
656
657    fn handle_sync_epoch(
658        &mut self,
659        sync_table_epochs: Vec<(HummockEpoch, HashSet<TableId>)>,
660        sync_result_sender: oneshot::Sender<HummockResult<SyncedData>>,
661    ) {
662        debug!(?sync_table_epochs, "awaiting for epoch to be synced",);
663        self.uploader
664            .start_sync_epoch(sync_result_sender, sync_table_epochs);
665    }
666
667    fn handle_clear(&mut self, notifier: oneshot::Sender<()>, table_ids: Option<HashSet<TableId>>) {
668        info!(
669            current_version_id = ?self.uploader.hummock_version().id(),
670            ?table_ids,
671            "handle clear event"
672        );
673
674        self.uploader.clear(table_ids.clone());
675
676        if table_ids.is_none() {
677            assert!(
678                self.local_read_versions.is_empty(),
679                "read version mapping not empty when clear. remaining tables: {:?}",
680                self.local_read_versions.table_ids().collect_vec()
681            );
682        }
683
684        // Notify completion of the Clear event.
685        let _ = notifier.send(()).inspect_err(|e| {
686            error!("failed to notify completion of clear event: {:?}", e);
687        });
688
689        info!("clear finished");
690    }
691
692    fn handle_observer_event(&mut self, event: HummockObserverEvent) {
693        match event {
694            HummockObserverEvent::VersionUpdate(version_update) => {
695                self.handle_version_update(version_update);
696            }
697            HummockObserverEvent::TableRefillRuntimeConfig(_, config) => {
698                let _timer = self
699                    .metrics
700                    .event_handler_on_apply_table_refill_runtime_config
701                    .start_timer();
702                Self::apply_table_refill_runtime_config(&mut self.refiller, config);
703            }
704        }
705    }
706
707    fn streaming_table_vnodes(&self, table_id: TableId) -> Option<Bitmap> {
708        self.local_read_versions.table_vnodes(table_id)
709    }
710
711    fn handle_version_update(&mut self, version_payload: HummockVersionUpdate) {
712        let _timer = self
713            .metrics
714            .event_handler_on_recv_version_update
715            .start_timer();
716        let pinned_version = self
717            .refiller
718            .last_new_pinned_version()
719            .cloned()
720            .unwrap_or_else(|| self.uploader.hummock_version().clone());
721
722        let mut sst_delta_infos = vec![];
723        if let Some(new_pinned_version) = Self::resolve_version_update_info(
724            &pinned_version,
725            version_payload,
726            Some(&mut sst_delta_infos),
727        ) {
728            self.refiller
729                .start_cache_refill(sst_delta_infos, pinned_version, new_pinned_version);
730        }
731    }
732
733    fn resolve_version_update_info(
734        pinned_version: &PinnedVersion,
735        version_payload: HummockVersionUpdate,
736        mut sst_delta_infos: Option<&mut Vec<SstDeltaInfo>>,
737    ) -> Option<PinnedVersion> {
738        match version_payload {
739            HummockVersionUpdate::VersionDeltas(version_deltas) => {
740                let mut version_to_apply = (**pinned_version).clone();
741                {
742                    for version_delta in version_deltas {
743                        assert_eq!(version_to_apply.id, version_delta.prev_id);
744                        if let Some(sst_delta_infos) = &mut sst_delta_infos {
745                            sst_delta_infos
746                                .extend(version_to_apply.build_sst_delta_infos(&version_delta));
747                        }
748
749                        version_to_apply.apply_version_delta(&version_delta);
750                    }
751                }
752
753                pinned_version.new_with_local_version(version_to_apply)
754            }
755            HummockVersionUpdate::PinnedVersion(version) => {
756                pinned_version.new_pin_version(*version)
757            }
758        }
759    }
760
761    fn apply_version_updates(&mut self, events: Vec<CacheRefillerEvent>) {
762        let Some(CacheRefillerEvent {
763            new_pinned_version: latest_pinned_version,
764            ..
765        }) = events.last()
766        else {
767            if cfg!(debug_assertions) {
768                panic!("empty events")
769            }
770            return;
771        };
772        if events.len() > 1 {
773            warn!(
774                count = events.len(),
775                "handle multiple version updates in batch"
776            );
777        }
778        let _timer = self
779            .metrics
780            .event_handler_on_apply_version_update
781            .start_timer();
782        self.recent_versions.rcu(|prev_recent_versions| {
783            let mut recent_versions = None;
784            for event in &events {
785                let CacheRefillerEvent {
786                    new_pinned_version, ..
787                } = event;
788                recent_versions = Some(
789                    recent_versions
790                        .as_ref()
791                        .unwrap_or(prev_recent_versions.as_ref())
792                        .with_new_version(new_pinned_version.clone()),
793                );
794            }
795            recent_versions.expect("non-empty events")
796        });
797
798        {
799            self.for_each_read_version(
800                self.local_read_versions.instance_ids(),
801                |_, read_version| {
802                    for CacheRefillerEvent {
803                        new_pinned_version, ..
804                    } in &events
805                    {
806                        read_version
807                            .update(VersionUpdate::CommittedSnapshot(new_pinned_version.clone()))
808                    }
809                },
810            );
811        }
812
813        self.version_update_notifier_tx.send_if_modified(|state| {
814            let mut modified = false;
815            for CacheRefillerEvent {
816                pinned_version,
817                new_pinned_version,
818            } in &events
819            {
820                assert_eq!(pinned_version.id(), state.id());
821                if state.id() == new_pinned_version.id() {
822                    continue;
823                }
824                assert!(new_pinned_version.id() > state.id());
825                *state = new_pinned_version.clone();
826                modified = true;
827            }
828            modified
829        });
830
831        debug!("update to hummock version: {}", latest_pinned_version.id(),);
832
833        self.uploader
834            .update_pinned_version(latest_pinned_version.clone());
835    }
836}
837
838impl HummockEventHandler {
839    pub async fn start_hummock_event_handler_worker(mut self) {
840        loop {
841            tokio::select! {
842                ssts = self.uploader.next_uploaded_ssts() => {
843                    self.handle_uploaded_ssts(ssts);
844                }
845                events = self.refiller.next_events() => {
846                    self.apply_version_updates(events);
847                }
848                event = pin!(self.hummock_event_rx.recv()) => {
849                    let Some(event) = event else { break };
850                    match event {
851                        HummockEvent::Shutdown => {
852                            info!("event handler shutdown");
853                            return;
854                        },
855                        event => {
856                            self.handle_hummock_event(event);
857                        }
858                    }
859                }
860                observer_event = pin!(self.observer_event_rx.recv()) => {
861                    let Some(observer_event) = observer_event else {
862                        warn!("observer event stream ends. event handle shutdown");
863                        return;
864                    };
865                    self.handle_observer_event(observer_event);
866                }
867            }
868        }
869    }
870
871    fn handle_uploaded_ssts(&mut self, ssts: Vec<Arc<StagingSstableInfo>>) {
872        let _timer = self
873            .metrics
874            .event_handler_on_upload_finish_latency
875            .start_timer();
876        self.handle_uploaded_ssts_inner(ssts);
877    }
878
879    /// Gracefully shutdown if returns `true`.
880    fn handle_hummock_event(&mut self, event: HummockEvent) {
881        match event {
882            HummockEvent::BufferMayFlush => {
883                self.uploader
884                    .may_flush(&self.metrics.event_handler_on_spiller);
885            }
886            HummockEvent::SyncEpoch {
887                sync_result_sender,
888                sync_table_epochs,
889            } => {
890                self.handle_sync_epoch(sync_table_epochs, sync_result_sender);
891            }
892            HummockEvent::Clear(notifier, table_ids) => {
893                self.handle_clear(notifier, table_ids);
894            }
895            HummockEvent::Shutdown => {
896                unreachable!("shutdown is handled specially")
897            }
898            HummockEvent::StartEpoch { epoch, table_ids } => {
899                self.uploader.start_epoch(epoch, table_ids);
900            }
901            HummockEvent::InitEpoch {
902                instance_id,
903                init_epoch,
904            } => {
905                let table_id = self
906                    .local_read_versions
907                    .get(&instance_id)
908                    .expect("should exist")
909                    .0;
910                self.uploader
911                    .init_instance(instance_id, table_id, init_epoch);
912            }
913            HummockEvent::ImmToUploader { instance_id, imms } => {
914                assert!(
915                    self.local_read_versions.contains_instance(&instance_id),
916                    "add imm from non-existing read version instance: instance_id: {}, table_id {:?}",
917                    instance_id,
918                    imms.first().map(|(imm, _)| imm.table_id),
919                );
920                self.uploader.add_imms(instance_id, imms);
921                self.uploader
922                    .may_flush(&self.metrics.event_handler_on_spiller);
923            }
924
925            HummockEvent::LocalSealEpoch {
926                next_epoch,
927                opts,
928                instance_id,
929            } => {
930                self.uploader
931                    .local_seal_epoch(instance_id, next_epoch, opts);
932            }
933
934            #[cfg(any(test, feature = "test"))]
935            HummockEvent::FlushEvent(sender) => {
936                let _ = sender.send(()).inspect_err(|e| {
937                    error!("unable to send flush result: {:?}", e);
938                });
939            }
940
941            HummockEvent::RegisterReadVersion {
942                table_id,
943                new_read_version_sender,
944                is_replicated,
945                vnodes,
946            } => {
947                let pinned_version = self.recent_versions.load().latest_version().clone();
948                let instance_id = self.generate_instance_id();
949                let basic_read_version = Arc::new(RwLock::new(
950                    HummockReadVersion::new_with_replication_option(
951                        table_id,
952                        instance_id,
953                        pinned_version,
954                        is_replicated,
955                        vnodes,
956                    ),
957                ));
958
959                debug!(
960                    "new read version registered: table_id: {}, instance_id: {}",
961                    table_id, instance_id
962                );
963
964                {
965                    self.local_read_versions.insert(
966                        instance_id,
967                        table_id,
968                        basic_read_version.clone(),
969                    );
970                    let mut read_version_mapping_guard = self.read_version_mapping.write();
971
972                    read_version_mapping_guard
973                        .entry(table_id)
974                        .or_default()
975                        .insert(instance_id, basic_read_version.clone());
976                }
977
978                match new_read_version_sender.send((
979                    basic_read_version,
980                    LocalInstanceGuard {
981                        table_id,
982                        instance_id,
983                        event_sender: Some(self.hummock_event_tx.clone()),
984                    },
985                )) {
986                    Ok(_) => {}
987                    Err((_, mut guard)) => {
988                        warn!(
989                            "RegisterReadVersion send fail table_id {:?} instance_is {:?}",
990                            table_id, instance_id
991                        );
992                        guard.event_sender.take().expect("sender is just set");
993                        self.destroy_read_version(instance_id);
994                    }
995                }
996
997                self.refiller
998                    .update_streaming_table_vnodes(table_id, self.streaming_table_vnodes(table_id));
999            }
1000
1001            HummockEvent::DestroyReadVersion { instance_id } => {
1002                let table_id = self
1003                    .local_read_versions
1004                    .get(&instance_id)
1005                    .unwrap_or_else(|| {
1006                        panic!("query nonexistent instance instance_id {instance_id}")
1007                    })
1008                    .0;
1009                self.uploader.may_destroy_instance(instance_id);
1010                self.destroy_read_version(instance_id);
1011                self.refiller
1012                    .update_streaming_table_vnodes(table_id, self.streaming_table_vnodes(table_id));
1013            }
1014            HummockEvent::GetMinUncommittedObjectId { result_tx } => {
1015                let _ = result_tx
1016                    .send(self.uploader.min_uncommitted_object_id())
1017                    .inspect_err(|e| {
1018                        error!("unable to send get_min_uncommitted_sst_id result: {:?}", e);
1019                    });
1020            }
1021            HummockEvent::GetTableCacheRefillMonitorSnapshot { result_tx } => {
1022                let _ = result_tx
1023                    .send(self.refiller.table_cache_refill_monitor_snapshot())
1024                    .inspect_err(|_| {
1025                        error!("unable to send table cache refill monitor snapshot result");
1026                    });
1027            }
1028            HummockEvent::RegisterVectorWriter {
1029                table_id,
1030                init_epoch,
1031            } => self.uploader.register_vector_writer(table_id, init_epoch),
1032            HummockEvent::VectorWriterSealEpoch {
1033                table_id,
1034                next_epoch,
1035                add,
1036            } => {
1037                self.uploader
1038                    .vector_writer_seal_epoch(table_id, next_epoch, add);
1039            }
1040            HummockEvent::DropVectorWriter { table_id } => {
1041                self.uploader.drop_vector_writer(table_id);
1042            }
1043        }
1044    }
1045
1046    fn destroy_read_version(&mut self, instance_id: LocalInstanceId) {
1047        {
1048            {
1049                debug!("read version deregister: instance_id: {}", instance_id);
1050                let (table_id, _) = self.local_read_versions.remove(instance_id);
1051                let mut read_version_mapping_guard = self.read_version_mapping.write();
1052                let entry = read_version_mapping_guard
1053                    .get_mut(&table_id)
1054                    .unwrap_or_else(|| {
1055                        panic!(
1056                            "DestroyHummockInstance table_id {} instance_id {} fail",
1057                            table_id, instance_id
1058                        )
1059                    });
1060                entry.remove(&instance_id).unwrap_or_else(|| {
1061                    panic!(
1062                        "DestroyHummockInstance nonexistent instance table_id {} instance_id {}",
1063                        table_id, instance_id
1064                    )
1065                });
1066                if entry.is_empty() {
1067                    read_version_mapping_guard.remove(&table_id);
1068                }
1069            }
1070        }
1071    }
1072
1073    fn generate_instance_id(&mut self) -> LocalInstanceId {
1074        self.last_instance_id += 1;
1075        self.last_instance_id
1076    }
1077}
1078
1079pub(super) fn send_sync_result(
1080    sender: oneshot::Sender<HummockResult<SyncedData>>,
1081    result: HummockResult<SyncedData>,
1082) {
1083    let _ = sender.send(result).inspect_err(|e| {
1084        error!("unable to send sync result. Err: {:?}", e);
1085    });
1086}
1087
1088impl SyncedData {
1089    pub fn into_sync_result(self) -> SyncResult {
1090        {
1091            let SyncedData {
1092                uploaded_ssts,
1093                table_watermarks,
1094                vector_index_adds,
1095            } = self;
1096            let mut sync_size = 0;
1097            let mut uncommitted_ssts = Vec::new();
1098            let mut old_value_ssts = Vec::new();
1099            // The newly uploaded `sstable_infos` contains newer data. Therefore,
1100            // `newly_upload_ssts` at the front
1101            for sst in uploaded_ssts {
1102                sync_size += sst.imm_size();
1103                uncommitted_ssts.extend(sst.sstable_infos().iter().cloned());
1104                old_value_ssts.extend(sst.old_value_sstable_infos().iter().cloned());
1105            }
1106            SyncResult {
1107                sync_size,
1108                uncommitted_ssts,
1109                table_watermarks,
1110                old_value_ssts,
1111                vector_index_adds,
1112            }
1113        }
1114    }
1115}
1116
1117#[cfg(test)]
1118mod tests {
1119    use std::collections::{HashMap, HashSet};
1120    use std::future::poll_fn;
1121    use std::sync::Arc;
1122    use std::task::Poll;
1123
1124    use futures::FutureExt;
1125    use parking_lot::Mutex;
1126    use risingwave_common::bitmap::Bitmap;
1127    use risingwave_common::catalog::TableId;
1128    use risingwave_common::config::Role;
1129    use risingwave_common::config::streaming::CacheRefillPolicy;
1130    use risingwave_common::hash::VirtualNode;
1131    use risingwave_common::util::epoch::{EpochExt, test_epoch};
1132    use risingwave_hummock_sdk::compaction_group::StaticCompactionGroupId;
1133    use risingwave_hummock_sdk::version::HummockVersion;
1134    use risingwave_pb::hummock::{PbHummockVersion, StateTableInfo};
1135    use risingwave_pb::meta::serving_table_vnode_mappings::PbServingTableVnodeMapping;
1136    use risingwave_pb::meta::table_cache_refill_policies::PbTableCacheRefillPolicy;
1137    use risingwave_pb::meta::table_cache_refill_policies::table_cache_refill_policy::PbCacheRefillPolicy;
1138    use risingwave_pb::meta::{
1139        PbServingTableVnodeMappings, PbTableRefillRuntimeConfig, TableCacheRefillPolicies,
1140    };
1141    use tokio::spawn;
1142    use tokio::sync::mpsc::unbounded_channel;
1143    use tokio::sync::oneshot;
1144
1145    use crate::hummock::HummockError;
1146    use crate::hummock::event_handler::hummock_event_handler::BufferTracker;
1147    use crate::hummock::event_handler::refiller::{CacheRefillConfig, CacheRefiller};
1148    use crate::hummock::event_handler::uploader::UploadTaskOutput;
1149    use crate::hummock::event_handler::uploader::test_utils::{
1150        TEST_TABLE_ID, gen_imm_inner, gen_imm_with_unlimit,
1151        prepare_uploader_order_test_spawn_task_fn,
1152    };
1153    use crate::hummock::event_handler::{
1154        HummockEvent, HummockEventHandler, HummockReadVersionRef, LocalInstanceGuard,
1155    };
1156    use crate::hummock::iterator::test_utils::mock_sstable_store;
1157    use crate::hummock::local_version::pinned_version::PinnedVersion;
1158    use crate::hummock::local_version::recent_versions::RecentVersions;
1159    use crate::hummock::test_utils::default_opts_for_test;
1160    use crate::mem_table::ImmutableMemtable;
1161    use crate::monitor::HummockStateStoreMetrics;
1162    use crate::store::SealCurrentEpochOptions;
1163
1164    #[tokio::test]
1165    async fn test_runtime_config_fields_are_independent_full_replacements() {
1166        let old_table_id = TableId::new(233);
1167        let internal_table_id = TableId::new(234);
1168        let new_table_id = TableId::new(235);
1169        let old_vnodes = Bitmap::ones(VirtualNode::COUNT_FOR_TEST);
1170        let new_vnodes = Bitmap::from_range(VirtualNode::COUNT_FOR_TEST, 0..8);
1171        let storage_opt = default_opts_for_test();
1172        let mut refiller = CacheRefiller::new(
1173            Role::Both,
1174            CacheRefillConfig::from_storage_opts(&storage_opt),
1175            mock_sstable_store().await,
1176            CacheRefiller::default_spawn_refill_task(),
1177        );
1178
1179        HummockEventHandler::apply_table_refill_runtime_config(
1180            &mut refiller,
1181            PbTableRefillRuntimeConfig {
1182                table_cache_refill_policies: Some(TableCacheRefillPolicies {
1183                    table_policies: vec![PbTableCacheRefillPolicy {
1184                        table_id: old_table_id.as_raw_id(),
1185                        policy: PbCacheRefillPolicy::Serving as i32,
1186                    }],
1187                    internal_table_policies: vec![PbTableCacheRefillPolicy {
1188                        table_id: internal_table_id.as_raw_id(),
1189                        policy: PbCacheRefillPolicy::Both as i32,
1190                    }],
1191                }),
1192                serving_table_vnode_mappings: Some(PbServingTableVnodeMappings {
1193                    mappings: vec![PbServingTableVnodeMapping {
1194                        table_id: old_table_id.as_raw_id(),
1195                        bitmap: Some(old_vnodes.to_protobuf()),
1196                    }],
1197                }),
1198                ..Default::default()
1199            },
1200        );
1201        let snapshot = refiller.table_cache_refill_monitor_snapshot();
1202        assert_eq!(
1203            snapshot.policies,
1204            HashMap::from([
1205                (old_table_id, CacheRefillPolicy::Serving),
1206                (internal_table_id, CacheRefillPolicy::Both),
1207            ])
1208        );
1209        assert_eq!(
1210            snapshot.serving_table_vnode_mapping,
1211            HashMap::from([(old_table_id, old_vnodes)])
1212        );
1213
1214        HummockEventHandler::apply_table_refill_runtime_config(
1215            &mut refiller,
1216            PbTableRefillRuntimeConfig {
1217                serving_table_vnode_mappings: Some(PbServingTableVnodeMappings {
1218                    mappings: vec![PbServingTableVnodeMapping {
1219                        table_id: new_table_id.as_raw_id(),
1220                        bitmap: Some(new_vnodes.to_protobuf()),
1221                    }],
1222                }),
1223                ..Default::default()
1224            },
1225        );
1226        let snapshot = refiller.table_cache_refill_monitor_snapshot();
1227        assert_eq!(
1228            snapshot.policies,
1229            HashMap::from([
1230                (old_table_id, CacheRefillPolicy::Serving),
1231                (internal_table_id, CacheRefillPolicy::Both),
1232            ])
1233        );
1234        assert_eq!(
1235            snapshot.serving_table_vnode_mapping,
1236            HashMap::from([(new_table_id, new_vnodes.clone())])
1237        );
1238
1239        HummockEventHandler::apply_table_refill_runtime_config(
1240            &mut refiller,
1241            PbTableRefillRuntimeConfig {
1242                table_cache_refill_policies: Some(TableCacheRefillPolicies::default()),
1243                ..Default::default()
1244            },
1245        );
1246        let snapshot = refiller.table_cache_refill_monitor_snapshot();
1247        assert!(snapshot.policies.is_empty());
1248        assert_eq!(
1249            snapshot.serving_table_vnode_mapping,
1250            HashMap::from([(new_table_id, new_vnodes)])
1251        );
1252
1253        HummockEventHandler::apply_table_refill_runtime_config(
1254            &mut refiller,
1255            PbTableRefillRuntimeConfig {
1256                serving_table_vnode_mappings: Some(PbServingTableVnodeMappings::default()),
1257                ..Default::default()
1258            },
1259        );
1260        let snapshot = refiller.table_cache_refill_monitor_snapshot();
1261        assert!(snapshot.policies.is_empty());
1262        assert!(snapshot.serving_table_vnode_mapping.is_empty());
1263    }
1264
1265    #[test]
1266    fn test_runtime_config_ignores_invalid_entries() {
1267        let table_id = TableId::new(233);
1268        assert!(
1269            super::table_cache_refill_policy_from_protobuf(PbTableCacheRefillPolicy {
1270                table_id: table_id.as_raw_id(),
1271                policy: i32::MAX,
1272            })
1273            .is_none()
1274        );
1275
1276        assert!(
1277            super::serving_table_vnode_mappings_to_map(PbServingTableVnodeMappings {
1278                mappings: vec![PbServingTableVnodeMapping {
1279                    table_id: table_id.as_raw_id(),
1280                    bitmap: None,
1281                }],
1282            })
1283            .is_empty()
1284        );
1285    }
1286
1287    #[tokio::test]
1288    async fn test_read_version_events_update_refiller_streaming_mapping() {
1289        fn assert_vnodes(vnodes: &Bitmap, expected: std::ops::Range<usize>) {
1290            assert_eq!(vnodes.count_ones(), expected.len());
1291            assert!(expected.into_iter().all(|vnode| vnodes.is_set(vnode)));
1292        }
1293
1294        let table_id = TableId::new(233);
1295        let other_table_id = TableId::new(234);
1296        let (_observer_event_tx, observer_event_rx) = unbounded_channel();
1297        let metrics = Arc::new(HummockStateStoreMetrics::unused());
1298        let storage_opt = default_opts_for_test();
1299        let (spawn_upload_task, _) = prepare_uploader_order_test_spawn_task_fn(false);
1300        let pinned_version = PinnedVersion::new(
1301            HummockVersion::from_rpc_protobuf(&PbHummockVersion {
1302                id: 1.into(),
1303                ..Default::default()
1304            }),
1305            unbounded_channel().0,
1306        );
1307        let mut event_handler = HummockEventHandler::new_inner(
1308            Role::Streaming,
1309            observer_event_rx,
1310            mock_sstable_store().await,
1311            metrics.clone(),
1312            CacheRefillConfig::from_storage_opts(&storage_opt),
1313            RecentVersions::new(pinned_version, 10, metrics.clone()),
1314            BufferTracker::from_storage_opts(&storage_opt, &metrics),
1315            spawn_upload_task,
1316            CacheRefiller::default_spawn_refill_task(),
1317        );
1318
1319        let (tx, rx) = oneshot::channel();
1320        event_handler.handle_hummock_event(HummockEvent::RegisterReadVersion {
1321            table_id,
1322            new_read_version_sender: tx,
1323            is_replicated: false,
1324            vnodes: Arc::new(Bitmap::from_range(VirtualNode::COUNT_FOR_TEST, 0..8)),
1325        });
1326        let (_read_version_1, mut guard_1) = rx.await.unwrap();
1327
1328        let (tx, rx) = oneshot::channel();
1329        event_handler.handle_hummock_event(HummockEvent::RegisterReadVersion {
1330            table_id,
1331            new_read_version_sender: tx,
1332            is_replicated: false,
1333            vnodes: Arc::new(Bitmap::from_range(VirtualNode::COUNT_FOR_TEST, 8..16)),
1334        });
1335        let (_read_version_2, mut guard_2) = rx.await.unwrap();
1336
1337        let (tx, rx) = oneshot::channel();
1338        event_handler.handle_hummock_event(HummockEvent::RegisterReadVersion {
1339            table_id: other_table_id,
1340            new_read_version_sender: tx,
1341            is_replicated: false,
1342            vnodes: Arc::new(Bitmap::from_range(VirtualNode::COUNT_FOR_TEST, 32..40)),
1343        });
1344        let (_other_read_version, mut other_guard) = rx.await.unwrap();
1345
1346        let snapshot = event_handler.refiller.table_cache_refill_monitor_snapshot();
1347        let vnodes = snapshot
1348            .streaming_table_vnode_mapping
1349            .get(&table_id)
1350            .unwrap();
1351        assert_vnodes(vnodes, 0..16);
1352        assert_vnodes(
1353            snapshot
1354                .streaming_table_vnode_mapping
1355                .get(&other_table_id)
1356                .unwrap(),
1357            32..40,
1358        );
1359
1360        guard_1.event_sender.take();
1361        event_handler.handle_hummock_event(HummockEvent::DestroyReadVersion {
1362            instance_id: guard_1.instance_id,
1363        });
1364        let snapshot = event_handler.refiller.table_cache_refill_monitor_snapshot();
1365        assert_vnodes(
1366            snapshot
1367                .streaming_table_vnode_mapping
1368                .get(&table_id)
1369                .unwrap(),
1370            8..16,
1371        );
1372
1373        guard_2.event_sender.take();
1374        event_handler.handle_hummock_event(HummockEvent::DestroyReadVersion {
1375            instance_id: guard_2.instance_id,
1376        });
1377        let snapshot = event_handler.refiller.table_cache_refill_monitor_snapshot();
1378        assert!(
1379            !snapshot
1380                .streaming_table_vnode_mapping
1381                .contains_key(&table_id)
1382        );
1383        assert_vnodes(
1384            snapshot
1385                .streaming_table_vnode_mapping
1386                .get(&other_table_id)
1387                .unwrap(),
1388            32..40,
1389        );
1390
1391        other_guard.event_sender.take();
1392        event_handler.handle_hummock_event(HummockEvent::DestroyReadVersion {
1393            instance_id: other_guard.instance_id,
1394        });
1395    }
1396
1397    #[tokio::test]
1398    async fn test_old_epoch_sync_fail() {
1399        let epoch0 = test_epoch(233);
1400
1401        let initial_version = PinnedVersion::new(
1402            HummockVersion::from_rpc_protobuf(&PbHummockVersion {
1403                id: 1.into(),
1404                state_table_info: HashMap::from_iter([(
1405                    TEST_TABLE_ID,
1406                    StateTableInfo {
1407                        committed_epoch: epoch0,
1408                        compaction_group_id: StaticCompactionGroupId::StateDefault,
1409                    },
1410                )]),
1411                ..Default::default()
1412            }),
1413            unbounded_channel().0,
1414        );
1415
1416        let (_observer_event_tx, observer_event_rx) = unbounded_channel();
1417
1418        let epoch1 = epoch0.next_epoch();
1419        let epoch2 = epoch1.next_epoch();
1420        let (tx, rx) = oneshot::channel();
1421        let rx = Arc::new(Mutex::new(Some(rx)));
1422
1423        let storage_opt = default_opts_for_test();
1424        let metrics = Arc::new(HummockStateStoreMetrics::unused());
1425
1426        let event_handler = HummockEventHandler::new_inner(
1427            Role::None,
1428            observer_event_rx,
1429            mock_sstable_store().await,
1430            metrics.clone(),
1431            CacheRefillConfig::from_storage_opts(&storage_opt),
1432            RecentVersions::new(initial_version.clone(), 10, metrics.clone()),
1433            BufferTracker::from_storage_opts(&storage_opt, &metrics),
1434            Arc::new(move |_, info| {
1435                assert_eq!(info.epochs.len(), 1);
1436                let epoch = info.epochs[0];
1437                match epoch {
1438                    epoch if epoch == epoch1 => {
1439                        let rx = rx.lock().take().unwrap();
1440                        spawn(async move {
1441                            rx.await.unwrap();
1442                            Err(HummockError::other("fail"))
1443                        })
1444                    }
1445                    epoch if epoch == epoch2 => spawn(async move {
1446                        Ok(UploadTaskOutput {
1447                            new_value_ssts: vec![],
1448                            old_value_ssts: vec![],
1449                            wait_poll_timer: None,
1450                        })
1451                    }),
1452                    _ => unreachable!(),
1453                }
1454            }),
1455            CacheRefiller::default_spawn_refill_task(),
1456        );
1457
1458        let event_tx = event_handler.event_sender();
1459
1460        let send_event = |event| event_tx.send(event).unwrap();
1461
1462        let join_handle = spawn(event_handler.start_hummock_event_handler_worker());
1463
1464        let (read_version, guard) = {
1465            let (tx, rx) = oneshot::channel();
1466            send_event(HummockEvent::RegisterReadVersion {
1467                table_id: TEST_TABLE_ID,
1468                new_read_version_sender: tx,
1469                is_replicated: false,
1470                vnodes: Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
1471            });
1472            rx.await.unwrap()
1473        };
1474
1475        send_event(HummockEvent::StartEpoch {
1476            epoch: epoch1,
1477            table_ids: HashSet::from_iter([TEST_TABLE_ID]),
1478        });
1479
1480        read_version.write().init();
1481
1482        send_event(HummockEvent::InitEpoch {
1483            instance_id: guard.instance_id,
1484            init_epoch: epoch1,
1485        });
1486
1487        let (imm1, tracker1) = gen_imm_with_unlimit(epoch1);
1488        read_version.write().add_pending_imm(imm1.clone(), tracker1);
1489
1490        send_event(HummockEvent::ImmToUploader {
1491            instance_id: guard.instance_id,
1492            imms: read_version.write().start_upload_pending_imms(),
1493        });
1494
1495        send_event(HummockEvent::StartEpoch {
1496            epoch: epoch2,
1497            table_ids: HashSet::from_iter([TEST_TABLE_ID]),
1498        });
1499
1500        send_event(HummockEvent::LocalSealEpoch {
1501            instance_id: guard.instance_id,
1502            next_epoch: epoch2,
1503            opts: SealCurrentEpochOptions::for_test(),
1504        });
1505
1506        {
1507            let (imm2, tracker2) = gen_imm_with_unlimit(epoch2);
1508            let mut read_version = read_version.write();
1509            read_version.add_pending_imm(imm2, tracker2);
1510
1511            send_event(HummockEvent::ImmToUploader {
1512                instance_id: guard.instance_id,
1513                imms: read_version.start_upload_pending_imms(),
1514            });
1515        }
1516
1517        let epoch3 = epoch2.next_epoch();
1518        send_event(HummockEvent::StartEpoch {
1519            epoch: epoch3,
1520            table_ids: HashSet::from_iter([TEST_TABLE_ID]),
1521        });
1522        send_event(HummockEvent::LocalSealEpoch {
1523            instance_id: guard.instance_id,
1524            next_epoch: epoch3,
1525            opts: SealCurrentEpochOptions::for_test(),
1526        });
1527
1528        let (tx1, mut rx1) = oneshot::channel();
1529        send_event(HummockEvent::SyncEpoch {
1530            sync_result_sender: tx1,
1531            sync_table_epochs: vec![(epoch1, HashSet::from_iter([TEST_TABLE_ID]))],
1532        });
1533        assert!(poll_fn(|cx| Poll::Ready(rx1.poll_unpin(cx).is_pending())).await);
1534        let (tx2, mut rx2) = oneshot::channel();
1535        send_event(HummockEvent::SyncEpoch {
1536            sync_result_sender: tx2,
1537            sync_table_epochs: vec![(epoch2, HashSet::from_iter([TEST_TABLE_ID]))],
1538        });
1539        assert!(poll_fn(|cx| Poll::Ready(rx2.poll_unpin(cx).is_pending())).await);
1540
1541        tx.send(()).unwrap();
1542        rx1.await.unwrap().unwrap_err();
1543        rx2.await.unwrap().unwrap_err();
1544
1545        send_event(HummockEvent::Shutdown);
1546        join_handle.await.unwrap();
1547    }
1548
1549    #[tokio::test]
1550    async fn test_clear_one_table_preserves_other_table_uploads() {
1551        let table_id1 = TableId::new(1);
1552        let table_id2 = TableId::new(2);
1553        let epoch0 = test_epoch(233);
1554
1555        let initial_version = PinnedVersion::new(
1556            HummockVersion::from_rpc_protobuf(&PbHummockVersion {
1557                id: 1.into(),
1558                state_table_info: HashMap::from_iter([
1559                    (
1560                        table_id1,
1561                        StateTableInfo {
1562                            committed_epoch: epoch0,
1563                            compaction_group_id: StaticCompactionGroupId::StateDefault,
1564                        },
1565                    ),
1566                    (
1567                        table_id2,
1568                        StateTableInfo {
1569                            committed_epoch: epoch0,
1570                            compaction_group_id: StaticCompactionGroupId::StateDefault,
1571                        },
1572                    ),
1573                ]),
1574                ..Default::default()
1575            }),
1576            unbounded_channel().0,
1577        );
1578
1579        let (_observer_event_tx, observer_event_rx) = unbounded_channel();
1580
1581        let epoch1 = epoch0.next_epoch();
1582        let epoch2 = epoch1.next_epoch();
1583        let epoch3 = epoch2.next_epoch();
1584
1585        let imm_size = gen_imm_inner(TEST_TABLE_ID, epoch1, 0).size();
1586
1587        // The buffer can hold at most 1 imm. When a new imm is added, the previous one will be spilled, and the newly added one will be retained.
1588        let buffer_tracker = BufferTracker::for_test_with_config(imm_size * 2 - 1, 1);
1589        let memory_limiter = buffer_tracker.get_memory_limiter().clone();
1590
1591        let gen_imm = |table_id, epoch, spill_offset| {
1592            let imm = gen_imm_inner(table_id, epoch, spill_offset);
1593            assert_eq!(imm.size(), imm_size);
1594            imm
1595        };
1596        let imm1_1 = gen_imm(table_id1, epoch1, 0);
1597        let imm1_2_1 = gen_imm(table_id1, epoch2, 0);
1598
1599        let storage_opt = default_opts_for_test();
1600        let metrics = Arc::new(HummockStateStoreMetrics::unused());
1601
1602        let (spawn_task, new_task_notifier) = prepare_uploader_order_test_spawn_task_fn(false);
1603
1604        let event_handler = HummockEventHandler::new_inner(
1605            Role::None,
1606            observer_event_rx,
1607            mock_sstable_store().await,
1608            metrics.clone(),
1609            CacheRefillConfig::from_storage_opts(&storage_opt),
1610            RecentVersions::new(initial_version.clone(), 10, metrics.clone()),
1611            buffer_tracker,
1612            spawn_task,
1613            CacheRefiller::default_spawn_refill_task(),
1614        );
1615
1616        let event_tx = event_handler.event_sender();
1617
1618        let send_event = |event| event_tx.send(event).unwrap();
1619        let flush_event = || async {
1620            let (tx, rx) = oneshot::channel();
1621            send_event(HummockEvent::FlushEvent(tx));
1622            rx.await.unwrap();
1623        };
1624        let start_epoch = |table_id, epoch| {
1625            send_event(HummockEvent::StartEpoch {
1626                epoch,
1627                table_ids: HashSet::from_iter([table_id]),
1628            })
1629        };
1630        let init_epoch = |instance: &LocalInstanceGuard, init_epoch| {
1631            send_event(HummockEvent::InitEpoch {
1632                instance_id: instance.instance_id,
1633                init_epoch,
1634            })
1635        };
1636        let event_tx_clone = event_tx.clone();
1637        let write_imm = {
1638            let memory_limiter = memory_limiter.clone();
1639            move |read_version: &HummockReadVersionRef,
1640                  instance: &LocalInstanceGuard,
1641                  imm: &ImmutableMemtable| {
1642                let memory_limiter = memory_limiter.clone();
1643                let event_tx = event_tx_clone.clone();
1644                let read_version = read_version.clone();
1645                let imm = imm.clone();
1646                let instance_id = instance.instance_id;
1647                async move {
1648                    let tracker = memory_limiter.require_memory(imm.size() as _).await;
1649                    let mut read_version = read_version.write();
1650                    read_version.add_pending_imm(imm.clone(), tracker);
1651
1652                    event_tx
1653                        .send(HummockEvent::ImmToUploader {
1654                            instance_id,
1655                            imms: read_version.start_upload_pending_imms(),
1656                        })
1657                        .unwrap();
1658                }
1659            }
1660        };
1661        let seal_epoch = |instance: &LocalInstanceGuard, next_epoch| {
1662            send_event(HummockEvent::LocalSealEpoch {
1663                instance_id: instance.instance_id,
1664                next_epoch,
1665                opts: SealCurrentEpochOptions::for_test(),
1666            })
1667        };
1668        let sync_epoch = |table_id, new_sync_epoch| {
1669            let (tx, rx) = oneshot::channel();
1670            send_event(HummockEvent::SyncEpoch {
1671                sync_result_sender: tx,
1672                sync_table_epochs: vec![(new_sync_epoch, HashSet::from_iter([table_id]))],
1673            });
1674            rx
1675        };
1676
1677        let join_handle = spawn(event_handler.start_hummock_event_handler_worker());
1678
1679        let (read_version1, guard1) = {
1680            let (tx, rx) = oneshot::channel();
1681            send_event(HummockEvent::RegisterReadVersion {
1682                table_id: table_id1,
1683                new_read_version_sender: tx,
1684                is_replicated: false,
1685                vnodes: Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
1686            });
1687            rx.await.unwrap()
1688        };
1689
1690        let (read_version2, guard2) = {
1691            let (tx, rx) = oneshot::channel();
1692            send_event(HummockEvent::RegisterReadVersion {
1693                table_id: table_id2,
1694                new_read_version_sender: tx,
1695                is_replicated: false,
1696                vnodes: Arc::new(Bitmap::ones(VirtualNode::COUNT_FOR_TEST)),
1697            });
1698            rx.await.unwrap()
1699        };
1700
1701        // prepare data of table1
1702        let (task1_1_finish_tx, task1_1_rx) = {
1703            start_epoch(table_id1, epoch1);
1704
1705            read_version1.write().init();
1706            init_epoch(&guard1, epoch1);
1707
1708            write_imm(&read_version1, &guard1, &imm1_1).await;
1709
1710            start_epoch(table_id1, epoch2);
1711
1712            seal_epoch(&guard1, epoch2);
1713
1714            let (wait_task_start, task_finish_tx) = new_task_notifier(HashMap::from_iter([(
1715                guard1.instance_id,
1716                vec![imm1_1.batch_id()],
1717            )]));
1718
1719            let mut rx = sync_epoch(table_id1, epoch1);
1720            wait_task_start.await;
1721            assert!(poll_fn(|cx| Poll::Ready(rx.poll_unpin(cx).is_pending())).await);
1722
1723            write_imm(&read_version1, &guard1, &imm1_2_1).await;
1724            flush_event().await;
1725
1726            (task_finish_tx, rx)
1727        };
1728        // by now, the state in uploader of table_id1
1729        // unsync:  epoch2 -> [imm1_2]
1730        // syncing: epoch1 -> [imm1_1]
1731
1732        let (task1_2_finish_tx, _finish_txs) = {
1733            let mut finish_txs = vec![];
1734            let imm2_1_1 = gen_imm(table_id2, epoch1, 0);
1735            start_epoch(table_id2, epoch1);
1736            read_version2.write().init();
1737            init_epoch(&guard2, epoch1);
1738            let (wait_task_start, task1_2_finish_tx) = new_task_notifier(HashMap::from_iter([(
1739                guard1.instance_id,
1740                vec![imm1_2_1.batch_id()],
1741            )]));
1742            write_imm(&read_version2, &guard2, &imm2_1_1).await;
1743            wait_task_start.await;
1744
1745            let imm2_1_2 = gen_imm(table_id2, epoch1, 1);
1746            let (wait_task_start, finish_tx) = new_task_notifier(HashMap::from_iter([(
1747                guard2.instance_id,
1748                vec![imm2_1_2.batch_id(), imm2_1_1.batch_id()],
1749            )]));
1750            finish_txs.push(finish_tx);
1751            write_imm(&read_version2, &guard2, &imm2_1_2).await;
1752            wait_task_start.await;
1753
1754            let imm2_1_3 = gen_imm(table_id2, epoch1, 2);
1755            write_imm(&read_version2, &guard2, &imm2_1_3).await;
1756            start_epoch(table_id2, epoch2);
1757            seal_epoch(&guard2, epoch2);
1758            let (wait_task_start, finish_tx) = new_task_notifier(HashMap::from_iter([(
1759                guard2.instance_id,
1760                vec![imm2_1_3.batch_id()],
1761            )]));
1762            finish_txs.push(finish_tx);
1763            let _sync_rx = sync_epoch(table_id2, epoch1);
1764            wait_task_start.await;
1765
1766            let imm2_2_1 = gen_imm(table_id2, epoch2, 0);
1767            write_imm(&read_version2, &guard2, &imm2_2_1).await;
1768            flush_event().await;
1769            let imm2_2_2 = gen_imm(table_id2, epoch2, 1);
1770            write_imm(&read_version2, &guard2, &imm2_2_2).await;
1771            let (wait_task_start, finish_tx) = new_task_notifier(HashMap::from_iter([(
1772                guard2.instance_id,
1773                vec![imm2_2_2.batch_id(), imm2_2_1.batch_id()],
1774            )]));
1775            finish_txs.push(finish_tx);
1776            wait_task_start.await;
1777
1778            let imm2_2_3 = gen_imm(table_id2, epoch2, 2);
1779            write_imm(&read_version2, &guard2, &imm2_2_3).await;
1780
1781            // by now, the state in uploader of table_id2
1782            // syncing: epoch1 -> spill: [imm2_1_2, imm2_1_1], sync: [imm2_1_3]
1783            // unsync: epoch2 -> spilling: [imm2_2_2, imm2_2_1], imm: [imm2_2_3]
1784            // the state in uploader of table_id1
1785            // unsync:  epoch2 -> spilling [imm1_2]
1786            // syncing: epoch1 -> [imm1_1]
1787
1788            drop(guard2);
1789            let (clear_tx, clear_rx) = oneshot::channel();
1790            send_event(HummockEvent::Clear(
1791                clear_tx,
1792                Some(HashSet::from_iter([table_id2])),
1793            ));
1794            clear_rx.await.unwrap();
1795            (task1_2_finish_tx, finish_txs)
1796        };
1797
1798        let imm1_2_2 = gen_imm(table_id1, epoch2, 1);
1799        write_imm(&read_version1, &guard1, &imm1_2_2).await;
1800        start_epoch(table_id1, epoch3);
1801        seal_epoch(&guard1, epoch3);
1802
1803        let (tx2, mut sync_rx2) = oneshot::channel();
1804        let (wait_task_start, task1_2_2_finish_tx) = new_task_notifier(HashMap::from_iter([(
1805            guard1.instance_id,
1806            vec![imm1_2_2.batch_id()],
1807        )]));
1808        send_event(HummockEvent::SyncEpoch {
1809            sync_result_sender: tx2,
1810            sync_table_epochs: vec![(epoch2, HashSet::from_iter([table_id1]))],
1811        });
1812        wait_task_start.await;
1813        assert!(poll_fn(|cx| Poll::Ready(sync_rx2.poll_unpin(cx).is_pending())).await);
1814
1815        task1_1_finish_tx.send(()).unwrap();
1816        let sync_data1 = task1_1_rx.await.unwrap().unwrap();
1817        assert!(!sync_data1.uploaded_ssts.is_empty());
1818        assert!(
1819            sync_data1
1820                .uploaded_ssts
1821                .iter()
1822                .all(|sst| sst.epochs().as_slice() == [epoch1])
1823        );
1824        task1_2_finish_tx.send(()).unwrap();
1825        assert!(poll_fn(|cx| Poll::Ready(sync_rx2.poll_unpin(cx).is_pending())).await);
1826        task1_2_2_finish_tx.send(()).unwrap();
1827        let sync_data2 = sync_rx2.await.unwrap().unwrap();
1828        assert!(!sync_data2.uploaded_ssts.is_empty());
1829        assert!(
1830            sync_data2
1831                .uploaded_ssts
1832                .iter()
1833                .all(|sst| sst.epochs().as_slice() == [epoch2])
1834        );
1835
1836        send_event(HummockEvent::Shutdown);
1837        join_handle.await.unwrap();
1838    }
1839}