1use std::collections::HashSet;
16use std::fmt::Debug;
17use std::sync::{Arc, LazyLock};
18use std::time::Duration;
19
20use enum_as_inner::EnumAsInner;
21use foyer::{
22 BlockEngineConfig, CacheBuilder, DeviceBuilder, FifoPicker, FsDeviceBuilder, HybridCacheBuilder,
23};
24use futures::FutureExt;
25use futures::future::BoxFuture;
26use mixtrics::registry::prometheus::PrometheusMetricsRegistry;
27use risingwave_common::catalog::TableId;
28use risingwave_common::config::Role;
29use risingwave_common::config::storage::FileCacheRuntimeConfig;
30use risingwave_common::license::Feature;
31use risingwave_common::monitor::GLOBAL_METRICS_REGISTRY;
32use risingwave_common_service::RpcNotificationClient;
33use risingwave_hummock_sdk::{HummockEpoch, HummockSstableObjectId, SyncResult};
34use risingwave_object_store::object::build_remote_object_store;
35use thiserror_ext::AsReport;
36
37use crate::StateStore;
38use crate::compaction_catalog_manager::{CompactionCatalogManager, RemoteTableAccessor};
39use crate::error::StorageResult;
40use crate::hummock::all::AllRecentFilter;
41use crate::hummock::hummock_meta_client::MonitoredHummockMetaClient;
42use crate::hummock::none::NoneRecentFilter;
43use crate::hummock::sharded::ShardedRecentFilter;
44use crate::hummock::simple::SimpleRecentFilter;
45use crate::hummock::{
46 Block, BlockCacheEventListener, HummockError, HummockStorage, Sstable, SstableBlockIndex,
47 SstableStore, SstableStoreConfig,
48};
49use crate::memory::MemoryStateStore;
50use crate::memory::sled::SledStateStore;
51use crate::monitor::{
52 CompactorMetrics, HummockStateStoreMetrics, MonitoredStateStore, MonitoredStorageMetrics,
53 ObjectStoreMetrics,
54};
55use crate::opts::StorageOpts;
56
57fn build_file_cache_spawner(
58 name: &str,
59 config: &FileCacheRuntimeConfig,
60) -> StorageResult<foyer::Spawner> {
61 let config = match config {
62 FileCacheRuntimeConfig::Disabled => return Ok(foyer::Spawner::current()),
63 FileCacheRuntimeConfig::Unified(config) => config.clone(),
64 FileCacheRuntimeConfig::Separated { .. } => {
65 return Err(HummockError::other(format!(
66 "{name} runtime_config.Separated is unsupported with Foyer 0.22; use runtime_config.Unified instead"
67 ))
68 .into());
69 }
70 };
71
72 let mut builder = tokio::runtime::Builder::new_multi_thread();
73 #[cfg(madsim)]
74 let _ = &config;
75 #[cfg(not(madsim))]
76 if config.worker_threads != 0 {
77 builder.worker_threads(config.worker_threads);
78 }
79 #[cfg(not(madsim))]
80 if config.max_blocking_threads != 0 {
81 builder.max_blocking_threads(config.max_blocking_threads);
82 }
83 builder.thread_name(format!("{name}-unified"));
84 let runtime = builder.enable_all().build().map_err(|error| {
85 HummockError::other(format!(
86 "failed to build {name} dedicated runtime: {}",
87 error.as_report()
88 ))
89 })?;
90 Ok(runtime.into())
91}
92
93#[cfg(test)]
94mod tests {
95 use risingwave_common::config::storage::FileCacheTokioRuntimeConfig;
96
97 use super::*;
98
99 #[test]
100 fn test_build_file_cache_spawner_rejects_separated_runtime() {
101 let runtime_options = FileCacheTokioRuntimeConfig {
102 worker_threads: 1,
103 max_blocking_threads: 1,
104 };
105 let error = build_file_cache_spawner(
106 "foyer.test",
107 &FileCacheRuntimeConfig::Separated {
108 read_runtime_options: runtime_options.clone(),
109 write_runtime_options: runtime_options,
110 },
111 )
112 .unwrap_err()
113 .to_string();
114
115 assert!(error.contains("foyer.test runtime_config.Separated"));
116 assert!(error.contains("use runtime_config.Unified instead"));
117 }
118
119 #[cfg(not(madsim))]
120 #[tokio::test]
121 async fn test_build_file_cache_spawner_runtime_modes() {
122 let current_runtime_id = tokio::runtime::Handle::current().id();
123
124 let disabled =
125 build_file_cache_spawner("foyer.test", &FileCacheRuntimeConfig::Disabled).unwrap();
126 let disabled_runtime_id = disabled
127 .spawn(async { tokio::runtime::Handle::current().id() })
128 .await
129 .unwrap();
130 assert_eq!(disabled_runtime_id, current_runtime_id);
131
132 let unified = build_file_cache_spawner(
133 "foyer.test",
134 &FileCacheRuntimeConfig::Unified(FileCacheTokioRuntimeConfig {
135 worker_threads: 1,
136 max_blocking_threads: 1,
137 }),
138 )
139 .unwrap();
140 let (unified_runtime_id, thread_name) = unified
141 .spawn(async {
142 (
143 tokio::runtime::Handle::current().id(),
144 std::thread::current().name().map(str::to_owned),
145 )
146 })
147 .await
148 .unwrap();
149
150 assert_ne!(unified_runtime_id, current_runtime_id);
151 assert_eq!(thread_name.as_deref(), Some("foyer.test-unified"));
152 }
153}
154
155static FOYER_METRICS_REGISTRY: LazyLock<Box<PrometheusMetricsRegistry>> = LazyLock::new(|| {
156 Box::new(PrometheusMetricsRegistry::new(
157 GLOBAL_METRICS_REGISTRY.clone(),
158 ))
159});
160
161mod opaque_type {
162 use super::*;
163
164 pub type HummockStorageType = impl StateStore + AsHummock;
165 pub type MemoryStateStoreType = impl StateStore + AsHummock;
166 pub type SledStateStoreType = impl StateStore + AsHummock;
167
168 #[define_opaque(MemoryStateStoreType)]
169 pub fn in_memory(state_store: MemoryStateStore) -> MemoryStateStoreType {
170 may_dynamic_dispatch(state_store)
171 }
172
173 #[define_opaque(HummockStorageType)]
174 pub fn hummock(state_store: HummockStorage) -> HummockStorageType {
175 may_dynamic_dispatch(may_verify(state_store))
176 }
177
178 #[define_opaque(SledStateStoreType)]
179 pub fn sled(state_store: SledStateStore) -> SledStateStoreType {
180 may_dynamic_dispatch(state_store)
181 }
182}
183pub use opaque_type::{HummockStorageType, MemoryStateStoreType, SledStateStoreType};
184use opaque_type::{hummock, in_memory, sled};
185
186#[cfg(feature = "hm-trace")]
187type Monitored<S> = MonitoredStateStore<crate::monitor::traced_store::TracedStateStore<S>>;
188
189#[cfg(not(feature = "hm-trace"))]
190type Monitored<S> = MonitoredStateStore<S>;
191
192fn monitored<S: StateStore>(
193 state_store: S,
194 storage_metrics: Arc<MonitoredStorageMetrics>,
195) -> Monitored<S> {
196 let inner = {
197 #[cfg(feature = "hm-trace")]
198 {
199 crate::monitor::traced_store::TracedStateStore::new_global(state_store)
200 }
201 #[cfg(not(feature = "hm-trace"))]
202 {
203 state_store
204 }
205 };
206 inner.monitored(storage_metrics)
207}
208
209fn inner<S>(state_store: &Monitored<S>) -> &S {
210 let inner = state_store.inner();
211 {
212 #[cfg(feature = "hm-trace")]
213 {
214 inner.inner()
215 }
216 #[cfg(not(feature = "hm-trace"))]
217 {
218 inner
219 }
220 }
221}
222
223#[derive(Clone, EnumAsInner)]
225#[expect(clippy::enum_variant_names)]
226pub enum StateStoreImpl {
227 HummockStateStore(Monitored<HummockStorageType>),
236 MemoryStateStore(Monitored<MemoryStateStoreType>),
241 SledStateStore(Monitored<SledStateStoreType>),
242}
243
244fn may_dynamic_dispatch(state_store: impl StateStore + AsHummock) -> impl StateStore + AsHummock {
245 #[cfg(not(debug_assertions))]
246 {
247 state_store
248 }
249 #[cfg(debug_assertions)]
250 {
251 use crate::store_impl::dyn_state_store::StateStorePointer;
252 StateStorePointer(Arc::new(state_store) as _)
253 }
254}
255
256fn may_verify(state_store: impl StateStore + AsHummock) -> impl StateStore + AsHummock {
257 #[cfg(not(debug_assertions))]
258 {
259 state_store
260 }
261 #[cfg(debug_assertions)]
262 {
263 use std::marker::PhantomData;
264
265 use risingwave_common::util::env_var::env_var_is_true;
266 use tracing::info;
267
268 use crate::store_impl::verify::VerifyStateStore;
269
270 let expected = if env_var_is_true("ENABLE_STATE_STORE_VERIFY") {
271 info!("enable verify state store");
272 Some(SledStateStore::new_temp())
273 } else {
274 info!("verify state store is not enabled");
275 None
276 };
277 VerifyStateStore {
278 actual: state_store,
279 expected,
280 _phantom: PhantomData::<()>,
281 }
282 }
283}
284
285impl StateStoreImpl {
286 fn in_memory(
287 state_store: MemoryStateStore,
288 storage_metrics: Arc<MonitoredStorageMetrics>,
289 ) -> Self {
290 Self::MemoryStateStore(monitored(in_memory(state_store), storage_metrics))
292 }
293
294 pub fn hummock(
295 state_store: HummockStorage,
296 storage_metrics: Arc<MonitoredStorageMetrics>,
297 ) -> Self {
298 Self::HummockStateStore(monitored(hummock(state_store), storage_metrics))
300 }
301
302 pub fn sled(
303 state_store: SledStateStore,
304 storage_metrics: Arc<MonitoredStorageMetrics>,
305 ) -> Self {
306 Self::SledStateStore(monitored(sled(state_store), storage_metrics))
307 }
308
309 pub fn shared_in_memory_store(storage_metrics: Arc<MonitoredStorageMetrics>) -> Self {
310 Self::in_memory(MemoryStateStore::shared(), storage_metrics)
311 }
312
313 pub fn for_test() -> Self {
314 Self::in_memory(
315 MemoryStateStore::new(),
316 Arc::new(MonitoredStorageMetrics::unused()),
317 )
318 }
319
320 pub fn as_hummock(&self) -> Option<&HummockStorage> {
321 match self {
322 StateStoreImpl::HummockStateStore(hummock) => {
323 Some(inner(hummock).as_hummock().expect("should be hummock"))
324 }
325 _ => None,
326 }
327 }
328}
329
330impl Debug for StateStoreImpl {
331 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
332 match self {
333 StateStoreImpl::HummockStateStore(_) => write!(f, "HummockStateStore"),
334 StateStoreImpl::MemoryStateStore(_) => write!(f, "MemoryStateStore"),
335 StateStoreImpl::SledStateStore(_) => write!(f, "SledStateStore"),
336 }
337 }
338}
339
340#[macro_export]
341macro_rules! dispatch_state_store {
342 ($impl:expr, $store:ident, $body:tt) => {{
343 use $crate::store_impl::StateStoreImpl;
344
345 match $impl {
346 StateStoreImpl::MemoryStateStore($store) => {
347 #[cfg(debug_assertions)]
350 {
351 $body
352 }
353 #[cfg(not(debug_assertions))]
354 {
355 let _store = $store;
356 unimplemented!("memory state store should never be used in release mode");
357 }
358 }
359
360 StateStoreImpl::SledStateStore($store) => {
361 #[cfg(debug_assertions)]
364 {
365 $body
366 }
367 #[cfg(not(debug_assertions))]
368 {
369 let _store = $store;
370 unimplemented!("sled state store should never be used in release mode");
371 }
372 }
373
374 StateStoreImpl::HummockStateStore($store) => $body,
375 }
376 }};
377}
378
379#[cfg(any(debug_assertions, test, feature = "test"))]
380pub mod verify {
381 use std::fmt::Debug;
382 use std::future::Future;
383 use std::marker::PhantomData;
384 use std::ops::Deref;
385 use std::sync::Arc;
386
387 use bytes::Bytes;
388 use risingwave_common::array::VectorRef;
389 use risingwave_common::bitmap::Bitmap;
390 use risingwave_common::hash::VirtualNode;
391 use risingwave_hummock_sdk::HummockReadEpoch;
392 use risingwave_hummock_sdk::key::{FullKey, TableKey, TableKeyRange};
393 use tracing::log::warn;
394
395 use crate::error::StorageResult;
396 use crate::hummock::HummockStorage;
397 use crate::store::*;
398 use crate::store_impl::AsHummock;
399
400 #[expect(dead_code)]
401 fn assert_result_eq<Item: PartialEq + Debug, E>(
402 first: &std::result::Result<Item, E>,
403 second: &std::result::Result<Item, E>,
404 ) {
405 match (first, second) {
406 (Ok(first), Ok(second)) => {
407 if first != second {
408 warn!("result different: {:?} {:?}", first, second);
409 }
410 assert_eq!(first, second);
411 }
412 (Err(_), Err(_)) => {}
413 _ => {
414 warn!("one success and one failed");
415 panic!("result not equal");
416 }
417 }
418 }
419
420 #[derive(Clone)]
421 pub struct VerifyStateStore<A, E, T = ()> {
422 pub actual: A,
423 pub expected: Option<E>,
424 pub _phantom: PhantomData<T>,
425 }
426
427 impl<A: AsHummock, E: AsHummock> AsHummock for VerifyStateStore<A, E> {
428 fn as_hummock(&self) -> Option<&HummockStorage> {
429 self.actual.as_hummock()
430 }
431 }
432
433 impl<A: StateStoreGet, E: StateStoreGet> StateStoreGet for VerifyStateStore<A, E> {
434 async fn on_key_value<'a, O: Send + 'a>(
435 &'a self,
436 key: TableKey<Bytes>,
437 read_options: ReadOptions,
438 on_key_value_fn: impl KeyValueFn<'a, O>,
439 ) -> StorageResult<Option<O>> {
440 let actual: Option<(FullKey<Bytes>, Bytes)> = self
441 .actual
442 .on_key_value(key.clone(), read_options.clone(), |key, value| {
443 Ok((key.copy_into(), Bytes::copy_from_slice(value)))
444 })
445 .await?;
446 if let Some(expected) = &self.expected {
447 let expected: Option<(FullKey<Bytes>, Bytes)> = expected
448 .on_key_value(key, read_options, |key, value| {
449 Ok((key.copy_into(), Bytes::copy_from_slice(value)))
450 })
451 .await?;
452 assert_eq!(
453 actual
454 .as_ref()
455 .map(|item| (item.0.epoch_with_gap.pure_epoch(), item)),
456 expected
457 .as_ref()
458 .map(|item| (item.0.epoch_with_gap.pure_epoch(), item))
459 );
460 }
461
462 actual
463 .map(|(key, value)| on_key_value_fn(key.to_ref(), value.as_ref()))
464 .transpose()
465 }
466 }
467
468 impl<A: StateStoreReadVector, E: StateStoreReadVector> StateStoreReadVector
469 for VerifyStateStore<A, E>
470 {
471 fn nearest<'a, O: Send + 'a>(
472 &'a self,
473 vec: VectorRef<'a>,
474 options: VectorNearestOptions,
475 on_nearest_item_fn: impl OnNearestItemFn<'a, O>,
476 ) -> impl StorageFuture<'a, Vec<O>> {
477 self.actual.nearest(vec, options, on_nearest_item_fn)
478 }
479 }
480
481 impl<A: StateStoreRead, E: StateStoreRead> StateStoreRead for VerifyStateStore<A, E> {
482 type Iter = impl StateStoreReadIter;
483 type RevIter = impl StateStoreReadIter;
484
485 #[expect(clippy::manual_async_fn)]
488 fn iter(
489 &self,
490 key_range: TableKeyRange,
491
492 read_options: ReadOptions,
493 ) -> impl Future<Output = StorageResult<Self::Iter>> + '_ {
494 async move {
495 let actual = self
496 .actual
497 .iter(key_range.clone(), read_options.clone())
498 .await?;
499 let expected = if let Some(expected) = &self.expected {
500 Some(expected.iter(key_range, read_options).await?)
501 } else {
502 None
503 };
504
505 Ok(verify_iter::<StateStoreKeyedRow>(actual, expected))
506 }
507 }
508
509 #[expect(clippy::manual_async_fn)]
510 fn rev_iter(
511 &self,
512 key_range: TableKeyRange,
513
514 read_options: ReadOptions,
515 ) -> impl Future<Output = StorageResult<Self::RevIter>> + '_ {
516 async move {
517 let actual = self
518 .actual
519 .rev_iter(key_range.clone(), read_options.clone())
520 .await?;
521 let expected = if let Some(expected) = &self.expected {
522 Some(expected.rev_iter(key_range, read_options).await?)
523 } else {
524 None
525 };
526
527 Ok(verify_iter::<StateStoreKeyedRow>(actual, expected))
528 }
529 }
530 }
531
532 impl<A: StateStoreReadLog, E: StateStoreReadLog> StateStoreReadLog for VerifyStateStore<A, E> {
533 type ChangeLogIter = impl StateStoreReadChangeLogIter;
534
535 async fn next_epoch(&self, epoch: u64, options: NextEpochOptions) -> StorageResult<u64> {
536 let actual = self.actual.next_epoch(epoch, options.clone()).await?;
537 if let Some(expected) = &self.expected {
538 assert_eq!(actual, expected.next_epoch(epoch, options).await?);
539 }
540 Ok(actual)
541 }
542
543 async fn iter_log(
544 &self,
545 epoch_range: (u64, u64),
546 key_range: TableKeyRange,
547 options: ReadLogOptions,
548 ) -> StorageResult<Self::ChangeLogIter> {
549 let actual = self
550 .actual
551 .iter_log(epoch_range, key_range.clone(), options.clone())
552 .await?;
553 let expected = if let Some(expected) = &self.expected {
554 Some(expected.iter_log(epoch_range, key_range, options).await?)
555 } else {
556 None
557 };
558
559 Ok(verify_iter::<StateStoreReadLogItem>(actual, expected))
560 }
561 }
562
563 impl<A: StateStoreIter<T>, E: StateStoreIter<T>, T: IterItem> StateStoreIter<T>
564 for VerifyStateStore<A, E, T>
565 where
566 for<'a> T::ItemRef<'a>: PartialEq + Debug,
567 {
568 async fn try_next(&mut self) -> StorageResult<Option<T::ItemRef<'_>>> {
569 let actual = self.actual.try_next().await?;
570 if let Some(expected) = self.expected.as_mut() {
571 let expected = expected.try_next().await?;
572 assert_eq!(actual, expected);
573 }
574 Ok(actual)
575 }
576 }
577
578 fn verify_iter<T: IterItem>(
579 actual: impl StateStoreIter<T>,
580 expected: Option<impl StateStoreIter<T>>,
581 ) -> impl StateStoreIter<T>
582 where
583 for<'a> T::ItemRef<'a>: PartialEq + Debug,
584 {
585 VerifyStateStore {
586 actual,
587 expected,
588 _phantom: PhantomData::<T>,
589 }
590 }
591
592 impl<A: LocalStateStore, E: LocalStateStore> LocalStateStore for VerifyStateStore<A, E> {
593 type FlushedSnapshotReader =
594 VerifyStateStore<A::FlushedSnapshotReader, E::FlushedSnapshotReader>;
595
596 type Iter<'a> = impl StateStoreIter + 'a;
597 type RevIter<'a> = impl StateStoreIter + 'a;
598
599 #[expect(clippy::manual_async_fn)]
600 fn iter(
601 &self,
602 key_range: TableKeyRange,
603 read_options: ReadOptions,
604 ) -> impl Future<Output = StorageResult<Self::Iter<'_>>> + Send + '_ {
605 async move {
606 let actual = self
607 .actual
608 .iter(key_range.clone(), read_options.clone())
609 .await?;
610 let expected = if let Some(expected) = &self.expected {
611 Some(expected.iter(key_range, read_options).await?)
612 } else {
613 None
614 };
615
616 Ok(verify_iter::<StateStoreKeyedRow>(actual, expected))
617 }
618 }
619
620 #[expect(clippy::manual_async_fn)]
621 fn rev_iter(
622 &self,
623 key_range: TableKeyRange,
624 read_options: ReadOptions,
625 ) -> impl Future<Output = StorageResult<Self::RevIter<'_>>> + Send + '_ {
626 async move {
627 let actual = self
628 .actual
629 .rev_iter(key_range.clone(), read_options.clone())
630 .await?;
631 let expected = if let Some(expected) = &self.expected {
632 Some(expected.rev_iter(key_range, read_options).await?)
633 } else {
634 None
635 };
636
637 Ok(verify_iter::<StateStoreKeyedRow>(actual, expected))
638 }
639 }
640
641 fn insert(
642 &mut self,
643 key: TableKey<Bytes>,
644 new_val: Bytes,
645 old_val: Option<Bytes>,
646 ) -> StorageResult<()> {
647 if let Some(expected) = &mut self.expected {
648 expected.insert(key.clone(), new_val.clone(), old_val.clone())?;
649 }
650 self.actual.insert(key, new_val, old_val)?;
651
652 Ok(())
653 }
654
655 fn delete(&mut self, key: TableKey<Bytes>, old_val: Bytes) -> StorageResult<()> {
656 if let Some(expected) = &mut self.expected {
657 expected.delete(key.clone(), old_val.clone())?;
658 }
659 self.actual.delete(key, old_val)?;
660 Ok(())
661 }
662
663 async fn update_vnode_bitmap(&mut self, vnodes: Arc<Bitmap>) -> StorageResult<Arc<Bitmap>> {
664 let ret = self.actual.update_vnode_bitmap(vnodes.clone()).await?;
665 if let Some(expected) = &mut self.expected {
666 assert_eq!(ret, expected.update_vnode_bitmap(vnodes).await?);
667 }
668 Ok(ret)
669 }
670
671 fn get_table_watermark(&self, vnode: VirtualNode) -> Option<Bytes> {
672 let ret = self.actual.get_table_watermark(vnode);
673 if let Some(expected) = &self.expected {
674 assert_eq!(ret, expected.get_table_watermark(vnode));
675 }
676 ret
677 }
678
679 fn new_flushed_snapshot_reader(&self) -> Self::FlushedSnapshotReader {
680 VerifyStateStore {
681 actual: self.actual.new_flushed_snapshot_reader(),
682 expected: self.expected.as_ref().map(E::new_flushed_snapshot_reader),
683 _phantom: Default::default(),
684 }
685 }
686 }
687
688 impl<A: StateStoreWriteEpochControl, E: StateStoreWriteEpochControl> StateStoreWriteEpochControl
689 for VerifyStateStore<A, E>
690 {
691 async fn flush(&mut self) -> StorageResult<usize> {
692 if let Some(expected) = &mut self.expected {
693 expected.flush().await?;
694 }
695 self.actual.flush().await
696 }
697
698 async fn try_flush(&mut self) -> StorageResult<()> {
699 if let Some(expected) = &mut self.expected {
700 expected.try_flush().await?;
701 }
702 self.actual.try_flush().await
703 }
704
705 async fn init(&mut self, options: InitOptions) -> StorageResult<()> {
706 self.actual.init(options.clone()).await?;
707 if let Some(expected) = &mut self.expected {
708 expected.init(options).await?;
709 }
710 Ok(())
711 }
712
713 fn seal_current_epoch(&mut self, next_epoch: u64, opts: SealCurrentEpochOptions) {
714 if let Some(expected) = &mut self.expected {
715 expected.seal_current_epoch(next_epoch, opts.clone());
716 }
717 self.actual.seal_current_epoch(next_epoch, opts);
718 }
719 }
720
721 impl<A: StateStore, E: StateStore> StateStore for VerifyStateStore<A, E> {
722 type Local = VerifyStateStore<A::Local, E::Local>;
723 type ReadSnapshot = VerifyStateStore<A::ReadSnapshot, E::ReadSnapshot>;
724 type VectorWriter = A::VectorWriter;
725
726 fn try_wait_epoch(
727 &self,
728 epoch: HummockReadEpoch,
729 options: TryWaitEpochOptions,
730 ) -> impl Future<Output = StorageResult<()>> + Send + '_ {
731 self.actual.try_wait_epoch(epoch, options)
732 }
733
734 async fn new_local(&self, option: NewLocalOptions) -> Self::Local {
735 let expected = if let Some(expected) = &self.expected {
736 Some(expected.new_local(option.clone()).await)
737 } else {
738 None
739 };
740 VerifyStateStore {
741 actual: self.actual.new_local(option).await,
742 expected,
743 _phantom: PhantomData::<()>,
744 }
745 }
746
747 async fn new_read_snapshot(
748 &self,
749 epoch: HummockReadEpoch,
750 options: NewReadSnapshotOptions,
751 ) -> StorageResult<Self::ReadSnapshot> {
752 let expected = if let Some(expected) = &self.expected {
753 Some(expected.new_read_snapshot(epoch, options).await?)
754 } else {
755 None
756 };
757 Ok(VerifyStateStore {
758 actual: self.actual.new_read_snapshot(epoch, options).await?,
759 expected,
760 _phantom: PhantomData::<()>,
761 })
762 }
763
764 fn new_vector_writer(
765 &self,
766 options: NewVectorWriterOptions,
767 ) -> impl Future<Output = Self::VectorWriter> + Send + '_ {
768 self.actual.new_vector_writer(options)
769 }
770 }
771
772 impl<A, E> Deref for VerifyStateStore<A, E> {
773 type Target = A;
774
775 fn deref(&self) -> &Self::Target {
776 &self.actual
777 }
778 }
779}
780
781impl StateStoreImpl {
782 #[cfg_attr(not(target_os = "linux"), allow(unused_variables))]
783 #[expect(clippy::borrowed_box)]
784 pub async fn new(
785 s: &str,
786 role: Role,
787 opts: Arc<StorageOpts>,
788 hummock_meta_client: Arc<MonitoredHummockMetaClient>,
789 state_store_metrics: Arc<HummockStateStoreMetrics>,
790 object_store_metrics: Arc<ObjectStoreMetrics>,
791 storage_metrics: Arc<MonitoredStorageMetrics>,
792 compactor_metrics: Arc<CompactorMetrics>,
793 await_tree_config: Option<await_tree::Config>,
794 use_new_object_prefix_strategy: bool,
795 ) -> StorageResult<Self> {
796 const KB: usize = 1 << 10;
797 const MB: usize = 1 << 20;
798
799 let meta_cache = {
800 let mut builder = HybridCacheBuilder::new()
801 .with_name("foyer.meta")
802 .with_metrics_registry(FOYER_METRICS_REGISTRY.clone())
803 .memory(opts.meta_cache_capacity_mb * MB)
804 .with_shards(opts.meta_cache_shard_num)
805 .with_eviction_config(opts.meta_cache_eviction_config.clone())
806 .with_weighter(|_: &HummockSstableObjectId, value: &Box<Sstable>| {
807 std::mem::size_of::<HummockSstableObjectId>()
808 + value.estimated_meta_cache_memory_weight()
809 })
810 .storage();
811
812 if !opts.meta_file_cache_dir.is_empty() {
813 if let Err(e) = Feature::ElasticDiskCache.check_available() {
814 tracing::warn!(error = %e.as_report(), "ElasticDiskCache is not available.");
815 } else {
816 let device_builder = FsDeviceBuilder::new(&opts.meta_file_cache_dir)
817 .with_capacity(opts.meta_file_cache_capacity_mb * MB)
818 .with_throttle(opts.meta_file_cache_throttle.clone());
819 #[cfg(target_os = "linux")]
820 let device_builder = device_builder.with_direct(opts.meta_file_cache_direct_io);
821 let device = device_builder.build().map_err(HummockError::foyer_error)?;
822 let engine_builder = BlockEngineConfig::new(device)
823 .with_block_size(opts.meta_file_cache_file_capacity_mb * MB)
824 .with_indexer_shards(opts.meta_file_cache_indexer_shards)
825 .with_flushers(opts.meta_file_cache_flushers)
826 .with_reclaimers(opts.meta_file_cache_reclaimers)
827 .with_buffer_pool_size(opts.meta_file_cache_flush_buffer_threshold_mb * MB)
828 .with_submit_queue_size_threshold(
829 opts.meta_file_cache_submit_queue_size_threshold_mb * MB,
830 )
831 .with_clean_block_threshold(
832 opts.meta_file_cache_reclaimers + opts.meta_file_cache_reclaimers / 2,
833 )
834 .with_recover_concurrency(opts.meta_file_cache_recover_concurrency)
835 .with_blob_index_size(opts.meta_file_cache_blob_index_size_kb * KB)
836 .with_eviction_pickers(vec![Box::new(FifoPicker::new(
837 opts.meta_file_cache_fifo_probation_ratio,
838 ))]);
839 builder = builder
840 .with_engine_config(engine_builder)
841 .with_recover_mode(opts.meta_file_cache_recover_mode)
842 .with_compression(opts.meta_file_cache_compression);
843 builder = builder.with_spawner(build_file_cache_spawner(
844 "foyer.meta",
845 &opts.meta_file_cache_runtime_config,
846 )?);
847 }
848 }
849
850 builder.build().await.map_err(HummockError::foyer_error)?
851 };
852
853 let block_cache = {
854 let mut builder = HybridCacheBuilder::new()
855 .with_name("foyer.data")
856 .with_metrics_registry(FOYER_METRICS_REGISTRY.clone())
857 .with_event_listener(Arc::new(BlockCacheEventListener::new(
858 state_store_metrics.clone(),
859 )))
860 .memory(opts.block_cache_capacity_mb * MB)
861 .with_shards(opts.block_cache_shard_num)
862 .with_eviction_config(opts.block_cache_eviction_config.clone())
863 .with_weighter(|_: &SstableBlockIndex, value: &Box<Block>| {
864 std::mem::size_of::<SstableBlockIndex>() + value.estimated_memory_weight()
865 })
866 .storage();
867
868 if !opts.data_file_cache_dir.is_empty() {
869 if let Err(e) = Feature::ElasticDiskCache.check_available() {
870 tracing::warn!(error = %e.as_report(), "ElasticDiskCache is not available.");
871 } else {
872 let device_builder = FsDeviceBuilder::new(&opts.data_file_cache_dir)
873 .with_capacity(opts.data_file_cache_capacity_mb * MB)
874 .with_throttle(opts.data_file_cache_throttle.clone());
875 #[cfg(target_os = "linux")]
876 let device_builder = device_builder.with_direct(opts.data_file_cache_direct_io);
877 let device = device_builder.build().map_err(HummockError::foyer_error)?;
878 let engine_builder = BlockEngineConfig::new(device)
879 .with_block_size(opts.data_file_cache_file_capacity_mb * MB)
880 .with_indexer_shards(opts.data_file_cache_indexer_shards)
881 .with_flushers(opts.data_file_cache_flushers)
882 .with_reclaimers(opts.data_file_cache_reclaimers)
883 .with_buffer_pool_size(opts.data_file_cache_flush_buffer_threshold_mb * MB)
884 .with_submit_queue_size_threshold(
885 opts.data_file_cache_submit_queue_size_threshold_mb * MB,
886 )
887 .with_clean_block_threshold(
888 opts.data_file_cache_reclaimers + opts.data_file_cache_reclaimers / 2,
889 )
890 .with_recover_concurrency(opts.data_file_cache_recover_concurrency)
891 .with_blob_index_size(opts.data_file_cache_blob_index_size_kb * KB)
892 .with_eviction_pickers(vec![Box::new(FifoPicker::new(
893 opts.data_file_cache_fifo_probation_ratio,
894 ))]);
895 builder = builder
896 .with_engine_config(engine_builder)
897 .with_recover_mode(opts.data_file_cache_recover_mode)
898 .with_compression(opts.data_file_cache_compression);
899 builder = builder.with_spawner(build_file_cache_spawner(
900 "foyer.data",
901 &opts.data_file_cache_runtime_config,
902 )?);
903 }
904 }
905
906 builder.build().await.map_err(HummockError::foyer_error)?
907 };
908
909 let vector_meta_cache = CacheBuilder::new(opts.vector_meta_cache_capacity_mb * MB)
910 .with_shards(opts.vector_meta_cache_shard_num)
911 .with_eviction_config(opts.vector_meta_cache_eviction_config.clone())
912 .build();
913
914 let vector_block_cache = CacheBuilder::new(opts.vector_block_cache_capacity_mb * MB)
915 .with_shards(opts.vector_block_cache_shard_num)
916 .with_eviction_config(opts.vector_block_cache_eviction_config.clone())
917 .build();
918
919 let recent_filter = if opts.data_file_cache_dir.is_empty() {
920 Arc::new(NoneRecentFilter::default().into())
921 } else if opts.cache_refill_recent_filter_shards == 1 {
922 Arc::new(
923 SimpleRecentFilter::new(
924 opts.cache_refill_recent_filter_layers,
925 Duration::from_millis(
926 opts.cache_refill_recent_filter_rotate_interval_ms as u64,
927 ),
928 )
929 .into(),
930 )
931 } else if opts.cache_refill_skip_recent_filter {
932 Arc::new(AllRecentFilter::default().into())
933 } else {
934 Arc::new(
935 ShardedRecentFilter::new(
936 opts.cache_refill_recent_filter_layers,
937 Duration::from_millis(
938 opts.cache_refill_recent_filter_rotate_interval_ms as u64,
939 ),
940 opts.cache_refill_recent_filter_shards,
941 )
942 .into(),
943 )
944 };
945
946 let store = match s {
947 hummock if hummock.starts_with("hummock+") => {
948 let object_store = build_remote_object_store(
949 hummock.strip_prefix("hummock+").unwrap(),
950 object_store_metrics.clone(),
951 "Hummock",
952 Arc::new(opts.object_store_config.clone()),
953 )
954 .await;
955
956 let sstable_store = Arc::new(SstableStore::new(SstableStoreConfig {
957 store: Arc::new(object_store),
958 path: opts.data_directory.clone(),
959 prefetch_buffer_capacity: opts.prefetch_buffer_capacity_mb * (1 << 20),
960 max_prefetch_block_number: opts.max_prefetch_block_number,
961 recent_filter,
962 state_store_metrics: state_store_metrics.clone(),
963 use_new_object_prefix_strategy,
964 skip_bloom_filter_in_serde: opts.sst_skip_bloom_filter_in_serde,
965
966 meta_cache,
967 block_cache,
968 vector_meta_cache,
969 vector_block_cache,
970 }));
971 let notification_client =
972 RpcNotificationClient::new(hummock_meta_client.get_inner().clone());
973 let compaction_catalog_manager_ref =
974 Arc::new(CompactionCatalogManager::new(Box::new(
975 RemoteTableAccessor::new(hummock_meta_client.get_inner().clone()),
976 )));
977
978 let inner = HummockStorage::new(
979 role,
980 opts.clone(),
981 sstable_store,
982 hummock_meta_client.clone(),
983 notification_client,
984 compaction_catalog_manager_ref,
985 state_store_metrics.clone(),
986 compactor_metrics.clone(),
987 await_tree_config,
988 )
989 .await?;
990
991 StateStoreImpl::hummock(inner, storage_metrics)
992 }
993
994 "in_memory" | "in-memory" => {
995 tracing::warn!(
996 "In-memory state store should never be used in end-to-end benchmarks or production environment. Scaling and recovery are not supported."
997 );
998 StateStoreImpl::shared_in_memory_store(storage_metrics.clone())
999 }
1000
1001 sled if sled.starts_with("sled://") => {
1002 tracing::warn!(
1003 "sled state store should never be used in end-to-end benchmarks or production environment. Scaling and recovery are not supported."
1004 );
1005 let path = sled.strip_prefix("sled://").unwrap();
1006 StateStoreImpl::sled(SledStateStore::new(path), storage_metrics.clone())
1007 }
1008
1009 other => unimplemented!("{} state store is not supported", other),
1010 };
1011
1012 Ok(store)
1013 }
1014}
1015
1016pub trait AsHummock: Send + Sync {
1017 fn as_hummock(&self) -> Option<&HummockStorage>;
1018
1019 fn sync(
1020 &self,
1021 sync_table_epochs: Vec<(HummockEpoch, HashSet<TableId>)>,
1022 ) -> BoxFuture<'_, StorageResult<SyncResult>> {
1023 async move {
1024 if let Some(hummock) = self.as_hummock() {
1025 hummock.sync(sync_table_epochs).await
1026 } else {
1027 Ok(SyncResult::default())
1028 }
1029 }
1030 .boxed()
1031 }
1032}
1033
1034impl AsHummock for HummockStorage {
1035 fn as_hummock(&self) -> Option<&HummockStorage> {
1036 Some(self)
1037 }
1038}
1039
1040impl AsHummock for MemoryStateStore {
1041 fn as_hummock(&self) -> Option<&HummockStorage> {
1042 None
1043 }
1044}
1045
1046impl AsHummock for SledStateStore {
1047 fn as_hummock(&self) -> Option<&HummockStorage> {
1048 None
1049 }
1050}
1051
1052#[cfg(debug_assertions)]
1053mod dyn_state_store {
1054 use std::future::Future;
1055 use std::ops::DerefMut;
1056 use std::sync::Arc;
1057
1058 use bytes::Bytes;
1059 use risingwave_common::array::VectorRef;
1060 use risingwave_common::bitmap::Bitmap;
1061 use risingwave_common::hash::VirtualNode;
1062 use risingwave_hummock_sdk::HummockReadEpoch;
1063 use risingwave_hummock_sdk::key::{TableKey, TableKeyRange};
1064
1065 use crate::error::StorageResult;
1066 use crate::hummock::HummockStorage;
1067 use crate::store::*;
1068 use crate::store_impl::AsHummock;
1069 use crate::vector::VectorDistance;
1070
1071 #[async_trait::async_trait]
1072 pub trait DynStateStoreIter<T: IterItem>: Send {
1073 async fn try_next(&mut self) -> StorageResult<Option<T::ItemRef<'_>>>;
1074 }
1075
1076 #[async_trait::async_trait]
1077 impl<T: IterItem, I: StateStoreIter<T>> DynStateStoreIter<T> for I {
1078 async fn try_next(&mut self) -> StorageResult<Option<T::ItemRef<'_>>> {
1079 self.try_next().await
1080 }
1081 }
1082
1083 pub type BoxStateStoreIter<'a, T> = Box<dyn DynStateStoreIter<T> + 'a>;
1084 impl<T: IterItem> StateStoreIter<T> for BoxStateStoreIter<'_, T> {
1085 fn try_next(
1086 &mut self,
1087 ) -> impl Future<Output = StorageResult<Option<T::ItemRef<'_>>>> + Send + '_ {
1088 self.deref_mut().try_next()
1089 }
1090 }
1091
1092 pub type BoxStateStoreReadIter = BoxStateStoreIter<'static, StateStoreKeyedRow>;
1095 pub type BoxStateStoreReadChangeLogIter = BoxStateStoreIter<'static, StateStoreReadLogItem>;
1096
1097 #[async_trait::async_trait]
1098 pub trait DynStateStoreGet: StaticSendSync {
1099 async fn get_keyed_row(
1100 &self,
1101 key: TableKey<Bytes>,
1102 read_options: ReadOptions,
1103 ) -> StorageResult<Option<StateStoreKeyedRow>>;
1104 }
1105
1106 #[async_trait::async_trait]
1107 pub trait DynStateStoreRead: DynStateStoreGet + StaticSendSync {
1108 async fn iter(
1109 &self,
1110 key_range: TableKeyRange,
1111
1112 read_options: ReadOptions,
1113 ) -> StorageResult<BoxStateStoreReadIter>;
1114
1115 async fn rev_iter(
1116 &self,
1117 key_range: TableKeyRange,
1118
1119 read_options: ReadOptions,
1120 ) -> StorageResult<BoxStateStoreReadIter>;
1121 }
1122
1123 #[async_trait::async_trait]
1124 pub trait DynStateStoreReadLog: StaticSendSync {
1125 async fn next_epoch(&self, epoch: u64, options: NextEpochOptions) -> StorageResult<u64>;
1126 async fn iter_log(
1127 &self,
1128 epoch_range: (u64, u64),
1129 key_range: TableKeyRange,
1130 options: ReadLogOptions,
1131 ) -> StorageResult<BoxStateStoreReadChangeLogIter>;
1132 }
1133
1134 pub type StateStoreReadDynRef = StateStorePointer<Arc<dyn DynStateStoreRead>>;
1135
1136 #[async_trait::async_trait]
1137 impl<S: StateStoreGet> DynStateStoreGet for S {
1138 async fn get_keyed_row(
1139 &self,
1140 key: TableKey<Bytes>,
1141 read_options: ReadOptions,
1142 ) -> StorageResult<Option<StateStoreKeyedRow>> {
1143 self.on_key_value(key, read_options, move |key, value| {
1144 Ok((key.copy_into(), Bytes::copy_from_slice(value)))
1145 })
1146 .await
1147 }
1148 }
1149
1150 #[async_trait::async_trait]
1151 impl<S: StateStoreRead> DynStateStoreRead for S {
1152 async fn iter(
1153 &self,
1154 key_range: TableKeyRange,
1155
1156 read_options: ReadOptions,
1157 ) -> StorageResult<BoxStateStoreReadIter> {
1158 Ok(Box::new(self.iter(key_range, read_options).await?))
1159 }
1160
1161 async fn rev_iter(
1162 &self,
1163 key_range: TableKeyRange,
1164
1165 read_options: ReadOptions,
1166 ) -> StorageResult<BoxStateStoreReadIter> {
1167 Ok(Box::new(self.rev_iter(key_range, read_options).await?))
1168 }
1169 }
1170
1171 #[async_trait::async_trait]
1172 impl<S: StateStoreReadLog> DynStateStoreReadLog for S {
1173 async fn next_epoch(&self, epoch: u64, options: NextEpochOptions) -> StorageResult<u64> {
1174 self.next_epoch(epoch, options).await
1175 }
1176
1177 async fn iter_log(
1178 &self,
1179 epoch_range: (u64, u64),
1180 key_range: TableKeyRange,
1181 options: ReadLogOptions,
1182 ) -> StorageResult<BoxStateStoreReadChangeLogIter> {
1183 Ok(Box::new(
1184 self.iter_log(epoch_range, key_range, options).await?,
1185 ))
1186 }
1187 }
1188
1189 pub type BoxLocalStateStoreIterStream<'a> = BoxStateStoreIter<'a, StateStoreKeyedRow>;
1191 #[async_trait::async_trait]
1192 pub trait DynLocalStateStore:
1193 DynStateStoreGet + DynStateStoreWriteEpochControl + StaticSendSync
1194 {
1195 async fn iter(
1196 &self,
1197 key_range: TableKeyRange,
1198 read_options: ReadOptions,
1199 ) -> StorageResult<BoxLocalStateStoreIterStream<'_>>;
1200
1201 async fn rev_iter(
1202 &self,
1203 key_range: TableKeyRange,
1204 read_options: ReadOptions,
1205 ) -> StorageResult<BoxLocalStateStoreIterStream<'_>>;
1206
1207 fn new_flushed_snapshot_reader(&self) -> StateStoreReadDynRef;
1208
1209 fn insert(
1210 &mut self,
1211 key: TableKey<Bytes>,
1212 new_val: Bytes,
1213 old_val: Option<Bytes>,
1214 ) -> StorageResult<()>;
1215
1216 fn delete(&mut self, key: TableKey<Bytes>, old_val: Bytes) -> StorageResult<()>;
1217
1218 async fn update_vnode_bitmap(&mut self, vnodes: Arc<Bitmap>) -> StorageResult<Arc<Bitmap>>;
1219
1220 fn get_table_watermark(&self, vnode: VirtualNode) -> Option<Bytes>;
1221 }
1222
1223 #[async_trait::async_trait]
1224 pub trait DynStateStoreWriteEpochControl: StaticSendSync {
1225 async fn flush(&mut self) -> StorageResult<usize>;
1226
1227 async fn try_flush(&mut self) -> StorageResult<()>;
1228
1229 async fn init(&mut self, epoch: InitOptions) -> StorageResult<()>;
1230
1231 fn seal_current_epoch(&mut self, next_epoch: u64, opts: SealCurrentEpochOptions);
1232 }
1233
1234 #[async_trait::async_trait]
1235 impl<S: LocalStateStore> DynLocalStateStore for S {
1236 async fn iter(
1237 &self,
1238 key_range: TableKeyRange,
1239 read_options: ReadOptions,
1240 ) -> StorageResult<BoxLocalStateStoreIterStream<'_>> {
1241 Ok(Box::new(self.iter(key_range, read_options).await?))
1242 }
1243
1244 async fn rev_iter(
1245 &self,
1246 key_range: TableKeyRange,
1247 read_options: ReadOptions,
1248 ) -> StorageResult<BoxLocalStateStoreIterStream<'_>> {
1249 Ok(Box::new(self.rev_iter(key_range, read_options).await?))
1250 }
1251
1252 fn new_flushed_snapshot_reader(&self) -> StateStoreReadDynRef {
1253 StateStorePointer(Arc::new(self.new_flushed_snapshot_reader()) as _)
1254 }
1255
1256 fn insert(
1257 &mut self,
1258 key: TableKey<Bytes>,
1259 new_val: Bytes,
1260 old_val: Option<Bytes>,
1261 ) -> StorageResult<()> {
1262 self.insert(key, new_val, old_val)
1263 }
1264
1265 fn delete(&mut self, key: TableKey<Bytes>, old_val: Bytes) -> StorageResult<()> {
1266 self.delete(key, old_val)
1267 }
1268
1269 async fn update_vnode_bitmap(&mut self, vnodes: Arc<Bitmap>) -> StorageResult<Arc<Bitmap>> {
1270 self.update_vnode_bitmap(vnodes).await
1271 }
1272
1273 fn get_table_watermark(&self, vnode: VirtualNode) -> Option<Bytes> {
1274 self.get_table_watermark(vnode)
1275 }
1276 }
1277
1278 #[async_trait::async_trait]
1279 impl<S: StateStoreWriteEpochControl> DynStateStoreWriteEpochControl for S {
1280 async fn flush(&mut self) -> StorageResult<usize> {
1281 self.flush().await
1282 }
1283
1284 async fn try_flush(&mut self) -> StorageResult<()> {
1285 self.try_flush().await
1286 }
1287
1288 async fn init(&mut self, options: InitOptions) -> StorageResult<()> {
1289 self.init(options).await
1290 }
1291
1292 fn seal_current_epoch(&mut self, next_epoch: u64, opts: SealCurrentEpochOptions) {
1293 self.seal_current_epoch(next_epoch, opts)
1294 }
1295 }
1296
1297 pub type BoxDynLocalStateStore = StateStorePointer<Box<dyn DynLocalStateStore>>;
1298
1299 impl LocalStateStore for BoxDynLocalStateStore {
1300 type FlushedSnapshotReader = StateStoreReadDynRef;
1301 type Iter<'a> = BoxLocalStateStoreIterStream<'a>;
1302 type RevIter<'a> = BoxLocalStateStoreIterStream<'a>;
1303
1304 fn iter(
1305 &self,
1306 key_range: TableKeyRange,
1307 read_options: ReadOptions,
1308 ) -> impl Future<Output = StorageResult<Self::Iter<'_>>> + Send + '_ {
1309 (*self.0).iter(key_range, read_options)
1310 }
1311
1312 fn rev_iter(
1313 &self,
1314 key_range: TableKeyRange,
1315 read_options: ReadOptions,
1316 ) -> impl Future<Output = StorageResult<Self::RevIter<'_>>> + Send + '_ {
1317 (*self.0).rev_iter(key_range, read_options)
1318 }
1319
1320 fn new_flushed_snapshot_reader(&self) -> Self::FlushedSnapshotReader {
1321 (*self.0).new_flushed_snapshot_reader()
1322 }
1323
1324 fn get_table_watermark(&self, vnode: VirtualNode) -> Option<Bytes> {
1325 (*self.0).get_table_watermark(vnode)
1326 }
1327
1328 fn insert(
1329 &mut self,
1330 key: TableKey<Bytes>,
1331 new_val: Bytes,
1332 old_val: Option<Bytes>,
1333 ) -> StorageResult<()> {
1334 (*self.0).insert(key, new_val, old_val)
1335 }
1336
1337 fn delete(&mut self, key: TableKey<Bytes>, old_val: Bytes) -> StorageResult<()> {
1338 (*self.0).delete(key, old_val)
1339 }
1340
1341 async fn update_vnode_bitmap(&mut self, vnodes: Arc<Bitmap>) -> StorageResult<Arc<Bitmap>> {
1342 (*self.0).update_vnode_bitmap(vnodes).await
1343 }
1344 }
1345
1346 impl<P> StateStoreWriteEpochControl for StateStorePointer<P>
1347 where
1348 StateStorePointer<P>: AsMut<dyn DynStateStoreWriteEpochControl> + StaticSendSync,
1349 {
1350 fn flush(&mut self) -> impl Future<Output = StorageResult<usize>> + Send + '_ {
1351 self.as_mut().flush()
1352 }
1353
1354 fn try_flush(&mut self) -> impl Future<Output = StorageResult<()>> + Send + '_ {
1355 self.as_mut().try_flush()
1356 }
1357
1358 fn init(
1359 &mut self,
1360 options: InitOptions,
1361 ) -> impl Future<Output = StorageResult<()>> + Send + '_ {
1362 self.as_mut().init(options)
1363 }
1364
1365 fn seal_current_epoch(&mut self, next_epoch: u64, opts: SealCurrentEpochOptions) {
1366 self.as_mut().seal_current_epoch(next_epoch, opts)
1367 }
1368 }
1369
1370 #[async_trait::async_trait]
1371 pub trait DynStateStoreWriteVector: DynStateStoreWriteEpochControl + StaticSendSync {
1372 fn insert(&mut self, vec: VectorRef<'_>, info: Bytes) -> StorageResult<()>;
1373 }
1374
1375 #[async_trait::async_trait]
1376 impl<S: StateStoreWriteVector> DynStateStoreWriteVector for S {
1377 fn insert(&mut self, vec: VectorRef<'_>, info: Bytes) -> StorageResult<()> {
1378 self.insert(vec, info)
1379 }
1380 }
1381
1382 pub type BoxDynStateStoreWriteVector = StateStorePointer<Box<dyn DynStateStoreWriteVector>>;
1383
1384 impl StateStoreWriteVector for BoxDynStateStoreWriteVector {
1385 fn insert(&mut self, vec: VectorRef<'_>, info: Bytes) -> StorageResult<()> {
1386 self.0.insert(vec, info)
1387 }
1388 }
1389
1390 #[async_trait::async_trait]
1393 pub trait DynStateStoreReadVector: StaticSendSync {
1394 async fn nearest(
1395 &self,
1396 vec: VectorRef<'_>,
1397 options: VectorNearestOptions,
1398 ) -> StorageResult<Vec<(Vector, VectorDistance, Bytes)>>;
1399 }
1400
1401 #[async_trait::async_trait]
1402 impl<S: StateStoreReadVector> DynStateStoreReadVector for S {
1403 async fn nearest(
1404 &self,
1405 vec: VectorRef<'_>,
1406 options: VectorNearestOptions,
1407 ) -> StorageResult<Vec<(Vector, VectorDistance, Bytes)>> {
1408 use risingwave_common::types::ScalarRef;
1409 self.nearest(vec, options, |vec, distance, info| {
1410 (
1411 vec.to_owned_scalar(),
1412 distance,
1413 Bytes::copy_from_slice(info),
1414 )
1415 })
1416 .await
1417 }
1418 }
1419
1420 impl<P> StateStoreReadVector for StateStorePointer<P>
1421 where
1422 StateStorePointer<P>: AsRef<dyn DynStateStoreReadVector> + StaticSendSync,
1423 {
1424 async fn nearest<'a, O: Send + 'a>(
1425 &'a self,
1426 vec: VectorRef<'a>,
1427 options: VectorNearestOptions,
1428 on_nearest_item_fn: impl OnNearestItemFn<'a, O>,
1429 ) -> StorageResult<Vec<O>> {
1430 let output = self.as_ref().nearest(vec, options).await?;
1431 Ok(output
1432 .into_iter()
1433 .map(|(vec, distance, info)| {
1434 on_nearest_item_fn(vec.to_ref(), distance, info.as_ref())
1435 })
1436 .collect())
1437 }
1438 }
1439
1440 pub trait DynStateStoreReadSnapshot:
1441 DynStateStoreRead + DynStateStoreReadVector + StaticSendSync
1442 {
1443 }
1444
1445 impl<S: DynStateStoreRead + DynStateStoreReadVector + StaticSendSync> DynStateStoreReadSnapshot
1446 for S
1447 {
1448 }
1449
1450 pub type StateStoreReadSnapshotDynRef = StateStorePointer<Arc<dyn DynStateStoreReadSnapshot>>;
1451 #[async_trait::async_trait]
1452 pub trait DynStateStoreExt: StaticSendSync {
1453 async fn try_wait_epoch(
1454 &self,
1455 epoch: HummockReadEpoch,
1456 options: TryWaitEpochOptions,
1457 ) -> StorageResult<()>;
1458
1459 async fn new_local(&self, option: NewLocalOptions) -> BoxDynLocalStateStore;
1460 async fn new_read_snapshot(
1461 &self,
1462 epoch: HummockReadEpoch,
1463 options: NewReadSnapshotOptions,
1464 ) -> StorageResult<StateStoreReadSnapshotDynRef>;
1465 async fn new_vector_writer(
1466 &self,
1467 options: NewVectorWriterOptions,
1468 ) -> BoxDynStateStoreWriteVector;
1469 }
1470
1471 #[async_trait::async_trait]
1472 impl<S: StateStore> DynStateStoreExt for S {
1473 async fn try_wait_epoch(
1474 &self,
1475 epoch: HummockReadEpoch,
1476 options: TryWaitEpochOptions,
1477 ) -> StorageResult<()> {
1478 self.try_wait_epoch(epoch, options).await
1479 }
1480
1481 async fn new_local(&self, option: NewLocalOptions) -> BoxDynLocalStateStore {
1482 StateStorePointer(Box::new(self.new_local(option).await))
1483 }
1484
1485 async fn new_read_snapshot(
1486 &self,
1487 epoch: HummockReadEpoch,
1488 options: NewReadSnapshotOptions,
1489 ) -> StorageResult<StateStoreReadSnapshotDynRef> {
1490 Ok(StateStorePointer(Arc::new(
1491 self.new_read_snapshot(epoch, options).await?,
1492 )))
1493 }
1494
1495 async fn new_vector_writer(
1496 &self,
1497 options: NewVectorWriterOptions,
1498 ) -> BoxDynStateStoreWriteVector {
1499 StateStorePointer(Box::new(self.new_vector_writer(options).await))
1500 }
1501 }
1502
1503 pub type StateStoreDynRef = StateStorePointer<Arc<dyn DynStateStore>>;
1504
1505 macro_rules! state_store_pointer_dyn_as_ref {
1506 ($pointer:ident < dyn $source_dyn_trait:ident > , $target_dyn_trait:ident) => {
1507 impl AsRef<dyn $target_dyn_trait>
1508 for StateStorePointer<$pointer<dyn $source_dyn_trait>>
1509 {
1510 fn as_ref(&self) -> &dyn $target_dyn_trait {
1511 (&*self.0) as _
1512 }
1513 }
1514 };
1515 }
1516
1517 state_store_pointer_dyn_as_ref!(Arc<dyn DynStateStoreReadSnapshot>, DynStateStoreRead);
1518 state_store_pointer_dyn_as_ref!(Arc<dyn DynStateStoreReadSnapshot>, DynStateStoreGet);
1519 state_store_pointer_dyn_as_ref!(Arc<dyn DynStateStoreReadSnapshot>, DynStateStoreReadVector);
1520 state_store_pointer_dyn_as_ref!(Arc<dyn DynStateStoreRead>, DynStateStoreRead);
1521 state_store_pointer_dyn_as_ref!(Arc<dyn DynStateStoreRead>, DynStateStoreGet);
1522 state_store_pointer_dyn_as_ref!(Box<dyn DynLocalStateStore>, DynStateStoreGet);
1523
1524 macro_rules! state_store_pointer_dyn_as_mut {
1525 ($pointer:ident < dyn $source_dyn_trait:ident > , $target_dyn_trait:ident) => {
1526 impl AsMut<dyn $target_dyn_trait>
1527 for StateStorePointer<$pointer<dyn $source_dyn_trait>>
1528 {
1529 fn as_mut(&mut self) -> &mut dyn $target_dyn_trait {
1530 (&mut *self.0) as _
1531 }
1532 }
1533 };
1534 }
1535
1536 state_store_pointer_dyn_as_mut!(Box<dyn DynLocalStateStore>, DynStateStoreWriteEpochControl);
1537 state_store_pointer_dyn_as_mut!(
1538 Box<dyn DynStateStoreWriteVector>,
1539 DynStateStoreWriteEpochControl
1540 );
1541
1542 #[derive(Clone)]
1543 pub struct StateStorePointer<P>(pub(crate) P);
1544
1545 impl<P> StateStoreGet for StateStorePointer<P>
1546 where
1547 StateStorePointer<P>: AsRef<dyn DynStateStoreGet> + StaticSendSync,
1548 {
1549 async fn on_key_value<'a, O: Send + 'a>(
1550 &'a self,
1551 key: TableKey<Bytes>,
1552 read_options: ReadOptions,
1553 on_key_value_fn: impl KeyValueFn<'a, O>,
1554 ) -> StorageResult<Option<O>> {
1555 let option = self.as_ref().get_keyed_row(key, read_options).await?;
1556 option
1557 .map(|(key, value)| on_key_value_fn(key.to_ref(), value.as_ref()))
1558 .transpose()
1559 }
1560 }
1561
1562 impl<P> StateStoreRead for StateStorePointer<P>
1563 where
1564 StateStorePointer<P>: AsRef<dyn DynStateStoreRead> + StateStoreGet + StaticSendSync,
1565 {
1566 type Iter = BoxStateStoreReadIter;
1567 type RevIter = BoxStateStoreReadIter;
1568
1569 fn iter(
1570 &self,
1571 key_range: TableKeyRange,
1572
1573 read_options: ReadOptions,
1574 ) -> impl Future<Output = StorageResult<Self::Iter>> + '_ {
1575 self.as_ref().iter(key_range, read_options)
1576 }
1577
1578 fn rev_iter(
1579 &self,
1580 key_range: TableKeyRange,
1581
1582 read_options: ReadOptions,
1583 ) -> impl Future<Output = StorageResult<Self::RevIter>> + '_ {
1584 self.as_ref().rev_iter(key_range, read_options)
1585 }
1586 }
1587
1588 impl StateStoreReadLog for StateStoreDynRef {
1589 type ChangeLogIter = BoxStateStoreReadChangeLogIter;
1590
1591 async fn next_epoch(&self, epoch: u64, options: NextEpochOptions) -> StorageResult<u64> {
1592 (*self.0).next_epoch(epoch, options).await
1593 }
1594
1595 fn iter_log(
1596 &self,
1597 epoch_range: (u64, u64),
1598 key_range: TableKeyRange,
1599 options: ReadLogOptions,
1600 ) -> impl Future<Output = StorageResult<Self::ChangeLogIter>> + Send + '_ {
1601 (*self.0).iter_log(epoch_range, key_range, options)
1602 }
1603 }
1604
1605 pub trait DynStateStore: DynStateStoreReadLog + DynStateStoreExt + AsHummock {}
1606
1607 impl AsHummock for StateStoreDynRef {
1608 fn as_hummock(&self) -> Option<&HummockStorage> {
1609 (*self.0).as_hummock()
1610 }
1611 }
1612
1613 impl<S: DynStateStoreReadLog + DynStateStoreExt + AsHummock> DynStateStore for S {}
1614
1615 impl StateStore for StateStoreDynRef {
1616 type Local = BoxDynLocalStateStore;
1617 type ReadSnapshot = StateStoreReadSnapshotDynRef;
1618 type VectorWriter = BoxDynStateStoreWriteVector;
1619
1620 fn try_wait_epoch(
1621 &self,
1622 epoch: HummockReadEpoch,
1623 options: TryWaitEpochOptions,
1624 ) -> impl Future<Output = StorageResult<()>> + Send + '_ {
1625 (*self.0).try_wait_epoch(epoch, options)
1626 }
1627
1628 fn new_local(
1629 &self,
1630 option: NewLocalOptions,
1631 ) -> impl Future<Output = Self::Local> + Send + '_ {
1632 (*self.0).new_local(option)
1633 }
1634
1635 async fn new_read_snapshot(
1636 &self,
1637 epoch: HummockReadEpoch,
1638 options: NewReadSnapshotOptions,
1639 ) -> StorageResult<Self::ReadSnapshot> {
1640 (*self.0).new_read_snapshot(epoch, options).await
1641 }
1642
1643 fn new_vector_writer(
1644 &self,
1645 options: NewVectorWriterOptions,
1646 ) -> impl Future<Output = Self::VectorWriter> + Send + '_ {
1647 (*self.0).new_vector_writer(options)
1648 }
1649 }
1650}