1use std::fmt::Display;
16use std::sync::Arc;
17
18use crate::logstore::LogStore;
19use crate::storage::RegionId;
20
21#[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 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#[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 pub fn type_name() -> &'static str {
57 "RaftEngineProvider"
58 }
59}
60
61#[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 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#[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 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 pub fn is_remote_wal(&self) -> bool {
145 matches!(self, Provider::Kafka(_) | Provider::ObjectStore(_))
146 }
147
148 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 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 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 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 #[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}