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            ),
250            _shutdown_guard: Arc::new(HummockStorageShutdownGuard {
251                shutdown_sender: event_tx,
252            }),
253            read_version_mapping: hummock_event_handler.read_version_mapping(),
254            backup_reader,
255            write_limiter,
256            compact_await_tree_reg: await_tree_reg,
257            hummock_meta_client,
258            simple_time_travel_version_cache: Arc::new(SimpleTimeTravelVersionCache::new(
259                options.time_travel_version_cache_capacity,
260            )),
261            table_change_log_manager,
262        };
263
264        tokio::spawn(hummock_event_handler.start_hummock_event_handler_worker());
265
266        Ok(instance)
267    }
268}
269
270impl HummockStorageReadSnapshot {
271    /// Gets the value of a specified `key` in the table specified in `read_options`.
272    /// The result is based on a snapshot corresponding to the given `epoch`.
273    /// if `key` has consistent hash virtual node value, then such value is stored in `value_meta`
274    ///
275    /// If `Ok(Some())` is returned, the key is found. If `Ok(None)` is returned,
276    /// the key is not found. If `Err()` is returned, the searching for the key
277    /// failed due to other non-EOF errors.
278    async fn get_inner<'a, O>(
279        &'a self,
280        key: TableKey<Bytes>,
281        read_options: ReadOptions,
282        on_key_value_fn: impl KeyValueFn<'a, O>,
283    ) -> StorageResult<Option<O>> {
284        let key_range = (Bound::Included(key.clone()), Bound::Included(key.clone()));
285
286        let (key_range, read_version_tuple) =
287            self.build_read_version_tuple(self.epoch, key_range).await?;
288
289        if is_empty_key_range(&key_range) {
290            return Ok(None);
291        }
292
293        self.hummock_version_reader
294            .get(
295                key,
296                self.epoch.get_epoch(),
297                self.table_id,
298                self.table_option,
299                read_options,
300                read_version_tuple,
301                on_key_value_fn,
302            )
303            .await
304    }
305
306    async fn iter_inner(
307        &self,
308        key_range: TableKeyRange,
309        read_options: ReadOptions,
310    ) -> StorageResult<HummockStorageIterator> {
311        let (key_range, read_version_tuple) =
312            self.build_read_version_tuple(self.epoch, key_range).await?;
313
314        self.hummock_version_reader
315            .iter(
316                key_range,
317                self.epoch.get_epoch(),
318                self.table_id,
319                self.table_option,
320                read_options,
321                read_version_tuple,
322            )
323            .await
324    }
325
326    async fn rev_iter_inner(
327        &self,
328        key_range: TableKeyRange,
329        read_options: ReadOptions,
330    ) -> StorageResult<HummockStorageRevIterator> {
331        let (key_range, read_version_tuple) =
332            self.build_read_version_tuple(self.epoch, key_range).await?;
333
334        self.hummock_version_reader
335            .rev_iter(
336                key_range,
337                self.epoch.get_epoch(),
338                self.table_id,
339                self.table_option,
340                read_options,
341                read_version_tuple,
342                None,
343            )
344            .await
345    }
346
347    async fn get_time_travel_version(
348        &self,
349        epoch: u64,
350        table_id: TableId,
351    ) -> StorageResult<PinnedVersion> {
352        let meta_client = self.hummock_meta_client.clone();
353        let fetch = async move {
354            let pb_version = meta_client
355                .get_version_by_epoch(epoch, table_id)
356                .await
357                .inspect_err(|e| tracing::error!("{}", e.to_report_string()))
358                .map_err(|e| match &e {
359                    RpcError::GrpcStatus(status)
360                        if status.inner().code() == tonic::Code::OutOfRange =>
361                    {
362                        HummockError::time_travel_version_expired(table_id, epoch)
363                    }
364                    _ => HummockError::meta_error(e.to_report_string()),
365                })?;
366            let version = HummockVersion::from_rpc_protobuf(&pb_version);
367            let (tx, _rx) = unbounded_channel();
368            Ok(PinnedVersion::new(version, tx))
369        };
370        let version = self
371            .simple_time_travel_version_cache
372            .get_or_insert(table_id, epoch, fetch)
373            .await?;
374        Ok(version)
375    }
376
377    async fn build_read_version_tuple(
378        &self,
379        epoch: HummockReadEpoch,
380        key_range: TableKeyRange,
381    ) -> StorageResult<(TableKeyRange, ReadVersionTuple)> {
382        match epoch {
383            HummockReadEpoch::Backup(epoch) => {
384                self.build_read_version_tuple_from_backup(epoch, self.table_id, key_range)
385                    .await
386            }
387            HummockReadEpoch::Committed(epoch) => {
388                let tuple = self
389                    .build_read_version_tuple_from_committed(epoch, self.table_id, key_range)
390                    .await?;
391                let (_, (_, _, version)) = &tuple;
392                let Some(committed_epoch) = version.table_committed_epoch(self.table_id) else {
393                    return Err(HummockError::other(format!(
394                        "table {} not found in version",
395                        self.table_id
396                    ))
397                    .into());
398                };
399                if committed_epoch != epoch {
400                    return Err(HummockError::committed_epoch_mismatch(
401                        self.table_id,
402                        committed_epoch,
403                        epoch,
404                    )
405                    .into());
406                }
407                Ok(tuple)
408            }
409            HummockReadEpoch::BatchQueryCommitted(epoch, _)
410            | HummockReadEpoch::TimeTravel(epoch) => {
411                self.build_read_version_tuple_from_committed(epoch, self.table_id, key_range)
412                    .await
413            }
414            HummockReadEpoch::NoWait(epoch) => {
415                self.build_read_version_tuple_from_all(epoch, self.table_id, key_range)
416                    .await
417            }
418        }
419    }
420
421    async fn build_read_version_tuple_from_backup(
422        &self,
423        epoch: u64,
424        table_id: TableId,
425        key_range: TableKeyRange,
426    ) -> StorageResult<(TableKeyRange, ReadVersionTuple)> {
427        match self
428            .backup_reader
429            .try_get_hummock_version(table_id, epoch)
430            .await
431        {
432            Ok(Some(backup_version)) => Ok(get_committed_read_version_tuple(
433                backup_version,
434                table_id,
435                key_range,
436                epoch,
437            )),
438            Ok(None) => Err(HummockError::read_backup_error(format!(
439                "backup include epoch {} not found",
440                epoch
441            ))
442            .into()),
443            Err(e) => Err(e),
444        }
445    }
446
447    async fn get_epoch_hummock_version(
448        &self,
449        epoch: u64,
450        table_id: TableId,
451    ) -> StorageResult<PinnedVersion> {
452        match self
453            .recent_versions
454            .load()
455            .get_safe_version(table_id, epoch)
456        {
457            Some(version) => Ok(version),
458            None => self.get_time_travel_version(epoch, table_id).await,
459        }
460    }
461
462    async fn build_read_version_tuple_from_committed(
463        &self,
464        epoch: u64,
465        table_id: TableId,
466        key_range: TableKeyRange,
467    ) -> StorageResult<(TableKeyRange, ReadVersionTuple)> {
468        let version = self.get_epoch_hummock_version(epoch, table_id).await?;
469        Ok(get_committed_read_version_tuple(
470            version, table_id, key_range, epoch,
471        ))
472    }
473
474    async fn build_read_version_tuple_from_all(
475        &self,
476        epoch: u64,
477        table_id: TableId,
478        key_range: TableKeyRange,
479    ) -> StorageResult<(TableKeyRange, ReadVersionTuple)> {
480        let pinned_version = self.recent_versions.load().latest_version().clone();
481        let info = pinned_version.state_table_info.info().get(&table_id);
482
483        // check epoch if lower mce
484        let ret = if let Some(info) = info
485            && epoch <= info.committed_epoch
486        {
487            let pinned_version = if epoch < info.committed_epoch {
488                pinned_version
489            } else {
490                self.get_epoch_hummock_version(epoch, table_id).await?
491            };
492            // read committed_version directly without build snapshot
493            get_committed_read_version_tuple(pinned_version, table_id, key_range, epoch)
494        } else {
495            let vnode = vnode(&key_range);
496            let mut matched_replicated_read_version_cnt = 0;
497            let read_version_vec = {
498                let read_guard = self.read_version_mapping.read();
499                read_guard
500                    .get(&table_id)
501                    .map(|v| {
502                        v.values()
503                            .filter(|v| {
504                                let read_version = v.read();
505                                if read_version.is_initialized() && read_version.contains(vnode) {
506                                    if read_version.is_replicated() {
507                                        matched_replicated_read_version_cnt += 1;
508                                        false
509                                    } else {
510                                        // Only non-replicated read version with matched vnode is considered
511                                        true
512                                    }
513                                } else {
514                                    false
515                                }
516                            })
517                            .cloned()
518                            .collect_vec()
519                    })
520                    .unwrap_or_default()
521            };
522
523            // When the system has just started and no state has been created, the memory state
524            // may be empty
525            if read_version_vec.is_empty() {
526                let table_committed_epoch = info.map(|info| info.committed_epoch);
527                if matched_replicated_read_version_cnt > 0 {
528                    tracing::warn!(
529                        "Read(table_id={} vnode={} epoch={}) is not allowed on replicated read version ({} found). Fall back to committed version (epoch={:?})",
530                        table_id,
531                        vnode.to_index(),
532                        epoch,
533                        matched_replicated_read_version_cnt,
534                        table_committed_epoch,
535                    );
536                } else {
537                    tracing::debug!(
538                        "No read version found for read(table_id={} vnode={} epoch={}). Fall back to committed version (epoch={:?})",
539                        table_id,
540                        vnode.to_index(),
541                        epoch,
542                        table_committed_epoch
543                    );
544                }
545                get_committed_read_version_tuple(pinned_version, table_id, key_range, epoch)
546            } else {
547                if read_version_vec.len() != 1 {
548                    let read_version_vnodes = read_version_vec
549                        .into_iter()
550                        .map(|v| {
551                            let v = v.read();
552                            v.vnodes().iter_ones().collect_vec()
553                        })
554                        .collect_vec();
555                    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());
556                }
557                read_filter_for_version(
558                    epoch,
559                    table_id,
560                    key_range,
561                    read_version_vec.first().unwrap(),
562                )?
563            }
564        };
565
566        Ok(ret)
567    }
568}
569
570impl HummockStorage {
571    async fn new_local_inner(&self, option: NewLocalOptions) -> LocalHummockStorage {
572        let (tx, rx) = tokio::sync::oneshot::channel();
573        self.hummock_event_sender
574            .send(HummockEvent::RegisterReadVersion {
575                table_id: option.table_id,
576                new_read_version_sender: tx,
577                is_replicated: option.is_replicated,
578                vnodes: option.vnodes.clone(),
579            })
580            .unwrap();
581
582        let (basic_read_version, instance_guard) = rx.await.unwrap();
583        let version_update_notifier_tx = self.version_update_notifier_tx.clone();
584        LocalHummockStorage::new(
585            instance_guard,
586            basic_read_version,
587            self.hummock_version_reader.clone(),
588            self.hummock_event_sender.clone(),
589            self.buffer_tracker.get_memory_limiter().clone(),
590            self.write_limiter.clone(),
591            option,
592            version_update_notifier_tx,
593            self.context.storage_opts.mem_table_spill_threshold,
594        )
595    }
596
597    pub async fn clear_shared_buffer(&self) {
598        let (tx, rx) = oneshot::channel();
599        self.hummock_event_sender
600            .send(HummockEvent::Clear(tx, None))
601            .expect("should send success");
602        rx.await.expect("should wait success");
603    }
604
605    pub async fn clear_tables(&self, table_ids: HashSet<TableId>) {
606        if !table_ids.is_empty() {
607            let (tx, rx) = oneshot::channel();
608            self.hummock_event_sender
609                .send(HummockEvent::Clear(tx, Some(table_ids)))
610                .expect("should send success");
611            rx.await.expect("should wait success");
612        }
613    }
614
615    /// Declare the start of an epoch. This information is provided for spill so that the spill task won't
616    /// include data of two or more syncs.
617    // TODO: remove this method when we support spill task that can include data of more two or more syncs
618    pub fn start_epoch(&self, epoch: HummockEpoch, table_ids: HashSet<TableId>) {
619        let _ = self
620            .hummock_event_sender
621            .send(HummockEvent::StartEpoch { epoch, table_ids });
622    }
623
624    pub fn sstable_store(&self) -> SstableStoreRef {
625        self.context.sstable_store.clone()
626    }
627
628    pub async fn table_cache_refill_monitor_snapshot(
629        &self,
630    ) -> HummockResult<TableCacheRefillMonitorSnapshot> {
631        let (tx, rx) = oneshot::channel();
632        self.hummock_event_sender
633            .send(HummockEvent::GetTableCacheRefillMonitorSnapshot { result_tx: tx })
634            .map_err(|_| HummockError::other("failed to send table cache refill monitor query"))?;
635        rx.await.map_err(|_| {
636            HummockError::other("failed to receive table cache refill monitor snapshot")
637        })
638    }
639
640    pub fn object_id_manager(&self) -> &ObjectIdManagerRef {
641        &self.object_id_manager
642    }
643
644    pub fn compaction_catalog_manager_ref(&self) -> CompactionCatalogManagerRef {
645        self.compaction_catalog_manager_ref.clone()
646    }
647
648    pub fn get_memory_limiter(&self) -> Arc<MemoryLimiter> {
649        self.buffer_tracker.get_memory_limiter().clone()
650    }
651
652    pub fn get_pinned_version(&self) -> PinnedVersion {
653        self.recent_versions.load().latest_version().clone()
654    }
655
656    pub fn backup_reader(&self) -> BackupReaderRef {
657        self.backup_reader.clone()
658    }
659
660    pub fn compaction_await_tree_reg(&self) -> Option<&await_tree::Registry> {
661        self.compact_await_tree_reg.as_ref()
662    }
663
664    pub async fn min_uncommitted_object_id(&self) -> Option<HummockRawObjectId> {
665        let (tx, rx) = oneshot::channel();
666        self.hummock_event_sender
667            .send(HummockEvent::GetMinUncommittedObjectId { result_tx: tx })
668            .expect("should send success");
669        rx.await.expect("should await success")
670    }
671
672    pub async fn sync(
673        &self,
674        sync_table_epochs: Vec<(HummockEpoch, HashSet<TableId>)>,
675    ) -> StorageResult<SyncResult> {
676        let (tx, rx) = oneshot::channel();
677        let _ = self.hummock_event_sender.send(HummockEvent::SyncEpoch {
678            sync_result_sender: tx,
679            sync_table_epochs,
680        });
681        let synced_data = rx
682            .await
683            .map_err(|_| HummockError::other("failed to receive sync result"))??;
684        Ok(synced_data.into_sync_result())
685    }
686}
687
688#[derive(Clone)]
689pub struct HummockStorageReadSnapshot {
690    epoch: HummockReadEpoch,
691    table_id: TableId,
692    table_option: TableOption,
693    recent_versions: Arc<ArcSwap<RecentVersions>>,
694    hummock_version_reader: HummockVersionReader,
695    read_version_mapping: ReadOnlyReadVersionMapping,
696    backup_reader: BackupReaderRef,
697    hummock_meta_client: Arc<dyn HummockMetaClient>,
698    simple_time_travel_version_cache: Arc<SimpleTimeTravelVersionCache>,
699}
700
701impl StateStoreGet for HummockStorageReadSnapshot {
702    fn on_key_value<'a, O: Send + 'a>(
703        &'a self,
704        key: TableKey<Bytes>,
705        read_options: ReadOptions,
706        on_key_value_fn: impl KeyValueFn<'a, O>,
707    ) -> impl StorageFuture<'a, Option<O>> {
708        self.get_inner(key, read_options, on_key_value_fn)
709    }
710}
711
712impl StateStoreRead for HummockStorageReadSnapshot {
713    type Iter = HummockStorageIterator;
714    type RevIter = HummockStorageRevIterator;
715
716    fn iter(
717        &self,
718        key_range: TableKeyRange,
719        read_options: ReadOptions,
720    ) -> impl Future<Output = StorageResult<Self::Iter>> + '_ {
721        let (l_vnode_inclusive, r_vnode_exclusive) = vnode_range(&key_range);
722        assert_eq!(
723            r_vnode_exclusive - l_vnode_inclusive,
724            1,
725            "read range {:?} for table {} iter contains more than one vnode",
726            key_range,
727            self.table_id
728        );
729        self.iter_inner(key_range, read_options)
730    }
731
732    fn rev_iter(
733        &self,
734        key_range: TableKeyRange,
735        read_options: ReadOptions,
736    ) -> impl Future<Output = StorageResult<Self::RevIter>> + '_ {
737        let (l_vnode_inclusive, r_vnode_exclusive) = vnode_range(&key_range);
738        assert_eq!(
739            r_vnode_exclusive - l_vnode_inclusive,
740            1,
741            "read range {:?} for table {} iter contains more than one vnode",
742            key_range,
743            self.table_id
744        );
745        self.rev_iter_inner(key_range, read_options)
746    }
747}
748
749impl StateStoreReadVector for HummockStorageReadSnapshot {
750    async fn nearest<'a, O: Send + 'a>(
751        &'a self,
752        vec: VectorRef<'a>,
753        options: VectorNearestOptions,
754        on_nearest_item_fn: impl OnNearestItemFn<'a, O>,
755    ) -> StorageResult<Vec<O>> {
756        let version = match self.epoch {
757            HummockReadEpoch::Committed(epoch)
758            | HummockReadEpoch::BatchQueryCommitted(epoch, _)
759            | HummockReadEpoch::TimeTravel(epoch) => {
760                self.get_epoch_hummock_version(epoch, self.table_id).await?
761            }
762            HummockReadEpoch::Backup(epoch) => self
763                .backup_reader
764                .try_get_hummock_version(self.table_id, epoch)
765                .await?
766                .ok_or_else(|| {
767                    HummockError::read_backup_error(format!(
768                        "backup include epoch {} not found",
769                        epoch
770                    ))
771                })?,
772            HummockReadEpoch::NoWait(_) => {
773                return Err(
774                    HummockError::other("nearest query does not support NoWait epoch").into(),
775                );
776            }
777        };
778        dispatch_distance_measurement!(options.measure, MeasurementType, {
779            Ok(self
780                .hummock_version_reader
781                .nearest::<MeasurementType, O>(
782                    version,
783                    self.table_id,
784                    vec,
785                    options,
786                    on_nearest_item_fn,
787                )
788                .await?)
789        })
790    }
791}
792
793impl StateStoreReadLog for HummockStorage {
794    type ChangeLogIter = ChangeLogIterator;
795
796    async fn next_epoch(&self, epoch: u64, options: NextEpochOptions) -> StorageResult<u64> {
797        fn next_epoch(
798            table_change_log: &TableChangeLogs,
799            epoch: u64,
800            table_id: TableId,
801        ) -> HummockResult<Option<u64>> {
802            let table_change_log = table_change_log.get(&table_id).ok_or_else(|| {
803                HummockError::next_epoch(format!("table {} has been dropped", table_id))
804            })?;
805            table_change_log
806                .next_epoch(epoch)
807                .map_err(|_| HummockError::change_log_retention_miss(table_id, epoch))
808        }
809        {
810            // fast path
811            if let Some(max_epoch) = self
812                .recent_versions
813                .load()
814                .latest_version()
815                .deref()
816                .state_table_info
817                .info()
818                .get(&options.table_id)
819                .map(|i| i.committed_epoch)
820                && max_epoch > epoch
821            {
822                // The next epoch exists either in the same `EpochNewChangeLog` or the next `EpochNewChangeLog`, so we fetch 2 `EpochNewChangeLogCommon`.
823                let table_change_log = self
824                    .table_change_log_manager
825                    .fetch_table_change_logs(options.table_id, (epoch, max_epoch), true, Some(2))
826                    .await?;
827                if let Some(next_epoch) = next_epoch(&table_change_log, epoch, options.table_id)? {
828                    return Ok(next_epoch);
829                }
830            }
831        }
832        let mut max_epoch = None;
833        wait_for_update(
834            &self.version_update_notifier_tx,
835            |version| {
836                let Some(mce) = version
837                    .state_table_info
838                    .info()
839                    .get(&options.table_id)
840                    .map(|i| i.committed_epoch)
841                else {
842                    return Ok(false);
843                };
844                max_epoch = Some(mce);
845                Ok(mce > epoch)
846            },
847            || format!("wait next_epoch: epoch: {} {}", epoch, options.table_id),
848        )
849        .await?;
850        // The next epoch exists either in the same `EpochNewChangeLog` or the next `EpochNewChangeLog`, so we fetch 2 `EpochNewChangeLogCommon`.
851        let table_change_log = self
852            .table_change_log_manager
853            .fetch_table_change_logs(
854                options.table_id,
855                (epoch, max_epoch.unwrap_or(u64::MAX)),
856                true,
857                Some(2),
858            )
859            .await?;
860        let next_epoch_ret = next_epoch(&table_change_log, epoch, options.table_id)?;
861        next_epoch_ret.ok_or_else(|| {
862            HummockError::next_epoch(format!(
863                "next_epoch for {} {} should be valid",
864                options.table_id, epoch
865            ))
866            .into()
867        })
868    }
869
870    async fn iter_log(
871        &self,
872        epoch_range: (u64, u64),
873        key_range: TableKeyRange,
874        options: ReadLogOptions,
875    ) -> StorageResult<Self::ChangeLogIter> {
876        let iter = self
877            .hummock_version_reader
878            .iter_log(
879                epoch_range,
880                key_range,
881                options,
882                self.table_change_log_manager.clone(),
883            )
884            .await?;
885        Ok(iter)
886    }
887}
888
889impl HummockStorage {
890    /// Waits until the local hummock version contains the epoch. If `wait_epoch` is `Current`,
891    /// we will only check whether it is le `sealed_epoch` and won't wait.
892    async fn try_wait_epoch_impl(
893        &self,
894        wait_epoch: HummockReadEpoch,
895        table_id: TableId,
896    ) -> StorageResult<()> {
897        tracing::debug!(
898            "try_wait_epoch: epoch: {:?}, table_id: {}",
899            wait_epoch,
900            table_id
901        );
902        match wait_epoch {
903            HummockReadEpoch::Committed(wait_epoch) => {
904                assert!(!is_max_epoch(wait_epoch), "epoch should not be MAX EPOCH");
905                wait_for_epoch(&self.version_update_notifier_tx, wait_epoch, table_id).await?;
906            }
907            HummockReadEpoch::BatchQueryCommitted(wait_epoch, wait_version_id) => {
908                assert!(!is_max_epoch(wait_epoch), "epoch should not be MAX EPOCH");
909                // fast path by checking recent_versions
910                {
911                    let recent_versions = self.recent_versions.load();
912                    let latest_version = recent_versions.latest_version();
913                    if latest_version.id >= wait_version_id
914                        && let Some(committed_epoch) =
915                            latest_version.table_committed_epoch(table_id)
916                        && committed_epoch >= wait_epoch
917                    {
918                        return Ok(());
919                    }
920                }
921                wait_for_update(
922                    &self.version_update_notifier_tx,
923                    |version| {
924                        if wait_version_id > version.id() {
925                            return Ok(false);
926                        }
927                        let committed_epoch =
928                            version.table_committed_epoch(table_id).ok_or_else(|| {
929                                // In batch query, since we have ensured that the current version must be after the
930                                // `wait_version_id`, when seeing that the table_id not exist in the latest version,
931                                // the table must have been dropped.
932                                HummockError::wait_epoch(format!(
933                                    "table id {} has been dropped",
934                                    table_id
935                                ))
936                            })?;
937                        Ok(committed_epoch >= wait_epoch)
938                    },
939                    || {
940                        format!(
941                            "try_wait_epoch: epoch: {}, version_id: {:?}",
942                            wait_epoch, wait_version_id
943                        )
944                    },
945                )
946                .await?;
947            }
948            _ => {}
949        };
950        Ok(())
951    }
952}
953
954impl StateStore for HummockStorage {
955    type Local = LocalHummockStorage;
956    type ReadSnapshot = HummockStorageReadSnapshot;
957    type VectorWriter = HummockVectorWriter;
958
959    /// Waits until the local hummock version contains the epoch. If `wait_epoch` is `Current`,
960    /// we will only check whether it is le `sealed_epoch` and won't wait.
961    async fn try_wait_epoch(
962        &self,
963        wait_epoch: HummockReadEpoch,
964        options: TryWaitEpochOptions,
965    ) -> StorageResult<()> {
966        self.try_wait_epoch_impl(wait_epoch, options.table_id).await
967    }
968
969    fn new_local(&self, option: NewLocalOptions) -> impl Future<Output = Self::Local> + Send + '_ {
970        self.new_local_inner(option)
971    }
972
973    async fn new_read_snapshot(
974        &self,
975        epoch: HummockReadEpoch,
976        options: NewReadSnapshotOptions,
977    ) -> StorageResult<Self::ReadSnapshot> {
978        self.try_wait_epoch_impl(epoch, options.table_id).await?;
979        Ok(HummockStorageReadSnapshot {
980            epoch,
981            table_id: options.table_id,
982            table_option: options.table_option,
983            recent_versions: self.recent_versions.clone(),
984            hummock_version_reader: self.hummock_version_reader.clone(),
985            read_version_mapping: self.read_version_mapping.clone(),
986            backup_reader: self.backup_reader.clone(),
987            hummock_meta_client: self.hummock_meta_client.clone(),
988            simple_time_travel_version_cache: self.simple_time_travel_version_cache.clone(),
989        })
990    }
991
992    async fn new_vector_writer(&self, options: NewVectorWriterOptions) -> Self::VectorWriter {
993        HummockVectorWriter::new(
994            options.table_id,
995            self.version_update_notifier_tx.clone(),
996            self.context.sstable_store.clone(),
997            self.object_id_manager.clone(),
998            self.hummock_event_sender.clone(),
999            self.hummock_version_reader.stats().clone(),
1000            self.context.storage_opts.clone(),
1001        )
1002    }
1003}
1004
1005#[cfg(any(test, feature = "test"))]
1006impl HummockStorage {
1007    pub async fn seal_and_sync_epoch(
1008        &self,
1009        epoch: u64,
1010        table_ids: HashSet<TableId>,
1011    ) -> StorageResult<risingwave_hummock_sdk::SyncResult> {
1012        self.sync(vec![(epoch, table_ids)]).await
1013    }
1014
1015    /// Used in the compaction test tool
1016    pub async fn update_version_and_wait(&self, version: HummockVersion) {
1017        use tokio::task::yield_now;
1018        let version_id = version.id;
1019        self._observer_event_sender
1020            .send(HummockObserverEvent::VersionUpdate(
1021                HummockVersionUpdate::PinnedVersion(Box::new(version)),
1022            ))
1023            .unwrap();
1024        loop {
1025            if self.recent_versions.load().latest_version().id() >= version_id {
1026                break;
1027            }
1028
1029            yield_now().await
1030        }
1031    }
1032
1033    pub async fn wait_version(&self, version: HummockVersion) {
1034        use tokio::task::yield_now;
1035        loop {
1036            if self.recent_versions.load().latest_version().id() >= version.id {
1037                break;
1038            }
1039
1040            yield_now().await
1041        }
1042    }
1043
1044    /// Creates a [`HummockStorage`] with default stats. Should only be used by tests.
1045    pub async fn for_test(
1046        options: Arc<StorageOpts>,
1047        sstable_store: SstableStoreRef,
1048        hummock_meta_client: Arc<dyn HummockMetaClient>,
1049        notification_client: impl NotificationClient,
1050    ) -> HummockResult<Self> {
1051        let compaction_catalog_manager = Arc::new(CompactionCatalogManager::new(Box::new(
1052            FakeRemoteTableAccessor {},
1053        )));
1054
1055        Self::new(
1056            Role::Both,
1057            options,
1058            sstable_store,
1059            hummock_meta_client,
1060            notification_client,
1061            compaction_catalog_manager,
1062            Arc::new(HummockStateStoreMetrics::unused()),
1063            Arc::new(CompactorMetrics::unused()),
1064            None,
1065        )
1066        .await
1067    }
1068
1069    pub fn storage_opts(&self) -> &Arc<StorageOpts> {
1070        &self.context.storage_opts
1071    }
1072
1073    pub fn version_reader(&self) -> &HummockVersionReader {
1074        &self.hummock_version_reader
1075    }
1076
1077    pub async fn wait_version_update(
1078        &self,
1079        old_id: risingwave_hummock_sdk::HummockVersionId,
1080    ) -> risingwave_hummock_sdk::HummockVersionId {
1081        use tokio::task::yield_now;
1082        loop {
1083            let cur_id = self.recent_versions.load().latest_version().id();
1084            if cur_id > old_id {
1085                return cur_id;
1086            }
1087            yield_now().await;
1088        }
1089    }
1090
1091    #[cfg(any(test, feature = "test"))]
1092    pub async fn flush_events_for_test(&self) {
1093        let (tx, rx) = oneshot::channel();
1094        self.hummock_event_sender
1095            .send(HummockEvent::FlushEvent(tx))
1096            .expect("flush event should succeed");
1097        rx.await.expect("flush event receiver dropped");
1098    }
1099}