1use std::fmt::Debug;
19use std::fmt::Formatter;
20use std::sync::Arc;
21
22use crate::raw::oio::Delete;
23use crate::raw::oio::List;
24use crate::raw::oio::PrefixLister;
25use crate::raw::*;
26use crate::*;
27
28#[derive(Debug, Clone)]
30pub struct SimulateLayer {
31 read_with_suffix: bool,
32 list_recursive: bool,
33 stat_dir: bool,
34 create_dir: bool,
35 delete_recursive: bool,
36}
37
38impl Default for SimulateLayer {
39 fn default() -> Self {
40 Self {
41 read_with_suffix: true,
42 list_recursive: true,
43 stat_dir: true,
44 create_dir: true,
45 delete_recursive: true,
46 }
47 }
48}
49
50impl SimulateLayer {
51 pub fn with_read_with_suffix(mut self, enabled: bool) -> Self {
53 self.read_with_suffix = enabled;
54 self
55 }
56
57 pub fn with_list_recursive(mut self, enabled: bool) -> Self {
59 self.list_recursive = enabled;
60 self
61 }
62
63 pub fn with_stat_dir(mut self, enabled: bool) -> Self {
65 self.stat_dir = enabled;
66 self
67 }
68
69 pub fn with_create_dir(mut self, enabled: bool) -> Self {
71 self.create_dir = enabled;
72 self
73 }
74
75 pub fn with_delete_recursive(mut self, enabled: bool) -> Self {
77 self.delete_recursive = enabled;
78 self
79 }
80}
81
82impl Layer for SimulateLayer {
83 fn apply_service(&self, srv: Servicer) -> Servicer {
84 Arc::new(self.layer(srv))
85 }
86}
87
88impl SimulateLayer {
89 fn layer(&self, srv: Servicer) -> SimulateService {
90 SimulateService {
91 srv,
92 config: self.clone(),
93 }
94 }
95}
96
97pub struct SimulateService {
99 srv: Servicer,
100 config: SimulateLayer,
101}
102
103impl Debug for SimulateService {
104 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
105 self.srv.fmt(f)
106 }
107}
108
109impl SimulateService {
110 fn simulate_capability(&self, mut cap: Capability) -> Capability {
111 if self.config.read_with_suffix && cap.read {
112 cap.read_with_suffix = true;
113 }
114 if self.config.create_dir && cap.list && cap.write_can_empty {
115 cap.create_dir = true;
116 }
117 if self.config.delete_recursive && cap.list && cap.delete {
118 cap.delete_with_recursive = true;
119 }
120 cap
121 }
122
123 async fn simulate_create_dir(
124 &self,
125 ctx: &OperationContext,
126 path: &str,
127 args: OpCreateDir,
128 ) -> Result<RpCreateDir> {
129 let capability = self.srv.capability();
130
131 if capability.create_dir || !self.config.create_dir {
132 return self.srv.create_dir(ctx, path, args).await;
133 }
134
135 if capability.write_can_empty && capability.list {
136 let mut w = self.srv.write(ctx, path, OpWrite::default())?;
137 oio::Write::close(&mut w).await?;
138 return Ok(RpCreateDir::default());
139 }
140
141 self.srv.create_dir(ctx, path, args).await
142 }
143
144 async fn simulate_stat(
145 &self,
146 ctx: &OperationContext,
147 path: &str,
148 args: OpStat,
149 ) -> Result<RpStat> {
150 let capability = self.srv.capability();
151
152 if path == "/" {
153 return Ok(RpStat::new(MetadataBuilder::dir().build()));
154 }
155
156 if path.ends_with('/') {
157 if capability.create_dir {
158 let meta = self
159 .srv
160 .stat(ctx, path, args.clone())
161 .await?
162 .into_metadata();
163
164 if meta.is_file() {
165 return Err(Error::new(
166 ErrorKind::NotFound,
167 "stat expected a directory, but found a file",
168 ));
169 }
170
171 return Ok(RpStat::new(meta));
172 }
173
174 if self.config.stat_dir && capability.list {
175 let mut l = self.srv.list(
176 ctx,
177 path,
178 options::ListOptions {
179 recursive: capability.list_with_recursive,
180 limit: Some(1),
181 ..Default::default()
182 }
183 .into(),
184 )?;
185
186 return if l.next().await?.is_some() {
187 Ok(RpStat::new(MetadataBuilder::dir().build()))
188 } else {
189 Err(Error::new(
190 ErrorKind::NotFound,
191 "the directory is not found",
192 ))
193 };
194 }
195 }
196
197 self.srv.stat(ctx, path, args).await
198 }
199
200 fn simulate_list(
201 &self,
202 ctx: &OperationContext,
203 path: &str,
204 args: OpList,
205 ) -> Result<SimulateLister> {
206 let cap = self.srv.capability();
207
208 let recursive = args.recursive();
209 let forward = args;
210
211 let lister = match (
212 recursive,
213 cap.list_with_recursive,
214 self.config.list_recursive,
215 ) {
216 (_, true, _) => {
218 let p = self.srv.list(ctx, path, forward)?;
219 SimulateLister::One(p)
220 }
221 (true, false, true) => {
223 if path.ends_with('/') {
224 let p = ServicerFlatLister::new(ctx.clone(), self.srv.clone(), path);
225 SimulateLister::Two(p)
226 } else {
227 let parent = get_parent(path);
228 let p = ServicerFlatLister::new(ctx.clone(), self.srv.clone(), parent);
229 let p = PrefixLister::new(p, path);
230 SimulateLister::Four(p)
231 }
232 }
233 (true, false, false) => {
235 let p = self.srv.list(ctx, path, forward)?;
236 SimulateLister::One(p)
237 }
238 (false, false, _) => {
240 if path.ends_with('/') {
241 let p = self.srv.list(ctx, path, forward)?;
242 SimulateLister::One(p)
243 } else {
244 let parent = get_parent(path);
245 let p = self.srv.list(ctx, parent, forward)?;
246 let p = PrefixLister::new(p, path);
247 SimulateLister::Three(p)
248 }
249 }
250 };
251
252 Ok(lister)
253 }
254
255 async fn simulate_delete_with_recursive(
256 &self,
257 ctx: &OperationContext,
258 path: &str,
259 args: OpDelete,
260 deleter: &mut oio::Deleter,
261 ) -> Result<()> {
262 if !self.capability().delete_with_recursive {
263 return Err(Error::new(
264 ErrorKind::Unsupported,
265 "recursive delete is not supported",
266 ));
267 }
268
269 let mut non_recursive = args;
270 non_recursive.set_recursive(false);
271
272 let mut lister = self.simulate_list(
273 ctx,
274 path,
275 options::ListOptions {
276 recursive: true,
277 ..Default::default()
278 }
279 .into(),
280 )?;
281
282 while let Some(entry) = lister.next().await? {
283 let entry = entry.into_entry();
284 let mut entry_args = non_recursive.clone();
285 if let Some(version) = entry.metadata().version() {
286 entry_args = entry_args.into_version(version);
287 }
288 deleter.delete(entry.path(), entry_args).await?;
289 }
290
291 Ok(())
292 }
293}
294
295impl Service for SimulateService {
296 type Reader = SimulateReader;
297 type Writer = oio::Writer;
298 type Lister = SimulateLister;
299 type Deleter = SimulateDeleter;
300 type Copier = oio::Copier;
301 type Composer = oio::Composer;
302
303 fn info(&self) -> ServiceInfo {
304 self.srv.info()
305 }
306
307 fn capability(&self) -> Capability {
308 self.simulate_capability(self.srv.capability())
309 }
310
311 async fn create_dir(
312 &self,
313 ctx: &OperationContext,
314 path: &str,
315 args: OpCreateDir,
316 ) -> Result<RpCreateDir> {
317 self.simulate_create_dir(ctx, path, args).await
318 }
319
320 fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
321 let capability = self.srv.capability();
322 let simulate_read_with_suffix =
323 self.config.read_with_suffix && capability.read && !capability.read_with_suffix;
324 let reader = self.srv.read(ctx, path, args.clone())?;
325 let reader = SimulateReader::new(
326 ctx.clone(),
327 self.srv.clone(),
328 path.to_string(),
329 args,
330 reader,
331 simulate_read_with_suffix,
332 );
333 Ok(reader)
334 }
335
336 fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
337 self.srv.write(ctx, path, args)
338 }
339
340 fn copy(
341 &self,
342 ctx: &OperationContext,
343 from: &str,
344 to: &str,
345 args: OpCopy,
346 ) -> Result<Self::Copier> {
347 self.srv.copy(ctx, from, to, args)
348 }
349
350 fn compose(&self, ctx: &OperationContext, to: &str, args: OpCompose) -> Result<Self::Composer> {
351 self.srv.compose(ctx, to, args)
352 }
353
354 async fn rename(
355 &self,
356 ctx: &OperationContext,
357 from: &str,
358 to: &str,
359 args: OpRename,
360 ) -> Result<RpRename> {
361 self.srv.rename(ctx, from, to, args).await
362 }
363
364 async fn restore(
365 &self,
366 ctx: &OperationContext,
367 path: &str,
368 args: OpRestore,
369 ) -> Result<RpRestore> {
370 self.srv.restore(ctx, path, args).await
371 }
372
373 async fn stat(&self, ctx: &OperationContext, path: &str, args: OpStat) -> Result<RpStat> {
374 self.simulate_stat(ctx, path, args).await
375 }
376
377 fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
378 let deleter = self.srv.delete(ctx)?;
379 Ok(SimulateDeleter::new(
380 ctx.clone(),
381 self.srv.clone(),
382 self.config.clone(),
383 deleter,
384 ))
385 }
386
387 fn list(&self, ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
388 self.simulate_list(ctx, path, args)
389 }
390
391 async fn presign(
392 &self,
393 ctx: &OperationContext,
394 path: &str,
395 args: OpPresign,
396 ) -> Result<RpPresign> {
397 self.srv.presign(ctx, path, args).await
398 }
399}
400
401pub type SimulateLister = FourWays<
402 oio::Lister,
403 ServicerFlatLister,
404 PrefixLister<oio::Lister>,
405 PrefixLister<ServicerFlatLister>,
406>;
407
408pub struct SimulateReader {
409 ctx: OperationContext,
410 srv: Servicer,
411 path: String,
412 args: OpRead,
413 inner: oio::Reader,
414 simulate_read_with_suffix: bool,
415}
416
417impl SimulateReader {
418 fn new(
419 ctx: OperationContext,
420 srv: Servicer,
421 path: String,
422 args: OpRead,
423 inner: oio::Reader,
424 simulate_read_with_suffix: bool,
425 ) -> Self {
426 Self {
427 ctx,
428 srv,
429 path,
430 args,
431 inner,
432 simulate_read_with_suffix,
433 }
434 }
435
436 async fn content_length(&self) -> Result<u64> {
437 if let Some(v) = self.args.content_length_hint() {
438 return Ok(v);
439 }
440
441 let op = options::StatOptions {
442 version: self.args.version().map(str::to_owned),
443 ..Default::default()
444 }
445 .into();
446
447 Ok(self
448 .srv
449 .stat(&self.ctx, &self.path, op)
450 .await?
451 .into_metadata()
452 .content_length())
453 }
454
455 async fn resolve_range(&self, range: BytesRange) -> Result<BytesRange> {
456 if !self.simulate_read_with_suffix || !range.is_suffix() {
457 return Ok(range);
458 }
459
460 let BytesRange::Suffix { size } = range else {
461 unreachable!("checked by BytesRange::is_suffix")
462 };
463
464 let content_length = self.content_length().await?;
465 let start = content_length.saturating_sub(size);
466 Ok(BytesRange::new(start, Some(content_length - start)))
467 }
468}
469
470impl oio::Read for SimulateReader {
471 async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
472 let range = self.resolve_range(range).await?;
473 self.inner.open(range).await
474 }
475
476 async fn read(&self, range: BytesRange) -> Result<(RpRead, Buffer)> {
477 let range = self.resolve_range(range).await?;
478 self.inner.read(range).await
479 }
480}
481
482pub struct ServicerFlatLister {
483 ctx: OperationContext,
484 srv: Servicer,
485 next_dir: Option<oio::Entry>,
486 active_lister: Vec<(Option<oio::Entry>, oio::Lister)>,
487}
488
489impl ServicerFlatLister {
490 fn new(ctx: OperationContext, srv: Servicer, path: &str) -> Self {
491 Self {
492 ctx,
493 srv,
494 next_dir: Some(oio::Entry::new(path, MetadataBuilder::dir().build())),
495 active_lister: vec![],
496 }
497 }
498}
499
500impl oio::List for ServicerFlatLister {
501 async fn next(&mut self) -> Result<Option<oio::Entry>> {
502 loop {
503 if let Some(de) = self.next_dir.take() {
504 let mut l = match self.srv.list(&self.ctx, de.path(), OpList::new()) {
505 Ok(v) => v,
506 Err(e) if e.kind() == ErrorKind::PermissionDenied => {
507 log::warn!(
508 "ServicerFlatLister skipping directory due to permission denied: {}",
509 de.path()
510 );
511 continue;
512 }
513 Err(e) if e.kind() == ErrorKind::NotFound => {
514 log::warn!(
515 "ServicerFlatLister skipping directory due to not found during listing: {}",
516 de.path()
517 );
518 continue;
519 }
520 Err(e) => return Err(e),
521 };
522 let first = loop {
523 match l.next().await {
524 Ok(v) => break v,
525 Err(e) if e.kind() == ErrorKind::NotFound => {
526 log::warn!(
527 "ServicerFlatLister skipping entry due to not found during listing: {}",
528 de.path()
529 );
530 continue;
531 }
532 Err(e) => return Err(e),
533 }
534 };
535 if let Some(v) = first {
536 self.active_lister.push((Some(de.clone()), l));
537
538 if v.mode().is_dir() {
539 if v.path() != de.path() {
540 self.next_dir = Some(v);
541 continue;
542 }
543 } else {
544 return Ok(Some(v));
545 }
546 }
547 }
548
549 if matches!(self.active_lister.last(), Some((None, _))) {
550 let _ = self.active_lister.pop();
551 continue;
552 }
553
554 let (de, lister) = match self.active_lister.last_mut() {
555 Some((de, lister)) => (de, lister),
556 None => return Ok(None),
557 };
558
559 match lister.next().await {
560 Err(e) if e.kind() == ErrorKind::NotFound => {
561 let path = de.as_ref().map(|entry| entry.path()).unwrap_or("<unknown>");
562 log::warn!(
563 "ServicerFlatLister skipping entry due to not found during recursive listing: {}",
564 path
565 );
566 continue;
567 }
568 Err(e) => return Err(e),
569 Ok(Some(v)) if v.mode().is_dir() => {
570 if v.path()
571 != de
572 .as_ref()
573 .expect("de must be present before listing")
574 .path()
575 {
576 self.next_dir = Some(v);
577 continue;
578 }
579 }
580 Ok(Some(v)) => return Ok(Some(v)),
581 Ok(None) => match de.take() {
582 Some(de) => return Ok(Some(de)),
583 None => {
584 let _ = self.active_lister.pop();
585 continue;
586 }
587 },
588 }
589 }
590 }
591}
592
593pub struct SimulateDeleter {
595 ctx: OperationContext,
596 srv: Servicer,
597 config: SimulateLayer,
598 inner: oio::Deleter,
599}
600
601impl SimulateDeleter {
602 pub fn new(
603 ctx: OperationContext,
604 srv: Servicer,
605 config: SimulateLayer,
606 inner: oio::Deleter,
607 ) -> Self {
608 Self {
609 ctx,
610 srv,
611 config,
612 inner,
613 }
614 }
615}
616
617impl oio::Delete for SimulateDeleter {
618 async fn delete(&mut self, path: &str, args: OpDelete) -> Result<()> {
619 if args.recursive() {
620 let cap = self.srv.capability();
621
622 if cap.delete_with_recursive {
623 return self.inner.delete(path, args).await;
624 }
625
626 if self.config.delete_recursive {
627 let service = SimulateService {
628 srv: self.srv.clone(),
629 config: self.config.clone(),
630 };
631 return service
632 .simulate_delete_with_recursive(&self.ctx, path, args, &mut self.inner)
633 .await;
634 }
635 }
636
637 self.inner.delete(path, args).await
638 }
639
640 async fn close(&mut self) -> Result<()> {
641 self.inner.close().await
642 }
643}
644
645#[cfg(test)]
646mod tests {
647 use std::sync::Mutex;
648
649 use super::*;
650
651 #[derive(Debug)]
652 struct MockService {
653 capability: Capability,
654 }
655
656 struct MockLister(bool);
657
658 impl oio::List for MockLister {
659 async fn next(&mut self) -> Result<Option<oio::Entry>> {
660 if self.0 {
661 return Ok(None);
662 }
663 self.0 = true;
664 Ok(Some(oio::Entry::new(
665 "parent/file",
666 MetadataBuilder::file(0).build(),
667 )))
668 }
669 }
670
671 impl Service for MockService {
672 type Reader = ();
673 type Writer = ();
674 type Lister = MockLister;
675 type Deleter = ();
676 type Copier = ();
677 type Composer = ();
678
679 fn info(&self) -> ServiceInfo {
680 ServiceInfo::with_scheme("mock")
681 }
682
683 fn capability(&self) -> Capability {
684 self.capability
685 }
686
687 async fn create_dir(
688 &self,
689 _: &OperationContext,
690 _: &str,
691 _: OpCreateDir,
692 ) -> Result<RpCreateDir> {
693 Err(Error::new(
694 ErrorKind::Unsupported,
695 "operation is not supported",
696 ))
697 }
698
699 async fn stat(&self, _: &OperationContext, _: &str, _: OpStat) -> Result<RpStat> {
700 Err(Error::new(
701 ErrorKind::Unsupported,
702 "operation is not supported",
703 ))
704 }
705
706 fn read(&self, _ctx: &OperationContext, _: &str, _: OpRead) -> Result<Self::Reader> {
707 Err(Error::new(
708 ErrorKind::Unsupported,
709 "operation is not supported",
710 ))
711 }
712
713 fn write(&self, _ctx: &OperationContext, _: &str, _: OpWrite) -> Result<Self::Writer> {
714 Err(Error::new(
715 ErrorKind::Unsupported,
716 "operation is not supported",
717 ))
718 }
719
720 fn delete(&self, _ctx: &OperationContext) -> Result<Self::Deleter> {
721 Err(Error::new(
722 ErrorKind::Unsupported,
723 "operation is not supported",
724 ))
725 }
726
727 fn list(&self, _ctx: &OperationContext, _: &str, _: OpList) -> Result<Self::Lister> {
728 if self.capability.list {
729 return Ok(MockLister(false));
730 }
731 Err(Error::new(ErrorKind::Unsupported, "list is not supported"))
732 }
733
734 fn copy(&self, _: &OperationContext, _: &str, _: &str, _: OpCopy) -> Result<Self::Copier> {
735 Err(Error::new(
736 ErrorKind::Unsupported,
737 "operation is not supported",
738 ))
739 }
740
741 async fn rename(
742 &self,
743 _: &OperationContext,
744 _: &str,
745 _: &str,
746 _: OpRename,
747 ) -> Result<RpRename> {
748 Err(Error::new(
749 ErrorKind::Unsupported,
750 "operation is not supported",
751 ))
752 }
753
754 async fn presign(&self, _: &OperationContext, _: &str, _: OpPresign) -> Result<RpPresign> {
755 Err(Error::new(
756 ErrorKind::Unsupported,
757 "operation is not supported",
758 ))
759 }
760 }
761
762 struct MockReader {
763 observed_range: Arc<Mutex<Option<BytesRange>>>,
764 }
765
766 impl oio::Read for MockReader {
767 async fn open(&self, range: BytesRange) -> Result<(RpRead, Box<dyn oio::ReadStreamDyn>)> {
768 *self.observed_range.lock().expect("mutex must not poison") = Some(range);
769 Ok((
770 RpRead::new({
771 let metadata = MetadataBuilder::file(0);
772 metadata.build()
773 }),
774 Box::new(Buffer::new()),
775 ))
776 }
777
778 async fn read(&self, range: BytesRange) -> Result<(RpRead, Buffer)> {
779 *self.observed_range.lock().expect("mutex must not poison") = Some(range);
780 Ok((
781 RpRead::new({
782 let metadata = MetadataBuilder::file(0);
783 metadata.build()
784 }),
785 Buffer::new(),
786 ))
787 }
788 }
789
790 #[test]
791 fn simulate_layer_exposes_read_with_suffix() {
792 let capability = Capability {
793 read: true,
794 read_with_suffix: false,
795 ..Default::default()
796 };
797 let srv = Arc::new(MockService { capability }) as Servicer;
798
799 let srv = SimulateLayer::default().apply_service(srv);
800
801 assert!(srv.capability().read_with_suffix);
802 }
803
804 #[test]
805 fn simulate_layer_can_disable_read_with_suffix() {
806 let capability = Capability {
807 read: true,
808 read_with_suffix: false,
809 ..Default::default()
810 };
811 let srv = Arc::new(MockService { capability }) as Servicer;
812
813 let srv = SimulateLayer::default()
814 .with_read_with_suffix(false)
815 .apply_service(srv);
816
817 assert!(!srv.capability().read_with_suffix);
818 }
819
820 #[tokio::test]
821 async fn simulate_stat_dir_without_recursive_list() -> Result<()> {
822 let capability = Capability {
823 stat: true,
824 list: true,
825 list_with_recursive: false,
826 ..Default::default()
827 };
828 let srv = Arc::new(MockService { capability }) as Servicer;
829 let srv = SimulateLayer::default().apply_service(srv);
830
831 let metadata = srv
832 .stat(&OperationContext::new(), "parent/", OpStat::default())
833 .await?
834 .into_metadata();
835
836 assert_eq!(metadata.mode(), EntryMode::DIR);
837 Ok(())
838 }
839
840 #[tokio::test]
841 async fn simulate_reader_uses_content_length_hint_for_suffix() -> Result<()> {
842 let observed_range = Arc::new(Mutex::new(None));
843 let (_, args, _) = options::ReadOptions {
844 content_length_hint: Some(42),
845 ..Default::default()
846 }
847 .into();
848 let reader = SimulateReader::new(
849 OperationContext::new(),
850 Arc::new(MockService {
851 capability: Capability::default(),
852 }),
853 "test".to_string(),
854 args,
855 Box::new(MockReader {
856 observed_range: observed_range.clone(),
857 }),
858 true,
859 );
860
861 oio::Read::read(&reader, BytesRange::suffix(10)).await?;
862
863 assert_eq!(
864 *observed_range.lock().expect("mutex must not poison"),
865 Some(BytesRange::new(32, Some(10)))
866 );
867
868 Ok(())
869 }
870}