Skip to main content

store_api/logstore/
provider.rs

1// Copyright 2023 Greptime Team
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::fmt::Display;
16use std::sync::Arc;
17
18use crate::logstore::LogStore;
19use crate::storage::RegionId;
20
21// The Provider of kafka log store
22#[derive(Debug, Clone, PartialEq, Eq, Hash)]
23pub struct KafkaProvider {
24    pub topic: String,
25}
26
27impl KafkaProvider {
28    pub fn new(topic: String) -> Self {
29        Self { topic }
30    }
31
32    /// Returns the type name.
33    pub fn type_name() -> &'static str {
34        "KafkaProvider"
35    }
36}
37
38impl Display for KafkaProvider {
39    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
40        write!(f, "{}", self.topic)
41    }
42}
43
44// The Provider of raft engine log store
45#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
46pub struct RaftEngineProvider {
47    pub id: u64,
48}
49
50impl RaftEngineProvider {
51    pub fn new(id: u64) -> Self {
52        Self { id }
53    }
54
55    /// Returns the type name.
56    pub fn type_name() -> &'static str {
57        "RaftEngineProvider"
58    }
59}
60
61/// The provider of the object store log store.
62///
63/// Objects under one prefix hold entries of many regions, so reads and obsoletes
64/// are scoped by `region_id`.
65#[derive(Debug, Clone, PartialEq, Eq, Hash)]
66pub struct ObjectStoreProvider {
67    pub region_id: RegionId,
68    pub prefix: String,
69}
70
71impl ObjectStoreProvider {
72    pub fn new(region_id: RegionId, prefix: String) -> Self {
73        Self { region_id, prefix }
74    }
75
76    /// Returns the type name.
77    pub fn type_name() -> &'static str {
78        "ObjectStoreProvider"
79    }
80}
81
82impl Display for ObjectStoreProvider {
83    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
84        write!(f, "{}/{}", self.prefix, self.region_id)
85    }
86}
87
88/// The Provider of LogStore
89#[derive(Debug, Clone, PartialEq, Eq, Hash)]
90pub enum Provider {
91    RaftEngine(RaftEngineProvider),
92    Kafka(Arc<KafkaProvider>),
93    ObjectStore(Arc<ObjectStoreProvider>),
94    Noop,
95}
96
97impl Display for Provider {
98    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
99        match &self {
100            Provider::RaftEngine(provider) => {
101                write!(f, "RaftEngine(region={})", RegionId::from_u64(provider.id))
102            }
103            Provider::Kafka(provider) => write!(f, "Kafka(topic={})", provider.topic),
104            Provider::ObjectStore(provider) => write!(
105                f,
106                "ObjectStore(prefix={}, region={})",
107                provider.prefix, provider.region_id
108            ),
109            Provider::Noop => write!(f, "Noop"),
110        }
111    }
112}
113
114impl Provider {
115    /// Returns the initial flushed entry id of the provider.
116    /// This is used to initialize the flushed entry id of the region when creating the region from scratch.
117    ///
118    /// Currently only used for remote WAL.
119    /// For local WAL, the initial flushed entry id is 0.
120    pub fn initial_flushed_entry_id<S: LogStore>(&self, wal: &S) -> u64 {
121        if self.is_remote_wal() {
122            return wal.latest_entry_id(self).unwrap_or(0);
123        }
124        0
125    }
126
127    pub fn raft_engine_provider(id: u64) -> Provider {
128        Provider::RaftEngine(RaftEngineProvider { id })
129    }
130
131    pub fn kafka_provider(topic: String) -> Provider {
132        Provider::Kafka(Arc::new(KafkaProvider { topic }))
133    }
134
135    pub fn object_store_provider(region_id: RegionId, prefix: String) -> Provider {
136        Provider::ObjectStore(Arc::new(ObjectStoreProvider { region_id, prefix }))
137    }
138
139    pub fn noop_provider() -> Provider {
140        Provider::Noop
141    }
142
143    /// Returns true if it's remote WAL.
144    pub fn is_remote_wal(&self) -> bool {
145        matches!(self, Provider::Kafka(_) | Provider::ObjectStore(_))
146    }
147
148    /// Returns the type name.
149    pub fn type_name(&self) -> &'static str {
150        match self {
151            Provider::RaftEngine(_) => RaftEngineProvider::type_name(),
152            Provider::Kafka(_) => KafkaProvider::type_name(),
153            Provider::ObjectStore(_) => ObjectStoreProvider::type_name(),
154            Provider::Noop => "Noop",
155        }
156    }
157
158    /// Returns the reference of [`RaftEngineProvider`] if it's the type of [`LogStoreProvider::RaftEngine`].
159    pub fn as_raft_engine_provider(&self) -> Option<&RaftEngineProvider> {
160        if let Provider::RaftEngine(ns) = self {
161            return Some(ns);
162        }
163        None
164    }
165
166    /// Returns the reference of [`KafkaProvider`] if it's the type of [`LogStoreProvider::Kafka`].
167    pub fn as_kafka_provider(&self) -> Option<&Arc<KafkaProvider>> {
168        if let Provider::Kafka(ns) = self {
169            return Some(ns);
170        }
171        None
172    }
173
174    /// Returns the reference of [`ObjectStoreProvider`] if it's the type of [`Provider::ObjectStore`].
175    pub fn as_object_store_provider(&self) -> Option<&Arc<ObjectStoreProvider>> {
176        if let Provider::ObjectStore(ns) = self {
177            return Some(ns);
178        }
179        None
180    }
181}
182
183#[cfg(test)]
184mod tests {
185    use common_error::mock::MockError;
186    use common_error::status_code::StatusCode;
187
188    use super::*;
189    use crate::logstore::entry::Entry;
190    use crate::logstore::{AppendBatchResponse, EntryId, SendableEntryStream, WalIndex};
191
192    /// A log store whose latest entry id is fixed.
193    #[derive(Debug)]
194    struct LatestEntryIdLogStore(EntryId);
195
196    #[async_trait::async_trait]
197    impl LogStore for LatestEntryIdLogStore {
198        type Error = MockError;
199
200        async fn stop(&self) -> Result<(), Self::Error> {
201            unreachable!()
202        }
203
204        async fn append_batch(
205            &self,
206            _entries: Vec<Entry>,
207        ) -> Result<AppendBatchResponse, Self::Error> {
208            unreachable!()
209        }
210
211        async fn read(
212            &self,
213            _provider: &Provider,
214            _id: EntryId,
215            _index: Option<WalIndex>,
216        ) -> Result<SendableEntryStream<'static, Entry, Self::Error>, Self::Error> {
217            unreachable!()
218        }
219
220        async fn create_namespace(&self, _ns: &Provider) -> Result<(), Self::Error> {
221            unreachable!()
222        }
223
224        async fn delete_namespace(&self, _ns: &Provider) -> Result<(), Self::Error> {
225            unreachable!()
226        }
227
228        async fn list_namespaces(&self) -> Result<Vec<Provider>, Self::Error> {
229            unreachable!()
230        }
231
232        async fn obsolete(
233            &self,
234            _provider: &Provider,
235            _region_id: RegionId,
236            _entry_id: EntryId,
237        ) -> Result<(), Self::Error> {
238            unreachable!()
239        }
240
241        async fn obsolete_all(
242            &self,
243            _provider: &Provider,
244            _region_id: RegionId,
245        ) -> Result<(), Self::Error> {
246            unreachable!()
247        }
248
249        fn entry(
250            &self,
251            _data: Vec<u8>,
252            _entry_id: EntryId,
253            _region_id: RegionId,
254            _provider: &Provider,
255        ) -> Result<Entry, Self::Error> {
256            unreachable!()
257        }
258
259        fn latest_entry_id(&self, provider: &Provider) -> Result<EntryId, Self::Error> {
260            if provider.is_remote_wal() {
261                Ok(self.0)
262            } else {
263                Err(MockError::new(StatusCode::Unexpected))
264            }
265        }
266    }
267
268    #[test]
269    fn test_object_store_provider() {
270        let region_id = RegionId::new(1, 2);
271        let provider = Provider::object_store_provider(region_id, "wal".to_string());
272
273        let object_store = provider.as_object_store_provider().unwrap();
274        assert_eq!(region_id, object_store.region_id);
275        assert_eq!("wal", object_store.prefix);
276        assert_eq!(ObjectStoreProvider::type_name(), provider.type_name());
277        assert_eq!(
278            format!("ObjectStore(prefix=wal, region={region_id})"),
279            provider.to_string()
280        );
281        assert!(provider.is_remote_wal());
282        assert!(provider.as_kafka_provider().is_none());
283        assert!(provider.as_raft_engine_provider().is_none());
284    }
285
286    #[test]
287    fn test_initial_flushed_entry_id_follows_remote_wal() {
288        let region_id = RegionId::new(1, 2);
289        let store = LatestEntryIdLogStore(42);
290
291        let object_store = Provider::object_store_provider(region_id, "wal".to_string());
292        assert_eq!(42, object_store.initial_flushed_entry_id(&store));
293
294        let kafka = Provider::kafka_provider("topic".to_string());
295        assert_eq!(42, kafka.initial_flushed_entry_id(&store));
296
297        let raft_engine = Provider::raft_engine_provider(region_id.as_u64());
298        assert_eq!(0, raft_engine.initial_flushed_entry_id(&store));
299        assert_eq!(
300            0,
301            Provider::noop_provider().initial_flushed_entry_id(&store)
302        );
303    }
304}