1use std::path::{Path, PathBuf};
2use std::sync::atomic::{AtomicU32, Ordering};
3use std::sync::Mutex;
4
5use crate::types::{ArtifactDescriptor, ArtifactFormat, ArtifactSubject};
6use crate::{BinocError, BinocResult, DataAccess, ItemRef};
7
8pub struct LocalDataAccess {
21 #[cfg(not(target_family = "wasm"))]
22 _session_dir: Option<tempfile::TempDir>,
23 #[cfg(target_family = "wasm")]
24 _session_dir: Option<PathBuf>,
25 data_root: PathBuf,
26 external_root: Option<PathBuf>,
27 workspace_counter: AtomicU32,
28 #[cfg(not(target_family = "wasm"))]
29 workspaces: Mutex<Vec<tempfile::TempDir>>,
30 #[cfg(target_family = "wasm")]
31 workspaces: Mutex<Vec<PathBuf>>,
32 #[cfg(not(target_family = "wasm"))]
33 provide_dir: Mutex<Option<tempfile::TempDir>>,
34 #[cfg(target_family = "wasm")]
35 provide_dir: Mutex<Option<PathBuf>>,
36 path_policy: PathPolicy,
37}
38
39enum PathPolicy {
40 Unrestricted,
41 Restricted {
42 snapshot_a: PathBuf,
43 snapshot_b: PathBuf,
44 extra_allowed: Mutex<Vec<PathBuf>>,
45 },
46}
47
48fn artifacts_dir(data_root: &Path) -> PathBuf {
49 data_root.join(".artifacts")
50}
51
52fn safe_name(s: &str) -> String {
53 s.bytes()
54 .map(|b| {
55 if b.is_ascii_alphanumeric() || b == b'-' || b == b'_' || b == b'.' {
56 (b as char).to_string()
57 } else {
58 format!("%{b:02x}")
59 }
60 })
61 .collect()
62}
63
64fn subject_dir_name(subject: ArtifactSubject) -> &'static str {
65 match subject {
66 ArtifactSubject::Left => "left",
67 ArtifactSubject::Right => "right",
68 ArtifactSubject::Pair => "pair",
69 }
70}
71
72#[cfg(target_family = "wasm")]
73fn policy_path(path: &Path) -> BinocResult<PathBuf> {
74 let mut out = PathBuf::new();
75 for component in path.components() {
76 match component {
77 std::path::Component::Prefix(_) => {}
78 std::path::Component::RootDir => {}
79 std::path::Component::CurDir => {}
80 std::path::Component::ParentDir => {
81 out.pop();
82 }
83 std::path::Component::Normal(part) => out.push(part),
84 }
85 }
86 Ok(out)
87}
88
89#[cfg(not(target_family = "wasm"))]
90fn policy_path(path: &Path) -> BinocResult<PathBuf> {
91 std::fs::canonicalize(path).map_err(BinocError::Io)
92}
93
94#[cfg(target_family = "wasm")]
95fn wasm_session_dir(prefix: &str) -> PathBuf {
96 let id: u64 = rand::random();
97 PathBuf::from(".binoc-tmp").join(format!("{prefix}-{id:016x}"))
98}
99
100#[cfg(target_family = "wasm")]
101fn temp_path(dir: &Path) -> &Path {
102 dir
103}
104
105#[cfg(not(target_family = "wasm"))]
106fn temp_path(dir: &tempfile::TempDir) -> &Path {
107 dir.path()
108}
109
110fn path_is_within(path: &Path, root: &Path) -> bool {
112 path.starts_with(root)
113}
114
115fn item_ref_from_physical(physical: &Path, logical: &str) -> ItemRef {
116 ItemRef {
117 logical_path: logical.to_string(),
118 is_dir: physical.is_dir(),
119 content_hash: None,
120 size: None,
121 media_type: None,
122 projection_hint: Default::default(),
123 tabular_parse: None,
124 handle: physical.to_string_lossy().to_string(),
125 }
126}
127
128impl LocalDataAccess {
129 #[cfg(not(target_family = "wasm"))]
130 pub fn new() -> Self {
131 let session = tempfile::tempdir().expect("failed to create session temp dir");
132 let data_root = session.path().to_path_buf();
133 Self {
134 _session_dir: Some(session),
135 data_root,
136 external_root: None,
137 workspace_counter: AtomicU32::new(0),
138 workspaces: Mutex::new(Vec::new()),
139 provide_dir: Mutex::new(None),
140 path_policy: PathPolicy::Unrestricted,
141 }
142 }
143
144 #[cfg(target_family = "wasm")]
145 pub fn new() -> Self {
146 let data_root = wasm_session_dir("session");
147 std::fs::create_dir_all(&data_root).expect("failed to create session temp dir");
148 Self {
149 _session_dir: Some(data_root.clone()),
150 data_root,
151 external_root: None,
152 workspace_counter: AtomicU32::new(0),
153 workspaces: Mutex::new(Vec::new()),
154 provide_dir: Mutex::new(None),
155 path_policy: PathPolicy::Unrestricted,
156 }
157 }
158
159 #[cfg(not(target_family = "wasm"))]
163 pub fn new_for_diff(snapshot_a: &Path, snapshot_b: &Path) -> BinocResult<Self> {
164 let session = tempfile::tempdir().map_err(BinocError::Io)?;
165 let data_root = session.path().to_path_buf();
166 let snap_a = std::fs::canonicalize(snapshot_a).map_err(BinocError::Io)?;
167 let snap_b = std::fs::canonicalize(snapshot_b).map_err(BinocError::Io)?;
168 let data_root_canon = std::fs::canonicalize(&data_root).map_err(BinocError::Io)?;
169 Ok(Self {
170 _session_dir: Some(session),
171 data_root,
172 external_root: None,
173 workspace_counter: AtomicU32::new(0),
174 workspaces: Mutex::new(Vec::new()),
175 provide_dir: Mutex::new(None),
176 path_policy: PathPolicy::Restricted {
177 snapshot_a: snap_a,
178 snapshot_b: snap_b,
179 extra_allowed: Mutex::new(vec![data_root_canon]),
180 },
181 })
182 }
183
184 #[cfg(target_family = "wasm")]
185 pub fn new_for_diff(snapshot_a: &Path, snapshot_b: &Path) -> BinocResult<Self> {
186 let data_root = wasm_session_dir("session");
187 std::fs::create_dir_all(&data_root).map_err(BinocError::Io)?;
188 let snap_a = policy_path(snapshot_a)?;
189 let snap_b = policy_path(snapshot_b)?;
190 let data_root_canon = policy_path(&data_root)?;
191 Ok(Self {
192 _session_dir: Some(data_root),
193 data_root: data_root_canon.clone(),
194 external_root: None,
195 workspace_counter: AtomicU32::new(0),
196 workspaces: Mutex::new(Vec::new()),
197 provide_dir: Mutex::new(None),
198 path_policy: PathPolicy::Restricted {
199 snapshot_a: snap_a,
200 snapshot_b: snap_b,
201 extra_allowed: Mutex::new(vec![data_root_canon]),
202 },
203 })
204 }
205
206 pub fn for_plugin(data_root: PathBuf, workspace: PathBuf) -> Self {
210 Self {
211 _session_dir: None,
212 data_root,
213 external_root: Some(workspace),
214 workspace_counter: AtomicU32::new(0),
215 workspaces: Mutex::new(Vec::new()),
216 provide_dir: Mutex::new(None),
217 path_policy: PathPolicy::Unrestricted,
218 }
219 }
220
221 pub fn with_data_root(data_root: PathBuf) -> Self {
224 Self {
225 _session_dir: None,
226 data_root,
227 external_root: None,
228 workspace_counter: AtomicU32::new(0),
229 workspaces: Mutex::new(Vec::new()),
230 provide_dir: Mutex::new(None),
231 path_policy: PathPolicy::Unrestricted,
232 }
233 }
234
235 fn record_allowed_if_restricted(&self, path: &Path) -> BinocResult<()> {
236 if let PathPolicy::Restricted { extra_allowed, .. } = &self.path_policy {
237 let c = policy_path(path)?;
238 extra_allowed.lock().unwrap().push(c);
239 }
240 Ok(())
241 }
242
243 fn enforce_path_policy_resolved(&self, resolved: &Path) -> BinocResult<()> {
244 match &self.path_policy {
245 PathPolicy::Unrestricted => Ok(()),
246 PathPolicy::Restricted {
247 snapshot_a,
248 snapshot_b,
249 extra_allowed,
250 } => {
251 if path_is_within(resolved, snapshot_a) || path_is_within(resolved, snapshot_b) {
252 return Ok(());
253 }
254 let roots = extra_allowed.lock().unwrap();
255 for root in roots.iter() {
256 if path_is_within(resolved, root) {
257 return Ok(());
258 }
259 }
260 Err(BinocError::PathPolicy(format!(
261 "path must stay under snapshot directories or session workspace: {}",
262 resolved.display()
263 )))
264 }
265 }
266 }
267
268 fn enforce_path_policy(&self, physical: &Path) -> BinocResult<()> {
270 let resolved = policy_path(physical)?;
271 self.enforce_path_policy_resolved(&resolved)
272 }
273
274 #[cfg(not(target_family = "wasm"))]
276 fn enforce_policy_for_read_path(&self, path: &Path) -> BinocResult<()> {
277 match &self.path_policy {
278 PathPolicy::Unrestricted => Ok(()),
279 PathPolicy::Restricted { .. } => {
280 if let Ok(c) = std::fs::canonicalize(path) {
281 return self.enforce_path_policy_resolved(&c);
282 }
283 let mut probe: Option<&Path> = Some(path);
284 while let Some(p) = probe {
285 if p.as_os_str().is_empty() {
286 break;
287 }
288 if p.exists() {
289 let base = std::fs::canonicalize(p).map_err(BinocError::Io)?;
290 self.enforce_path_policy_resolved(&base)?;
291 return Ok(());
292 }
293 probe = p.parent();
294 }
295 Err(BinocError::PathPolicy(format!(
296 "cannot resolve path under session: {}",
297 path.display()
298 )))
299 }
300 }
301 }
302
303 #[cfg(target_family = "wasm")]
304 fn enforce_policy_for_read_path(&self, path: &Path) -> BinocResult<()> {
305 match &self.path_policy {
306 PathPolicy::Unrestricted => Ok(()),
307 PathPolicy::Restricted { .. } => self.enforce_path_policy_resolved(&policy_path(path)?),
308 }
309 }
310
311 fn ensure_provide_dir(&self) -> BinocResult<PathBuf> {
312 if let Some(root) = &self.external_root {
313 let d = root.join("_provide");
314 std::fs::create_dir_all(&d).map_err(BinocError::Io)?;
315 self.record_allowed_if_restricted(&d)?;
316 return Ok(d);
317 }
318 let mut guard = self.provide_dir.lock().unwrap();
319 if guard.is_none() {
320 #[cfg(not(target_family = "wasm"))]
321 let dir = tempfile::tempdir().map_err(BinocError::Io)?;
322 #[cfg(target_family = "wasm")]
323 let dir = {
324 let dir = wasm_session_dir("provide");
325 std::fs::create_dir_all(&dir).map_err(BinocError::Io)?;
326 dir
327 };
328 self.record_allowed_if_restricted(temp_path(&dir))?;
329 *guard = Some(dir);
330 }
331 Ok(temp_path(guard.as_ref().unwrap()).to_path_buf())
332 }
333}
334
335impl Default for LocalDataAccess {
336 fn default() -> Self {
337 Self::new()
338 }
339}
340
341impl DataAccess for LocalDataAccess {
342 fn read_bytes(&self, item: &ItemRef) -> BinocResult<Vec<u8>> {
343 let p = Path::new(&item.handle);
344 self.enforce_policy_for_read_path(p)?;
345 std::fs::read(p).map_err(BinocError::Io)
346 }
347
348 fn open_read(&self, item: &ItemRef) -> BinocResult<Box<dyn std::io::Read + Send>> {
349 let p = Path::new(&item.handle);
350 self.enforce_policy_for_read_path(p)?;
351 let file = std::fs::File::open(p).map_err(BinocError::Io)?;
352 Ok(Box::new(file))
353 }
354
355 fn local_path(&self, item: &ItemRef) -> BinocResult<PathBuf> {
356 let p = PathBuf::from(&item.handle);
357 self.enforce_policy_for_read_path(&p)?;
358 Ok(p)
359 }
360
361 fn provide(&self, logical_path: &str, content: &[u8]) -> BinocResult<ItemRef> {
362 let dir = self.ensure_provide_dir()?;
363 let safe_name = logical_path.replace(['/', '\\'], "_");
364 let file_path = dir.join(&safe_name);
365 std::fs::write(&file_path, content).map_err(BinocError::Io)?;
366 self.enforce_path_policy(&file_path)?;
367 Ok(item_ref_from_physical(&file_path, logical_path))
368 }
369
370 fn workspace(&self) -> BinocResult<PathBuf> {
371 if let Some(root) = &self.external_root {
372 let n = self.workspace_counter.fetch_add(1, Ordering::Relaxed);
373 let subdir = root.join(format!("ws-{n}"));
374 std::fs::create_dir_all(&subdir).map_err(BinocError::Io)?;
375 self.record_allowed_if_restricted(&subdir)?;
376 return Ok(subdir);
377 }
378 #[cfg(not(target_family = "wasm"))]
379 let dir = tempfile::tempdir().map_err(BinocError::Io)?;
380 #[cfg(target_family = "wasm")]
381 let dir = {
382 let dir = wasm_session_dir("workspace");
383 std::fs::create_dir_all(&dir).map_err(BinocError::Io)?;
384 dir
385 };
386 let path = temp_path(&dir).to_path_buf();
387 self.record_allowed_if_restricted(&path)?;
388 self.workspaces.lock().unwrap().push(dir);
389 Ok(path)
390 }
391
392 fn register_local(&self, physical: &Path, logical: &str) -> BinocResult<ItemRef> {
393 self.enforce_path_policy(physical)?;
394 Ok(item_ref_from_physical(physical, logical))
395 }
396
397 fn publish_artifact(
398 &self,
399 format: &ArtifactFormat,
400 subject: ArtifactSubject,
401 producer: &str,
402 data: &[u8],
403 ) -> BinocResult<ArtifactDescriptor> {
404 let id: u64 = rand::random();
405 let dir = artifacts_dir(&self.data_root)
406 .join(safe_name(&format.package))
407 .join(safe_name(&format.name))
408 .join(format!("v{}", format.version))
409 .join(subject_dir_name(subject));
410 std::fs::create_dir_all(&dir).map_err(BinocError::Io)?;
411 let filename = format!("{}-{id:016x}", safe_name(producer));
412 let handle = dir.join(filename).to_string_lossy().to_string();
413 std::fs::write(&handle, data).map_err(BinocError::Io)?;
414 Ok(ArtifactDescriptor {
415 format: format.clone(),
416 subject,
417 producer: producer.to_string(),
418 handle,
419 })
420 }
421
422 fn get_artifact(&self, descriptor: &ArtifactDescriptor) -> BinocResult<Option<Vec<u8>>> {
423 let path = PathBuf::from(&descriptor.handle);
424 self.enforce_policy_for_read_path(&path)?;
425 match std::fs::read(&path) {
426 Ok(data) => Ok(Some(data)),
427 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
428 Err(e) => Err(BinocError::Io(e)),
429 }
430 }
431
432 fn data_root(&self) -> BinocResult<PathBuf> {
433 Ok(self.data_root.clone())
434 }
435}
436
437#[cfg(test)]
438mod tests {
439 use super::*;
440
441 #[test]
442 fn publish_and_get_artifact_round_trip() {
443 let da = LocalDataAccess::new();
444 let fmt = ArtifactFormat::new("binoc", "tabular", 1);
445 let desc = da
446 .publish_artifact(&fmt, ArtifactSubject::Left, "binoc.csv", b"hello world")
447 .unwrap();
448 assert_eq!(desc.format, fmt);
449 assert_eq!(desc.subject, ArtifactSubject::Left);
450 assert_eq!(desc.producer, "binoc.csv");
451 let loaded = da.get_artifact(&desc).unwrap();
452 assert_eq!(loaded, Some(b"hello world".to_vec()));
453 }
454
455 #[test]
456 fn get_artifact_missing_returns_none() {
457 let da = LocalDataAccess::new();
458 let desc = ArtifactDescriptor {
459 format: ArtifactFormat::new("nonexistent", "thing", 1),
460 subject: ArtifactSubject::Pair,
461 producer: "test".into(),
462 handle: "/tmp/does-not-exist-binoc-test".into(),
463 };
464 assert_eq!(da.get_artifact(&desc).unwrap(), None);
465 }
466
467 #[test]
468 fn cross_instance_artifact_visibility() {
469 let da = LocalDataAccess::new();
470 let fmt = ArtifactFormat::new("binoc", "tabular", 1);
471 let desc = da
472 .publish_artifact(&fmt, ArtifactSubject::Right, "binoc.csv", b"shared-value")
473 .unwrap();
474 let data_root = da.data_root().unwrap();
475
476 let plugin_da = LocalDataAccess::with_data_root(data_root);
477 let loaded = plugin_da.get_artifact(&desc).unwrap();
478 assert_eq!(loaded, Some(b"shared-value".to_vec()));
479 }
480
481 #[test]
482 fn for_plugin_shares_artifacts() {
483 let da = LocalDataAccess::new();
484 let data_root = da.data_root().unwrap();
485 let ws = da.workspace().unwrap();
486
487 let plugin_da = LocalDataAccess::for_plugin(data_root, ws);
488 let fmt = ArtifactFormat::new("myplugin", "schema", 1);
489 let desc = plugin_da
490 .publish_artifact(&fmt, ArtifactSubject::Pair, "myplugin", b"plugin-data")
491 .unwrap();
492
493 let loaded = da.get_artifact(&desc).unwrap();
494 assert_eq!(loaded, Some(b"plugin-data".to_vec()));
495 }
496
497 #[test]
498 fn data_root_returns_valid_path() {
499 let da = LocalDataAccess::new();
500 let root = da.data_root().unwrap();
501 assert!(root.exists());
502 }
503
504 #[test]
505 fn restricted_rejects_register_outside_snapshots() {
506 let tmp_a = tempfile::tempdir().unwrap();
507 let tmp_b = tempfile::tempdir().unwrap();
508 let outside = tempfile::tempdir().unwrap();
509 std::fs::write(outside.path().join("x.txt"), b"x").unwrap();
510
511 let da = LocalDataAccess::new_for_diff(tmp_a.path(), tmp_b.path()).unwrap();
512 let p = outside.path().join("x.txt");
513 let err = da.register_local(&p, "x.txt").unwrap_err();
514 assert!(matches!(err, BinocError::PathPolicy(_)));
515 }
516
517 #[test]
518 fn restricted_allows_register_under_snapshot() {
519 let tmp_a = tempfile::tempdir().unwrap();
520 let tmp_b = tempfile::tempdir().unwrap();
521 let f = tmp_a.path().join("f.txt");
522 std::fs::write(&f, b"ok").unwrap();
523
524 let da = LocalDataAccess::new_for_diff(tmp_a.path(), tmp_b.path()).unwrap();
525 da.register_local(&f, "f.txt").unwrap();
526 }
527}