Skip to main content

risingwave_storage/
store_impl.rs

1// Copyright 2022 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::HashSet;
16use std::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/// The type erased [`StateStore`].
224#[derive(Clone, EnumAsInner)]
225#[expect(clippy::enum_variant_names)]
226pub enum StateStoreImpl {
227    /// The Hummock state store, which operates on an S3-like service. URLs beginning with
228    /// `hummock` will be automatically recognized as Hummock state store.
229    ///
230    /// Example URLs:
231    ///
232    /// * `hummock+s3://bucket`
233    /// * `hummock+minio://KEY:SECRET@minio-ip:port`
234    /// * `hummock+memory` (should only be used in 1 compute node mode)
235    HummockStateStore(Monitored<HummockStorageType>),
236    /// In-memory B-Tree state store. Should only be used in unit and integration tests. If you
237    /// want speed up e2e test, you should use Hummock in-memory mode instead. Also, this state
238    /// store misses some critical implementation to ensure the correctness of persisting streaming
239    /// state. (e.g., no `read_epoch` support, no async checkpoint)
240    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        // The specific type of MemoryStateStoreType in deducted here.
291        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        // The specific type of HummockStateStoreType in deducted here.
299        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                // WARNING: don't change this. Enabling memory backend will cause monomorphization
348                // explosion and thus slow compile time in release mode.
349                #[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                // WARNING: don't change this. Enabling memory backend will cause monomorphization
362                // explosion and thus slow compile time in release mode.
363                #[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        // TODO: may avoid manual async fn when the bug of rust compiler is fixed. Currently it will
486        // fail to compile.
487        #[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    // For StateStoreRead
1093
1094    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    // For LocalStateStore
1190    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    // For global StateStore
1391
1392    #[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}