1use std::fmt::Debug;
19use std::sync::Arc;
20
21use super::PERSY_SCHEME;
22use super::config::PersyConfig;
23use super::core::*;
24use super::deleter::PersyDeleter;
25use super::reader::*;
26use super::writer::PersyWriter;
27use opendal_core::raw::*;
28use opendal_core::*;
29
30#[doc = include_str!("docs.md")]
32#[derive(Debug, Default)]
33pub struct PersyBuilder {
34 pub(super) config: PersyConfig,
35}
36
37impl PersyBuilder {
38 pub fn datafile(mut self, path: &str) -> Self {
40 self.config.datafile = Some(path.into());
41 self
42 }
43
44 pub fn segment(mut self, path: &str) -> Self {
46 self.config.segment = Some(path.into());
47 self
48 }
49
50 pub fn index(mut self, path: &str) -> Self {
52 self.config.index = Some(path.into());
53 self
54 }
55}
56
57impl Builder for PersyBuilder {
58 type Config = PersyConfig;
59
60 fn build(self) -> Result<impl Service> {
61 let datafile_path = self.config.datafile.ok_or_else(|| {
62 Error::new(ErrorKind::ConfigInvalid, "datafile is required but not set")
63 .with_context("service", PERSY_SCHEME)
64 })?;
65
66 let segment_name = self.config.segment.ok_or_else(|| {
67 Error::new(ErrorKind::ConfigInvalid, "segment is required but not set")
68 .with_context("service", PERSY_SCHEME)
69 })?;
70
71 let segment = segment_name.clone();
72
73 let index_name = self.config.index.ok_or_else(|| {
74 Error::new(ErrorKind::ConfigInvalid, "index is required but not set")
75 .with_context("service", PERSY_SCHEME)
76 })?;
77
78 let index = index_name.clone();
79
80 let persy = persy::OpenOptions::new()
81 .create(true)
82 .prepare_with(move |p| init(p, &segment_name, &index_name))
83 .open(&datafile_path)
84 .map_err(|e| {
85 Error::new(ErrorKind::ConfigInvalid, "open db")
86 .with_context("service", PERSY_SCHEME)
87 .with_context("datafile", datafile_path.clone())
88 .set_source(e)
89 })?;
90
91 fn init(
93 persy: &persy::Persy,
94 segment_name: &str,
95 index_name: &str,
96 ) -> Result<(), Box<dyn std::error::Error>> {
97 let mut tx = persy.begin()?;
98
99 if !tx.exists_segment(segment_name)? {
100 tx.create_segment(segment_name)?;
101 }
102 if !tx.exists_index(index_name)? {
103 tx.create_index::<String, persy::PersyId>(index_name, persy::ValueMode::Replace)?;
104 }
105
106 let prepared = tx.prepare()?;
107 prepared.commit()?;
108
109 Ok(())
110 }
111
112 Ok(PersyBackend::new(PersyCore {
113 datafile: datafile_path,
114 segment,
115 index,
116 persy,
117 }))
118 }
119}
120
121#[derive(Clone, Debug)]
123pub struct PersyBackend {
124 pub(crate) core: Arc<PersyCore>,
125 pub(crate) root: String,
126 pub(crate) info: ServiceInfo,
127 pub(crate) capability: Capability,
128}
129
130impl PersyBackend {
131 pub fn new(core: PersyCore) -> Self {
132 let info = ServiceInfo::new(PERSY_SCHEME, "/", &core.datafile);
133 let capability = Capability {
134 read: true,
135 stat: true,
136 write: true,
137 write_can_empty: true,
138 delete: true,
139 ..Default::default()
140 };
141
142 Self {
143 core: Arc::new(core),
144 root: "/".to_string(),
145 info,
146 capability,
147 }
148 }
149}
150
151impl Service for PersyBackend {
152 type Reader = oio::StreamReader<PersyReader>;
153 type Writer = PersyWriter;
154 type Lister = ();
155 type Deleter = oio::OneShotDeleter<PersyDeleter>;
156 type Copier = ();
157 type Composer = ();
158
159 fn info(&self) -> ServiceInfo {
160 self.info.clone()
161 }
162
163 fn capability(&self) -> Capability {
164 self.capability
165 }
166
167 async fn create_dir(
168 &self,
169 _ctx: &OperationContext,
170 _path: &str,
171 _args: OpCreateDir,
172 ) -> Result<RpCreateDir> {
173 Err(Error::new(
174 ErrorKind::Unsupported,
175 "operation is not supported",
176 ))
177 }
178
179 async fn stat(&self, _ctx: &OperationContext, path: &str, _: OpStat) -> Result<RpStat> {
180 let p = build_abs_path(&self.root, path);
181
182 if p == build_abs_path(&self.root, "") {
183 Ok(RpStat::new(MetadataBuilder::dir().build()))
184 } else {
185 let bs = self.core.get(&p)?;
186 match bs {
187 Some(bs) => Ok(RpStat::new({
188 let metadata = MetadataBuilder::file(bs.len() as u64);
189 metadata.build()
190 })),
191 None => Err(Error::new(ErrorKind::NotFound, "kv not found in persy")),
192 }
193 }
194 }
195 fn read(&self, _ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
196 let output: oio::StreamReader<PersyReader> = {
197 Ok(oio::StreamReader::new(PersyReader::new(
198 self.clone(),
199 path,
200 args,
201 )))
202 }?;
203
204 Ok(output)
205 }
206
207 fn write(&self, _ctx: &OperationContext, path: &str, _: OpWrite) -> Result<Self::Writer> {
208 let output: PersyWriter = {
209 let p = build_abs_path(&self.root, path);
210 Ok(PersyWriter::new(self.core.clone(), p))
211 }?;
212
213 Ok(output)
214 }
215
216 fn delete(&self, _ctx: &OperationContext) -> Result<Self::Deleter> {
217 let output: oio::OneShotDeleter<PersyDeleter> = {
218 Ok(oio::OneShotDeleter::new(PersyDeleter::new(
219 self.core.clone(),
220 self.root.clone(),
221 )))
222 }?;
223
224 Ok(output)
225 }
226
227 fn list(&self, _ctx: &OperationContext, _path: &str, _args: OpList) -> Result<Self::Lister> {
228 Err(Error::new(
229 ErrorKind::Unsupported,
230 "operation is not supported",
231 ))
232 }
233
234 fn copy(
235 &self,
236 _ctx: &OperationContext,
237 _from: &str,
238 _to: &str,
239 _args: OpCopy,
240 ) -> Result<Self::Copier> {
241 Err(Error::new(
242 ErrorKind::Unsupported,
243 "operation is not supported",
244 ))
245 }
246
247 async fn rename(
248 &self,
249 _ctx: &OperationContext,
250 _from: &str,
251 _to: &str,
252 _args: OpRename,
253 ) -> Result<RpRename> {
254 Err(Error::new(
255 ErrorKind::Unsupported,
256 "operation is not supported",
257 ))
258 }
259
260 async fn presign(
261 &self,
262 _ctx: &OperationContext,
263 _path: &str,
264 _args: OpPresign,
265 ) -> Result<RpPresign> {
266 Err(Error::new(
267 ErrorKind::Unsupported,
268 "operation is not supported",
269 ))
270 }
271}