#![cfg_attr(docsrs, feature(doc_cfg))]
#![deny(missing_docs)]
use std::collections::HashSet;
use std::sync::Arc;
use std::vec::IntoIter;
use opendal_core::raw::*;
use opendal_core::*;
#[derive(Clone, Debug, Default)]
pub struct ImmutableIndexLayer {
vec: Vec<String>,
}
impl ImmutableIndexLayer {
pub fn new() -> Self {
Self::default()
}
}
impl ImmutableIndexLayer {
pub fn insert(&mut self, key: String) {
self.vec.push(key);
}
pub fn extend_iter<I>(&mut self, iter: I)
where
I: IntoIterator<Item = String>,
{
self.vec.extend(iter);
}
}
impl Layer for ImmutableIndexLayer {
fn apply_service(&self, inner: Servicer) -> Servicer {
Arc::new(self.layer(inner))
}
}
impl ImmutableIndexLayer {
fn layer(&self, inner: Servicer) -> ImmutableIndexService {
ImmutableIndexService {
inner,
vec: self.vec.clone(),
}
}
}
#[doc(hidden)]
#[derive(Debug)]
pub struct ImmutableIndexService {
inner: Servicer,
vec: Vec<String>,
}
impl ImmutableIndexService {
fn children_flat(&self, path: &str) -> Vec<String> {
self.vec
.iter()
.filter(|v| v.starts_with(path) && v.as_str() != path)
.cloned()
.collect()
}
fn children_hierarchy(&self, path: &str) -> Vec<String> {
let mut res = HashSet::new();
for i in self.vec.iter() {
if !i.starts_with(path) {
continue;
}
if i == path {
continue;
}
match i[path.len()..].find('/') {
None => {
res.insert(i.to_string());
}
Some(idx) => {
let dir_idx = idx + 1 + path.len();
if dir_idx == i.len() {
res.insert(i.to_string());
} else {
res.insert(i[..dir_idx].to_string());
}
}
}
}
res.into_iter().collect()
}
}
impl Service for ImmutableIndexService {
type Reader = oio::Reader;
type Writer = oio::Writer;
type Lister = ImmutableDir;
type Deleter = oio::Deleter;
type Copier = oio::Copier;
fn info(&self) -> ServiceInfo {
self.inner.info()
}
fn capability(&self) -> Capability {
let mut capability = self.inner.capability();
capability.list = true;
capability.list_with_recursive = true;
capability
}
async fn create_dir(
&self,
ctx: &OperationContext,
path: &str,
args: OpCreateDir,
) -> Result<RpCreateDir> {
self.inner.create_dir(ctx, path, args).await
}
async fn stat(&self, ctx: &OperationContext, path: &str, args: OpStat) -> Result<RpStat> {
self.inner.stat(ctx, path, args).await
}
fn read(&self, ctx: &OperationContext, path: &str, args: OpRead) -> Result<Self::Reader> {
self.inner.read(ctx, path, args)
}
fn write(&self, ctx: &OperationContext, path: &str, args: OpWrite) -> Result<Self::Writer> {
self.inner.write(ctx, path, args)
}
fn copy(
&self,
ctx: &OperationContext,
from: &str,
to: &str,
args: OpCopy,
opts: OpCopier,
) -> Result<Self::Copier> {
self.inner.copy(ctx, from, to, args, opts)
}
fn list(&self, _ctx: &OperationContext, path: &str, args: OpList) -> Result<Self::Lister> {
let mut path = path;
if path == "/" {
path = ""
}
let idx = if args.recursive() {
self.children_flat(path)
} else {
self.children_hierarchy(path)
};
Ok(ImmutableDir::new(idx))
}
fn delete(&self, ctx: &OperationContext) -> Result<Self::Deleter> {
self.inner.delete(ctx)
}
async fn rename(
&self,
ctx: &OperationContext,
from: &str,
to: &str,
args: OpRename,
) -> Result<RpRename> {
self.inner.rename(ctx, from, to, args).await
}
async fn presign(
&self,
ctx: &OperationContext,
path: &str,
args: OpPresign,
) -> Result<RpPresign> {
self.inner.presign(ctx, path, args).await
}
}
#[doc(hidden)]
pub struct ImmutableDir {
idx: IntoIter<String>,
}
impl ImmutableDir {
fn new(idx: Vec<String>) -> Self {
Self {
idx: idx.into_iter(),
}
}
fn inner_next(&mut self) -> Option<oio::Entry> {
self.idx.next().map(|v| {
let mode = if v.ends_with('/') {
EntryMode::DIR
} else {
EntryMode::FILE
};
let meta = Metadata::new(mode);
oio::Entry::with(v, meta)
})
}
}
impl oio::List for ImmutableDir {
async fn next(&mut self) -> Result<Option<oio::Entry>> {
Ok(self.inner_next())
}
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use std::sync::Arc;
use super::*;
use futures::TryStreamExt;
use log::debug;
use logforth::append::Testing;
use logforth::filter::rustlog::RustLogFilterBuilder;
use logforth::layout::TextLayout;
#[derive(Debug)]
struct MockService;
impl Service for MockService {
type Reader = ();
type Writer = ();
type Lister = ();
type Deleter = ();
type Copier = ();
fn info(&self) -> ServiceInfo {
ServiceInfo::with_scheme("mock")
}
fn capability(&self) -> Capability {
Capability {
list: true,
list_with_recursive: true,
..Default::default()
}
}
async fn create_dir(
&self,
_: &OperationContext,
_: &str,
_: OpCreateDir,
) -> Result<RpCreateDir> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
}
async fn stat(&self, _: &OperationContext, _: &str, _: OpStat) -> Result<RpStat> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
}
fn read(&self, _ctx: &OperationContext, _: &str, _: OpRead) -> Result<Self::Reader> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
}
fn write(&self, _ctx: &OperationContext, _: &str, _: OpWrite) -> Result<Self::Writer> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
}
fn delete(&self, _ctx: &OperationContext) -> Result<Self::Deleter> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
}
fn list(&self, _ctx: &OperationContext, _: &str, _: OpList) -> Result<Self::Lister> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
}
fn copy(
&self,
_: &OperationContext,
_: &str,
_: &str,
_: OpCopy,
_: OpCopier,
) -> Result<Self::Copier> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
}
async fn rename(
&self,
_: &OperationContext,
_: &str,
_: &str,
_: OpRename,
) -> Result<RpRename> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
}
async fn presign(&self, _: &OperationContext, _: &str, _: OpPresign) -> Result<RpPresign> {
Err(Error::new(
ErrorKind::Unsupported,
"operation is not supported",
))
}
}
fn build_operator(layer: ImmutableIndexLayer) -> Operator {
Operator::from_parts(OperationContext::default(), Arc::new(MockService)).layer(layer)
}
fn setup() {
let _ = logforth::starter_log::builder()
.dispatch(|d| {
d.filter(RustLogFilterBuilder::from_default_env().build())
.append(Testing::default().with_layout(TextLayout::default()))
})
.try_apply();
}
#[tokio::test]
async fn test_list() -> Result<()> {
setup();
let mut iil = ImmutableIndexLayer::default();
for i in ["file", "dir/", "dir/file", "dir_without_prefix/file"] {
iil.insert(i.to_string())
}
let op = build_operator(iil);
let mut map = HashMap::new();
let mut set = HashSet::new();
let mut ds = op.lister("").await?;
while let Some(entry) = ds.try_next().await? {
debug!("got entry: {}", entry.path());
assert!(
set.insert(entry.path().to_string()),
"duplicated value: {}",
entry.path()
);
map.insert(entry.path().to_string(), entry.metadata().mode());
}
assert_eq!(map["file"], EntryMode::FILE);
assert_eq!(map["dir/"], EntryMode::DIR);
assert_eq!(map["dir_without_prefix/"], EntryMode::DIR);
Ok(())
}
#[tokio::test]
async fn test_scan() -> Result<()> {
setup();
let mut iil = ImmutableIndexLayer::default();
for i in ["file", "dir/", "dir/file", "dir_without_prefix/file"] {
iil.insert(i.to_string())
}
let op = build_operator(iil);
let mut ds = op.lister_with("/").recursive(true).await?;
let mut set = HashSet::new();
let mut map = HashMap::new();
while let Some(entry) = ds.try_next().await? {
debug!("got entry: {}", entry.path());
assert!(
set.insert(entry.path().to_string()),
"duplicated value: {}",
entry.path()
);
map.insert(entry.path().to_string(), entry.metadata().mode());
}
debug!("current files: {map:?}");
assert_eq!(map["file"], EntryMode::FILE);
assert_eq!(map["dir/"], EntryMode::DIR);
assert_eq!(map["dir_without_prefix/file"], EntryMode::FILE);
Ok(())
}
#[tokio::test]
async fn test_list_dir() -> Result<()> {
setup();
let mut iil = ImmutableIndexLayer::default();
for i in [
"dataset/stateful/ontime_2007_200.csv",
"dataset/stateful/ontime_2008_200.csv",
"dataset/stateful/ontime_2009_200.csv",
] {
iil.insert(i.to_string())
}
let op = build_operator(iil);
let mut map = HashMap::new();
let mut set = HashSet::new();
let mut ds = op.lister("/").await?;
while let Some(entry) = ds.try_next().await? {
assert!(
set.insert(entry.path().to_string()),
"duplicated value: {}",
entry.path()
);
map.insert(entry.path().to_string(), entry.metadata().mode());
}
assert_eq!(map.len(), 1);
assert_eq!(map["dataset/"], EntryMode::DIR);
let mut map = HashMap::new();
let mut set = HashSet::new();
let mut ds = op.lister("dataset/stateful/").await?;
while let Some(entry) = ds.try_next().await? {
assert!(
set.insert(entry.path().to_string()),
"duplicated value: {}",
entry.path()
);
map.insert(entry.path().to_string(), entry.metadata().mode());
}
assert_eq!(map["dataset/stateful/ontime_2007_200.csv"], EntryMode::FILE);
assert_eq!(map["dataset/stateful/ontime_2008_200.csv"], EntryMode::FILE);
assert_eq!(map["dataset/stateful/ontime_2009_200.csv"], EntryMode::FILE);
Ok(())
}
#[tokio::test]
async fn test_walk_top_down_dir() -> Result<()> {
setup();
let mut iil = ImmutableIndexLayer::default();
for i in [
"dataset/stateful/ontime_2007_200.csv",
"dataset/stateful/ontime_2008_200.csv",
"dataset/stateful/ontime_2009_200.csv",
] {
iil.insert(i.to_string())
}
let op = build_operator(iil);
let mut ds = op.lister_with("/").recursive(true).await?;
let mut map = HashMap::new();
let mut set = HashSet::new();
while let Some(entry) = ds.try_next().await? {
assert!(
set.insert(entry.path().to_string()),
"duplicated value: {}",
entry.path()
);
map.insert(entry.path().to_string(), entry.metadata().mode());
}
debug!("current files: {map:?}");
assert_eq!(map["dataset/stateful/ontime_2007_200.csv"], EntryMode::FILE);
assert_eq!(map["dataset/stateful/ontime_2008_200.csv"], EntryMode::FILE);
assert_eq!(map["dataset/stateful/ontime_2009_200.csv"], EntryMode::FILE);
Ok(())
}
}