Skip to main content

risingwave_storage/hummock/store/
hummock_storage.rs

1// Copyright 2023 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::collections::HashSet;
16use std::future::Future;
17use std::ops::{Bound, Deref};
18use std::sync::Arc;
19
20use arc_swap::ArcSwap;
21use bytes::Bytes;
22use itertools::Itertools;
23use risingwave_common::array::VectorRef;
24use risingwave_common::catalog::{TableId, TableOption};
25use risingwave_common::config::Role;
26use risingwave_common::dispatch_distance_measurement;
27use risingwave_common::util::epoch::is_max_epoch;
28use risingwave_common_service::{NotificationClient, ObserverManager};
29use risingwave_hummock_sdk::change_log::TableChangeLogs;
30use risingwave_hummock_sdk::key::{
31    TableKey, TableKeyRange, is_empty_key_range, vnode, vnode_range,
32};
33use risingwave_hummock_sdk::sstable_info::SstableInfo;
34use risingwave_hummock_sdk::table_watermark::TableWatermarksIndex;
35use risingwave_hummock_sdk::version::HummockVersion;
36use risingwave_hummock_sdk::{HummockRawObjectId, HummockReadEpoch, SyncResult};
37use risingwave_rpc_client::HummockMetaClient;
38use risingwave_rpc_client::error::RpcError;
39use thiserror_ext::AsReport;
40use tokio::sync::mpsc::{UnboundedSender, unbounded_channel};
41use tokio::sync::oneshot;
42
43use super::local_hummock_storage::LocalHummockStorage;
44use super::version::{CommittedVersion, HummockVersionReader, read_filter_for_version};
45use crate::compaction_catalog_manager::CompactionCatalogManagerRef;
46#[cfg(any(test, feature = "test"))]
47use crate::compaction_catalog_manager::{CompactionCatalogManager, FakeRemoteTableAccessor};
48use crate::error::StorageResult;
49use crate::hummock::backup_reader::{BackupReader, BackupReaderRef};
50use crate::hummock::compactor::{
51    CompactionAwaitTreeRegRef, CompactorContext, new_compaction_await_tree_reg_ref,
52};
53use crate::hummock::event_handler::hummock_event_handler::{BufferTracker, HummockEventSender};
54use crate::hummock::event_handler::refiller::TableCacheRefillMonitorSnapshot;
55use crate::hummock::event_handler::{
56    HummockEvent, HummockEventHandler, HummockObserverEvent, HummockVersionUpdate,
57    ReadOnlyReadVersionMapping,
58};
59use crate::hummock::iterator::change_log::ChangeLogIterator;
60use crate::hummock::local_version::pinned_version::{PinnedVersion, start_pinned_version_worker};
61use crate::hummock::local_version::recent_versions::RecentVersions;
62use crate::hummock::observer_manager::HummockObserverNode;
63use crate::hummock::store::vector_writer::HummockVectorWriter;
64use crate::hummock::table_change_log_manager::TableChangeLogManager;
65use crate::hummock::time_travel_version_cache::SimpleTimeTravelVersionCache;
66use crate::hummock::utils::{wait_for_epoch, wait_for_update};
67use crate::hummock::write_limiter::{WriteLimiter, WriteLimiterRef};
68use crate::hummock::{
69    HummockEpoch, HummockError, HummockResult, HummockStorageIterator, HummockStorageRevIterator,
70    MemoryLimiter, ObjectIdManager, ObjectIdManagerRef, SstableStoreRef,
71};
72use crate::mem_table::ImmutableMemtable;
73use crate::monitor::{CompactorMetrics, HummockStateStoreMetrics};
74use crate::opts::StorageOpts;
75use crate::store::*;
76
77struct HummockStorageShutdownGuard {
78    shutdown_sender: HummockEventSender,
79}
80
81impl Drop for HummockStorageShutdownGuard {
82    fn drop(&mut self) {
83        let _ = self
84            .shutdown_sender
85            .send(HummockEvent::Shutdown)
86            .inspect_err(|e| tracing::debug!(event = ?e.0, "unable to send shutdown"));
87    }
88}
89
90/// `HummockStorage` is the entry point of the Hummock state store backend.
91/// It implements the `StateStore` and `StateStoreRead` traits but without any write method
92/// since all writes should be done via `LocalHummockStorage` to ensure the single writer property
93/// of hummock. `LocalHummockStorage` instance can be created via `new_local` call.
94/// Hummock is the state store backend.
95#[derive(Clone)]
96pub struct HummockStorage {
97    hummock_event_sender: HummockEventSender,
98    // Only used in tests to inject observer events such as version updates.
99    _observer_event_sender: UnboundedSender<HummockObserverEvent>,
100
101    context: CompactorContext,
102
103    compaction_catalog_manager_ref: CompactionCatalogManagerRef,
104
105    object_id_manager: ObjectIdManagerRef,
106
107    buffer_tracker: BufferTracker,
108
109    version_update_notifier_tx: Arc<tokio::sync::watch::Sender<PinnedVersion>>,
110
111    recent_versions: Arc<ArcSwap<RecentVersions>>,
112
113    hummock_version_reader: HummockVersionReader,
114
115    _shutdown_guard: Arc<HummockStorageShutdownGuard>,
116
117    read_version_mapping: ReadOnlyReadVersionMapping,
118
119    backup_reader: BackupReaderRef,
120
121    write_limiter: WriteLimiterRef,
122
123    compact_await_tree_reg: Option<CompactionAwaitTreeRegRef>,
124
125    hummock_meta_client: Arc<dyn HummockMetaClient>,
126
127    simple_time_travel_version_cache: Arc<SimpleTimeTravelVersionCache>,
128
129    table_change_log_manager: Arc<TableChangeLogManager>,
130}
131
132pub type ReadVersionTuple = (Vec<ImmutableMemtable>, Vec<SstableInfo>, CommittedVersion);
133
134pub fn get_committed_read_version_tuple(
135    version: PinnedVersion,
136    table_id: TableId,
137    mut key_range: TableKeyRange,
138    epoch: HummockEpoch,
139) -> (TableKeyRange, ReadVersionTuple) {
140    if let Some(table_watermarks) = version.table_watermarks.get(&table_id) {
141        TableWatermarksIndex::new_committed(
142            table_watermarks.clone(),
143            version
144                .state_table_info
145                .info()
146                .get(&table_id)
147                .expect("should exist when having table watermark")
148                .committed_epoch,
149            table_watermarks.watermark_type,
150        )
151        .rewrite_range_with_table_watermark(epoch, &mut key_range)
152    }
153    (key_range, (vec![], vec![], version))
154}
155
156impl HummockStorage {
157    /// Creates a [`HummockStorage`].
158    #[allow(clippy::too_many_arguments)]
159    pub async fn new(
160        role: Role,
161        options: Arc<StorageOpts>,
162        sstable_store: SstableStoreRef,
163        hummock_meta_client: Arc<dyn HummockMetaClient>,
164        notification_client: impl NotificationClient,
165        compaction_catalog_manager_ref: CompactionCatalogManagerRef,
166        state_store_metrics: Arc<HummockStateStoreMetrics>,
167        compactor_metrics: Arc<CompactorMetrics>,
168        await_tree_config: Option<await_tree::Config>,
169    ) -> HummockResult<Self> {
170        let object_id_manager = Arc::new(ObjectIdManager::new(
171            hummock_meta_client.clone(),
172            options.sstable_id_remote_fetch_number,
173        ));
174        let backup_reader = BackupReader::new(
175            &options.backup_storage_url,
176            &options.backup_storage_directory,
177            &options.object_store_config,
178        )
179        .await
180        .map_err(HummockError::read_backup_error)?;
181        let write_limiter = Arc::new(WriteLimiter::default());
182        let (observer_event_tx, mut observer_event_rx) = unbounded_channel();
183        let observer_manager = ObserverManager::new(
184            notification_client,
185            HummockObserverNode::new(
186                role,
187                compaction_catalog_manager_ref.clone(),
188                backup_reader.clone(),
189                observer_event_tx.clone(),
190                write_limiter.clone(),
191            ),
192        )
193        .await;
194        observer_manager.start().await;
195
196        let hummock_version = match observer_event_rx.recv().await {
197            Some(HummockObserverEvent::VersionUpdate(HummockVersionUpdate::PinnedVersion(
198                version,
199            ))) => *version,
200            _ => unreachable!(
201                "the hummock observer manager is the first one to take the event tx. Should be full hummock version"
202            ),
203        };
204        let (pin_version_tx, pin_version_rx) = unbounded_channel();
205        let pinned_version = PinnedVersion::new(hummock_version, pin_version_tx);
206        tokio::spawn(start_pinned_version_worker(
207            pin_version_rx,
208            hummock_meta_client.clone(),
209            options.max_version_pinning_duration_sec,
210        ));
211
212        let await_tree_reg = await_tree_config.map(new_compaction_await_tree_reg_ref);
213
214        let compactor_context = CompactorContext::new_local_compact_context(
215            options.clone(),
216            sstable_store.clone(),
217            compactor_metrics.clone(),
218            await_tree_reg.clone(),
219        );
220
221        let hummock_event_handler = HummockEventHandler::new(
222            role,
223            observer_event_rx,
224            pinned_version,
225            compactor_context.clone(),
226            compaction_catalog_manager_ref.clone(),
227            object_id_manager.clone(),
228            state_store_metrics.clone(),
229        );
230
231        let event_tx = hummock_event_handler.event_sender();
232        let table_change_log_manager = Arc::new(TableChangeLogManager::new(
233            options.table_change_log_cache_capacity,
234            hummock_meta_client.clone(),
235            state_store_metrics.clone(),
236        ));
237        let instance = Self {
238            context: compactor_context,
239            compaction_catalog_manager_ref: compaction_catalog_manager_ref.clone(),
240            object_id_manager,
241            buffer_tracker: hummock_event_handler.buffer_tracker().clone(),
242            version_update_notifier_tx: hummock_event_handler.version_update_notifier_tx(),
243            hummock_event_sender: event_tx.clone(),
244            _observer_event_sender: observer_event_tx,
245            recent_versions: hummock_event_handler.recent_versions(),
246            hummock_version_reader: HummockVersionReader::new(
247                sstable_store,
248                state_store_metrics.clone(),
249                options.max_preload_io_retry_times,
250            ),
251            _shutdown_guard: Arc::new(HummockStorageShutdownGuard {
252                shutdown_sender: event_tx,
253            }),
254            read_version_mapping: hummock_event_handler.read_version_mapping(),
255            backup_reader,
256            write_limiter,
257            compact_await_tree_reg: await_tree_reg,
258            hummock_meta_client,
259            simple_time_travel_version_cache: Arc::new(SimpleTimeTravelVersionCache::new(
260                options.time_travel_version_cache_capacity,
261            )),
262            table_change_log_manager,
263        };
264
265        tokio::spawn(hummock_event_handler.start_hummock_event_handler_worker());
266
267        Ok(instance)
268    }
269}
270
271impl HummockStorageReadSnapshot {
272    /// Gets the value of a specified `key` in the table specified in `read_options`.
273    /// The result is based on a snapshot corresponding to the given `epoch`.
274    /// if `key` has consistent hash virtual node value, then such value is stored in `value_meta`
275    ///
276    /// If `Ok(Some())` is returned, the key is found. If `Ok(None)` is returned,
277    /// the key is not found. If `Err()` is returned, the searching for the key
278    /// failed due to other non-EOF errors.
279    async fn get_inner<'a, O>(
280        &'a self,
281        key: TableKey<Bytes>,
282        read_options: ReadOptions,
283        on_key_value_fn: impl KeyValueFn<'a, O>,
284    ) -> StorageResult<Option<O>> {
285        let key_range = (Bound::Included(key.clone()), Bound::Included(key.clone()));
286
287        let (key_range, read_version_tuple) =
288            self.build_read_version_tuple(self.epoch, key_range).await?;
289
290        if is_empty_key_range(&key_range) {
291            return Ok(None);
292        }
293
294        self.hummock_version_reader
295            .get(
296                key,
297                self.epoch.get_epoch(),
298                self.table_id,
299                self.table_option,
300                read_options,
301                read_version_tuple,
302                on_key_value_fn,
303            )
304            .await
305    }
306
307    async fn iter_inner(
308        &self,
309        key_range: TableKeyRange,
310        read_options: ReadOptions,
311    ) -> StorageResult<HummockStorageIterator> {
312        let (key_range, read_version_tuple) =
313            self.build_read_version_tuple(self.epoch, key_range).await?;
314
315        self.hummock_version_reader
316            .iter(
317                key_range,
318                self.epoch.get_epoch(),
319                self.table_id,
320                self.table_option,
321                read_options,
322                read_version_tuple,
323            )
324            .await
325    }
326
327    async fn rev_iter_inner(
328        &self,
329        key_range: TableKeyRange,
330        read_options: ReadOptions,
331    ) -> StorageResult<HummockStorageRevIterator> {
332        let (key_range, read_version_tuple) =
333            self.build_read_version_tuple(self.epoch, key_range).await?;
334
335        self.hummock_version_reader
336            .rev_iter(
337                key_range,
338                self.epoch.get_epoch(),
339                self.table_id,
340                self.table_option,
341                read_options,
342                read_version_tuple,
343                None,
344            )
345            .await
346    }
347
348    async fn get_time_travel_version(
349        &self,
350        epoch: u64,
351        table_id: TableId,
352    ) -> StorageResult<PinnedVersion> {
353        let meta_client = self.hummock_meta_client.clone();
354        let fetch = async move {
355            let pb_version = meta_client
356                .get_version_by_epoch(epoch, table_id)
357                .await
358                .inspect_err(|e| tracing::error!("{}", e.to_report_string()))
359                .map_err(|e| match &e {
360                    RpcError::GrpcStatus(status)
361                        if status.inner().code() == tonic::Code::OutOfRange =>
362                    {
363                        HummockError::time_travel_version_expired(table_id, epoch)
364                    }
365                    _ => HummockError::meta_error(e.to_report_string()),
366                })?;
367            let version = HummockVersion::from_rpc_protobuf(&pb_version);
368            let (tx, _rx) = unbounded_channel();
369            Ok(PinnedVersion::new(version, tx))
370        };
371        let version = self
372            .simple_time_travel_version_cache
373            .get_or_insert(table_id, epoch, fetch)
374            .await?;
375        Ok(version)
376    }
377
378    async fn build_read_version_tuple(
379        &self,
380        epoch: HummockReadEpoch,
381        key_range: TableKeyRange,
382    ) -> StorageResult<(TableKeyRange, ReadVersionTuple)> {
383        match epoch {
384            HummockReadEpoch::Backup(epoch) => {
385                self.build_read_version_tuple_from_backup(epoch, self.table_id, key_range)
386                    .await
387            }
388            HummockReadEpoch::Committed(epoch) => {
389                let tuple = self
390                    .build_read_version_tuple_from_committed(epoch, self.table_id, key_range)
391                    .await?;
392                let (_, (_, _, version)) = &tuple;
393                let Some(committed_epoch) = version.table_committed_epoch(self.table_id) else {
394                    return Err(HummockError::other(format!(
395                        "table {} not found in version",
396                        self.table_id
397                    ))
398                    .into());
399                };
400                if committed_epoch != epoch {
401                    return Err(HummockError::committed_epoch_mismatch(
402                        self.table_id,
403                        committed_epoch,
404                        epoch,
405                    )
406                    .into());
407                }
408                Ok(tuple)
409            }
410            HummockReadEpoch::BatchQueryCommitted(epoch, _)
411            | HummockReadEpoch::TimeTravel(epoch) => {
412                self.build_read_version_tuple_from_committed(epoch, self.table_id, key_range)
413                    .await
414            }
415            HummockReadEpoch::NoWait(epoch) => {
416                self.build_read_version_tuple_from_all(epoch, self.table_id, key_range)
417                    .await
418            }
419        }
420    }
421
422    async fn build_read_version_tuple_from_backup(
423        &self,
424        epoch: u64,
425        table_id: TableId,
426        key_range: TableKeyRange,
427    ) -> StorageResult<(TableKeyRange, ReadVersionTuple)> {
428        match self
429            .backup_reader
430            .try_get_hummock_version(table_id, epoch)
431            .await
432        {
433            Ok(Some(backup_version)) => Ok(get_committed_read_version_tuple(
434                backup_version,
435                table_id,
436                key_range,
437                epoch,
438            )),
439            Ok(None) => Err(HummockError::read_backup_error(format!(
440                "backup include epoch {} not found",
441                epoch
442            ))
443            .into()),
444            Err(e) => Err(e),
445        }
446    }
447
448    async fn get_epoch_hummock_version(
449        &self,
450        epoch: u64,
451        table_id: TableId,
452    ) -> StorageResult<PinnedVersion> {
453        match self
454            .recent_versions
455            .load()
456            .get_safe_version(table_id, epoch)
457        {
458            Some(version) => Ok(version),
459            None => self.get_time_travel_version(epoch, table_id).await,
460        }
461    }
462
463    async fn build_read_version_tuple_from_committed(
464        &self,
465        epoch: u64,
466        table_id: TableId,
467        key_range: TableKeyRange,
468    ) -> StorageResult<(TableKeyRange, ReadVersionTuple)> {
469        let version = self.get_epoch_hummock_version(epoch, table_id).await?;
470        Ok(get_committed_read_version_tuple(
471            version, table_id, key_range, epoch,
472        ))
473    }
474
475    async fn build_read_version_tuple_from_all(
476        &self,
477        epoch: u64,
478        table_id: TableId,
479        key_range: TableKeyRange,
480    ) -> StorageResult<(TableKeyRange, ReadVersionTuple)> {
481        let pinned_version = self.recent_versions.load().latest_version().clone();
482        let info = pinned_version.state_table_info.info().get(&table_id);
483
484        // check epoch if lower mce
485        let ret = if let Some(info) = info
486            && epoch <= info.committed_epoch
487        {
488            let pinned_version = if epoch < info.committed_epoch {
489                pinned_version
490            } else {
491                self.get_epoch_hummock_version(epoch, table_id).await?
492            };
493            // read committed_version directly without build snapshot
494            get_committed_read_version_tuple(pinned_version, table_id, key_range, epoch)
495        } else {
496            let vnode = vnode(&key_range);
497            let mut matched_replicated_read_version_cnt = 0;
498            let read_version_vec = {
499                let read_guard = self.read_version_mapping.read();
500                read_guard
501                    .get(&table_id)
502                    .map(|v| {
503                        v.values()
504                            .filter(|v| {
505                                let read_version = v.read();
506                                if read_version.is_initialized() && read_version.contains(vnode) {
507                                    if read_version.is_replicated() {
508                                        matched_replicated_read_version_cnt += 1;
509                                        false
510                                    } else {
511                                        // Only non-replicated read version with matched vnode is considered
512                                        true
513                                    }
514                                } else {
515                                    false
516                                }
517                            })
518                            .cloned()
519                            .collect_vec()
520                    })
521                    .unwrap_or_default()
522            };
523
524            // When the system has just started and no state has been created, the memory state
525            // may be empty
526            if read_version_vec.is_empty() {
527                let table_committed_epoch = info.map(|info| info.committed_epoch);
528                if matched_replicated_read_version_cnt > 0 {
529                    tracing::warn!(
530                        "Read(table_id={} vnode={} epoch={}) is not allowed on replicated read version ({} found). Fall back to committed version (epoch={:?})",
531                        table_id,
532                        vnode.to_index(),
533                        epoch,
534                        matched_replicated_read_version_cnt,
535                        table_committed_epoch,
536                    );
537                } else {
538                    tracing::debug!(
539                        "No read version found for read(table_id={} vnode={} epoch={}). Fall back to committed version (epoch={:?})",
540                        table_id,
541                        vnode.to_index(),
542                        epoch,
543                        table_committed_epoch
544                    );
545                }
546                get_committed_read_version_tuple(pinned_version, table_id, key_range, epoch)
547            } else {
548                if read_version_vec.len() != 1 {
549                    let read_version_vnodes = read_version_vec
550                        .into_iter()
551                        .map(|v| {
552                            let v = v.read();
553                            v.vnodes().iter_ones().collect_vec()
554                        })
555                        .collect_vec();
556                    return Err(HummockError::other(format!("There are {} read version associated with vnode {}. read_version_vnodes={:?}", read_version_vnodes.len(), vnode.to_index(), read_version_vnodes)).into());
557                }
558                read_filter_for_version(
559                    epoch,
560                    table_id,
561                    key_range,
562                    read_version_vec.first().unwrap(),
563                )?
564            }
565        };
566
567        Ok(ret)
568    }
569}
570
571impl HummockStorage {
572    async fn new_local_inner(&self, option: NewLocalOptions) -> LocalHummockStorage {
573        let (tx, rx) = tokio::sync::oneshot::channel();
574        self.hummock_event_sender
575            .send(HummockEvent::RegisterReadVersion {
576                table_id: option.table_id,
577                new_read_version_sender: tx,
578                is_replicated: option.is_replicated,
579                vnodes: option.vnodes.clone(),
580            })
581            .unwrap();
582
583        let (basic_read_version, instance_guard) = rx.await.unwrap();
584        let version_update_notifier_tx = self.version_update_notifier_tx.clone();
585        LocalHummockStorage::new(
586            instance_guard,
587            basic_read_version,
588            self.hummock_version_reader.clone(),
589            self.hummock_event_sender.clone(),
590            self.buffer_tracker.get_memory_limiter().clone(),
591            self.write_limiter.clone(),
592            option,
593            version_update_notifier_tx,
594            self.context.storage_opts.mem_table_spill_threshold,
595        )
596    }
597
598    pub async fn clear_shared_buffer(&self) {
599        let (tx, rx) = oneshot::channel();
600        self.hummock_event_sender
601            .send(HummockEvent::Clear(tx, None))
602            .expect("should send success");
603        rx.await.expect("should wait success");
604    }
605
606    pub async fn clear_tables(&self, table_ids: HashSet<TableId>) {
607        if !table_ids.is_empty() {
608            let (tx, rx) = oneshot::channel();
609            self.hummock_event_sender
610                .send(HummockEvent::Clear(tx, Some(table_ids)))
611                .expect("should send success");
612            rx.await.expect("should wait success");
613        }
614    }
615
616    /// Declare the start of an epoch. This information is provided for spill so that the spill task won't
617    /// include data of two or more syncs.
618    // TODO: remove this method when we support spill task that can include data of more two or more syncs
619    pub fn start_epoch(&self, epoch: HummockEpoch, table_ids: HashSet<TableId>) {
620        let _ = self
621            .hummock_event_sender
622            .send(HummockEvent::StartEpoch { epoch, table_ids });
623    }
624
625    pub fn sstable_store(&self) -> SstableStoreRef {
626        self.context.sstable_store.clone()
627    }
628
629    pub async fn table_cache_refill_monitor_snapshot(
630        &self,
631    ) -> HummockResult<TableCacheRefillMonitorSnapshot> {
632        let (tx, rx) = oneshot::channel();
633        self.hummock_event_sender
634            .send(HummockEvent::GetTableCacheRefillMonitorSnapshot { result_tx: tx })
635            .map_err(|_| HummockError::other("failed to send table cache refill monitor query"))?;
636        rx.await.map_err(|_| {
637            HummockError::other("failed to receive table cache refill monitor snapshot")
638        })
639    }
640
641    pub fn object_id_manager(&self) -> &ObjectIdManagerRef {
642        &self.object_id_manager
643    }
644
645    pub fn compaction_catalog_manager_ref(&self) -> CompactionCatalogManagerRef {
646        self.compaction_catalog_manager_ref.clone()
647    }
648
649    pub fn get_memory_limiter(&self) -> Arc<MemoryLimiter> {
650        self.buffer_tracker.get_memory_limiter().clone()
651    }
652
653    pub fn get_pinned_version(&self) -> PinnedVersion {
654        self.recent_versions.load().latest_version().clone()
655    }
656
657    pub fn backup_reader(&self) -> BackupReaderRef {
658        self.backup_reader.clone()
659    }
660
661    pub fn compaction_await_tree_reg(&self) -> Option<&await_tree::Registry> {
662        self.compact_await_tree_reg.as_ref()
663    }
664
665    pub async fn min_uncommitted_object_id(&self) -> Option<HummockRawObjectId> {
666        let (tx, rx) = oneshot::channel();
667        self.hummock_event_sender
668            .send(HummockEvent::GetMinUncommittedObjectId { result_tx: tx })
669            .expect("should send success");
670        rx.await.expect("should await success")
671    }
672
673    pub async fn sync(
674        &self,
675        sync_table_epochs: Vec<(HummockEpoch, HashSet<TableId>)>,
676    ) -> StorageResult<SyncResult> {
677        let (tx, rx) = oneshot::channel();
678        let _ = self.hummock_event_sender.send(HummockEvent::SyncEpoch {
679            sync_result_sender: tx,
680            sync_table_epochs,
681        });
682        let synced_data = rx
683            .await
684            .map_err(|_| HummockError::other("failed to receive sync result"))??;
685        Ok(synced_data.into_sync_result())
686    }
687}
688
689#[derive(Clone)]
690pub struct HummockStorageReadSnapshot {
691    epoch: HummockReadEpoch,
692    table_id: TableId,
693    table_option: TableOption,
694    recent_versions: Arc<ArcSwap<RecentVersions>>,
695    hummock_version_reader: HummockVersionReader,
696    read_version_mapping: ReadOnlyReadVersionMapping,
697    backup_reader: BackupReaderRef,
698    hummock_meta_client: Arc<dyn HummockMetaClient>,
699    simple_time_travel_version_cache: Arc<SimpleTimeTravelVersionCache>,
700}
701
702impl StateStoreGet for HummockStorageReadSnapshot {
703    fn on_key_value<'a, O: Send + 'a>(
704        &'a self,
705        key: TableKey<Bytes>,
706        read_options: ReadOptions,
707        on_key_value_fn: impl KeyValueFn<'a, O>,
708    ) -> impl StorageFuture<'a, Option<O>> {
709        self.get_inner(key, read_options, on_key_value_fn)
710    }
711}
712
713impl StateStoreRead for HummockStorageReadSnapshot {
714    type Iter = HummockStorageIterator;
715    type RevIter = HummockStorageRevIterator;
716
717    fn iter(
718        &self,
719        key_range: TableKeyRange,
720        read_options: ReadOptions,
721    ) -> impl Future<Output = StorageResult<Self::Iter>> + '_ {
722        let (l_vnode_inclusive, r_vnode_exclusive) = vnode_range(&key_range);
723        assert_eq!(
724            r_vnode_exclusive - l_vnode_inclusive,
725            1,
726            "read range {:?} for table {} iter contains more than one vnode",
727            key_range,
728            self.table_id
729        );
730        self.iter_inner(key_range, read_options)
731    }
732
733    fn rev_iter(
734        &self,
735        key_range: TableKeyRange,
736        read_options: ReadOptions,
737    ) -> impl Future<Output = StorageResult<Self::RevIter>> + '_ {
738        let (l_vnode_inclusive, r_vnode_exclusive) = vnode_range(&key_range);
739        assert_eq!(
740            r_vnode_exclusive - l_vnode_inclusive,
741            1,
742            "read range {:?} for table {} iter contains more than one vnode",
743            key_range,
744            self.table_id
745        );
746        self.rev_iter_inner(key_range, read_options)
747    }
748}
749
750impl StateStoreReadVector for HummockStorageReadSnapshot {
751    async fn nearest<'a, O: Send + 'a>(
752        &'a self,
753        vec: VectorRef<'a>,
754        options: VectorNearestOptions,
755        on_nearest_item_fn: impl OnNearestItemFn<'a, O>,
756    ) -> StorageResult<Vec<O>> {
757        let version = match self.epoch {
758            HummockReadEpoch::Committed(epoch)
759            | HummockReadEpoch::BatchQueryCommitted(epoch, _)
760            | HummockReadEpoch::TimeTravel(epoch) => {
761                self.get_epoch_hummock_version(epoch, self.table_id).await?
762            }
763            HummockReadEpoch::Backup(epoch) => self
764                .backup_reader
765                .try_get_hummock_version(self.table_id, epoch)
766                .await?
767                .ok_or_else(|| {
768                    HummockError::read_backup_error(format!(
769                        "backup include epoch {} not found",
770                        epoch
771                    ))
772                })?,
773            HummockReadEpoch::NoWait(_) => {
774                return Err(
775                    HummockError::other("nearest query does not support NoWait epoch").into(),
776                );
777            }
778        };
779        dispatch_distance_measurement!(options.measure, MeasurementType, {
780            Ok(self
781                .hummock_version_reader
782                .nearest::<MeasurementType, O>(
783                    version,
784                    self.table_id,
785                    vec,
786                    options,
787                    on_nearest_item_fn,
788                )
789                .await?)
790        })
791    }
792}
793
794impl StateStoreReadLog for HummockStorage {
795    type ChangeLogIter = ChangeLogIterator;
796
797    async fn next_epoch(&self, epoch: u64, options: NextEpochOptions) -> StorageResult<u64> {
798        fn next_epoch(
799            table_change_log: &TableChangeLogs,
800            epoch: u64,
801            table_id: TableId,
802        ) -> HummockResult<Option<u64>> {
803            let table_change_log = table_change_log.get(&table_id).ok_or_else(|| {
804                HummockError::next_epoch(format!("table {} has been dropped", table_id))
805            })?;
806            table_change_log
807                .next_epoch(epoch)
808                .map_err(|_| HummockError::change_log_retention_miss(table_id, epoch))
809        }
810        {
811            // fast path
812            if let Some(max_epoch) = self
813                .recent_versions
814                .load()
815                .latest_version()
816                .deref()
817                .state_table_info
818                .info()
819                .get(&options.table_id)
820                .map(|i| i.committed_epoch)
821                && max_epoch > epoch
822            {
823                // The next epoch exists either in the same `EpochNewChangeLog` or the next `EpochNewChangeLog`, so we fetch 2 `EpochNewChangeLogCommon`.
824                let table_change_log = self
825                    .table_change_log_manager
826                    .fetch_table_change_logs(options.table_id, (epoch, max_epoch), true, Some(2))
827                    .await?;
828                if let Some(next_epoch) = next_epoch(&table_change_log, epoch, options.table_id)? {
829                    return Ok(next_epoch);
830                }
831            }
832        }
833        let mut max_epoch = None;
834        wait_for_update(
835            &self.version_update_notifier_tx,
836            |version| {
837                let Some(mce) = version
838                    .state_table_info
839                    .info()
840                    .get(&options.table_id)
841                    .map(|i| i.committed_epoch)
842                else {
843                    return Ok(false);
844                };
845                max_epoch = Some(mce);
846                Ok(mce > epoch)
847            },
848            || format!("wait next_epoch: epoch: {} {}", epoch, options.table_id),
849        )
850        .await?;
851        // The next epoch exists either in the same `EpochNewChangeLog` or the next `EpochNewChangeLog`, so we fetch 2 `EpochNewChangeLogCommon`.
852        let table_change_log = self
853            .table_change_log_manager
854            .fetch_table_change_logs(
855                options.table_id,
856                (epoch, max_epoch.unwrap_or(u64::MAX)),
857                true,
858                Some(2),
859            )
860            .await?;
861        let next_epoch_ret = next_epoch(&table_change_log, epoch, options.table_id)?;
862        next_epoch_ret.ok_or_else(|| {
863            HummockError::next_epoch(format!(
864                "next_epoch for {} {} should be valid",
865                options.table_id, epoch
866            ))
867            .into()
868        })
869    }
870
871    async fn iter_log(
872        &self,
873        epoch_range: (u64, u64),
874        key_range: TableKeyRange,
875        options: ReadLogOptions,
876    ) -> StorageResult<Self::ChangeLogIter> {
877        let iter = self
878            .hummock_version_reader
879            .iter_log(
880                epoch_range,
881                key_range,
882                options,
883                self.table_change_log_manager.clone(),
884            )
885            .await?;
886        Ok(iter)
887    }
888}
889
890impl HummockStorage {
891    /// Waits until the local hummock version contains the epoch. If `wait_epoch` is `Current`,
892    /// we will only check whether it is le `sealed_epoch` and won't wait.
893    async fn try_wait_epoch_impl(
894        &self,
895        wait_epoch: HummockReadEpoch,
896        table_id: TableId,
897    ) -> StorageResult<()> {
898        tracing::debug!(
899            "try_wait_epoch: epoch: {:?}, table_id: {}",
900            wait_epoch,
901            table_id
902        );
903        match wait_epoch {
904            HummockReadEpoch::Committed(wait_epoch) => {
905                assert!(!is_max_epoch(wait_epoch), "epoch should not be MAX EPOCH");
906                wait_for_epoch(&self.version_update_notifier_tx, wait_epoch, table_id).await?;
907            }
908            HummockReadEpoch::BatchQueryCommitted(wait_epoch, wait_version_id) => {
909                assert!(!is_max_epoch(wait_epoch), "epoch should not be MAX EPOCH");
910                // fast path by checking recent_versions
911                {
912                    let recent_versions = self.recent_versions.load();
913                    let latest_version = recent_versions.latest_version();
914                    if latest_version.id >= wait_version_id
915                        && let Some(committed_epoch) =
916                            latest_version.table_committed_epoch(table_id)
917                        && committed_epoch >= wait_epoch
918                    {
919                        return Ok(());
920                    }
921                }
922                wait_for_update(
923                    &self.version_update_notifier_tx,
924                    |version| {
925                        if wait_version_id > version.id() {
926                            return Ok(false);
927                        }
928                        let committed_epoch =
929                            version.table_committed_epoch(table_id).ok_or_else(|| {
930                                // In batch query, since we have ensured that the current version must be after the
931                                // `wait_version_id`, when seeing that the table_id not exist in the latest version,
932                                // the table must have been dropped.
933                                HummockError::wait_epoch(format!(
934                                    "table id {} has been dropped",
935                                    table_id
936                                ))
937                            })?;
938                        Ok(committed_epoch >= wait_epoch)
939                    },
940                    || {
941                        format!(
942                            "try_wait_epoch: epoch: {}, version_id: {:?}",
943                            wait_epoch, wait_version_id
944                        )
945                    },
946                )
947                .await?;
948            }
949            _ => {}
950        };
951        Ok(())
952    }
953}
954
955impl StateStore for HummockStorage {
956    type Local = LocalHummockStorage;
957    type ReadSnapshot = HummockStorageReadSnapshot;
958    type VectorWriter = HummockVectorWriter;
959
960    /// Waits until the local hummock version contains the epoch. If `wait_epoch` is `Current`,
961    /// we will only check whether it is le `sealed_epoch` and won't wait.
962    async fn try_wait_epoch(
963        &self,
964        wait_epoch: HummockReadEpoch,
965        options: TryWaitEpochOptions,
966    ) -> StorageResult<()> {
967        self.try_wait_epoch_impl(wait_epoch, options.table_id).await
968    }
969
970    fn new_local(&self, option: NewLocalOptions) -> impl Future<Output = Self::Local> + Send + '_ {
971        self.new_local_inner(option)
972    }
973
974    async fn new_read_snapshot(
975        &self,
976        epoch: HummockReadEpoch,
977        options: NewReadSnapshotOptions,
978    ) -> StorageResult<Self::ReadSnapshot> {
979        self.try_wait_epoch_impl(epoch, options.table_id).await?;
980        Ok(HummockStorageReadSnapshot {
981            epoch,
982            table_id: options.table_id,
983            table_option: options.table_option,
984            recent_versions: self.recent_versions.clone(),
985            hummock_version_reader: self.hummock_version_reader.clone(),
986            read_version_mapping: self.read_version_mapping.clone(),
987            backup_reader: self.backup_reader.clone(),
988            hummock_meta_client: self.hummock_meta_client.clone(),
989            simple_time_travel_version_cache: self.simple_time_travel_version_cache.clone(),
990        })
991    }
992
993    async fn new_vector_writer(&self, options: NewVectorWriterOptions) -> Self::VectorWriter {
994        HummockVectorWriter::new(
995            options.table_id,
996            self.version_update_notifier_tx.clone(),
997            self.context.sstable_store.clone(),
998            self.object_id_manager.clone(),
999            self.hummock_event_sender.clone(),
1000            self.hummock_version_reader.stats().clone(),
1001            self.context.storage_opts.clone(),
1002        )
1003    }
1004}
1005
1006#[cfg(any(test, feature = "test"))]
1007impl HummockStorage {
1008    pub async fn seal_and_sync_epoch(
1009        &self,
1010        epoch: u64,
1011        table_ids: HashSet<TableId>,
1012    ) -> StorageResult<risingwave_hummock_sdk::SyncResult> {
1013        self.sync(vec![(epoch, table_ids)]).await
1014    }
1015
1016    /// Used in the compaction test tool
1017    pub async fn update_version_and_wait(&self, version: HummockVersion) {
1018        use tokio::task::yield_now;
1019        let version_id = version.id;
1020        self._observer_event_sender
1021            .send(HummockObserverEvent::VersionUpdate(
1022                HummockVersionUpdate::PinnedVersion(Box::new(version)),
1023            ))
1024            .unwrap();
1025        loop {
1026            if self.recent_versions.load().latest_version().id() >= version_id {
1027                break;
1028            }
1029
1030            yield_now().await
1031        }
1032    }
1033
1034    pub async fn wait_version(&self, version: HummockVersion) {
1035        use tokio::task::yield_now;
1036        loop {
1037            if self.recent_versions.load().latest_version().id() >= version.id {
1038                break;
1039            }
1040
1041            yield_now().await
1042        }
1043    }
1044
1045    /// Creates a [`HummockStorage`] with default stats. Should only be used by tests.
1046    pub async fn for_test(
1047        options: Arc<StorageOpts>,
1048        sstable_store: SstableStoreRef,
1049        hummock_meta_client: Arc<dyn HummockMetaClient>,
1050        notification_client: impl NotificationClient,
1051    ) -> HummockResult<Self> {
1052        let compaction_catalog_manager = Arc::new(CompactionCatalogManager::new(Box::new(
1053            FakeRemoteTableAccessor {},
1054        )));
1055
1056        Self::new(
1057            Role::Both,
1058            options,
1059            sstable_store,
1060            hummock_meta_client,
1061            notification_client,
1062            compaction_catalog_manager,
1063            Arc::new(HummockStateStoreMetrics::unused()),
1064            Arc::new(CompactorMetrics::unused()),
1065            None,
1066        )
1067        .await
1068    }
1069
1070    pub fn storage_opts(&self) -> &Arc<StorageOpts> {
1071        &self.context.storage_opts
1072    }
1073
1074    pub fn version_reader(&self) -> &HummockVersionReader {
1075        &self.hummock_version_reader
1076    }
1077
1078    pub async fn wait_version_update(
1079        &self,
1080        old_id: risingwave_hummock_sdk::HummockVersionId,
1081    ) -> risingwave_hummock_sdk::HummockVersionId {
1082        use tokio::task::yield_now;
1083        loop {
1084            let cur_id = self.recent_versions.load().latest_version().id();
1085            if cur_id > old_id {
1086                return cur_id;
1087            }
1088            yield_now().await;
1089        }
1090    }
1091
1092    #[cfg(any(test, feature = "test"))]
1093    pub async fn flush_events_for_test(&self) {
1094        let (tx, rx) = oneshot::channel();
1095        self.hummock_event_sender
1096            .send(HummockEvent::FlushEvent(tx))
1097            .expect("flush event should succeed");
1098        rx.await.expect("flush event receiver dropped");
1099    }
1100}