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 ),
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 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 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 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 true
512 }
513 } else {
514 false
515 }
516 })
517 .cloned()
518 .collect_vec()
519 })
520 .unwrap_or_default()
521 };
522
523 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 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 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 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 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 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 {
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 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 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 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 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}