use std::path::PathBuf;
use crate::domain::error::{WireError, WireResult};
use crate::infrastructure::wire_uri::WireUri;
use async_trait::async_trait;
#[async_trait]
pub trait Adapter: Send + Sync {
fn scheme(&self) -> &'static str;
async fn fetch(&self, uri: &WireUri) -> WireResult<serde_json::Value>;
}
pub struct FileAdapter;
impl FileAdapter {
pub async fn fetch_file(&self, raw_path: &str) -> WireResult<serde_json::Value> {
let resolved = resolve_file_path(raw_path)?;
let meta = std::fs::metadata(&resolved)
.map_err(|e| WireError::Storage(format!("file adapter: stat: {e}")))?;
if meta.is_dir() {
let newest = newest_child(&resolved)?;
let body = std::fs::read_to_string(&newest)
.map_err(|e| WireError::Storage(format!("file adapter: read: {e}")))?;
Ok(serde_json::json!({
"scheme": "file",
"kind": "newest_in_dir",
"dir": resolved.display().to_string(),
"path": newest.display().to_string(),
"body": body,
}))
} else {
let body = std::fs::read_to_string(&resolved)
.map_err(|e| WireError::Storage(format!("file adapter: read: {e}")))?;
Ok(serde_json::json!({
"scheme": "file",
"kind": "file",
"path": resolved.display().to_string(),
"body": body,
}))
}
}
}
#[async_trait]
impl Adapter for FileAdapter {
fn scheme(&self) -> &'static str {
"file"
}
async fn fetch(&self, uri: &WireUri) -> WireResult<serde_json::Value> {
let source_uri = uri.as_raw();
let rest = source_uri
.strip_prefix("file://")
.or_else(|| source_uri.strip_prefix("file:"))
.ok_or_else(|| WireError::Storage(format!("file adapter: bad uri: {source_uri}")))?;
self.fetch_file(rest).await
}
}
fn resolve_file_path(raw: &str) -> WireResult<PathBuf> {
let stripped = raw.split('#').next().unwrap_or(raw);
let expanded = if let Some(rest) = stripped.strip_prefix("~/") {
let home = std::env::var("HOME")
.map_err(|_| WireError::Storage("file adapter: HOME unset".to_string()))?;
PathBuf::from(home).join(rest)
} else {
PathBuf::from(stripped)
};
Ok(expanded)
}
fn newest_child(dir: &std::path::Path) -> WireResult<PathBuf> {
let mut entries: Vec<_> = std::fs::read_dir(dir)
.map_err(|e| WireError::Storage(format!("file adapter: read_dir: {e}")))?
.filter_map(|r| r.ok())
.filter(|e| e.path().is_file())
.collect();
if entries.is_empty() {
return Err(WireError::Storage(format!(
"file adapter: empty dir: {}",
dir.display()
)));
}
entries.sort_by_key(|e| {
e.metadata()
.and_then(|m| m.modified())
.ok()
.unwrap_or(std::time::SystemTime::UNIX_EPOCH)
});
Ok(entries
.last()
.map(|e| e.path())
.expect("non-empty sorted entries"))
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn file_adapter_reads_existing_file() {
let me = file!();
let abs = std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.parent()
.unwrap()
.parent()
.unwrap()
.join(me);
let uri = WireUri::parse(&format!("file://{}", abs.display())).unwrap();
let a = FileAdapter;
let v = a.fetch(&uri).await.unwrap();
let body = v["body"].as_str().unwrap();
assert!(body.contains("Layer 6 Adapter"));
}
#[tokio::test]
async fn file_adapter_rejects_non_file_uri() {
let a = FileAdapter;
let uri = WireUri::parse("ssh://nope/x").unwrap();
let r = a.fetch(&uri).await;
assert!(r.is_err());
}
}