zen_engine/loader/
filesystem.rs1use std::fs::File;
2use std::future::Future;
3use std::io::BufReader;
4use std::path::{Path, PathBuf};
5use std::pin::Pin;
6use std::sync::Arc;
7
8use serde::{Deserialize, Serialize};
9
10use crate::loader::{DecisionLoader, LoaderError, LoaderResponse};
11use crate::model::DecisionContent;
12
13#[derive(Debug)]
15pub struct FilesystemLoader {
16 root: String,
17}
18
19#[derive(Serialize, Deserialize)]
20pub struct FilesystemLoaderOptions<R: Into<String>> {
21 pub root: R,
22}
23
24impl FilesystemLoader {
25 pub fn new<R>(options: FilesystemLoaderOptions<R>) -> Self
26 where
27 R: Into<String>,
28 {
29 Self {
30 root: options.root.into(),
31 }
32 }
33
34 fn key_to_path<K: AsRef<str>>(&self, key: K) -> PathBuf {
35 Path::new(&self.root).join(key.as_ref())
36 }
37
38 fn read_content<K: AsRef<str>>(&self, key: K) -> LoaderResponse {
39 let path = self.key_to_path(key.as_ref());
40 if !Path::exists(&path) {
41 return Err(LoaderError::NotFound(String::from(key.as_ref())));
42 }
43
44 let file = File::open(path).map_err(|e| LoaderError::Internal {
45 key: String::from(key.as_ref()),
46 source: e.into(),
47 })?;
48
49 let reader = BufReader::new(file);
50 let result: DecisionContent =
51 serde_json::from_reader(reader).map_err(|e| LoaderError::Internal {
52 key: String::from(key.as_ref()),
53 source: e.into(),
54 })?;
55
56 Ok(Arc::new(result))
57 }
58}
59
60impl DecisionLoader for FilesystemLoader {
61 fn load<'a>(
62 &'a self,
63 key: &'a str,
64 ) -> Pin<Box<dyn Future<Output = LoaderResponse> + 'a + Send>> {
65 Box::pin(async move { self.read_content(key) })
66 }
67
68 fn load_sync(&self, key: &str) -> Option<LoaderResponse> {
69 Some(self.read_content(key))
70 }
71
72 fn keys(&self) -> Option<Vec<Arc<str>>> {
73 let root = Path::new(&self.root);
74 let mut keys = Vec::new();
75 let mut stack = vec![root.to_path_buf()];
76 while let Some(dir) = stack.pop() {
77 let Ok(entries) = std::fs::read_dir(&dir) else {
78 continue;
79 };
80 for entry in entries.flatten() {
81 let path = entry.path();
82 if path.is_dir() {
83 stack.push(path);
84 } else if path.extension().and_then(|e| e.to_str()) == Some("json") {
85 let key = path.strip_prefix(root).ok().and_then(|rel| {
86 rel.components()
87 .map(|component| component.as_os_str().to_str())
88 .collect::<Option<Vec<_>>>()
89 .map(|segments| segments.join("/"))
90 });
91 if let Some(key) = key {
92 keys.push(Arc::from(key));
93 }
94 }
95 }
96 }
97 Some(keys)
98 }
99}
100
101#[cfg(test)]
102mod tests {
103 use super::*;
104
105 fn test_loader() -> FilesystemLoader {
106 let root = Path::new(env!("CARGO_MANIFEST_DIR"))
107 .parent()
108 .unwrap()
109 .parent()
110 .unwrap()
111 .join("test-data");
112 FilesystemLoader::new(FilesystemLoaderOptions {
113 root: root.to_string_lossy().to_string(),
114 })
115 }
116
117 #[tokio::test]
118 async fn load_and_load_sync_resolve_existing_key() {
119 let loader = test_loader();
120
121 assert!(loader.load("table.json").await.is_ok());
122 assert!(loader.load_sync("table.json").unwrap().is_ok());
123 }
124
125 #[tokio::test]
126 async fn load_reports_missing_key() {
127 let loader = test_loader();
128
129 assert!(loader.load("missing.json").await.is_err());
130 assert!(loader.load_sync("missing.json").unwrap().is_err());
131 }
132}