opendal_service_rocksdb/
backend.rs1use std::sync::Arc;
19
20use opendal_core::raw::*;
21use opendal_core::*;
22use rocksdb::DB;
23
24use super::ROCKSDB_SCHEME;
25use super::config::RocksdbConfig;
26use super::core::*;
27use super::deleter::RocksdbDeleter;
28use super::lister::RocksdbLister;
29use super::reader::*;
30use super::writer::RocksdbWriter;
31
32#[doc = include_str!("docs.md")]
34#[derive(Debug, Default)]
35pub struct RocksdbBuilder {
36 pub(super) config: RocksdbConfig,
37}
38
39impl RocksdbBuilder {
40 pub fn datadir(mut self, path: &str) -> Self {
42 self.config.datadir = Some(path.into());
43 self
44 }
45
46 pub fn root(mut self, root: &str) -> Self {
50 self.config.root = if root.is_empty() {
51 None
52 } else {
53 Some(root.to_string())
54 };
55
56 self
57 }
58}
59
60impl Builder for RocksdbBuilder {
61 type Config = RocksdbConfig;
62
63 fn build(self) -> Result<impl Service> {
64 let path = self.config.datadir.ok_or_else(|| {
65 Error::new(ErrorKind::ConfigInvalid, "datadir is required but not set")
66 .with_context("service", ROCKSDB_SCHEME)
67 })?;
68 let db = DB::open_default(&path).map_err(|e| {
69 Error::new(ErrorKind::ConfigInvalid, "open default transaction db")
70 .with_context("service", ROCKSDB_SCHEME)
71 .with_context("datadir", path)
72 .set_source(e)
73 })?;
74
75 let root = normalize_root(&self.config.root.unwrap_or_default());
76
77 Ok(RocksdbBackend::new(RocksdbCore { db: Arc::new(db) }).with_normalized_root(root))
78 }
79}
80
81#[derive(Clone, Debug)]
83pub struct RocksdbBackend {
84 pub(crate) core: Arc<RocksdbCore>,
85 pub(crate) root: String,
86 pub(crate) info: ServiceInfo,
87 pub(crate) capability: Capability,
88}
89
90impl RocksdbBackend {
91 pub fn new(core: RocksdbCore) -> Self {
92 let info = ServiceInfo::new(ROCKSDB_SCHEME, "/", core.db.path().to_string_lossy());
93 let capability = Capability {
94 read: true,
95 stat: true,
96 write: true,
97 write_can_empty: true,
98 delete: true,
99 list: true,
100 list_with_recursive: true,
101 ..Default::default()
102 };
103
104 Self {
105 core: Arc::new(core),
106 root: "/".to_string(),
107 info,
108 capability,
109 }
110 }
111
112 fn with_normalized_root(mut self, root: String) -> Self {
113 self.info = self.info.with_root(&root);
114 self.root = root;
115 self
116 }
117}
118
119impl Service for RocksdbBackend {
120 type Reader = oio::StreamReader<RocksdbReader>;
121 type Writer = RocksdbWriter;
122 type Lister = oio::HierarchyLister<RocksdbLister>;
123 type Deleter = oio::OneShotDeleter<RocksdbDeleter>;
124 type Copier = ();
125 type Composer = ();
126
127 fn info(&self) -> ServiceInfo {
128 self.info.clone()
129 }
130
131 fn capability(&self) -> Capability {
132 self.capability
133 }
134
135 async fn create_dir(
136 &self,
137 _ctx: &OperationContext,
138 _path: &str,
139 _args: OpCreateDir,
140 ) -> Result<RpCreateDir> {
141 Err(Error::new(
142 ErrorKind::Unsupported,
143 "operation is not supported",
144 ))
145 }
146
147 async fn stat(&self, _ctx: &OperationContext, path: &str, _: OpStat) -> Result<RpStat> {
148 let p = build_abs_path(&self.root, path);
149
150 if p == build_abs_path(&self.root, "") {
151 Ok(RpStat::new(MetadataBuilder::dir().build()))
152 } else {
153 let bs = self.core.get(&p)?;
154 match bs {
155 Some(bs) => Ok(RpStat::new({
156 let metadata = MetadataBuilder::file(bs.len() as u64);
157 metadata.build()
158 })),
159 None => Err(Error::new(ErrorKind::NotFound, "kv not found in rocksdb")),
160 }
161 }
162 }
163 fn read(&self, _ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
164 let output: oio::StreamReader<RocksdbReader> = {
165 Ok(oio::StreamReader::new(RocksdbReader::new(
166 self.clone(),
167 path,
168 args,
169 )))
170 }?;
171
172 Ok(output)
173 }
174
175 fn write(&self, _ctx: &OperationContext, path: &str, _: OpWrite) -> Result<Self::Writer> {
176 let output: RocksdbWriter = {
177 let p = build_abs_path(&self.root, path);
178 let writer = RocksdbWriter::new(self.core.clone(), p);
179 Ok(writer)
180 }?;
181
182 Ok(output)
183 }
184
185 fn delete(&self, _ctx: &OperationContext) -> Result<Self::Deleter> {
186 let output: oio::OneShotDeleter<RocksdbDeleter> = {
187 let deleter = RocksdbDeleter::new(self.core.clone(), self.root.clone());
188 Ok(oio::OneShotDeleter::new(deleter))
189 }?;
190
191 Ok(output)
192 }
193
194 fn list(&self, _ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
195 let output: oio::HierarchyLister<RocksdbLister> = {
196 let p = build_abs_path(&self.root, path);
197 let lister = RocksdbLister::new(self.core.clone(), self.root.clone(), p)?;
198 Ok(oio::HierarchyLister::new(lister, path, args.recursive()))
199 }?;
200
201 Ok(output)
202 }
203
204 fn copy(
205 &self,
206 _ctx: &OperationContext,
207 _from: &str,
208 _to: &str,
209 _args: OpCopy,
210 ) -> Result<Self::Copier> {
211 Err(Error::new(
212 ErrorKind::Unsupported,
213 "operation is not supported",
214 ))
215 }
216
217 async fn rename(
218 &self,
219 _ctx: &OperationContext,
220 _from: &str,
221 _to: &str,
222 _args: OpRename,
223 ) -> Result<RpRename> {
224 Err(Error::new(
225 ErrorKind::Unsupported,
226 "operation is not supported",
227 ))
228 }
229
230 async fn presign(
231 &self,
232 _ctx: &OperationContext,
233 _path: &str,
234 _args: OpPresign,
235 ) -> Result<RpPresign> {
236 Err(Error::new(
237 ErrorKind::Unsupported,
238 "operation is not supported",
239 ))
240 }
241}