Skip to main content

binoc_sdk/
data_access.rs

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
8/// In-process DataAccess backed by the local filesystem, temp directories,
9/// and a filesystem-backed artifact store under `data_root/.artifacts/`.
10///
11/// Construction modes:
12///
13/// - [`Self::new`] — unrestricted paths (tests, ad-hoc tooling).
14/// - [`Self::new_for_diff`] — paths must stay under the two snapshot trees or
15///   session workspace (used by the controller).
16/// - [`Self::for_plugin`] — shares the host's `data_root` for artifact access,
17///   plus a pre-allocated workspace for expansion (C ABI plugins).
18/// - [`Self::with_data_root`] — shares an existing `data_root` for artifact
19///   reads only (no expansion workspace).
20pub 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
110/// True when `path` is `root` or a descendant (component-wise).
111fn 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    /// Session-backed access with path confinement: filesystem reads and
160    /// `register_local` targets must lie under the snapshot roots, the session
161    /// `data_root`, or a workspace / provide directory created by this instance.
162    #[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    /// Create a LocalDataAccess for a plugin running across the C ABI.
207    /// Shares the host's `data_root` for cache access and uses `workspace`
208    /// for expansion (provide, workspace calls).
209    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    /// Create a LocalDataAccess that can only read from an existing data_root
222    /// cache. No workspace for expansion. Used during extract-only access.
223    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    /// Enforce policy for a path that must already exist on disk (e.g. `register_local`).
269    fn enforce_path_policy(&self, physical: &Path) -> BinocResult<()> {
270        let resolved = policy_path(physical)?;
271        self.enforce_path_policy_resolved(&resolved)
272    }
273
274    /// Enforce policy before reading; allows missing leaf paths under an allowed directory.
275    #[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}