1use 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#[derive(Clone)]
96pub struct HummockStorage {
97 hummock_event_sender: HummockEventSender,
98 _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 #[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 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 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 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 true
513 }
514 } else {
515 false
516 }
517 })
518 .cloned()
519 .collect_vec()
520 })
521 .unwrap_or_default()
522 };
523
524 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 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 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 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 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 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 {
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 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 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 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 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}