use delta_kernel::object_store::path::Path;
use delta_kernel::object_store::{ObjectMeta, Result as ObjectStoreResult};
use super::generic_error;
pub(crate) struct RestListPage {
pub(crate) objects: Vec<ObjectMeta>,
pub(crate) next_page_token: Option<String>,
}
#[derive(Debug, Clone)]
pub struct RestEndpointConfig {
pub files_prefix: String,
pub directories_prefix: String,
pub page_token_param: String,
pub start_from_param: String,
pub recursive_param: String,
pub overwrite_param: String,
pub contents_field: String,
pub next_page_token_field: String,
pub entry_path_field: String,
pub entry_size_field: String,
pub entry_is_directory_field: String,
pub entry_last_modified_field: String,
pub entry_strip_prefix: Option<String>,
}
fn join_url(base_url: &str, prefix: &str, path: &str) -> String {
let mut segments = vec![base_url.trim_end_matches('/')];
let prefix = prefix.trim_matches('/');
if !prefix.is_empty() {
segments.push(prefix);
}
let path = path.trim_start_matches('/');
if !path.is_empty() {
segments.push(path);
}
segments.join("/")
}
impl RestEndpointConfig {
pub(crate) fn file_url(&self, base_url: &str, path: &str) -> String {
join_url(base_url, &self.files_prefix, path)
}
pub(crate) fn directory_url(&self, base_url: &str, path: &str) -> String {
join_url(base_url, &self.directories_prefix, path)
}
pub(crate) fn list_query(
&self,
page_token: Option<&str>,
start_from: Option<&str>,
recursive: bool,
) -> Vec<(String, String)> {
let mut q = Vec::new();
if let Some(t) = page_token {
q.push((self.page_token_param.clone(), t.to_string()));
}
if let Some(s) = start_from {
q.push((self.start_from_param.clone(), s.to_string()));
}
if recursive {
q.push((self.recursive_param.clone(), "true".to_string()));
}
q
}
pub(crate) fn parse_list(&self, body: &[u8]) -> ObjectStoreResult<RestListPage> {
let root: serde_json::Value = serde_json::from_slice(body).map_err(generic_error)?;
let mut objects = Vec::new();
if let Some(entries) = root.get(&self.contents_field).and_then(|v| v.as_array()) {
for entry in entries {
let is_dir = entry
.get(&self.entry_is_directory_field)
.and_then(|v| v.as_bool())
.unwrap_or(false);
if is_dir {
continue;
}
let Some(path) = entry.get(&self.entry_path_field).and_then(|v| v.as_str()) else {
return Err(generic_error(format!(
"list entry is missing the `{}` field",
self.entry_path_field
)));
};
let path = match &self.entry_strip_prefix {
Some(prefix) => match path.strip_prefix(prefix.as_str()) {
Some(rest) => rest.trim_start_matches('/'),
None => {
return Err(generic_error(format!(
"list entry `{path}` is outside the configured prefix `{prefix}`"
)))
}
},
None => path,
};
let size = match entry.get(&self.entry_size_field) {
None => 0,
Some(v) => v.as_u64().ok_or_else(|| {
generic_error(format!(
"list entry field `{}` is not a non-negative integer",
self.entry_size_field
))
})?,
};
let last_modified = match entry.get(&self.entry_last_modified_field) {
None => chrono::DateTime::UNIX_EPOCH,
Some(v) => {
let ms = v.as_u64().ok_or_else(|| {
generic_error(format!(
"list entry field `{}` is not a non-negative integer",
self.entry_last_modified_field
))
})?;
(std::time::SystemTime::UNIX_EPOCH + std::time::Duration::from_millis(ms))
.into()
}
};
let location = Path::parse(path).map_err(|e| {
generic_error(format!("list entry path `{path}` is not a valid path: {e}"))
})?;
objects.push(ObjectMeta {
location,
last_modified,
size,
e_tag: None,
version: None,
});
}
}
let next_page_token = root
.get(&self.next_page_token_field)
.and_then(|v| v.as_str())
.filter(|s| !s.is_empty())
.map(|s| s.to_string());
Ok(RestListPage {
objects,
next_page_token,
})
}
pub(crate) fn put_query(&self, overwrite: bool) -> Vec<(String, String)> {
vec![(self.overwrite_param.clone(), overwrite.to_string())]
}
}